diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 4405d0996..30546d00c 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -4163,30 +4163,58 @@ export class AgentSession { return queued; } + /** + * Orders subscriber fan-out across concurrent `#emitSessionEvent` calls. + * Extension emits only await when the event type has handlers, so an event + * with no handlers could otherwise overtake an earlier event still inside + * its extension emit — an instant refusal delivered its assistant + * `message_end` to the TUI before its own `message_start`, skipping the + * turn-ending error render entirely. + */ + #subscriberEmitGate: Promise = Promise.resolve(); + async #emitSessionEvent(event: AgentSessionEvent): Promise { if (event.type === "message_update") { this.#emit(event); void this.#queueExtensionEvent(event); return; } - await this.#emitExtensionEvent(event); - // Hold the wire-level agent_end until in-flight prompts unwind. Subscribers - // (rpc-mode, ACP, Cursor) treat agent_end as the "session is idle" signal; - // emitting while #promptInFlightCount > 0 lets a client fire its next - // `prompt` into a session that still reports isStreaming === true. Flush - // happens in #endInFlight / #resetInFlight. A later agent_end (e.g. from - // an auto-compaction turn that starts before the original prompt unwinds) - // supersedes the pending one, which is what subscribers want — they only - // care about the final settle. - if (event.type === "agent_end" && this.#promptInFlightCount > 0) { - this.#pendingAgentEndEmit = event; - return; + // Take a FIFO ticket before the extension emit: extension deliveries for + // consecutive events still run concurrently, but subscriber fan-out waits + // for every earlier event's fan-out (or deferral) to happen first. + const previousGate = this.#subscriberEmitGate; + const { promise: gate, resolve: releaseGate } = Promise.withResolvers(); + this.#subscriberEmitGate = gate; + try { + await this.#emitExtensionEvent(event); + await previousGate; + // Hold the wire-level agent_end until in-flight prompts unwind. Subscribers + // (rpc-mode, ACP, Cursor) treat agent_end as the "session is idle" signal; + // emitting while #promptInFlightCount > 0 lets a client fire its next + // `prompt` into a session that still reports isStreaming === true. Flush + // happens in #endInFlight / #resetInFlight. A later agent_end (e.g. from + // an auto-compaction turn that starts before the original prompt unwinds) + // supersedes the pending one, which is what subscribers want — they only + // care about the final settle. + if (event.type === "agent_end" && this.#promptInFlightCount > 0) { + this.#pendingAgentEndEmit = event; + return; + } + this.#emit(event); + } finally { + releaseGate(); } - this.#emit(event); } // Track last assistant message for auto-compaction check #lastAssistantMessage: AssistantMessage | undefined = undefined; + /** + * Classifier-refusal turn pruned from active context at settle (#3591). + * Retained until the next run starts so post-settle readers + * ({@link getLastAssistantMessage}: print mode, task executor) still see + * the terminal error instead of a silently successful-looking state. + */ + #prunedTerminalRefusal: AssistantMessage | undefined = undefined; /** Internal handler for agent events - shared by subscribe and reconnect. * diff --git a/packages/coding-agent/test/agent-session-event-order.test.ts b/packages/coding-agent/test/agent-session-event-order.test.ts new file mode 100644 index 000000000..52f4c9343 --- /dev/null +++ b/packages/coding-agent/test/agent-session-event-order.test.ts @@ -0,0 +1,127 @@ +/** + * Subscriber event-order contract for `AgentSession`. + * + * Extension emits inside the session's event pipeline only await when the + * event type has registered handlers. For a turn whose provider events all + * land in one tick (e.g. an instant Anthropic classifier refusal: empty + * `message_start` + terminal `message_delta` + `message_stop` in a single SSE + * flush), an event type WITHOUT extension handlers used to overtake an earlier + * event WITH handlers — the TUI received the assistant `message_end` before + * its own `message_start`, so no streaming component existed yet and the + * turn-ending error (pinned banner + inline `Error:` line) was never rendered. + * The session must deliver events to subscribers in emission order regardless + * of which event types extensions subscribe to. + */ + +import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "bun:test"; +import * as path from "node:path"; +import { Agent } from "@oh-my-pi/pi-agent-core"; +import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock"; +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 type { ExtensionRunner } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/runner"; +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 { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; +import { TempDir } from "@oh-my-pi/pi-utils"; + +describe("AgentSession subscriber event order", () => { + let tempDir: TempDir; + let authStorage: AuthStorage; + let modelRegistry: ModelRegistry; + let session: AgentSession | undefined; + + beforeAll(async () => { + tempDir = TempDir.createSync("@pi-event-order-"); + authStorage = await AuthStorage.create(path.join(tempDir.path(), "testauth.db")); + authStorage.setRuntimeApiKey("anthropic", "anthropic-test-key"); + modelRegistry = new ModelRegistry(authStorage, path.join(tempDir.path(), "models.yml")); + }); + + afterAll(() => { + authStorage.close(); + tempDir.removeSync(); + }); + + afterEach(async () => { + if (session) { + await session.dispose(); + session = undefined; + } + vi.restoreAllMocks(); + }); + + it("delivers message_start before message_end when extension handlers are asymmetric", async () => { + const model = getBundledModel("anthropic", "claude-sonnet-4-5"); + if (!model) throw new Error("Expected bundled test model to exist"); + + // Instant terminal error turn: start and end reach the session pipeline + // in the same tick, exactly like a real classifier refusal. + const mock = createMockModel({ + responses: [ + { + stopReason: "error", + stopDetails: { type: "refusal", category: "cyber", explanation: "Declined." }, + errorMessage: "Refusal (cyber): Declined.", + }, + ], + }); + const agent = new Agent({ + getApiKey: agentModel => `${agentModel.provider}-test-key`, + initialState: { + model, + systemPrompt: ["Test"], + tools: [], + messages: [], + }, + streamFn: (streamModel, context, options) => mock.stream(streamModel, context, options), + }); + + // Handlers exist for message_start only, and each emit burns a burst of + // microtasks — so the handler-less message_end path would skip its + // extension await and overtake the start without subscriber-order + // serialization. + const extensionRunner = { + emit: vi.fn(async () => { + for (let hop = 0; hop < 25; hop++) await Promise.resolve(); + }), + emitBeforeAgentStart: vi.fn().mockResolvedValue(undefined), + hasHandlers: vi.fn((eventType: string) => eventType === "message_start"), + emitSessionStop: vi.fn().mockResolvedValue(undefined), + } as unknown as ExtensionRunner; + + const settings = Settings.isolated({ + "compaction.enabled": false, + "retry.baseDelayMs": 5, + "retry.maxRetries": 1, + "retry.modelFallback": false, + }); + settings.setModelRole("default", `${model.provider}/${model.id}`); + + session = new AgentSession({ + agent, + sessionManager: SessionManager.inMemory(), + settings, + modelRegistry, + extensionRunner, + }); + + const order: string[] = []; + session.subscribe(event => { + if (event.type === "message_start" || event.type === "message_end") { + order.push(`${event.type}:${event.message.role}`); + } + }); + + await session.prompt("Trigger instant refusal"); + await session.waitForIdle(); + + for (const role of ["user", "assistant"]) { + const startIndex = order.indexOf(`message_start:${role}`); + const endIndex = order.indexOf(`message_end:${role}`); + expect(startIndex).toBeGreaterThanOrEqual(0); + expect(endIndex).toBeGreaterThan(startIndex); + } + }); +}); diff --git a/packages/coding-agent/test/transcript-history-exactly-once.test.ts b/packages/coding-agent/test/transcript-history-exactly-once.test.ts new file mode 100644 index 000000000..1575467a7 --- /dev/null +++ b/packages/coding-agent/test/transcript-history-exactly-once.test.ts @@ -0,0 +1,103 @@ +import { describe, expect, it } from "bun:test"; +import { TranscriptContainer } from "@oh-my-pi/pi-coding-agent/modes/components/transcript-container"; +import { type Component, TUI } from "@oh-my-pi/pi-tui"; +import { StressRenderScheduler } from "../../tui/test/render-stress-scheduler"; +import { VirtualTerminal } from "../../tui/test/virtual-terminal"; + +/** + * Finalized history block. With `tracked`, reports a post-finalize content + * version like `AssistantMessageComponent`; otherwise it is version-untracked + * like most tool blocks. + */ +class HistoryBlock implements Component { + #lines: readonly string[]; + getTranscriptBlockVersion?: () => number; + constructor(lines: readonly string[], tracked: boolean) { + this.#lines = lines; + if (tracked) this.getTranscriptBlockVersion = () => 1; + } + render(width: number): readonly string[] { + return this.#lines.map(line => line.slice(0, width)); + } + isTranscriptBlockFinalized(): boolean { + return true; + } +} + +/** Streaming live block with a settled prefix, like a streaming assistant reply. */ +class LiveBlock implements Component { + lines: string[] = ["live-000"]; + settled = 0; + render(width: number): readonly string[] { + return this.lines.map(line => line.slice(0, width)); + } + isTranscriptBlockFinalized(): boolean { + return false; + } + getTranscriptBlockSettledRows(): number { + return this.settled; + } +} + +// Streams a live block behind a run of small finalized history blocks until the +// history fully commits to native scrollback, then verifies exactly-once history +// on the terminal tape. Regression guard: transcript-side committed-prefix +// compaction (dropping committed rows from the local frame) shifted the frame +// under the engine's committed-prefix ledger, the audit re-anchored, and +// already-taped rows were recommitted below their first copy — visibly +// duplicated blocks. The transcript now always keeps its full local frame. +async function streamPastCommit(tracked: boolean): Promise> { + const term = new VirtualTerminal(40, 6); + Object.defineProperty(term, "isNativeViewportAtBottom", { configurable: true, value: () => undefined }); + const scheduler = new StressRenderScheduler(); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); + const chat = new TranscriptContainer(); + const historyRows: string[] = []; + for (let i = 0; i < 6; i++) { + const rows = [`box-${i}-alpha`, `box-${i}-beta`]; + historyRows.push(...rows); + chat.addChild(new HistoryBlock(rows, tracked)); + } + const live = new LiveBlock(); + chat.addChild(live); + tui.addChild(chat); + + try { + tui.start(); + await scheduler.drain(term); + // Grow the live block one row per frame with the settled prefix trailing + // by one, pushing the finalized history through commit and compaction. + for (let i = 1; i < 40; i++) { + live.lines.push(`live-${String(i).padStart(3, "0")}`); + live.settled = live.lines.length - 1; + tui.requestRender(); + await scheduler.drain(term); + } + } finally { + tui.stop(); + await term.flush(); + } + + const counts = new Map(); + for (const row of term.getScrollBuffer()) { + const text = Bun.stripANSI(row).trimEnd(); + if (text.length === 0) continue; + counts.set(text, (counts.get(text) ?? 0) + 1); + } + // Loss check alongside the duplication check: every history row must have + // reached the tape exactly once. + for (const row of historyRows) expect(counts.get(row) ?? 0).toBe(1); + return counts; +} + +describe("transcript committed history", () => { + it("keeps version-tracked committed history exactly once on the tape", async () => { + const counts = await streamPastCommit(true); + expect([...counts.entries()].filter(([, count]) => count > 1)).toEqual([]); + }); + + it("keeps version-untracked committed history exactly once on the tape", async () => { + const counts = await streamPastCommit(false); + expect([...counts.entries()].filter(([, count]) => count > 1)).toEqual([]); + }); +});