From dbc3b36f6d94d719cb1a07074e5d74ce22d2fad3 Mon Sep 17 00:00:00 2001 From: Adam Malczewski Date: Mon, 1 Jun 2026 07:51:23 +0900 Subject: fix(queue): consume queued messages after a turn ends (start a new turn) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A message queued while the agent was mid-turn was only handled if it arrived DURING a tool batch (injected as a [USER INTERRUPT]). If it landed after the last tool call — or the turn had no tools — the agent silently appended it to history and ended the turn with no response, so it sat there unanswered. This affected both user-queued messages and agent-queued ones (send_to_tab). - agent.ts: stop the end-of-turn drain that swallowed trailing queued messages into history. They now stay on the queue. - agent-manager: after a CLEAN turn settles, continueFromQueue() drains the queue and starts a fresh turn to answer it. Skipped on a user-stopped or errored turn (queue preserved for the next send). - Loop safety: continuation draws from the existing autoWakeBudget, so a runaway agent<->agent chain is bounded; human sends refill it, so human conversations are never throttled. - dequeueMessages now tags message-consumed with reason "interrupt" | "continuation"; the frontend collapses continuation- consumed queued bubbles into the next turn's initiator row (avoids the linger/dup traps documented in queue-interrupt-reconcile-edge-cases.md). - Tests: agent (no-swallow + interrupt regression), agent-manager (continuation, no-op when empty, user-stop preserves queue, bounded loop), frontend (continuation bubble becomes next initiator). - wishlist: remove the now-fixed item. --- packages/core/src/agent/agent.ts | Bin 57628 -> 57720 bytes packages/core/src/types/index.ts | 16 ++++++- packages/core/tests/agent/agent.test.ts | 73 ++++++++++++++++++++++++++++++++ 3 files changed, 88 insertions(+), 1 deletion(-) (limited to 'packages/core') diff --git a/packages/core/src/agent/agent.ts b/packages/core/src/agent/agent.ts index c2dfef1..c6f1322 100644 Binary files a/packages/core/src/agent/agent.ts and b/packages/core/src/agent/agent.ts differ diff --git a/packages/core/src/types/index.ts b/packages/core/src/types/index.ts index a4230ca..ced3dc2 100644 --- a/packages/core/src/types/index.ts +++ b/packages/core/src/types/index.ts @@ -275,7 +275,21 @@ export type AgentEvent = agentModels?: Array<{ key_id: string; model_id: string }> | null; } | { type: "message-queued"; tabId: string; messageId: string; message: string } - | { type: "message-consumed"; tabId: string; messageIds: string[] } + | { + type: "message-consumed"; + tabId: string; + messageIds: string[]; + /** + * Why the queue was drained: + * - "interrupt": consumed mid-turn, folded into a running turn's tool + * result as a [USER INTERRUPT]. The optimistic bubble collapses into + * that sealed turn. + * - "continuation": consumed between turns to START a new turn. The + * optimistic bubble becomes that new turn's initiating user row. + * Absent ⇒ treat as "interrupt" (back-compat). + */ + reason?: "interrupt" | "continuation"; + } | { type: "message-cancelled"; tabId: string; messageId: string }; // ─── Tool Types ────────────────────────────────────────────────── diff --git a/packages/core/tests/agent/agent.test.ts b/packages/core/tests/agent/agent.test.ts index 7560443..188129c 100644 --- a/packages/core/tests/agent/agent.test.ts +++ b/packages/core/tests/agent/agent.test.ts @@ -219,6 +219,79 @@ describe("Agent", () => { }); }); + it("does NOT swallow trailing queued messages into history at turn end", async () => { + // Regression for the "queue not consumed after the turn ends" bug. A + // message that lands on the queue after the last tool call (here: a + // no-tool turn) must be LEFT on the queue for the orchestrator to start + // a new turn — not silently appended to history with no response. + vi.mocked(streamText).mockReturnValue( + makeMockStreamResult([{ type: "text-delta", id: "t0", text: "done" }, finishStop]), + ); + + const queue = [{ id: "q1", message: "answer me next", timestamp: 1 }]; + const dequeueMessages = vi.fn(() => queue.splice(0, queue.length)); + const agent = new Agent(makeConfig(), { + dequeueMessages, + waitForQueuedMessage: () => ({ promise: Promise.resolve(), cancel: () => {} }), + }); + + const before = agent.messages.length; + for await (const _ of agent.run("hello")) { + // consume + } + + // The agent appended exactly the user turn + its own assistant reply; + // it did NOT drain the queue or append a trailing user message for it. + expect(dequeueMessages).not.toHaveBeenCalled(); + expect(queue).toHaveLength(1); + const added = agent.messages.slice(before); + expect(added.map((m) => m.role)).toEqual(["user", "assistant"]); + expect( + added.some((m) => m.chunks.some((c) => c.type === "text" && c.text === "answer me next")), + ).toBe(false); + }); + + it("still injects a mid-turn queued message into the last tool result", async () => { + // The interrupt path (site 1) must be untouched by the turn-end fix: a + // message present DURING a tool batch is folded into that batch's last + // tool result as a [USER INTERRUPT], and the agent loops back to the LLM. + vi.mocked(streamText) + .mockReturnValueOnce( + makeMockStreamResult([ + { type: "tool-call", toolCallId: "tc1", toolName: "read_file", input: { path: "a.txt" } }, + finishToolCalls, + ]), + ) + .mockReturnValueOnce( + makeMockStreamResult([{ type: "text-delta", id: "t0", text: "ok" }, finishStop]), + ); + + const queue = [{ id: "q1", message: "stop and do X", timestamp: 1 }]; + const dequeueMessages = vi.fn(() => queue.splice(0, queue.length)); + const toolDef = { + name: "read_file", + description: "reads a file", + parameters: z.object({ path: z.string() }), + execute: async () => "file contents", + }; + const agent = new Agent(makeConfig({ tools: [toolDef] }), { + dequeueMessages, + waitForQueuedMessage: () => ({ promise: Promise.resolve(), cancel: () => {} }), + }); + + const events: AgentEvent[] = []; + for await (const event of agent.run("read it")) { + events.push(event); + } + + expect(dequeueMessages).toHaveBeenCalled(); + const toolResult = events.find((e) => e.type === "tool-result") as + | (AgentEvent & { toolResult: { result: string } }) + | undefined; + expect(toolResult?.toolResult.result).toContain("[USER INTERRUPT]"); + expect(toolResult?.toolResult.result).toContain("stop and do X"); + }); + it("yields reasoning-delta events", async () => { vi.mocked(streamText).mockReturnValue( makeMockStreamResult([ -- cgit v1.2.3