From d4d89fadb94a93dc46e44f3ec4c211c883fc41fa Mon Sep 17 00:00:00 2001 From: can1357 Date: Sat, 13 Jun 2026 16:17:56 +0200 Subject: [PATCH] fix: resolved queued-message flush flow by coalescing in-flight calls - Coalesced concurrent interruptAndFlushQueuedMessages() calls through one in-flight promise. - Replaced continue()-retry logic with agent.prompt() to flush queued messages from empty contexts. - Skipped queued-message flush replay while compacting or streaming to avoid turn overlap. - Updated flush path to consume queued steering first, then follow-ups, via dequeuing helper. - Added a regression test for empty-state interrupt-and-flush delivering queued steers safely. --- packages/agent/src/agent.ts | 12 ++ .../coding-agent/src/session/agent-session.ts | 76 ++++------ ...interrupt-and-flush-empty-messages.test.ts | 142 ++++++++++++++++++ .../event-controller-interrupt.test.ts | 1 + 4 files changed, 188 insertions(+), 43 deletions(-) create mode 100644 packages/coding-agent/test/issue-interrupt-and-flush-empty-messages.test.ts diff --git a/packages/agent/src/agent.ts b/packages/agent/src/agent.ts index ee81971e6..821320b57 100644 --- a/packages/agent/src/agent.ts +++ b/packages/agent/src/agent.ts @@ -751,6 +751,18 @@ export class Agent { return this.#steeringQueue.length > 0 || this.#followUpQueue.length > 0; } + /** + * Drain queued messages for an immediate resume: all steering messages if any + * (honoring steeringMode), otherwise follow-ups. Mirrors continue()'s dequeue + * precedence so an empty-Enter flush and the natural turn-end resume agree on + * what to deliver next. + */ + takeQueuedMessages(): AgentMessage[] { + const steering = this.#dequeueSteeringMessages(); + if (steering.length > 0) return steering; + return this.#dequeueFollowUpMessages(); + } + #dequeueSteeringMessages(): AgentMessage[] { if (this.#steeringMode === "one-at-a-time") { if (this.#steeringQueue.length > 0) { diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 398c1bdca..1099629a1 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -1084,7 +1084,7 @@ export class AgentSession { // both #endInFlight (normal) and #resetInFlight (abort). #pendingAgentEndEmit: AgentSessionEvent | undefined; #resumingQueuedMessages = false; - #queuedMessageFlushStartPromise: Promise<{ continuation: Promise | undefined }> | undefined; + #queuedFlushInterrupt: Promise | undefined; #obfuscator: SecretObfuscator | undefined; #checkpointState: CheckpointState | undefined = undefined; #pendingRewindReport: string | undefined = undefined; @@ -5509,61 +5509,51 @@ export class AgentSession { /** * Abort active work, then immediately resume the agent so queued steer/follow-up * messages drain instead of waiting for another natural turn boundary. + * + * The drained queue is re-run via `agent.prompt()`, which appends and runs + * regardless of the trailing message role. Earlier this called + * `agent.continue()`, which only dequeues steers after an assistant message + * and throws "No messages to continue from" on an empty context — both states + * a flushed empty-Enter can land in, which stranded the steer in the queue and + * surfaced that error in the TUI. + * + * Concurrent calls (e.g. a double empty-Enter) coalesce: the first owns the + * abort→resume handoff while the rest await it and return. The guard clears + * once the resumed turn has started (not at turn end), so a genuinely-later + * flush of a freshly queued steer still works. */ async interruptAndFlushQueuedMessages(options?: { reason?: string }): Promise { - const existingStart = this.#queuedMessageFlushStartPromise; - if (existingStart) { - const { continuation } = await existingStart; - await continuation; + const inFlight = this.#queuedFlushInterrupt; + if (inFlight) { + await inFlight; return; } if (!this.agent.hasQueuedMessages()) return; - const start = this.#beginQueuedMessageFlush(options); - this.#queuedMessageFlushStartPromise = start; - let continuation: Promise | undefined; - try { - ({ continuation } = await start); - } finally { - if (this.#queuedMessageFlushStartPromise === start) { - this.#queuedMessageFlushStartPromise = undefined; - } - } - await continuation; - } - - async #beginQueuedMessageFlush(options?: { reason?: string }): Promise<{ continuation: Promise | undefined }> { - if (!this.agent.hasQueuedMessages()) return { continuation: undefined }; + const { promise: handoff, resolve: settleHandoff } = Promise.withResolvers(); + this.#queuedFlushInterrupt = handoff; this.#resumingQueuedMessages = true; + let turn: Promise | undefined; try { await this.abort({ reason: options?.reason }); - if (!this.agent.hasQueuedMessages()) return { continuation: undefined }; + if (this.isCompacting || this.isGeneratingHandoff) return; await this.#maybeRestoreRetryFallbackPrimary(); - if (!this.agent.hasQueuedMessages()) return { continuation: undefined }; - const continuation = this.#continueQueuedMessagesWithIdleRetry(); - return { continuation }; + // A turn slipped in while we were settling (e.g. a fresh user prompt): + // leave the queue for its next steering boundary rather than draining + // and double-running, which would throw AgentBusyError. Checking + // isStreaming and prompting below is synchronous, so no turn can start + // between the drain and the resume. + if (this.agent.state.isStreaming) return; + const queued = this.agent.takeQueuedMessages(); + if (queued.length === 0) return; + this.#resumingQueuedMessages = false; + turn = this.agent.prompt(queued); } finally { this.#resumingQueuedMessages = false; + this.#queuedFlushInterrupt = undefined; + settleHandoff(); } - } - - async #continueQueuedMessagesWithIdleRetry(): Promise { - const deadline = Date.now() + 30_000; - for (;;) { - if (!this.agent.hasQueuedMessages()) return; - try { - await this.agent.continue(); - return; - } catch (err) { - if (!(err instanceof AgentBusyError)) { - throw err; - } - if (Date.now() >= deadline) { - throw new Error("Timed out waiting for prior agent run to finish before continuing queued messages."); - } - await this.agent.waitForIdle(); - } - } + await turn; } /** diff --git a/packages/coding-agent/test/issue-interrupt-and-flush-empty-messages.test.ts b/packages/coding-agent/test/issue-interrupt-and-flush-empty-messages.test.ts new file mode 100644 index 000000000..55c5fd4fe --- /dev/null +++ b/packages/coding-agent/test/issue-interrupt-and-flush-empty-messages.test.ts @@ -0,0 +1,142 @@ +/** + * Regression test for the empty-messages interrupt-and-flush bug. + * + * Scenario the user hit: + * 1. The agent is already streaming (a hidden turn: context promotion, + * auto-retry, auto-compaction, an extension-emitted turn, etc.). + * 2. The user types "hi" and presses Enter. Because isStreaming is true, + * the input controller routes the message to `session.steer("hi")` + * (see the streamingBehavior branch in `AgentSession.prompt`). + * Critically, the user message is NEVER appended to `#state.messages` + * in the agent — it lives only in the steering queue. + * 3. The user presses Enter again on an empty editor to inject the steer + * immediately. The input controller calls + * `session.interruptAndFlushQueuedMessages()`. + * 4. `interruptAndFlushQueuedMessages` aborts the streaming turn, then + * calls `agent.continue()` to drain the queued steer. `continue()` is: + * + * if (messages.length === 0) throw "No messages to continue from"; + * if (last role === "assistant") { drainSteering(); runLoop(drained); } + * else runLoop(undefined); + * + * When `#state.messages` is empty, step 4 throws and: + * - `agent.continue()` propagates the throw out of + * `interruptAndFlushQueuedMessages`. + * - The input controller catches it and surfaces the error string + * ("Error: No messages to continue from") — the visible + * "Error:" line in the user's first screenshot. + * - The steering queue is never drained, so the session's + * `#steeringMessages` mirror is never cleared by the + * message_start handler. The "Steer: hi" chip stays visible + * forever (the second screenshot) until the user manually + * dequeues with Alt+Up or starts a new session. + * + * Contract this test defends: + * - `interruptAndFlushQueuedMessages` MUST deliver a queued steer even + * when the agent has no prior messages in `#state.messages`. + * - The flush MUST NOT reject with any error. + * - It MUST clear both the agent queue and the session mirror in a + * single atomic step, so the visible chip disappears the moment the + * user presses Enter, not after a downstream message_start event + * fires. + * - The delivered messages must contain the original steer text. + */ + +import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { Agent } from "@oh-my-pi/pi-agent-core"; +import type { Message } from "@oh-my-pi/pi-ai"; +import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; +import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; +import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; +import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session"; +import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage"; +import { convertToLlm } from "@oh-my-pi/pi-coding-agent/session/messages"; +import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; +import { Snowflake } from "@oh-my-pi/pi-utils"; +import { createAssistantMessage } from "./helpers/agent-session-setup"; + +describe("interrupt-and-flush must not strand a queued steer when the agent has no prior messages", () => { + let session: AgentSession; + let tempDir: string; + const authStorages: AuthStorage[] = []; + + beforeEach(() => { + tempDir = path.join(os.tmpdir(), `pi-flush-empty-messages-${Snowflake.next()}`); + fs.mkdirSync(tempDir, { recursive: true }); + }); + + afterEach(async () => { + if (session) { + await session.dispose(); + } + for (const authStorage of authStorages.splice(0)) { + authStorage.close(); + } + if (tempDir && fs.existsSync(tempDir)) { + fs.rmSync(tempDir, { recursive: true }); + } + }); + + it("delivers a steer queued on a session whose #state.messages is empty", async () => { + const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; + const callMessages: Message[][] = []; + + const agent = new Agent({ + getApiKey: () => "test-key", + initialState: { model, systemPrompt: ["Test"], tools: [] }, + convertToLlm, + streamFn(_model, context, _options) { + callMessages.push([...context.messages]); + const stream = new AssistantMessageEventStream(); + queueMicrotask(() => { + stream.push({ type: "start", partial: createAssistantMessage("") }); + stream.push({ type: "done", reason: "stop", message: createAssistantMessage("ack") }); + }); + return stream; + }, + }); + + const sessionManager = SessionManager.inMemory(); + const settings = Settings.isolated(); + const authStorage = await AuthStorage.create(path.join(tempDir, "auth.db")); + authStorages.push(authStorage); + const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml")); + authStorage.setRuntimeApiKey("anthropic", "test-key"); + + session = new AgentSession({ agent, sessionManager, settings, modelRegistry }); + + // Reproduce the user flow: a hidden turn is already streaming, so + // the input controller routed the user's "hi" to session.steer + // instead of session.prompt. The session mirror and agent queue + // are populated; #state.messages stays empty. + await session.steer("hi"); + expect(session.agent.hasQueuedMessages()).toBe(true); + expect(session.getQueuedMessages().steering).toEqual(["hi"]); + + // Empty-Enter → interruptAndFlushQueuedMessages. Must NOT throw. + // Awaiting the flush is enough: it resolves only after the resumed + // turn's prompt() has run and the assistant message has been + // emitted. The call to callMessages.push() happens synchronously + // inside streamFn, so by the time the flush resolves the model + // call has been recorded. + await expect( + session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" }), + ).resolves.toBeUndefined(); + + // The steer must be delivered as a fresh turn. + expect(callMessages.length).toBeGreaterThanOrEqual(1); + const delivered = callMessages[0]?.some((message: Message) => { + if (typeof message.content === "string") return message.content.includes("hi"); + return message.content.some(c => c.type === "text" && c.text.includes("hi")); + }); + expect(delivered).toBe(true); + + // Both queues must clear atomically — no stranded chip. + expect(session.agent.hasQueuedMessages()).toBe(false); + expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] }); + }); +}); diff --git a/packages/coding-agent/test/modes/controllers/event-controller-interrupt.test.ts b/packages/coding-agent/test/modes/controllers/event-controller-interrupt.test.ts index 943b76e23..ef66ad274 100644 --- a/packages/coding-agent/test/modes/controllers/event-controller-interrupt.test.ts +++ b/packages/coding-agent/test/modes/controllers/event-controller-interrupt.test.ts @@ -132,4 +132,5 @@ describe("EventController user interrupt acknowledgement", () => { expect(setWorkingMessage).not.toHaveBeenCalled(); }); + });