Skip to content

Commit 920990c

Browse files
Merge pull request #1063 from corbitsdev/cl-7988-queue-overlapping-subagent-steers-instead-of-dropping-them
fix(subagent): queue overlapping steers instead of dropping them
2 parents a7a5082 + 30b0485 commit 920990c

2 files changed

Lines changed: 223 additions & 41 deletions

File tree

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

Lines changed: 56 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1839,17 +1839,31 @@ describe("CL-7344 follow-up stash", () => {
18391839
test("complete drops a stashed follow-up and keeps the original report", async () => {
18401840
const store = createSubAgentSessionStore();
18411841
const started: string[] = [];
1842+
const failures: unknown[] = [];
18421843
const session = runningRetained(store, async (message) => {
18431844
started.push(message);
18441845
return "should not run";
18451846
});
1846-
store.sendInputOne(session.id, "steer now", { interrupt: true });
1847+
const onFail = (err: unknown): void => {
1848+
failures.push(err);
1849+
};
1850+
store.sendInputOne(session.id, "steer one", { interrupt: true, onFail });
1851+
store.sendInputOne(session.id, "steer two", { interrupt: true, onFail });
18471852
store.complete(session.id, "## Summary\nOriginal done.");
18481853
await Promise.resolve();
18491854
expect(started).toEqual([]);
18501855
expect(store.get(session.id)?.report).toBe("## Summary\nOriginal done.");
18511856
expect(store.get(session.id)?.lifecycleStatus).toBe("completed");
18521857
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");
1861+
expect(store.get(session.id)?.entries).toContainEqual(
1862+
expect.objectContaining({
1863+
kind: "report",
1864+
content: expect.stringContaining("steer two"),
1865+
}),
1866+
);
18531867
});
18541868

18551869
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", () => {
18781892
expect(store.get(session.id)?.lifecycleStatus).toBe("running");
18791893
});
18801894

1881-
test("last send_input interrupt overwrites the stash", async () => {
1895+
test("back-to-back send_input interrupts queue and deliver in order", async () => {
18821896
const store = createSubAgentSessionStore();
18831897
const started: string[] = [];
1884-
const session = runningRetained(store, async (message) => {
1885-
started.push(message);
1886-
return "followup";
1887-
});
1898+
const resolvers: ((reply: string) => void)[] = [];
1899+
const session = runningRetained(
1900+
store,
1901+
(message) =>
1902+
new Promise<string>((resolve) => {
1903+
started.push(message);
1904+
resolvers.push(resolve);
1905+
}),
1906+
);
18881907
store.sendInputOne(session.id, "first", { interrupt: true });
18891908
store.sendInputOne(session.id, "second", { interrupt: true });
18901909
store.attachReport(session.id, "salvage", { stopReason: "interrupted" });
18911910
await Promise.resolve();
1892-
expect(started).toEqual(["second"]);
1911+
expect(started).toEqual(["first"]);
1912+
defined(resolvers[0])("reply one");
1913+
await new Promise((resolve) => setTimeout(resolve, 0));
1914+
expect(started).toEqual(["first", "second"]);
1915+
defined(resolvers[1])("reply two");
1916+
await new Promise((resolve) => setTimeout(resolve, 0));
1917+
expect(store.get(session.id)?.lifecycle.state).toBe("completed");
1918+
expect(store.get(session.id)?.report).toBe("reply two");
1919+
});
1920+
1921+
test("dropped steers report which message was lost and why", async () => {
1922+
const store = createSubAgentSessionStore();
1923+
const failures: unknown[] = [];
1924+
const session = runningRetained(store, async () => "x");
1925+
store.sendInputOne(session.id, "steer now", {
1926+
interrupt: true,
1927+
onFail: (err) => {
1928+
failures.push(err);
1929+
},
1930+
});
1931+
store.complete(session.id, "## Summary\nOriginal done.");
1932+
await Promise.resolve();
1933+
expect(failures).toHaveLength(1);
1934+
expect(String(defined(failures[0]))).toContain("steer now");
1935+
expect(String(defined(failures[0]))).toContain("completed");
1936+
expect(store.get(session.id)?.entries).toContainEqual(
1937+
expect.objectContaining({
1938+
kind: "report",
1939+
content: expect.stringContaining("steer now"),
1940+
}),
1941+
);
18931942
});
18941943

18951944
test("fail, cancel, close, interrupt_agent, and settleRun drop the stash", async () => {

0 commit comments

Comments
 (0)