diff --git a/src/subagent/lifecycle-tools.test.ts b/src/subagent/lifecycle-tools.test.ts index 91f64145c..69a15e756 100644 --- a/src/subagent/lifecycle-tools.test.ts +++ b/src/subagent/lifecycle-tools.test.ts @@ -950,7 +950,7 @@ describe("send_input", () => { expect(sessions.get(missing.id)?.lifecycleStatus).toBe("running"); }); - test("completion during interrupt keeps the original terminal report", async () => { + test("completion during interrupt delivers the stashed steer as a follow-up", async () => { const sessions = createSubAgentSessionStore(); const fleetRecords = createFleetMailbox(sessions); const worker = sessions.start({ @@ -978,7 +978,7 @@ describe("send_input", () => { message: "late interrupt", interrupt: true, }); - expect(result).toEqual({ agent_id: worker.id, status: "completed" }); + expect(result).toEqual({ agent_id: worker.id, status: "interrupted" }); const collected = await callTool(wait, { targets: [worker.id], @@ -989,10 +989,16 @@ describe("send_input", () => { expect.objectContaining({ agent_id: worker.id, status: "done", - report: "original report", + report: "follow-up report", }), ]); - expect(followupStarted).toBe(false); + expect(followupStarted).toBe(true); + expect(sessions.get(worker.id)?.entries).toContainEqual( + expect.objectContaining({ + kind: "report", + content: expect.stringContaining("original report"), + }), + ); }); test("CL-7344: interrupt:true stashes until attachReport; resume stays fail-closed", async () => { diff --git a/src/subagent/session-store.test.ts b/src/subagent/session-store.test.ts index f24852a75..727ea1228 100644 --- a/src/subagent/session-store.test.ts +++ b/src/subagent/session-store.test.ts @@ -1836,13 +1836,13 @@ describe("CL-7344 follow-up stash", () => { expect(store.isRunInFlight(session.id)).toBe(true); }); - test("complete drops a stashed follow-up and keeps the original report", async () => { + test("complete on a resumable session delivers the stash and preserves the original report", async () => { const store = createSubAgentSessionStore(); const started: string[] = []; const failures: unknown[] = []; const session = runningRetained(store, async (message) => { started.push(message); - return "should not run"; + return `reply to ${message}`; }); const onFail = (err: unknown): void => { failures.push(err); @@ -1850,18 +1850,16 @@ describe("CL-7344 follow-up stash", () => { store.sendInputOne(session.id, "steer one", { interrupt: true, onFail }); store.sendInputOne(session.id, "steer two", { interrupt: true, onFail }); 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); - expect(failures).toHaveLength(2); - expect(String(defined(failures[0]))).toContain("steer one"); - expect(String(defined(failures[1]))).toContain("steer two"); + await new Promise((resolve) => setTimeout(resolve, 0)); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(started).toEqual(["steer one", "steer two"]); + expect(failures).toEqual([]); + expect(store.get(session.id)?.lifecycle.state).toBe("completed"); + expect(store.get(session.id)?.report).toBe("reply to steer two"); expect(store.get(session.id)?.entries).toContainEqual( expect.objectContaining({ kind: "report", - content: expect.stringContaining("steer two"), + content: expect.stringContaining("Original done"), }), ); }); @@ -1921,7 +1919,17 @@ describe("CL-7344 follow-up stash", () => { test("dropped steers report which message was lost and why", async () => { const store = createSubAgentSessionStore(); const failures: unknown[] = []; - const session = runningRetained(store, async () => "x"); + // No retained:true: the session cannot resume, so a run-completion win + // surfaces the queued steer as a loss instead of delivering it. + const session = store.start({ + description: "d", + agentId: "a", + brief: "b", + }); + store.markRunning(session.id); + store.markRunInFlight(session.id); + store.registerInterrupt(session.id, () => undefined); + store.registerFollowup(session.id, async () => "x"); store.sendInputOne(session.id, "steer now", { interrupt: true, onFail: (err) => { @@ -2052,4 +2060,175 @@ describe("CL-7344 follow-up stash", () => { await new Promise((resolve) => setTimeout(resolve, 0)); expect(store.get(session.id)?.lifecycle.state).toBe("shutdown"); }); + + test("CL-7989 interrupt wins the race: stashed steer launches from attachReport", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const failures: unknown[] = []; + const replies: string[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return "followup reply"; + }); + store.sendInputOne(session.id, "steer now", { + interrupt: true, + onFail: (err: unknown) => { + failures.push(err); + }, + onFollowupReply: (reply: string) => { + replies.push(reply); + }, + }); + store.attachReport(session.id, "salvage", { stopReason: "interrupted" }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(started).toEqual(["steer now"]); + expect(failures).toEqual([]); + expect(replies).toEqual(["followup reply"]); + expect(store.get(session.id)?.lifecycle.state).toBe("completed"); + expect(store.get(session.id)?.report).toBe("followup reply"); + }); + + test("CL-7989 run-completion wins the race: stashed steer delivers as a fresh follow-up", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const failures: unknown[] = []; + const replies: string[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return "followup reply"; + }); + store.sendInputOne(session.id, "steer now", { + interrupt: true, + onFail: (err: unknown) => { + failures.push(err); + }, + onFollowupReply: (reply: string) => { + replies.push(reply); + }, + }); + store.complete(session.id, "## Summary\nOriginal done."); + expect(store.get(session.id)?.lifecycleStatus).toBe("running"); + expect(store.isRunInFlight(session.id)).toBe(true); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(started).toEqual(["steer now"]); + expect(failures).toEqual([]); + expect(replies).toEqual(["followup reply"]); + expect(store.get(session.id)?.lifecycle.state).toBe("completed"); + expect(store.get(session.id)?.report).toBe("followup reply"); + expect(store.get(session.id)?.entries).toContainEqual( + expect.objectContaining({ + kind: "report", + content: expect.stringContaining("Original done"), + }), + ); + }); + + test("CL-7989 run-completion wins with queued steers: all deliver in order", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const failures: unknown[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return `reply to ${message}`; + }); + const onFail = (err: unknown): void => { + failures.push(err); + }; + store.sendInputOne(session.id, "steer one", { interrupt: true, onFail }); + store.sendInputOne(session.id, "steer two", { interrupt: true, onFail }); + store.complete(session.id, "## Summary\nOriginal done."); + await new Promise((resolve) => setTimeout(resolve, 0)); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(started).toEqual(["steer one", "steer two"]); + expect(failures).toEqual([]); + expect(store.get(session.id)?.lifecycle.state).toBe("completed"); + expect(store.get(session.id)?.report).toBe("reply to steer two"); + }); + + test("CL-7989 run-completion wins on a non-retained session: loss is surfaced", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const failures: unknown[] = []; + const session = store.start({ + description: "d", + agentId: "a", + brief: "b", + }); + store.markRunning(session.id); + store.markRunInFlight(session.id); + store.registerInterrupt(session.id, () => undefined); + store.registerFollowup(session.id, async (message) => { + started.push(message); + return "should not run"; + }); + store.sendInputOne(session.id, "steer now", { + interrupt: true, + onFail: (err: unknown) => { + failures.push(err); + }, + }); + 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(failures).toHaveLength(1); + expect(String(defined(failures[0]))).toContain("steer now"); + expect(store.get(session.id)?.entries).toContainEqual( + expect.objectContaining({ + kind: "report", + content: expect.stringContaining("steer now"), + }), + ); + }); + + test("a throwing handoff onReply does not stall the queued steers behind it", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return `reply to ${message}`; + }); + store.sendInputOne(session.id, "steer one", { + interrupt: true, + onFollowupReply: () => { + throw new Error("observer blew up"); + }, + }); + store.sendInputOne(session.id, "steer two", { interrupt: true }); + store.complete(session.id, "## Summary\nOriginal done."); + await new Promise((resolve) => setTimeout(resolve, 0)); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(started).toEqual(["steer one", "steer two"]); + expect(store.get(session.id)?.lifecycle.state).toBe("completed"); + expect(store.get(session.id)?.report).toBe("reply to steer two"); + }); + + test("CL-7989 run-failure wins the race: loss is surfaced, never silent", async () => { + const store = createSubAgentSessionStore(); + const started: string[] = []; + const failures: unknown[] = []; + const session = runningRetained(store, async (message) => { + started.push(message); + return "should not run"; + }); + store.sendInputOne(session.id, "steer now", { + interrupt: true, + onFail: (err: unknown) => { + failures.push(err); + }, + }); + store.fail(session.id, "provider 500"); + await Promise.resolve(); + expect(started).toEqual([]); + expect(store.get(session.id)?.lifecycle.state).toBe("failed"); + expect(failures).toHaveLength(1); + expect(String(defined(failures[0]))).toContain("steer now"); + expect(store.get(session.id)?.entries).toContainEqual( + expect.objectContaining({ + kind: "report", + content: expect.stringContaining("steer now"), + }), + ); + }); }); diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index 50eba073d..9d99ed79a 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -1012,7 +1012,11 @@ export function createSubAgentSessionStore( content: capText(reply, maxEntryChars), }); }); - onReply?.(reply); + try { + onReply?.(reply); + } catch { + // A throwing onReply must not break the handoff to the next steer. + } launchNextStashedFollowup(id); return; } @@ -1027,7 +1031,11 @@ export function createSubAgentSessionStore( }); }); runInFlight.delete(id); - onReply?.(reply); + try { + onReply?.(reply); + } catch { + // A throwing onReply must not break session settlement. + } pruneRetained(); }; const settleFollowupFailure = ( @@ -1043,7 +1051,11 @@ export function createSubAgentSessionStore( !(err instanceof AgentClosedError) && (stashedFollowups.get(id)?.length ?? 0) > 0 ) { - onFail?.(err); + try { + onFail?.(err); + } catch { + // A throwing onFail must not break the handoff to the next steer. + } log.error("followup turn failed for {id}: {error}", { id, error: err instanceof Error ? err.message : String(err), @@ -1052,7 +1064,11 @@ export function createSubAgentSessionStore( return; } runInFlight.delete(id); - onFail?.(err); + try { + onFail?.(err); + } catch { + // A throwing onFail must not break session settlement. + } // 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 @@ -1091,7 +1107,11 @@ export function createSubAgentSessionStore( const takesSlot = queue !== undefined && !queue.occupied(id); const start = (): void => { beginFollowupTurn(id); - opts?.onStart?.(); + try { + opts?.onStart?.(); + } catch { + // A throwing onStart must not break the follow-up turn. + } void followup(message) .then((reply) => { settleFollowupReply(id, reply, opts?.onReply); @@ -1159,7 +1179,11 @@ export function createSubAgentSessionStore( runInFlight.delete(id); return false; } - next.onStart?.(); + try { + next.onStart?.(); + } catch { + // A throwing onStart must not strand the lane on a phantom turn. + } const pending = followup(next.message); void Promise.resolve().then(() => { void pending.then( @@ -1451,6 +1475,13 @@ export function createSubAgentSessionStore( // spawn_agent path ever has a salvage to report, and it // always passes this flag explicitly (see its call site). const agentRetained = opts?.agentRetained ?? true; + // CL-7989: when the original run wins the race against a stashed steer + // and the session stays open and resumable, deliver the queue as a + // fresh follow-up instead of dropping it. The lane flips to running + // inside this same mutation so observers never see a completed session + // with a pending steer; the hand-off below reuses the stash launcher + // so the rest of the queue chains in FIFO order. + let deliverStash = false; mutate(id, (session) => { // Cancel and interrupt_agent win races: a late complete must not // resurrect the session as done. Interrupted is still strip-live @@ -1479,14 +1510,33 @@ export function createSubAgentSessionStore( // release it now rather than leaving a stale reference around. cancelHandles.delete(id); if (!agentRetained) closeHandles.delete(id); - runInFlight.delete(id); + const pending = stashedFollowups.get(id); + if ( + pending !== undefined && + pending.length > 0 && + session.retained === true && + followupHandles.has(id) + ) { + deliverStash = true; + session.lifecycle = { state: "running" }; + delete session.finishedAt; + runInFlight.add(id); + } else { + runInFlight.delete(id); + } pruneCompleted(); pruneRetained(); }); // CL-7988: a completed turn supersedes any steer still queued for this // session — surface it. Runs after the mutate so the completion lands - // first even when the queue is non-empty. - dropStashedFollowups(id, "session completed"); + // first even when the queue is non-empty. CL-7989: when the run won the + // race but the session stays open and resumable, the queue launches as + // a fresh follow-up above instead of being dropped here. + if (deliverStash && sessions.has(id)) { + launchNextStashedFollowup(id); + } else { + dropStashedFollowups(id, "session completed"); + } }, fail(id: string, error: string): void {