Skip to content

Commit f9948ca

Browse files
Merge pull request #1066 from corbitsdev/cl-7989-deliver-stashed-steer-when-interrupt-loses-the-race-with-run
Deliver a stashed steer as a follow-up when the run wins the race
2 parents 920990c + 663aff0 commit f9948ca

3 files changed

Lines changed: 260 additions & 25 deletions

File tree

‎src/subagent/lifecycle-tools.test.ts‎

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -950,7 +950,7 @@ describe("send_input", () => {
950950
expect(sessions.get(missing.id)?.lifecycleStatus).toBe("running");
951951
});
952952

953-
test("completion during interrupt keeps the original terminal report", async () => {
953+
test("completion during interrupt delivers the stashed steer as a follow-up", async () => {
954954
const sessions = createSubAgentSessionStore();
955955
const fleetRecords = createFleetMailbox(sessions);
956956
const worker = sessions.start({
@@ -978,7 +978,7 @@ describe("send_input", () => {
978978
message: "late interrupt",
979979
interrupt: true,
980980
});
981-
expect(result).toEqual({ agent_id: worker.id, status: "completed" });
981+
expect(result).toEqual({ agent_id: worker.id, status: "interrupted" });
982982

983983
const collected = await callTool(wait, {
984984
targets: [worker.id],
@@ -989,10 +989,16 @@ describe("send_input", () => {
989989
expect.objectContaining({
990990
agent_id: worker.id,
991991
status: "done",
992-
report: "original report",
992+
report: "follow-up report",
993993
}),
994994
]);
995-
expect(followupStarted).toBe(false);
995+
expect(followupStarted).toBe(true);
996+
expect(sessions.get(worker.id)?.entries).toContainEqual(
997+
expect.objectContaining({
998+
kind: "report",
999+
content: expect.stringContaining("original report"),
1000+
}),
1001+
);
9961002
});
9971003

9981004
test("CL-7344: interrupt:true stashes until attachReport; resume stays fail-closed", async () => {

‎src/subagent/session-store.test.ts‎

Lines changed: 191 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1836,32 +1836,30 @@ describe("CL-7344 follow-up stash", () => {
18361836
expect(store.isRunInFlight(session.id)).toBe(true);
18371837
});
18381838

1839-
test("complete drops a stashed follow-up and keeps the original report", async () => {
1839+
test("complete on a resumable session delivers the stash and preserves the original report", async () => {
18401840
const store = createSubAgentSessionStore();
18411841
const started: string[] = [];
18421842
const failures: unknown[] = [];
18431843
const session = runningRetained(store, async (message) => {
18441844
started.push(message);
1845-
return "should not run";
1845+
return `reply to ${message}`;
18461846
});
18471847
const onFail = (err: unknown): void => {
18481848
failures.push(err);
18491849
};
18501850
store.sendInputOne(session.id, "steer one", { interrupt: true, onFail });
18511851
store.sendInputOne(session.id, "steer two", { interrupt: true, onFail });
18521852
store.complete(session.id, "## Summary\nOriginal done.");
1853-
await Promise.resolve();
1854-
expect(started).toEqual([]);
1855-
expect(store.get(session.id)?.report).toBe("## Summary\nOriginal done.");
1856-
expect(store.get(session.id)?.lifecycleStatus).toBe("completed");
1857-
expect(store.isRunInFlight(session.id)).toBe(false);
1858-
expect(failures).toHaveLength(2);
1859-
expect(String(defined(failures[0]))).toContain("steer one");
1860-
expect(String(defined(failures[1]))).toContain("steer two");
1853+
await new Promise((resolve) => setTimeout(resolve, 0));
1854+
await new Promise((resolve) => setTimeout(resolve, 0));
1855+
expect(started).toEqual(["steer one", "steer two"]);
1856+
expect(failures).toEqual([]);
1857+
expect(store.get(session.id)?.lifecycle.state).toBe("completed");
1858+
expect(store.get(session.id)?.report).toBe("reply to steer two");
18611859
expect(store.get(session.id)?.entries).toContainEqual(
18621860
expect.objectContaining({
18631861
kind: "report",
1864-
content: expect.stringContaining("steer two"),
1862+
content: expect.stringContaining("Original done"),
18651863
}),
18661864
);
18671865
});
@@ -1921,7 +1919,17 @@ describe("CL-7344 follow-up stash", () => {
19211919
test("dropped steers report which message was lost and why", async () => {
19221920
const store = createSubAgentSessionStore();
19231921
const failures: unknown[] = [];
1924-
const session = runningRetained(store, async () => "x");
1922+
// No retained:true: the session cannot resume, so a run-completion win
1923+
// surfaces the queued steer as a loss instead of delivering it.
1924+
const session = store.start({
1925+
description: "d",
1926+
agentId: "a",
1927+
brief: "b",
1928+
});
1929+
store.markRunning(session.id);
1930+
store.markRunInFlight(session.id);
1931+
store.registerInterrupt(session.id, () => undefined);
1932+
store.registerFollowup(session.id, async () => "x");
19251933
store.sendInputOne(session.id, "steer now", {
19261934
interrupt: true,
19271935
onFail: (err) => {
@@ -2052,4 +2060,175 @@ describe("CL-7344 follow-up stash", () => {
20522060
await new Promise((resolve) => setTimeout(resolve, 0));
20532061
expect(store.get(session.id)?.lifecycle.state).toBe("shutdown");
20542062
});
2063+
2064+
test("CL-7989 interrupt wins the race: stashed steer launches from attachReport", async () => {
2065+
const store = createSubAgentSessionStore();
2066+
const started: string[] = [];
2067+
const failures: unknown[] = [];
2068+
const replies: string[] = [];
2069+
const session = runningRetained(store, async (message) => {
2070+
started.push(message);
2071+
return "followup reply";
2072+
});
2073+
store.sendInputOne(session.id, "steer now", {
2074+
interrupt: true,
2075+
onFail: (err: unknown) => {
2076+
failures.push(err);
2077+
},
2078+
onFollowupReply: (reply: string) => {
2079+
replies.push(reply);
2080+
},
2081+
});
2082+
store.attachReport(session.id, "salvage", { stopReason: "interrupted" });
2083+
await new Promise((resolve) => setTimeout(resolve, 0));
2084+
expect(started).toEqual(["steer now"]);
2085+
expect(failures).toEqual([]);
2086+
expect(replies).toEqual(["followup reply"]);
2087+
expect(store.get(session.id)?.lifecycle.state).toBe("completed");
2088+
expect(store.get(session.id)?.report).toBe("followup reply");
2089+
});
2090+
2091+
test("CL-7989 run-completion wins the race: stashed steer delivers as a fresh follow-up", async () => {
2092+
const store = createSubAgentSessionStore();
2093+
const started: string[] = [];
2094+
const failures: unknown[] = [];
2095+
const replies: string[] = [];
2096+
const session = runningRetained(store, async (message) => {
2097+
started.push(message);
2098+
return "followup reply";
2099+
});
2100+
store.sendInputOne(session.id, "steer now", {
2101+
interrupt: true,
2102+
onFail: (err: unknown) => {
2103+
failures.push(err);
2104+
},
2105+
onFollowupReply: (reply: string) => {
2106+
replies.push(reply);
2107+
},
2108+
});
2109+
store.complete(session.id, "## Summary\nOriginal done.");
2110+
expect(store.get(session.id)?.lifecycleStatus).toBe("running");
2111+
expect(store.isRunInFlight(session.id)).toBe(true);
2112+
await new Promise((resolve) => setTimeout(resolve, 0));
2113+
expect(started).toEqual(["steer now"]);
2114+
expect(failures).toEqual([]);
2115+
expect(replies).toEqual(["followup reply"]);
2116+
expect(store.get(session.id)?.lifecycle.state).toBe("completed");
2117+
expect(store.get(session.id)?.report).toBe("followup reply");
2118+
expect(store.get(session.id)?.entries).toContainEqual(
2119+
expect.objectContaining({
2120+
kind: "report",
2121+
content: expect.stringContaining("Original done"),
2122+
}),
2123+
);
2124+
});
2125+
2126+
test("CL-7989 run-completion wins with queued steers: all deliver in order", async () => {
2127+
const store = createSubAgentSessionStore();
2128+
const started: string[] = [];
2129+
const failures: unknown[] = [];
2130+
const session = runningRetained(store, async (message) => {
2131+
started.push(message);
2132+
return `reply to ${message}`;
2133+
});
2134+
const onFail = (err: unknown): void => {
2135+
failures.push(err);
2136+
};
2137+
store.sendInputOne(session.id, "steer one", { interrupt: true, onFail });
2138+
store.sendInputOne(session.id, "steer two", { interrupt: true, onFail });
2139+
store.complete(session.id, "## Summary\nOriginal done.");
2140+
await new Promise((resolve) => setTimeout(resolve, 0));
2141+
await new Promise((resolve) => setTimeout(resolve, 0));
2142+
expect(started).toEqual(["steer one", "steer two"]);
2143+
expect(failures).toEqual([]);
2144+
expect(store.get(session.id)?.lifecycle.state).toBe("completed");
2145+
expect(store.get(session.id)?.report).toBe("reply to steer two");
2146+
});
2147+
2148+
test("CL-7989 run-completion wins on a non-retained session: loss is surfaced", async () => {
2149+
const store = createSubAgentSessionStore();
2150+
const started: string[] = [];
2151+
const failures: unknown[] = [];
2152+
const session = store.start({
2153+
description: "d",
2154+
agentId: "a",
2155+
brief: "b",
2156+
});
2157+
store.markRunning(session.id);
2158+
store.markRunInFlight(session.id);
2159+
store.registerInterrupt(session.id, () => undefined);
2160+
store.registerFollowup(session.id, async (message) => {
2161+
started.push(message);
2162+
return "should not run";
2163+
});
2164+
store.sendInputOne(session.id, "steer now", {
2165+
interrupt: true,
2166+
onFail: (err: unknown) => {
2167+
failures.push(err);
2168+
},
2169+
});
2170+
store.complete(session.id, "## Summary\nOriginal done.");
2171+
await Promise.resolve();
2172+
expect(started).toEqual([]);
2173+
expect(store.get(session.id)?.report).toBe("## Summary\nOriginal done.");
2174+
expect(store.get(session.id)?.lifecycleStatus).toBe("completed");
2175+
expect(failures).toHaveLength(1);
2176+
expect(String(defined(failures[0]))).toContain("steer now");
2177+
expect(store.get(session.id)?.entries).toContainEqual(
2178+
expect.objectContaining({
2179+
kind: "report",
2180+
content: expect.stringContaining("steer now"),
2181+
}),
2182+
);
2183+
});
2184+
2185+
test("a throwing handoff onReply does not stall the queued steers behind it", async () => {
2186+
const store = createSubAgentSessionStore();
2187+
const started: string[] = [];
2188+
const session = runningRetained(store, async (message) => {
2189+
started.push(message);
2190+
return `reply to ${message}`;
2191+
});
2192+
store.sendInputOne(session.id, "steer one", {
2193+
interrupt: true,
2194+
onFollowupReply: () => {
2195+
throw new Error("observer blew up");
2196+
},
2197+
});
2198+
store.sendInputOne(session.id, "steer two", { interrupt: true });
2199+
store.complete(session.id, "## Summary\nOriginal done.");
2200+
await new Promise((resolve) => setTimeout(resolve, 0));
2201+
await new Promise((resolve) => setTimeout(resolve, 0));
2202+
expect(started).toEqual(["steer one", "steer two"]);
2203+
expect(store.get(session.id)?.lifecycle.state).toBe("completed");
2204+
expect(store.get(session.id)?.report).toBe("reply to steer two");
2205+
});
2206+
2207+
test("CL-7989 run-failure wins the race: loss is surfaced, never silent", async () => {
2208+
const store = createSubAgentSessionStore();
2209+
const started: string[] = [];
2210+
const failures: unknown[] = [];
2211+
const session = runningRetained(store, async (message) => {
2212+
started.push(message);
2213+
return "should not run";
2214+
});
2215+
store.sendInputOne(session.id, "steer now", {
2216+
interrupt: true,
2217+
onFail: (err: unknown) => {
2218+
failures.push(err);
2219+
},
2220+
});
2221+
store.fail(session.id, "provider 500");
2222+
await Promise.resolve();
2223+
expect(started).toEqual([]);
2224+
expect(store.get(session.id)?.lifecycle.state).toBe("failed");
2225+
expect(failures).toHaveLength(1);
2226+
expect(String(defined(failures[0]))).toContain("steer now");
2227+
expect(store.get(session.id)?.entries).toContainEqual(
2228+
expect.objectContaining({
2229+
kind: "report",
2230+
content: expect.stringContaining("steer now"),
2231+
}),
2232+
);
2233+
});
20552234
});

0 commit comments

Comments
 (0)