diff --git a/src/subagent/agent-fleet.test.ts b/src/subagent/agent-fleet.test.ts index a5d96f0f0..0a126afeb 100644 --- a/src/subagent/agent-fleet.test.ts +++ b/src/subagent/agent-fleet.test.ts @@ -1502,13 +1502,16 @@ describe("interrupt_agent unblocks wait_agents", () => { test("close overlay without in-flight wait stays interrupted after followup complete", async () => { const gate = deferred(); - const followupGate = deferred(); const closeHold = deferred(); + const followupCalls: string[] = []; const deps = makeDeps(async (params) => { params.onAgentReady?.({ close: async () => closeHold.promise, interrupt: () => undefined, - followup: async () => followupGate.promise, + followup: async (message: string) => { + followupCalls.push(message); + return "should never run"; + }, deliver: () => undefined, }); return gate.promise; @@ -1552,25 +1555,16 @@ describe("interrupt_agent unblocks wait_agents", () => { new AbortController().signal, ); - followupGate.resolve("followup after close overlay"); - await new Promise((resolve) => { - const done = (): boolean => - deps.sessions.get(id)?.lifecycle.state === "completed"; - if (done()) { - resolve(); - return; - } - const unsub = deps.sessions.subscribe(() => { - if (done()) { - unsub(); - resolve(); - } - }); - if (done()) { - unsub(); - resolve(); - } - }); + closeHold.resolve(undefined); + await closing; + // CL-7344: close drops the stashed follow-up — it never runs against + // the closed agent, and the close overlay survives the run settling. + expect(followupCalls).toEqual([]); + gate.resolve({ + report: "original interrupted", + interrupted: true, + } as RunSubAgentResult); + await new Promise((resolve) => setTimeout(resolve, 20)); expect(deps.fleetRecords.peek(id)?.status).toBe("interrupted"); const listed = await callTool(list, {}); @@ -1583,13 +1577,6 @@ describe("interrupt_agent unblocks wait_agents", () => { expect(waited.timed_out).toBe(false); const results = waited.results as { status: string }[]; expect(defined(results[0]).status).toBe("interrupted"); - - closeHold.resolve(undefined); - await closing; - gate.resolve({ - report: "original interrupted", - interrupted: true, - } as RunSubAgentResult); }); test("completeAfterInterrupt does not clear a close overlay", () => { @@ -1642,14 +1629,17 @@ describe("interrupt_agent unblocks wait_agents", () => { message: "stop that", interrupt: true, }); - followupGate.reject(new Error("followup failed")); - await new Promise((resolve) => setTimeout(resolve, 20)); + // CL-7344: the follow-up is stashed until the original run settles; the + // salvage handoff launches it, and its rejection restores the + // interrupted lane so wait collects the salvage. gate.resolve({ report: "## Summary\nStopped.\n## Findings\nsalvage\n## Blockers\ninterrupted\n## Paths\n", interrupted: true, } as RunSubAgentResult); await new Promise((resolve) => setTimeout(resolve, 20)); + followupGate.reject(new Error("followup failed")); + await new Promise((resolve) => setTimeout(resolve, 20)); const waited = await callTool(wait, { targets: [id], timeout_ms: 5000 }); expect(waited.timed_out).toBe(false); @@ -1658,21 +1648,14 @@ describe("interrupt_agent unblocks wait_agents", () => { expect(defined(results[0]).report).toContain("salvage"); }); - test("send_input interrupt queued overlay clears when the followup is admitted", async () => { - const admission = createAdmissionQueue({ capacity: 1 }); - const sessions = createSubAgentSessionStore({ admission }); + test("send_input interrupt on a live lane stashes the followup until the run settles", async () => { + const sessions = createSubAgentSessionStore(); const fleetRecords = createFleetMailbox(sessions); - admission.enqueue({ - id: "holder", - provider: "p", - start: () => undefined, - }); const worker = sessions.start({ description: "looping", agentId: "explorer", brief: "b", retained: true, - provider: "p", }); sessions.markRunning(worker.id); fleetRecords.register(worker.id); @@ -1690,51 +1673,44 @@ describe("interrupt_agent unblocks wait_agents", () => { interrupt: true, }); expect(sent.status).toBe("interrupted"); - expect(sessions.get(worker.id)?.lifecycleStatus).toBe("pending_init"); + // CL-7344: the interrupt stashes the follow-up instead of starting it; + // the lane is still the live run, so the session stays running while + // wait/list stay live until the run settles and hands off. + expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running"); - const queuedWait = await callTool(wait, { + const liveWait = await callTool(wait, { targets: [worker.id], timeout_ms: 50, }); - expect(queuedWait.timed_out).toBe(true); - expect( - defined((queuedWait.results as { status: string }[])[0]).status, - ).toBe("queued"); - const queuedList = await callTool(list, {}); - const queuedEntry = ( - queuedList.agents as { + expect(liveWait.timed_out).toBe(true); + expect(defined((liveWait.results as { status: string }[])[0]).status).toBe( + "running", + ); + const liveList = await callTool(list, {}); + const liveEntry = ( + liveList.agents as { agent_id: string; status: string; lifecycle: string; }[] ).find((a) => a.agent_id === worker.id); - expect(queuedEntry?.status).toBe("queued"); - expect(queuedEntry?.lifecycle).toBe("pending_init"); - - admission.release("holder"); - await new Promise((resolve) => setTimeout(resolve, 20)); + expect(liveEntry?.status).toBe("running"); + expect(liveEntry?.lifecycle).toBe("running"); + // Settling the original run launches the stashed follow-up. + sessions.attachReport(worker.id, "interrupted salvage", { + stopReason: "interrupted", + }); expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running"); - const runningWait = await callTool(wait, { + followupGate.resolve("later"); + const done = await callTool(wait, { targets: [worker.id], - timeout_ms: 50, + timeout_ms: 5000, }); - expect(runningWait.timed_out).toBe(true); - expect( - defined((runningWait.results as { status: string }[])[0]).status, - ).toBe("running"); - const runningList = await callTool(list, {}); - const runningEntry = ( - runningList.agents as { - agent_id: string; - status: string; - lifecycle: string; - }[] - ).find((a) => a.agent_id === worker.id); - expect(runningEntry?.status).toBe("running"); - expect(runningEntry?.lifecycle).toBe("running"); - - followupGate.resolve("later"); + expect(done.timed_out).toBe(false); + const results = done.results as { status: string; report?: string }[]; + expect(defined(results[0]).status).toBe("done"); + expect(defined(results[0]).report).toBe("later"); }); test("soft-interrupt wait path collects so omitted re-wait does not re-deliver", async () => { @@ -1910,8 +1886,13 @@ describe("interrupt_agent unblocks wait_agents", () => { message: "stop that", interrupt: true, }); + // CL-7344: the follow-up is stashed until the original run settles; the + // salvage handoff launches it, so the run must settle first. + gate.resolve({ + report: "original interrupted", + interrupted: true, + } as RunSubAgentResult); followupGate.resolve("followup report"); - await new Promise((resolve) => setTimeout(resolve, 20)); const waited = await callTool(wait, { targets: [id], timeout_ms: 5000 }); expect(waited.timed_out).toBe(false); const results = waited.results as { status: string; report?: string }[]; diff --git a/src/subagent/agent-fleet.ts b/src/subagent/agent-fleet.ts index 936ecbe01..709853562 100644 --- a/src/subagent/agent-fleet.ts +++ b/src/subagent/agent-fleet.ts @@ -249,11 +249,8 @@ class FleetMailbox { /** * CL-7331: mark that a send_input interrupt:true followup owns this lane. * Suppresses any interrupt overlay so wait stays live (running/queued) - * until the followup settles, and tells the spawn settlement to swallow - * the original run's interrupted result instead of attaching salvage over - * the live followup. No-op on an unknown id; safe to call on a collected - * mailbox (frozen status still wins for projection, but the settlement - * swallow still applies). + * until the followup settles. No-op on an unknown id; safe to call on a + * collected mailbox (frozen status still wins for projection). */ noteFollowup(id: string): void { const existing = this.records.get(id); @@ -263,11 +260,6 @@ class FleetMailbox { this.sessions?.wake(); } - /** True while a send_input interrupt:true followup owns this lane. */ - hasLiveFollowup(id: string): boolean { - return this.records.get(id)?.followupLive === true; - } - /** * send_input interrupt:true followup finished. Clear the followup lane flag * (and any admission queued overlay) so wait can project the settled session. @@ -1387,19 +1379,17 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool { if (result.interrupted === true) { keepWorktreeAlive = true; runInterrupted = true; - // CL-7331: a followup started via send_input interrupt owns - // this lane now — the settling original turn must not attach - // salvage over it. The mailbox flag (set by send_input, not - // by interrupt_agent or a bare settle) identifies that lane; - // lifecycle alone cannot, since a never-started run and a - // queued followup both read pending_init. - if (!deps.fleetRecords.hasLiveFollowup(session.id)) { - deps.sessions.attachReport(session.id, result.report, { - ...(result.stopReason !== undefined - ? { stopReason: result.stopReason } - : {}), - }); - } + // CL-7344: a followup stashed via send_input interrupt launches + // from attachReport, so the settling original turn must + // attach its salvage here — attachReport records the salvage + // and launches the stashed followup in one atomic step. + // Skipping the attach would strand the stash and hang + // wait_agents on a phantom turn. + deps.sessions.attachReport(session.id, result.report, { + ...(result.stopReason !== undefined + ? { stopReason: result.stopReason } + : {}), + }); return; } const alreadyCancelled = diff --git a/src/subagent/lifecycle-tools.test.ts b/src/subagent/lifecycle-tools.test.ts index d70b38ad1..91f64145c 100644 --- a/src/subagent/lifecycle-tools.test.ts +++ b/src/subagent/lifecycle-tools.test.ts @@ -7,7 +7,11 @@ import { createSendInputTool, resumeAgentToolDefinition, } from "./lifecycle-tools.js"; -import { createFleetMailbox, createWaitAgentsTool } from "./agent-fleet.js"; +import { + createFleetMailbox, + createListAgentsTool, + createWaitAgentsTool, +} from "./agent-fleet.js"; import { createSubAgentSessionStore, DEFAULT_MAX_ENTRY_CHARS, @@ -28,6 +32,7 @@ async function callTool( | ReturnType | ReturnType | ReturnType + | ReturnType | ReturnType, args: Record, ): Promise> { @@ -795,6 +800,9 @@ describe("send_input", () => { expect(sessions.get(worker.id)?.stopReason).toBe("interrupted"); expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running"); + // CL-7344: the interrupt stashes the follow-up until the original run + // settles; the salvage handoff launches it. + sessions.attachReport(worker.id, "salvage", { stopReason: "interrupted" }); finish("followup report"); await new Promise((resolve) => setTimeout(resolve, 0)); const collected = await callTool(wait, { @@ -836,6 +844,10 @@ describe("send_input", () => { message: "stop that", interrupt: true, }); + // CL-7344: the interrupt stashes the follow-up until the original run + // settles; the salvage handoff launches it, it rejects, and the session + // restamps interrupted. + sessions.attachReport(worker.id, "salvage", { stopReason: "interrupted" }); await new Promise((resolve) => setTimeout(resolve, 0)); const collected = await callTool(wait, { targets: [worker.id], @@ -907,6 +919,12 @@ describe("send_input", () => { }); expect(result).toEqual({ agent_id: worker.id, status: "interrupted" }); expect(interrupted).toBe(true); + // CL-7344: the interrupt stashes the follow-up until the original run + // settles; the salvage handoff launches it. + expect(followupStarted).toBe(false); + expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running"); + sessions.attachReport(worker.id, "salvage", { stopReason: "interrupted" }); + await new Promise((resolve) => setTimeout(resolve, 0)); expect(followupStarted).toBe(true); expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running"); expect(sessions.get(worker.id)?.finishedAt).toBeUndefined(); @@ -932,6 +950,164 @@ describe("send_input", () => { expect(sessions.get(missing.id)?.lifecycleStatus).toBe("running"); }); + test("completion during interrupt keeps the original terminal report", async () => { + const sessions = createSubAgentSessionStore(); + const fleetRecords = createFleetMailbox(sessions); + const worker = sessions.start({ + description: "worker", + agentId: "a", + brief: "b", + retained: true, + }); + sessions.markRunning(worker.id); + sessions.markRunInFlight(worker.id); + sessions.registerInterrupt(worker.id, () => { + sessions.complete(worker.id, "original report"); + }); + let followupStarted = false; + sessions.registerFollowup(worker.id, async () => { + followupStarted = true; + return "follow-up report"; + }); + fleetRecords.register(worker.id); + + const sendInput = createSendInputTool({ sessions, fleetRecords }); + const wait = createWaitAgentsTool({ sessions, fleetRecords }); + const result = await callTool(sendInput, { + target: worker.id, + message: "late interrupt", + interrupt: true, + }); + expect(result).toEqual({ agent_id: worker.id, status: "completed" }); + + const collected = await callTool(wait, { + targets: [worker.id], + timeout_ms: 1000, + }); + expect(collected.timed_out).toBe(false); + expect(collected.results).toEqual([ + expect.objectContaining({ + agent_id: worker.id, + status: "done", + report: "original report", + }), + ]); + expect(followupStarted).toBe(false); + }); + + test("CL-7344: interrupt:true stashes until attachReport; resume stays fail-closed", async () => { + const sessions = createSubAgentSessionStore(); + const fleetRecords = createFleetMailbox(sessions); + const worker = sessions.start({ + description: "worker", + agentId: "a", + brief: "b", + retained: true, + }); + sessions.markRunning(worker.id); + sessions.markRunInFlight(worker.id); + sessions.registerInterrupt(worker.id, () => undefined); + let followupStarted = false; + sessions.registerFollowup(worker.id, async (message) => { + followupStarted = true; + expect(message).toBe("patch only the test"); + return "queued turn finished"; + }); + fleetRecords.register(worker.id); + + const sendInput = createSendInputTool({ sessions, fleetRecords }); + const resume = createResumeAgentTool({ sessions, fleetRecords }); + const wait = createWaitAgentsTool({ sessions, fleetRecords }); + const result = await callTool(sendInput, { + target: worker.id, + message: "patch only the test", + interrupt: true, + }); + expect(result).toEqual({ agent_id: worker.id, status: "interrupted" }); + expect(followupStarted).toBe(false); + expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running"); + if (resume.kind !== "full") throw new Error("expected full tool"); + const resumed = await resume.handler( + { + id: "resume-during-stash", + name: "resume_agent", + arguments: { target: worker.id, message: "x" }, + }, + new AbortController().signal, + ); + expect(resumed.isError).toBe(true); + expect(String(resumed.content)).toContain("status: running"); + const pending = await callTool(wait, { + targets: [worker.id], + timeout_ms: 50, + }); + expect(pending.timed_out).toBe(true); + expect(defined((pending.results as { status: string }[])[0]).status).toBe( + "running", + ); + + sessions.attachReport(worker.id, "salvage", { stopReason: "interrupted" }); + await Promise.resolve(); + expect(followupStarted).toBe(true); + }); + + test("CL-7344: AgentClosedError follow-up wait collects failed", async () => { + const { AgentClosedError } = await import("@intx/agent"); + const sessions = createSubAgentSessionStore(); + const fleetRecords = createFleetMailbox(sessions); + const worker = sessions.start({ + description: "worker", + agentId: "a", + brief: "b", + retained: true, + }); + sessions.markRunning(worker.id); + sessions.markRunInFlight(worker.id); + sessions.registerInterrupt(worker.id, () => undefined); + sessions.registerFollowup(worker.id, async () => { + throw new AgentClosedError(); + }); + fleetRecords.register(worker.id); + const sendInput = createSendInputTool({ sessions, fleetRecords }); + const wait = createWaitAgentsTool({ sessions, fleetRecords }); + const list = createListAgentsTool({ sessions, fleetRecords }); + const resume = createResumeAgentTool({ sessions, fleetRecords }); + await callTool(sendInput, { + target: worker.id, + message: "stop that", + interrupt: true, + }); + sessions.attachReport(worker.id, "salvage", { stopReason: "interrupted" }); + await new Promise((resolve) => setTimeout(resolve, 0)); + const collected = await callTool(wait, { + targets: [worker.id], + timeout_ms: 1000, + }); + expect(collected.timed_out).toBe(false); + const results = collected.results as { status: string; error?: string }[]; + expect(defined(results[0]).status).toBe("failed"); + expect(defined(results[0]).error).toContain("closed"); + + const listed = await callTool(list, {}); + const listedWorker = ( + listed.agents as { agent_id: string; status: string; lifecycle: string }[] + ).find((agent) => agent.agent_id === worker.id); + expect(listedWorker?.status).toBe("failed"); + expect(listedWorker?.lifecycle).toBe("shutdown"); + + if (resume.kind !== "full") throw new Error("expected full tool"); + const resumed = await resume.handler( + { + id: "resume-closed-followup", + name: "resume_agent", + arguments: { target: worker.id, message: "retry" }, + }, + new AbortController().signal, + ); + expect(resumed.isError).toBe(true); + expect(String(resumed.content)).toContain("status: shutdown"); + }); + test("rejects completed, interrupted, and closed sessions — steering is in-flight only", async () => { const sessions = createSubAgentSessionStore(); const fleetRecords = createFleetMailbox(sessions); diff --git a/src/subagent/lifecycle-tools.ts b/src/subagent/lifecycle-tools.ts index 6d38dd3eb..19ab71041 100644 --- a/src/subagent/lifecycle-tools.ts +++ b/src/subagent/lifecycle-tools.ts @@ -470,7 +470,11 @@ export function createSendInputTool(deps: LifecycleToolDeps): AgentTool { // wait_agents unblocks as interrupted; a queued followup must instead // stay wait-live (running/queued) so the followup reply surfaces via // wait_agents instead of freezing as an already-collected interrupt. - if (interrupt && deps.fleetRecords !== undefined) { + if ( + interrupt && + outcome.status === "interrupted" && + deps.fleetRecords !== undefined + ) { deps.fleetRecords.noteFollowup(target); const after = deps.sessions.get(target); if (after?.lifecycle.state === "pending_init") diff --git a/src/subagent/session-store.test.ts b/src/subagent/session-store.test.ts index 0c92368d0..244ee206e 100644 --- a/src/subagent/session-store.test.ts +++ b/src/subagent/session-store.test.ts @@ -594,6 +594,12 @@ describe("CL-6943 reusable worker sessions", () => { ok: true, status: "interrupted", }); + // CL-7344: the interrupt stashes the follow-up until the original run + // settles; the salvage handoff launches it, it rejects, and the session + // restamps interrupted. + store.attachReport(session.id, "interrupted salvage", { + stopReason: "interrupted", + }); await new Promise((resolve) => setTimeout(resolve, 0)); const after = store.get(session.id); expect(after?.lifecycle.state).toBe("interrupted"); @@ -1465,6 +1471,11 @@ describe("pending ask_director", () => { expect(rejected).toBeInstanceOf(Error); expect(String(rejected)).toContain("cancelled by send_input interrupt"); expect(delivered).toEqual([]); + // CL-7344: the interrupt stashes the follow-up until the original run + // settles; the salvage handoff launches it. + store.attachReport(session.id, "interrupted salvage", { + stopReason: "interrupted", + }); await new Promise((resolve) => setTimeout(resolve, 0)); expect(followups).toEqual(["stop that"]); }); @@ -1771,3 +1782,215 @@ describe("pending ask_director", () => { expect(String(rejected)).toContain("session completed"); }); }); + +describe("CL-7344 follow-up stash", () => { + function runningRetained( + store: ReturnType, + followup: (message: string) => Promise, + ) { + const session = store.start({ + description: "d", + agentId: "a", + brief: "b", + retained: true, + }); + store.markRunning(session.id); + store.markRunInFlight(session.id); + store.registerInterrupt(session.id, () => undefined); + store.registerFollowup(session.id, followup); + return session; + } + + test("sendInputOne interrupt stashes and does not start follow-up until attachReport", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return "followup"; + }); + + expect( + store.sendInputOne(session.id, "steer now", { interrupt: true }), + ).toEqual({ ok: true, status: "interrupted" }); + await Promise.resolve(); + expect(started).toEqual([]); + expect(store.get(session.id)?.lifecycleStatus).toBe("running"); + expect(store.resumeOne(session.id, "later").ok).toBe(false); + + store.attachReport(session.id, "interrupted salvage", { + stopReason: "interrupted", + }); + await Promise.resolve(); + expect(started).toEqual(["steer now"]); + expect(store.get(session.id)?.lifecycleStatus).toBe("running"); + expect(store.isRunInFlight(session.id)).toBe(true); + }); + + test("complete drops a stashed follow-up and keeps the original report", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return "should not run"; + }); + store.sendInputOne(session.id, "steer now", { interrupt: true }); + store.complete(session.id, "## Summary\nOriginal done."); + await Promise.resolve(); + expect(started).toEqual([]); + expect(store.get(session.id)?.report).toBe("## Summary\nOriginal done."); + expect(store.get(session.id)?.lifecycleStatus).toBe("completed"); + expect(store.isRunInFlight(session.id)).toBe(false); + }); + + test("attachReport interrupted starts follow-up in the same notify as clearing the original run", async () => { + const store = createSubAgentSessionStore(); + let started = 0; + const session = runningRetained(store, async () => { + started += 1; + return "followup"; + }); + const gaps: string[] = []; + store.subscribe(() => { + const status = store.get(session.id)?.lifecycleStatus ?? "missing"; + if ( + (!store.isRunInFlight(session.id) || status !== "running") && + started === 0 + ) { + gaps.push(status); + } + }); + store.sendInputOne(session.id, "steer now", { interrupt: true }); + expect(started).toBe(0); + store.attachReport(session.id, "salvage", { stopReason: "interrupted" }); + expect(started).toBe(1); + expect(gaps).toEqual([]); + expect(store.isRunInFlight(session.id)).toBe(true); + expect(store.get(session.id)?.lifecycleStatus).toBe("running"); + }); + + test("last send_input interrupt overwrites the stash", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return "followup"; + }); + store.sendInputOne(session.id, "first", { interrupt: true }); + store.sendInputOne(session.id, "second", { interrupt: true }); + store.attachReport(session.id, "salvage", { stopReason: "interrupted" }); + await Promise.resolve(); + expect(started).toEqual(["second"]); + }); + + test("fail, cancel, close, interrupt_agent, and settleRun drop the stash", async () => { + const make = ( + followup: (message: string) => Promise, + ): ReturnType => { + const store = createSubAgentSessionStore(); + runningRetained(store, followup); + return store; + }; + + const failStarted: string[] = []; + const failStore = make(async (m) => { + failStarted.push(m); + return "x"; + }); + const failId = defined(failStore.list()[0]).id; + failStore.sendInputOne(failId, "steer", { interrupt: true }); + failStore.fail(failId, "boom"); + failStore.attachReport(failId, "salvage", { stopReason: "interrupted" }); + await Promise.resolve(); + expect(failStarted).toEqual([]); + + const cancelStarted: string[] = []; + const cancelStore = make(async (m) => { + cancelStarted.push(m); + return "x"; + }); + const cancelId = defined(cancelStore.list()[0]).id; + cancelStore.sendInputOne(cancelId, "steer", { interrupt: true }); + cancelStore.cancel(cancelId); + cancelStore.attachReport(cancelId, "salvage", { + stopReason: "interrupted", + }); + await Promise.resolve(); + expect(cancelStarted).toEqual([]); + + const closeStarted: string[] = []; + const closeStore = make(async (m) => { + closeStarted.push(m); + return "x"; + }); + const closeId = defined(closeStore.list()[0]).id; + closeStore.registerClose(closeId, async () => undefined); + closeStore.sendInputOne(closeId, "steer", { interrupt: true }); + await closeStore.closeOne(closeId, 1000); + closeStore.attachReport(closeId, "salvage", { stopReason: "interrupted" }); + await Promise.resolve(); + expect(closeStarted).toEqual([]); + + const interruptStarted: string[] = []; + const interruptStore = make(async (m) => { + interruptStarted.push(m); + return "x"; + }); + const interruptId = defined(interruptStore.list()[0]).id; + interruptStore.sendInputOne(interruptId, "steer", { interrupt: true }); + expect(interruptStore.interruptOne(interruptId).ok).toBe(true); + interruptStore.attachReport(interruptId, "salvage", { + stopReason: "interrupted", + }); + await Promise.resolve(); + expect(interruptStarted).toEqual([]); + + const settleStarted: string[] = []; + const settleStore = make(async (m) => { + settleStarted.push(m); + return "x"; + }); + const settleId = defined(settleStore.list()[0]).id; + settleStore.sendInputOne(settleId, "steer", { interrupt: true }); + settleStore.settleRun(settleId); + settleStore.attachReport(settleId, "salvage", { + stopReason: "interrupted", + }); + await Promise.resolve(); + expect(settleStarted).toEqual([]); + }); + + test("AgentClosedError on an interrupt-won follow-up fails terminal", async () => { + const { AgentClosedError } = await import("@intx/agent"); + const store = createSubAgentSessionStore(); + const session = runningRetained(store, async () => { + throw new AgentClosedError(); + }); + store.sendInputOne(session.id, "steer now", { interrupt: true }); + store.attachReport(session.id, "salvage", { stopReason: "interrupted" }); + await new Promise((resolve) => setTimeout(resolve, 0)); + const after = store.get(session.id); + expect(after?.lifecycle.state).toBe("failed"); + expect(after?.error).toContain("closed"); + expect(store.isRunInFlight(session.id)).toBe(false); + }); + + test("AgentClosedError after close_agent shutdown does not rewrite close", async () => { + const { AgentClosedError } = await import("@intx/agent"); + const store = createSubAgentSessionStore(); + let followup: (message: string) => Promise = async () => "x"; + const session = runningRetained(store, (message) => followup(message)); + store.registerClose(session.id, async () => undefined); + let rejectFollowup: (err: unknown) => void = () => undefined; + followup = () => + new Promise((_, reject) => { + rejectFollowup = reject; + }); + store.sendInputOne(session.id, "steer now", { interrupt: true }); + store.attachReport(session.id, "salvage", { stopReason: "interrupted" }); + await store.closeOne(session.id, 1000); + expect(store.get(session.id)?.lifecycle.state).toBe("shutdown"); + rejectFollowup(new AgentClosedError()); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(store.get(session.id)?.lifecycle.state).toBe("shutdown"); + }); +}); diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index 016f97a63..1441c8a23 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -5,6 +5,7 @@ // this store is the dedicated child record the enter-session UI reads. import type { ReactorEmittedEvent } from "@intx/inference"; +import { AgentClosedError } from "@intx/agent"; import { getLogger } from "@intx/log"; import { LOG_NAMESPACE_ROOT } from "../branding.js"; import { awaitBoundedTeardown, DEFAULT_CLOSE_DEADLINE_MS } from "./dispose.js"; @@ -510,6 +511,21 @@ export function createSubAgentSessionStore( (message: string) => Promise >(); const deliverHandles = new Map void>(); + // CL-7344: a send_input interrupt that lands while the original run is still + // in flight must not start its follow-up against a run that is about to + // settle. The message is stashed here and launched atomically from + // attachReport (same mutation/notify as the salvage handoff) when the run's + // report arrives; a completing run drops it. Any terminal transition — + // interrupt, close, cancel, fail, eviction — drops it too, so a queued + // follow-up can never run against a closed agent. + interface StashedFollowup { + message: string; + failLifecycle: "completed" | "interrupted"; + onStart?: () => void; + onReply?: (reply: string) => void; + onFail?: (error: unknown) => void; + } + const stashedFollowups = new Map(); const pendingAsks = new Map< string, { @@ -639,6 +655,7 @@ export function createSubAgentSessionStore( const cancelSession = (id: string, reason: string): boolean => { settleCancelsAsks(id, reason); + stashedFollowups.delete(id); const session = sessions.get(id); if (session === undefined || !isLiveStrip(session.lifecycle)) return false; const abort = cancelHandles.get(id); @@ -686,6 +703,7 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); + stashedFollowups.delete(id); }; // An open retained session (spawn_agent's reusable-session contract: @@ -849,6 +867,97 @@ export function createSubAgentSessionStore( s.finishedAt = now(); }); }; + // CL-7344: shared terminal transition used when a queued follow-up rejects + // because the agent closed mid-invocation. Session and fleet move together: + // endFollowupTurn restores the visible fleet state while fail records the + // actionable error, so wait_agents, list_agents, and resume_agent agree. + const failClosedSession = (id: string, error: string): void => { + endFollowupTurn(id, "interrupted"); + settleCancelsAsks(id, "session failed"); + mutate(id, (session) => { + if ( + !isLiveStrip(session.lifecycle) || + session.lifecycle.state === "cancelled" + ) { + return; + } + session.lifecycle = { state: "failed", error }; + session.finishedAt = now(); + session.lastActivityAt = now(); + session.error = error; + session.stopReason = "error"; + clearToolCalls(session); + pushEntry(session, { + kind: "report", + content: capText(`Error: ${error}`, maxEntryChars), + }); + }); + releaseHandles(id); + pruneCompleted(); + }; + // CL-7344: shared follow-up settlement, used both by queueFollowupTurn's + // immediate/queued start and by the stashed-interrupt launch in + // attachReport, so every follow-up completes, fails, and releases its + // admission slot the same way. + const settleFollowupReply = ( + id: string, + reply: string, + onReply?: (reply: string) => void, + ): void => { + const still = sessions.get(id); + if (still === undefined) { + runInFlight.delete(id); + return; + } + if ( + still.lifecycle.state === "shutdown" || + still.lifecycle.state === "cancelled" || + still.lifecycle.state === "failed" || + still.lifecycle.state === "interrupted" + ) { + runInFlight.delete(id); + return; + } + mutate(id, (s) => { + s.lifecycle = { state: "completed", report: reply }; + s.finishedAt = now(); + s.report = reply; + delete s.stopReason; + pushEntry(s, { + kind: "report", + content: capText(reply, maxEntryChars), + }); + }); + runInFlight.delete(id); + onReply?.(reply); + pruneRetained(); + }; + const settleFollowupFailure = ( + id: string, + err: unknown, + failLifecycle: "completed" | "interrupted", + onFail?: (error: unknown) => void, + ): void => { + runInFlight.delete(id); + onFail?.(err); + // CL-7344: the agent closed between queueing and invocation, so the + // follow-up can never run. Move session and fleet records to the + // same terminal state with an actionable error instead of silently + // restoring the stale interrupted snapshot. + if (err instanceof AgentClosedError) { + failClosedSession( + id, + `Follow-up rejected: agent ${id} closed before the message could be delivered.`, + ); + log.error("followup rejected for closed agent {id}", { id }); + return; + } + endFollowupTurn(id, failLifecycle); + log.error("followup turn failed for {id}: {error}", { + id, + error: err instanceof Error ? err.message : String(err), + }); + }; const queueFollowupTurn = ( id: string, message: string, @@ -872,42 +981,10 @@ export function createSubAgentSessionStore( opts?.onStart?.(); void followup(message) .then((reply) => { - const still = sessions.get(id); - if (still === undefined) { - runInFlight.delete(id); - return; - } - if ( - still.lifecycle.state === "shutdown" || - still.lifecycle.state === "cancelled" || - still.lifecycle.state === "failed" || - still.lifecycle.state === "interrupted" - ) { - runInFlight.delete(id); - return; - } - mutate(id, (s) => { - s.lifecycle = { state: "completed", report: reply }; - s.finishedAt = now(); - s.report = reply; - delete s.stopReason; - pushEntry(s, { - kind: "report", - content: capText(reply, maxEntryChars), - }); - }); - runInFlight.delete(id); - opts?.onReply?.(reply); - pruneRetained(); + settleFollowupReply(id, reply, opts?.onReply); }) .catch((err: unknown) => { - runInFlight.delete(id); - opts?.onFail?.(err); - endFollowupTurn(id, failLifecycle); - log.error("followup turn failed for {id}: {error}", { - id, - error: err instanceof Error ? err.message : String(err), - }); + settleFollowupFailure(id, err, failLifecycle, opts?.onFail); }) .finally(() => { if (takesSlot) queue?.release(id); @@ -934,6 +1011,33 @@ export function createSubAgentSessionStore( } return status; }; + // attachReport moves the lifecycle to running before calling this, keeping + // the interrupted run and follow-up handoff atomic to observers. + const launchStashedFollowup = ( + id: string, + stashed: StashedFollowup, + ): void => { + const followup = followupHandles.get(id); + // The follow-up handle can only be gone if teardown raced the handoff; + // then there is no follow-up to inherit the run, so settle it instead of + // leaving wait_agents stuck on a phantom turn. + if (followup === undefined || sessions.get(id) === undefined) { + runInFlight.delete(id); + return; + } + stashed.onStart?.(); + const pending = followup(stashed.message); + void Promise.resolve().then(() => { + void pending.then( + (reply) => { + settleFollowupReply(id, reply, stashed.onReply); + }, + (err: unknown) => { + settleFollowupFailure(id, err, stashed.failLifecycle, stashed.onFail); + }, + ); + }); + }; return { list(): readonly SubAgentSession[] { @@ -969,6 +1073,7 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); + stashedFollowups.delete(id); pinCounts.delete(id); runInFlight.delete(id); forgetRevision(id); @@ -1239,6 +1344,7 @@ export function createSubAgentSessionStore( cancelHandles.delete(id); if (!agentRetained) closeHandles.delete(id); runInFlight.delete(id); + stashedFollowups.delete(id); pruneCompleted(); pruneRetained(); }); @@ -1318,6 +1424,10 @@ export function createSubAgentSessionStore( return "not_found"; } settleCancelsAsks(id, "session closed"); + // CL-7344: a stashed send_input follow-up must never launch against a + // closing agent; dropped again in each terminal path below in case the + // stash lands during the setup-window wait. + stashedFollowups.delete(id); let close = closeHandles.get(id); const alreadyClosed = isAlreadyClosed(session.lifecycle); if (alreadyClosed && close === undefined) { @@ -1345,6 +1455,7 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); + stashedFollowups.delete(id); runInFlight.delete(id); pruneCompleted(); return "shutdown"; @@ -1386,6 +1497,7 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); + stashedFollowups.delete(id); runInFlight.delete(id); pruneCompleted(); if (closeError !== undefined) throw closeError; @@ -1413,6 +1525,7 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); + stashedFollowups.delete(id); runInFlight.delete(id); pruneCompleted(); if (closeError !== undefined) throw closeError; @@ -1465,17 +1578,31 @@ export function createSubAgentSessionStore( }; } settleCancelsAsks(id, "cancelled by send_input interrupt"); - interrupt(); - queueFollowupTurn(id, message, "interrupted", { + const stashed: StashedFollowup = { + message, + failLifecycle: "interrupted", ...(opts.onStart !== undefined ? { onStart: opts.onStart } : {}), ...(opts.onFollowupReply !== undefined ? { onReply: opts.onFollowupReply } : {}), ...(opts.onFail !== undefined ? { onFail: opts.onFail } : {}), - }); - // After beginFollowupTurn, which clears leftover stopReason. Stamp - // here so an in-flight wait_agents overlay can project interrupted - // without flipping lifecycle off the live follow-up. + }; + // Stash before interrupting because interrupt callbacks may settle the + // original run synchronously. That terminal transition consumes the + // stash before control returns here, so it cannot be resurrected. + stashedFollowups.set(id, stashed); + interrupt(); + if (stashedFollowups.get(id) !== stashed) { + const settled = sessions.get(id); + const status = + settled === undefined + ? "not_found" + : projectLifecycleStatus(settled.lifecycle); + return { + ok: true, + status: status === "running" ? "interrupted" : status, + }; + } mutate(id, (s) => { s.stopReason = "interrupted"; }); @@ -1550,6 +1677,9 @@ export function createSubAgentSessionStore( if (!isLiveStrip(session.lifecycle)) { return { ok: false, status: projectLifecycleStatus(session.lifecycle) }; } + // CL-7344: interrupt_agent settles the run itself, so a stashed + // send_input follow-up must not launch from a later attachReport. + stashedFollowups.delete(id); const interrupt = interruptHandles.get(id); if (interrupt === undefined) { if (session.lifecycle.state === "pending_init") { @@ -1673,6 +1803,7 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); + stashedFollowups.delete(id); mutate(id, (s) => { s.lifecycle = { state: "shutdown", @@ -1716,6 +1847,12 @@ export function createSubAgentSessionStore( report: string, opts?: { stopReason?: ForcedStopReason }, ): void { + // Consume the stashed interrupt follow-up before mutating. Interrupted + // salvage moves directly to the next running turn; any terminal outcome + // drops the stash so no follow-up runs against a closed agent. + const stashed = stashedFollowups.get(id); + stashedFollowups.delete(id); + let toLaunch: StashedFollowup | undefined; mutate(id, (session) => { const state = session.lifecycle.state; if (state === "completed" || state === "failed") { @@ -1732,6 +1869,14 @@ export function createSubAgentSessionStore( kind: "report", content: capText(report, maxEntryChars), }); + if (stashed !== undefined) { + // Move directly into the follow-up lifecycle before mutate notifies + // subscribers. No observer can resume the interrupted handoff. + toLaunch = stashed; + session.lifecycle = { state: "running" }; + delete session.finishedAt; + delete session.stopReason; + } } else if ( (state === "cancelled" || state === "interrupted" || @@ -1747,10 +1892,14 @@ export function createSubAgentSessionStore( content: capText(report, maxEntryChars), }); } - runInFlight.delete(id); + // A launched follow-up inherits the original run's in-flight marker. + if (toLaunch === undefined) runInFlight.delete(id); pruneCompleted(); pruneRetained(); }); + if (toLaunch !== undefined && sessions.has(id)) { + launchStashedFollowup(id, toLaunch); + } }, isRunInFlight(id: string): boolean { @@ -1759,6 +1908,7 @@ export function createSubAgentSessionStore( settleRun(id: string): void { cancelAskInternal(id, "run settled"); + stashedFollowups.delete(id); if (!runInFlight.delete(id)) return; notify(); }, @@ -1786,6 +1936,7 @@ export function createSubAgentSessionStore( interruptHandles.clear(); followupHandles.clear(); deliverHandles.clear(); + stashedFollowups.clear(); sessions.clear(); pinCounts.clear(); runInFlight.clear();