From 30b04851661c7321c19a36318d52f6f471e8840e Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 14 Sep 2026 14:21:35 -0700 Subject: [PATCH] fix(subagent): queue overlapping steers instead of dropping them A second interrupt-steer overwrote the first, losing it silently. Steers now queue FIFO per session and any steer that never launches fails loudly with which message was lost and why. --- src/subagent/session-store.test.ts | 63 ++++++++- src/subagent/session-store.ts | 201 ++++++++++++++++++++++++----- 2 files changed, 223 insertions(+), 41 deletions(-) diff --git a/src/subagent/session-store.test.ts b/src/subagent/session-store.test.ts index af2caf996..f24852a75 100644 --- a/src/subagent/session-store.test.ts +++ b/src/subagent/session-store.test.ts @@ -1839,17 +1839,31 @@ describe("CL-7344 follow-up stash", () => { test("complete drops a stashed follow-up and keeps 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"; }); - store.sendInputOne(session.id, "steer now", { interrupt: true }); + 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 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"); + expect(store.get(session.id)?.entries).toContainEqual( + expect.objectContaining({ + kind: "report", + content: expect.stringContaining("steer two"), + }), + ); }); test("attachReport interrupted starts follow-up in the same notify as clearing the original run", async () => { @@ -1878,18 +1892,53 @@ describe("CL-7344 follow-up stash", () => { expect(store.get(session.id)?.lifecycleStatus).toBe("running"); }); - test("last send_input interrupt overwrites the stash", async () => { + test("back-to-back send_input interrupts queue and deliver in order", async () => { const store = createSubAgentSessionStore(); const started: string[] = []; - const session = runningRetained(store, async (message) => { - started.push(message); - return "followup"; - }); + const resolvers: ((reply: string) => void)[] = []; + const session = runningRetained( + store, + (message) => + new Promise((resolve) => { + started.push(message); + resolvers.push(resolve); + }), + ); 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"]); + expect(started).toEqual(["first"]); + defined(resolvers[0])("reply one"); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(started).toEqual(["first", "second"]); + defined(resolvers[1])("reply two"); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(store.get(session.id)?.lifecycle.state).toBe("completed"); + expect(store.get(session.id)?.report).toBe("reply two"); + }); + + test("dropped steers report which message was lost and why", async () => { + const store = createSubAgentSessionStore(); + const failures: unknown[] = []; + const session = runningRetained(store, async () => "x"); + store.sendInputOne(session.id, "steer now", { + interrupt: true, + onFail: (err) => { + failures.push(err); + }, + }); + store.complete(session.id, "## Summary\nOriginal done."); + await Promise.resolve(); + expect(failures).toHaveLength(1); + expect(String(defined(failures[0]))).toContain("steer now"); + expect(String(defined(failures[0]))).toContain("completed"); + expect(store.get(session.id)?.entries).toContainEqual( + expect.objectContaining({ + kind: "report", + content: expect.stringContaining("steer now"), + }), + ); }); test("fail, cancel, close, interrupt_agent, and settleRun drop the stash", async () => { diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index cccea823c..50eba073d 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -525,7 +525,50 @@ export function createSubAgentSessionStore( onReply?: (reply: string) => void; onFail?: (error: unknown) => void; } - const stashedFollowups = new Map(); + // CL-7988: overlapping interrupt-steers queue FIFO per session instead of + // overwriting each other. The head launches when the live run settles via + // attachReport; each settled follow-up turn launches the next in order. + const stashedFollowups = new Map(); + + const steerPreview = (message: string): string => { + const firstLine = message.split("\n", 1)[0] ?? ""; + return firstLine.length > 120 ? `${firstLine.slice(0, 117)}...` : firstLine; + }; + + // CL-7988: a steer that never launches is superseded, never silent. Each + // queued steer fails with which message was lost and why, and the loss is + // recorded on the session transcript so the operator can see it. + const dropStashedFollowups = (id: string, reason: string): void => { + const queue = stashedFollowups.get(id); + if (queue === undefined || queue.length === 0) { + stashedFollowups.delete(id); + return; + } + stashedFollowups.delete(id); + const lost = queue + .map((stashed) => `"${steerPreview(stashed.message)}"`) + .join(", "); + for (const stashed of queue) { + try { + stashed.onFail?.( + new Error( + `send_input steer "${steerPreview(stashed.message)}" dropped (${reason})`, + ), + ); + } catch { + // A throwing onFail must not break session settlement. + } + } + mutate(id, (session) => { + pushEntry(session, { + kind: "report", + content: capText( + `Steer ${queue.length === 1 ? "dropped" : `${queue.length} steers dropped`} (${reason}): ${lost}`, + maxEntryChars, + ), + }); + }); + }; const pendingAsks = new Map< string, { @@ -686,7 +729,8 @@ export function createSubAgentSessionStore( const cancelSession = (id: string, reason: string): boolean => { settleCancelsAsks(id, reason); - stashedFollowups.delete(id); + // CL-7988: cancelled steers are superseded — surface which were dropped. + dropStashedFollowups(id, "session cancelled"); const session = sessions.get(id); if (session === undefined || !isLiveStrip(session.lifecycle)) return false; const abort = cancelHandles.get(id); @@ -734,7 +778,9 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); - stashedFollowups.delete(id); + // CL-7988 backstop: every teardown path funnels through here, so a queue + // that somehow survived its semantic drop point is surfaced, never silent. + dropStashedFollowups(id, "session ended"); }; // An open retained session (spawn_agent's reusable-session contract: @@ -906,6 +952,8 @@ export function createSubAgentSessionStore( const failClosedSession = (id: string, error: string): void => { endFollowupTurn(id, "interrupted"); settleCancelsAsks(id, "session failed"); + // CL-7988: the failed turn supersedes steers still queued behind it. + dropStashedFollowups(id, "session failed"); mutate(id, (session) => { if ( !isLiveStrip(session.lifecycle) || @@ -939,6 +987,8 @@ export function createSubAgentSessionStore( const still = sessions.get(id); if (still === undefined) { runInFlight.delete(id); + // CL-7988: the session vanished with steers still queued — surface them. + dropStashedFollowups(id, "session ended"); return; } if ( @@ -948,6 +998,22 @@ export function createSubAgentSessionStore( still.lifecycle.state === "interrupted" ) { runInFlight.delete(id); + // CL-7988: the lane died with steers still queued — surface them. + dropStashedFollowups(id, "session settled"); + return; + } + // CL-7988: more steers queued behind this one — record its reply and hand + // off to the next steer in order instead of completing the session. + if ((stashedFollowups.get(id)?.length ?? 0) > 0) { + mutate(id, (s) => { + s.report = reply; + pushEntry(s, { + kind: "report", + content: capText(reply, maxEntryChars), + }); + }); + onReply?.(reply); + launchNextStashedFollowup(id); return; } mutate(id, (s) => { @@ -970,6 +1036,21 @@ export function createSubAgentSessionStore( failLifecycle: "completed" | "interrupted", onFail?: (error: unknown) => void, ): void => { + // CL-7988: a failed steer hands off to the next queued steer in order. + // The lane stays live across the handoff so no observer sees a gap + // between the two turns; only the last settlement restores the lane. + if ( + !(err instanceof AgentClosedError) && + (stashedFollowups.get(id)?.length ?? 0) > 0 + ) { + onFail?.(err); + log.error("followup turn failed for {id}: {error}", { + id, + error: err instanceof Error ? err.message : String(err), + }); + launchNextStashedFollowup(id); + return; + } runInFlight.delete(id); onFail?.(err); // CL-7344: the agent closed between queueing and invocation, so the @@ -1045,30 +1126,52 @@ export function createSubAgentSessionStore( }; // 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 => { + // CL-7988: shifts the head steer off the session queue and wires its + // settlement to launch the next queued steer in order. Returns false when + // nothing launches (queue empty, or the agent tore down mid-handoff). + const launchNextStashedFollowup = (id: string): boolean => { + // A follow-up must only launch into a live lane. If the session settled + // or vanished while steers were queued, the queue is superseded — drop + // it loudly instead of running against a closed agent. + const live = sessions.get(id); + if ( + live === undefined || + (live.lifecycle.state !== "running" && + live.lifecycle.state !== "pending_init") + ) { + dropStashedFollowups(id, "session settled"); + runInFlight.delete(id); + return false; + } + const queue = stashedFollowups.get(id); + const next = queue?.shift(); + if (queue !== undefined && queue.length === 0) stashedFollowups.delete(id); + if (next === undefined) return false; 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) { + // then no follow-up can inherit the run, so drop the queue loudly and + // settle instead of leaving wait_agents stuck on a phantom turn. + if (followup === undefined) { + const rest = stashedFollowups.get(id); + if (rest !== undefined) rest.unshift(next); + else stashedFollowups.set(id, [next]); + dropStashedFollowups(id, "agent tore down before delivery"); runInFlight.delete(id); - return; + return false; } - stashed.onStart?.(); - const pending = followup(stashed.message); + next.onStart?.(); + const pending = followup(next.message); void Promise.resolve().then(() => { void pending.then( (reply) => { - settleFollowupReply(id, reply, stashed.onReply); + settleFollowupReply(id, reply, next.onReply); }, (err: unknown) => { - settleFollowupFailure(id, err, stashed.failLifecycle, stashed.onFail); + settleFollowupFailure(id, err, next.failLifecycle, next.onFail); }, ); }); + return true; }; return { @@ -1105,7 +1208,8 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); - stashedFollowups.delete(id); + // CL-7988: the old session's queued steers are superseded — surface them. + dropStashedFollowups(id, "session replaced"); pinCounts.delete(id); runInFlight.delete(id); forgetRevision(id); @@ -1376,10 +1480,13 @@ export function createSubAgentSessionStore( cancelHandles.delete(id); if (!agentRetained) closeHandles.delete(id); runInFlight.delete(id); - stashedFollowups.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"); }, fail(id: string, error: string): void { @@ -1405,6 +1512,8 @@ export function createSubAgentSessionStore( content: capText(`Error: ${error}`, maxEntryChars), }); runInFlight.delete(id); + // CL-7988: the failure supersedes steers queued behind the failed turn. + dropStashedFollowups(id, "session failed"); releaseHandles(id); pruneCompleted(); }); @@ -1459,7 +1568,8 @@ export function createSubAgentSessionStore( // 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); + // CL-7988: closing supersedes queued steers — surface them, never silent. + dropStashedFollowups(id, "session closed"); let close = closeHandles.get(id); const alreadyClosed = isAlreadyClosed(session.lifecycle); if (alreadyClosed && close === undefined) { @@ -1487,7 +1597,8 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); - stashedFollowups.delete(id); + // CL-7988: closing supersedes queued steers — surface them. + dropStashedFollowups(id, "session closed"); runInFlight.delete(id); pruneCompleted(); return "shutdown"; @@ -1529,7 +1640,8 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); - stashedFollowups.delete(id); + // CL-7988: closing supersedes queued steers — surface them. + dropStashedFollowups(id, "session closed"); runInFlight.delete(id); pruneCompleted(); if (closeError !== undefined) throw closeError; @@ -1557,7 +1669,8 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); - stashedFollowups.delete(id); + // CL-7988: closing supersedes queued steers — surface them. + dropStashedFollowups(id, "session closed"); runInFlight.delete(id); pruneCompleted(); if (closeError !== undefined) throw closeError; @@ -1622,9 +1735,17 @@ export function createSubAgentSessionStore( // 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) { + // CL-7988: overlapping steers queue FIFO. Only the first steer fires + // the interrupt; later ones ride on the already-signalled handoff. + const queued = stashedFollowups.get(id); + const opening = queued === undefined; + if (opening) { + stashedFollowups.set(id, [stashed]); + } else { + queued.push(stashed); + } + if (opening) interrupt(); + if (!(stashedFollowups.get(id)?.includes(stashed) ?? false)) { const settled = sessions.get(id); const status = settled === undefined @@ -1711,7 +1832,8 @@ export function createSubAgentSessionStore( } // 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); + // CL-7988: the dropped steers fail loudly with which message was lost. + dropStashedFollowups(id, "interrupted by interrupt_agent"); const interrupt = interruptHandles.get(id); if (interrupt === undefined) { if (session.lifecycle.state === "pending_init") { @@ -1840,7 +1962,8 @@ export function createSubAgentSessionStore( interruptHandles.delete(id); followupHandles.delete(id); deliverHandles.delete(id); - stashedFollowups.delete(id); + // CL-7988: releasing supersedes queued steers — surface them. + dropStashedFollowups(id, "session handles released"); mutate(id, (s) => { s.lifecycle = { state: "shutdown", @@ -1884,11 +2007,13 @@ export function createSubAgentSessionStore( report: string, opts?: { stopReason?: ForcedStopReason }, ): void { - // Consume the stashed interrupt follow-up before mutating. Interrupted + // Consume the head 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); + // drops the queue loudly so no steer is lost silently and no follow-up + // runs against a closed agent. Steers behind the head stay queued until + // each follow-up turn settles and launches the next in order (CL-7988). + // Peek here; launchNextStashedFollowup shifts after the mutate below. + const head = stashedFollowups.get(id)?.[0]; let toLaunch: StashedFollowup | undefined; mutate(id, (session) => { const state = session.lifecycle.state; @@ -1906,10 +2031,10 @@ export function createSubAgentSessionStore( kind: "report", content: capText(report, maxEntryChars), }); - if (stashed !== undefined) { + if (head !== undefined) { // Move directly into the follow-up lifecycle before mutate notifies // subscribers. No observer can resume the interrupted handoff. - toLaunch = stashed; + toLaunch = head; session.lifecycle = { state: "running" }; delete session.finishedAt; delete session.stopReason; @@ -1935,7 +2060,11 @@ export function createSubAgentSessionStore( pruneRetained(); }); if (toLaunch !== undefined && sessions.has(id)) { - launchStashedFollowup(id, toLaunch); + launchNextStashedFollowup(id); + } else { + // Terminal outcome (or a session that vanished mid-handoff): nothing + // launches, so the queued steers are superseded — surface them. + dropStashedFollowups(id, "session settled"); } }, @@ -1945,7 +2074,8 @@ export function createSubAgentSessionStore( settleRun(id: string): void { cancelAskInternal(id, "run settled"); - stashedFollowups.delete(id); + // CL-7988: the run settled with steers still queued — surface them. + dropStashedFollowups(id, "run settled"); if (!runInFlight.delete(id)) return; notify(); }, @@ -1973,6 +2103,9 @@ export function createSubAgentSessionStore( interruptHandles.clear(); followupHandles.clear(); deliverHandles.clear(); + // CL-7988: surface every queued steer before wiping the map. + for (const id of stashedFollowups.keys()) + dropStashedFollowups(id, "store cleared"); stashedFollowups.clear(); sessions.clear(); pinCounts.clear();