diff --git a/CHANGELOG.md b/CHANGELOG.md index 2b7e86a..dc01a33 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,9 @@ ## Unreleased +- Gateway: a failed turn no longer drops a conversation's history when a new + chat evicted it while that turn was still in flight. The history loaded for + the run is written back. - Gateway (#144): sender allowlists and the gateway toolset are one `GatewayPolicy`, built from config at startup (`src/gateway/access.ts`). Platform adapters declare `capabilities` (`{ kind: "text", max_reply_chars }`); diff --git a/src/gateway/bus.ts b/src/gateway/bus.ts index 9447650..29f010c 100644 --- a/src/gateway/bus.ts +++ b/src/gateway/bus.ts @@ -92,6 +92,10 @@ export class GatewayBus { const kept = history_after_run_error(error); if (kept !== undefined) { this.store_history(key, cap_history(kept, this.history_cap)); + } else if (this.histories.has(key) === false && history.length > 0) { + // A new chat may have evicted this key while the run was in flight. + // Nothing was stored, so put the loaded history back. + this.store_history(key, history); } return { text: sanitize_agent_error(error), failed: true }; } diff --git a/test/gateway_conversation_bounds.test.ts b/test/gateway_conversation_bounds.test.ts index 0ca46b1..e868b95 100644 --- a/test/gateway_conversation_bounds.test.ts +++ b/test/gateway_conversation_bounds.test.ts @@ -188,4 +188,55 @@ describe("gateway max_conversations", () => { expect(await third).toBe("echo:a3"); expect(await other).toBe("echo:b1"); }); + + it("restores history when a new chat evicts an in-flight conversation and that run fails", async () => { + const work_dir = make_temp_dir(); + const records: RunRecord[] = []; + const started: HeldRun[] = []; + const bus = new GatewayBus( + { + config: parse_agent_config({ + providers: [{ kind: "openai_compat", name: "main", model: "mock-model" }], + work_dir, + log_level: "error", + }), + agent_factory: () => failing_held_agent(records, started, "a2"), + }, + { max_conversations: 2, history_cap: 40 }, + ); + + expect(await bus.handle("webhook", "a", "u", "a1")).toBe("echo:a1"); + expect(await bus.handle("webhook", "b", "u", "b1")).toBe("echo:b1"); + + const failed = bus.handle("webhook", "a", "u", "a2"); + await wait_for_start(started, "a2"); + // c arrives while a2 is in flight, so history_for evicts a (the LRU slot). + expect(await bus.handle("webhook", "c", "u", "c1")).toBe("echo:c1"); + started.find((run) => run.input === "a2")?.release(); + expect(await failed).toContain("agent error:"); + + expect(await bus.handle("webhook", "a", "u", "a3")).toBe("echo:a3"); + const a3 = records.find((record) => record.input === "a3"); + expect(a3?.history_len).toBeGreaterThan(0); + }); }); + +/** Holds `fail_input` until release, then throws; every other input echoes. */ +function failing_held_agent(records: RunRecord[], started: HeldRun[], fail_input: string): Agent { + return { + run: async (options: { input: string; history?: readonly Message[] }): Promise => { + const history = options.history ?? []; + if (options.input === fail_input) { + let release: () => void = () => undefined; + const gate = new Promise((resolve) => { + release = resolve; + }); + started.push({ input: options.input, history_len: history.length, release }); + await gate; + throw new Error("provider down"); + } + records.push({ input: options.input, history_len: history.length }); + return echo_result(options.input, history); + }, + } as unknown as Agent; +}