From a51ac2decb032587c9f8915902c7e49e4623f3e3 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Sun, 27 Sep 2026 23:07:58 +0100 Subject: [PATCH 01/12] fix(web): close wedged websockets with an app.version heartbeat A half-open socket answers nothing and never fires close, so every recovery mechanism (reconnect, hydration, subscription replay) waited forever while the UI showed stale streaming state. A 30s interval now issues the existing app.version RPC with a 10s deadline; an unanswered probe closes the socket locally and the established onclose ladder owns recovery. agent.send gains a 20s timeout so a composer submit cannot park forever on a dead socket. The wedged-socket test now answers the liveness probe, since a socket that ignores it is exactly what the watchdog closes. --- apps/web/src/__tests__/connection.test.ts | 48 +++++++++++++++++++ .../__tests__/rpc-timeout-recovery.test.ts | 5 ++ apps/web/src/transport/ws-transport.ts | 43 ++++++++++++++++- 3 files changed, 95 insertions(+), 1 deletion(-) diff --git a/apps/web/src/__tests__/connection.test.ts b/apps/web/src/__tests__/connection.test.ts index 5f8a69d02..2f7a063d8 100644 --- a/apps/web/src/__tests__/connection.test.ts +++ b/apps/web/src/__tests__/connection.test.ts @@ -234,6 +234,54 @@ describe("4001 auth failure handling", () => { }); }); +describe("liveness heartbeat", () => { + it("closes a silent socket when the heartbeat goes unanswered", async () => { + vi.useFakeTimers(); + const statusSpy = vi.fn(); + const transport = createWsTransport("ws://localhost:1234", { + onStatusChange: statusSpy, + }); + mockWsInstance.simulateOpen(); + + // 30s interval fires the probe; the mock never responds; the 10s + // heartbeat timeout closes the socket and drives the reconnect path. + await vi.advanceTimersByTimeAsync(30_000); + statusSpy.mockClear(); + await vi.advanceTimersByTimeAsync(10_000); + + expect(statusSpy).toHaveBeenCalledWith("reconnecting"); + transport.close(); + vi.useRealTimers(); + }); + + it("keeps the socket open when the heartbeat is answered", async () => { + vi.useFakeTimers(); + const statusSpy = vi.fn(); + const transport = createWsTransport("ws://localhost:1234", { + onStatusChange: statusSpy, + }); + mockWsInstance.simulateOpen(); + const sendSpy = vi.spyOn(mockWsInstance, "send"); + + await vi.advanceTimersByTimeAsync(30_000); + const heartbeat = sendSpy.mock.calls + .map(([raw]) => JSON.parse(raw as string) as { id: string; method?: string }) + .find((message) => message.method === "app.version"); + expect(heartbeat).toBeDefined(); + mockWsInstance.onmessage?.({ + data: JSON.stringify({ id: heartbeat!.id, result: "0.0.1" }), + }); + + await vi.advanceTimersByTimeAsync(10_000); + expect(mockWsInstance.readyState).toBe(1); + statusSpy.mockClear(); + await vi.advanceTimersByTimeAsync(10_000); + expect(statusSpy).not.toHaveBeenCalledWith("reconnecting"); + transport.close(); + vi.useRealTimers(); + }); +}); + /** * Browser-faithful socket: send() throws while CONNECTING and silently * discards while CLOSING/CLOSED, matching the WebSocket spec. The silent diff --git a/apps/web/src/transport/__tests__/rpc-timeout-recovery.test.ts b/apps/web/src/transport/__tests__/rpc-timeout-recovery.test.ts index af0d5cc3c..28ffc90f0 100644 --- a/apps/web/src/transport/__tests__/rpc-timeout-recovery.test.ts +++ b/apps/web/src/transport/__tests__/rpc-timeout-recovery.test.ts @@ -49,6 +49,11 @@ class TimeoutSocket { const parsed: unknown = JSON.parse(raw); if (!isRpcRequest(parsed)) throw new Error("Expected an RPC request"); this.requests.push(parsed); + // A socket that ignores the liveness probe is closed by the watchdog; + // these tests exercise RPC timeouts on a live connection instead. + if (parsed.method === "app.version") { + queueMicrotask(() => this.respond(parsed, "0.0.1")); + } } open(): void { diff --git a/apps/web/src/transport/ws-transport.ts b/apps/web/src/transport/ws-transport.ts index 5329652c7..925d3caf4 100644 --- a/apps/web/src/transport/ws-transport.ts +++ b/apps/web/src/transport/ws-transport.ts @@ -107,6 +107,13 @@ const TERMINAL_LATE_CREATE_CLEANUP_TIMEOUT_MS = 10_000; /** Maximum Terminal creates whose late responses still need exact cleanup. */ const MAX_PENDING_TERMINAL_CREATE_CLEANUPS = 8; +/** Interval between liveness probes on an open socket. */ +const HEARTBEAT_INTERVAL_MS = 30_000; +/** How long a heartbeat may go unanswered before the socket is treated as dead. */ +const HEARTBEAT_TIMEOUT_MS = 10_000; +/** Deadline for `agent.send` so a composer submit cannot park forever on a dead socket. */ +const SEND_MESSAGE_TIMEOUT_MS = 20_000; + /** Last thread-list refresh timestamp per workspace, triggered on WS reconnect. */ const lastLoadThreadsAtByWorkspace = new Map(); /** Minimum interval between reconnect-triggered thread-list fetches to avoid rapid-reconnect storms. */ @@ -386,6 +393,8 @@ export function createWsTransport( let closed = false; let reconnectDelay = MIN_RECONNECT_MS; let reconnectTimer: ReturnType | null = null; + let heartbeatTimer: ReturnType | null = null; + let heartbeatInFlight = false; let terminalSelectionPromise: Promise | null = null; // Track consecutive auth failures so we apply backoff after 3 immediate // retries, preventing a tight loop when the token is persistently wrong. @@ -554,6 +563,7 @@ export function createWsTransport( setAttachmentTransportWsUrl(url); resolveReady(); options?.onStatusChange?.("connected"); + startHeartbeat(); invalidateLiveTurnDiff(); void selectTerminalClientWithRecovery(); @@ -614,6 +624,7 @@ export function createWsTransport( ws.onclose = (event: CloseEvent) => { freshTurnDiffThreads.clear(); + stopHeartbeat(); rejectPending("WebSocket disconnected"); lateResponseHandlers = new Map(); pendingTerminalCreateCleanups.clear(); @@ -640,6 +651,35 @@ export function createWsTransport( pending = new Map(); } + // A half-open socket answers nothing and never fires `close`; without a + // probe, reconnect and runtime resync wait forever while the UI shows stale + // streaming state. The probe is a cheap existing RPC bounded by timeoutMs; + // on timeout we close locally so onclose drives the recovery ladder. + function heartbeatTick(): void { + if (heartbeatInFlight || closed || ws.readyState !== WebSocket.OPEN) return; + heartbeatInFlight = true; + rpc("app.version", {}, { timeoutMs: HEARTBEAT_TIMEOUT_MS }) + .catch(() => { + if (ws.readyState === WebSocket.OPEN) ws.close(); + }) + .finally(() => { + heartbeatInFlight = false; + }); + } + + function startHeartbeat(): void { + stopHeartbeat(); + heartbeatTimer = setInterval(heartbeatTick, HEARTBEAT_INTERVAL_MS); + } + + function stopHeartbeat(): void { + if (heartbeatTimer) { + clearInterval(heartbeatTimer); + heartbeatTimer = null; + } + heartbeatInFlight = false; + } + function invalidateLiveTurnDiff(): void { void Promise.all([import("@/features/projects/state/workspaceStore"), import("@/stores/diffStore")]).then(([workspace, diff]) => { const threadId = workspace.useWorkspaceStore.getState().activeThreadId; @@ -1153,7 +1193,7 @@ export function createWsTransport( ...(replyToMessageId && { replyToMessageId }), ...(quotedText && { quotedText }), ...guardrails, - }); + }, { timeoutMs: SEND_MESSAGE_TIMEOUT_MS }); }, getRecoveryIncident: () => rpc("agent.recoveryIncident", {}), @@ -1543,6 +1583,7 @@ export function createWsTransport( // Lifecycle close: () => { closed = true; + stopHeartbeat(); if (reconnectTimer) { clearTimeout(reconnectTimer); reconnectTimer = null; From 64d4cdd8359b97623fcc3950b56dd04f9a381c30 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Sun, 27 Sep 2026 23:08:29 +0100 Subject: [PATCH 02/12] feat(agent-model): stamp turn records with their execution identity The reducer now writes routing.executionId onto the stored turn so clients can correlate a canonical turn to the local runtime without threading the event envelope through state. A payload-supplied executionId that disagrees with routing is rejected as a routing conflict. --- packages/agent-model/src/records.ts | 4 ++++ packages/agent-model/src/reducer.ts | 8 +++++++- 2 files changed, 11 insertions(+), 1 deletion(-) diff --git a/packages/agent-model/src/records.ts b/packages/agent-model/src/records.ts index 16f103648..c5a2d244b 100644 --- a/packages/agent-model/src/records.ts +++ b/packages/agent-model/src/records.ts @@ -2,6 +2,7 @@ import { z } from "zod"; import { AgentItemIdSchema, AgentThreadIdSchema, + AgentTurnExecutionIdSchema, AgentTurnIdSchema, CanonicalTimestampSchema, CollaborationActionIdSchema, @@ -86,6 +87,9 @@ export const AgentTurnSchema = z threadId: AgentThreadIdSchema, status: AgentTurnStatusSchema, trigger: AgentTurnTriggerSchema, + // Stamped by the reducer from routing so clients can correlate the turn to + // the local runtime without threading the event envelope through state. + executionId: AgentTurnExecutionIdSchema.optional(), permissionMode: z.enum(["supervised", "full"]), approvalReviewMode: z.enum(["manual", "automatic"]), approvalReviewReason: z.string().min(1).max(128), diff --git a/packages/agent-model/src/reducer.ts b/packages/agent-model/src/reducer.ts index 83ab5ff85..29880fa9e 100644 --- a/packages/agent-model/src/reducer.ts +++ b/packages/agent-model/src/reducer.ts @@ -243,12 +243,18 @@ function reduceTurnCreated( if (event.routing.threadId !== turn.threadId || event.routing.turnId !== turn.id) { return { state, outcome: "routing-conflict" }; } + if (turn.executionId !== undefined && turn.executionId !== event.routing.executionId) { + return { state, outcome: "routing-conflict" }; + } const currentTurn = state.turns[turn.id]; if (currentTurn && TERMINAL_TURN_STATUSES.has(currentTurn.status)) { return { state: { ...state, ...acceptedInputState }, outcome: "terminal-outcome-confirmed" }; } + const stored = turn.executionId === event.routing.executionId + ? turn + : { ...turn, executionId: event.routing.executionId }; return { - state: { ...state, turns: { ...state.turns, [turn.id]: turn }, ...acceptedInputState }, + state: { ...state, turns: { ...state.turns, [turn.id]: stored }, ...acceptedInputState }, outcome: "applied", }; } From 5059a3fffaf2b4de723c24c9122fd54940beda97 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Sun, 27 Sep 2026 23:08:29 +0100 Subject: [PATCH 03/12] fix(web): derive thread runtime phase from canonical turn state Composer, project tree, and running-thread indicators now follow the canonical replica once a turn correlates by execution identity, so a dropped legacy terminal event can no longer strand runningThreadIds. The optimistic window before turnStarted and the finalizing phase stay on the legacy path until a canonical turn proves correlation. The narrative projection keeps its existing child-only lifecycle gate. --- .../canonical-runtime-reconcile.test.ts | 147 ++++++++++++++++++ apps/web/src/stores/thread-lifecycle.ts | 54 +++++-- apps/web/src/stores/threadStore.ts | 40 ++++- 3 files changed, 228 insertions(+), 13 deletions(-) create mode 100644 apps/web/src/__tests__/canonical-runtime-reconcile.test.ts diff --git a/apps/web/src/__tests__/canonical-runtime-reconcile.test.ts b/apps/web/src/__tests__/canonical-runtime-reconcile.test.ts new file mode 100644 index 000000000..c84d53093 --- /dev/null +++ b/apps/web/src/__tests__/canonical-runtime-reconcile.test.ts @@ -0,0 +1,147 @@ +import { describe, expect, it, beforeEach, vi } from "vitest"; +import type { CanonicalAgentEventEnvelope } from "@mcode/contracts"; +import { useThreadStore } from "@/stores/threadStore"; +import { + resetThreadStoreForTests, + seedThreadRecord, + readThreadField, +} from "@/stores/thread-store-test-utils"; +import { clearRecordCache } from "@/features/conversation/hydration/record-cache"; +import { mockTransport } from "./mocks/transport"; + +vi.mock("@/transport", async () => ({ + ...(await vi.importActual("@/transport")), + getTransport: () => mockTransport, +})); + +const THREAD_ID = "thread-1"; +const TURN_ID = "turn-1"; +const EXECUTION_ID = "00000000-0000-4000-8000-000000000001"; +const OTHER_EXECUTION_ID = "00000000-0000-4000-8000-000000000099"; +const NOW = "2026-08-11T12:00:00.000Z"; + +function envelope( + eventId: string, + acceptedSequence: number, + executionId: string, + payload: CanonicalAgentEventEnvelope["payload"], +): CanonicalAgentEventEnvelope { + return { + eventId, + routing: { + threadId: THREAD_ID, + turnId: payload.type === "thread.recorded" ? undefined : TURN_ID, + executionId, + }, + sourceProviderId: "codex", + sourceIdentities: [], + acceptedSequence, + durableRevision: acceptedSequence, + serverTimestamps: { acceptedAt: NOW, persistedAt: NOW }, + payload, + }; +} + +function turnEvents(executionId: string, terminal?: "completed" | "interrupted"): CanonicalAgentEventEnvelope[] { + const events: CanonicalAgentEventEnvelope[] = [ + envelope("e1", 1, executionId, { + type: "thread.recorded", + thread: { + id: THREAD_ID, + workspaceId: "workspace-1", + rootThreadId: THREAD_ID, + providerId: "codex", + providerIdentities: [], + activityState: "Active", + conversationRevision: 1, + rosterRevision: 0, + createdAt: NOW, + updatedAt: NOW, + }, + }), + envelope("e2", 2, executionId, { + type: "turn.created", + turn: { + id: TURN_ID, + threadId: THREAD_ID, + status: "Pending", + trigger: { kind: "user" }, + permissionMode: "full", + approvalReviewMode: "manual", + approvalReviewReason: "manual-requested", + providerIdentities: [], + startedAt: null, + endedAt: null, + createdAt: NOW, + updatedAt: NOW, + }, + }), + envelope("e3", 3, executionId, { type: "turn.started", startedAt: NOW }), + ]; + if (terminal === "completed") { + events.push(envelope("e4", 4, executionId, { type: "turn.completed", endedAt: NOW })); + } else if (terminal === "interrupted") { + events.push(envelope("e4", 4, executionId, { type: "turn.interrupted", endedAt: NOW, reason: "restart" })); + } + return events; +} + +function seedRuntime(phase: "running" | "idle", turnExecutionId: string | null, running: boolean) { + useThreadStore.setState({ + records: seedThreadRecord(THREAD_ID, { runtimePhase: phase, turnExecutionId }), + runningThreadIds: running ? new Set([THREAD_ID]) : new Set(), + }); +} + +describe("canonical runtime reconciliation", () => { + beforeEach(() => { + clearRecordCache(); + resetThreadStoreForTests(); + vi.clearAllMocks(); + }); + + it("clears a stale running flag when the correlated canonical turn terminates", () => { + seedRuntime("running", EXECUTION_ID, true); + + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, turnEvents(EXECUTION_ID, "completed")); + + expect(readThreadField(THREAD_ID, (r) => r.runtimePhase)).toBe("completed"); + expect(useThreadStore.getState().runningThreadIds.has(THREAD_ID)).toBe(false); + }); + + it("claims an idle record when canonical reports a running turn", () => { + seedRuntime("idle", null, false); + + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, turnEvents(EXECUTION_ID)); + + expect(readThreadField(THREAD_ID, (r) => r.runtimePhase)).toBe("running"); + expect(useThreadStore.getState().runningThreadIds.has(THREAD_ID)).toBe(true); + }); + + it("keeps the optimistic running window while no canonical turn correlates", () => { + seedRuntime("running", null, true); + + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, turnEvents(EXECUTION_ID, "completed")); + + expect(readThreadField(THREAD_ID, (r) => r.runtimePhase)).toBe("running"); + expect(useThreadStore.getState().runningThreadIds.has(THREAD_ID)).toBe(true); + }); + + it("ignores a canonical turn carrying a different execution identity", () => { + seedRuntime("running", OTHER_EXECUTION_ID, true); + + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, turnEvents(EXECUTION_ID, "completed")); + + expect(readThreadField(THREAD_ID, (r) => r.runtimePhase)).toBe("running"); + expect(useThreadStore.getState().runningThreadIds.has(THREAD_ID)).toBe(true); + }); + + it("maps an interrupted canonical turn to the interrupted phase", () => { + seedRuntime("running", EXECUTION_ID, true); + + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, turnEvents(EXECUTION_ID, "interrupted")); + + expect(readThreadField(THREAD_ID, (r) => r.runtimePhase)).toBe("interrupted"); + expect(useThreadStore.getState().runningThreadIds.has(THREAD_ID)).toBe(false); + }); +}); diff --git a/apps/web/src/stores/thread-lifecycle.ts b/apps/web/src/stores/thread-lifecycle.ts index cd30db2f3..6d7cc57cc 100644 --- a/apps/web/src/stores/thread-lifecycle.ts +++ b/apps/web/src/stores/thread-lifecycle.ts @@ -1,24 +1,47 @@ -import type { AgentTurn, TurnRuntimePhase } from "@mcode/contracts"; +import type { AgentTurn, AgentTurnStatus, TurnRuntimePhase } from "@mcode/contracts"; import type { ThreadRecord } from "./thread-record"; -/** Selects the provider-owned child lifecycle when no local execution is active. */ -export function getCanonicalLifecycleTurn(threadId: string, record: ThreadRecord): AgentTurn | undefined { - if (record.runtimePhase === "running" || record.runtimePhase === "finalizing") return undefined; - const latest = Object.values(record.canonicalAgent.state.turns) +function latestCanonicalTurn(threadId: string, record: ThreadRecord): AgentTurn | undefined { + return Object.values(record.canonicalAgent.state.turns) .filter((turn) => turn.threadId === threadId) .sort((left, right) => Date.parse(left.startedAt ?? left.createdAt) - Date.parse(right.startedAt ?? right.createdAt) || left.id.localeCompare(right.id)) .at(-1); +} + +/** Selects the provider-owned child lifecycle when no local execution is active. */ +export function getCanonicalLifecycleTurn(threadId: string, record: ThreadRecord): AgentTurn | undefined { + if (record.runtimePhase === "running" || record.runtimePhase === "finalizing") return undefined; + const latest = latestCanonicalTurn(threadId, record); if (latest?.trigger.kind !== "child") return undefined; if (record.runtimePhase !== "idle" && (latest.status === "Pending" || latest.status === "Running")) return undefined; return latest; } -/** Resolves the lifecycle shared by transcript, Composer, and running-thread indicators. */ -export function getThreadRuntimePhase(threadId: string, record: ThreadRecord): TurnRuntimePhase { - const turn = getCanonicalLifecycleTurn(threadId, record); - if (!turn) return record.runtimePhase; - switch (turn.status) { +/** + * Canonical turn that owns this record's runtime phase, once correlation is + * proven. A tracked local execution only follows the canonical turn carrying + * the same execution identity, so an older persisted turn cannot clear or + * resurrect a live run while a matching terminal can clear a stale one. + */ +export function getCanonicalRuntimeTurn(threadId: string, record: ThreadRecord): AgentTurn | undefined { + if (record.runtimePhase === "finalizing") return undefined; + const latest = latestCanonicalTurn(threadId, record); + if (!latest) return undefined; + if (record.turnExecutionId !== null) { + return latest.executionId === record.turnExecutionId ? latest : undefined; + } + // Optimistic sends run without an identity until turnStarted lands; a prior + // terminal turn must not cancel that window. + if (record.runtimePhase === "running") return undefined; + // Without an identity, canonical Pending or Running claims only idle + // records; terminal truth may claim any non-busy record. + if ((latest.status === "Pending" || latest.status === "Running") && record.runtimePhase !== "idle") return undefined; + return latest; +} + +function phaseForTurnStatus(status: AgentTurnStatus): TurnRuntimePhase { + switch (status) { case "Pending": case "Running": return "running"; case "Completed": return "completed"; @@ -28,6 +51,17 @@ export function getThreadRuntimePhase(threadId: string, record: ThreadRecord): T } } +/** Canonical-owned runtime phase for this record, or null when legacy owns it. */ +export function getCanonicalRuntimePhase(threadId: string, record: ThreadRecord): TurnRuntimePhase | null { + const turn = getCanonicalRuntimeTurn(threadId, record); + return turn ? phaseForTurnStatus(turn.status) : null; +} + +/** Resolves the lifecycle shared by transcript, Composer, and running-thread indicators. */ +export function getThreadRuntimePhase(threadId: string, record: ThreadRecord): TurnRuntimePhase { + return getCanonicalRuntimePhase(threadId, record) ?? record.runtimePhase; +} + /** Whether the current lifecycle still owns execution, including final persistence. */ export function isThreadRuntimeActive(threadId: string, record: ThreadRecord): boolean { const phase = getThreadRuntimePhase(threadId, record); diff --git a/apps/web/src/stores/threadStore.ts b/apps/web/src/stores/threadStore.ts index e4616371e..19f5bf881 100644 --- a/apps/web/src/stores/threadStore.ts +++ b/apps/web/src/stores/threadStore.ts @@ -105,7 +105,7 @@ import { } from "./thread-store/usage"; export { mergeProviderUsageSnapshot } from "./thread-store/usage"; -import { isThreadRuntimeActive } from "./thread-lifecycle"; +import { getCanonicalRuntimePhase, isThreadRuntimeActive } from "./thread-lifecycle"; function deriveRunningThreadIds(records: Map): Set { return new Set( @@ -2715,6 +2715,31 @@ export const useThreadStore = create((zustandSet, get) => { return false; }; + // Canonical lifecycle owns runtime state once its turn correlates to the + // record, so a dropped legacy terminal event cannot strand the running set + // or the composer gate. + const reconcileCanonicalRuntime = ( + records: Map, + runningThreadIds: Set, + threadId: string, + ): { records: Map; runningThreadIds: Set } => { + const record = getThreadRecord(records, threadId); + const phase = getCanonicalRuntimePhase(threadId, record); + if (phase === null) return { records, runningThreadIds }; + let nextRecords = records; + if (record.runtimePhase !== phase) { + nextRecords = patchThreadRecord(records, threadId, { runtimePhase: phase }); + } + const running = phase === "running"; + if (runningThreadIds.has(threadId) === running) { + return { records: nextRecords, runningThreadIds }; + } + const nextRunning = new Set(runningThreadIds); + if (running) nextRunning.add(threadId); + else nextRunning.delete(threadId); + return { records: nextRecords, runningThreadIds: nextRunning }; + }; + return { records: new Map(), currentThreadId: null, @@ -2728,6 +2753,7 @@ export const useThreadStore = create((zustandSet, get) => { flushPendingTextDeltas(); set((state) => { let records = state.records; + let runningThreadIds = state.runningThreadIds; let changed = false; for (const recovery of recoveries) { const current = getThreadRecord(records, recovery.threadId); @@ -2738,9 +2764,14 @@ export const useThreadStore = create((zustandSet, get) => { canonicalAgent: update.replica, ...(update.installedSnapshot ? recoverParentNarrative(recovery.threadId, update.replica.state) : {}), }); + const reconciled = reconcileCanonicalRuntime(records, runningThreadIds, recovery.threadId); + records = reconciled.records; + runningThreadIds = reconciled.runningThreadIds; changed = true; } - return changed ? { records } : {}; + return changed + ? { records, ...(runningThreadIds === state.runningThreadIds ? {} : { runningThreadIds }) } + : {}; }); }, @@ -2757,8 +2788,11 @@ export const useThreadStore = create((zustandSet, get) => { const update = applyCanonicalPushEvents(current.canonicalAgent, threadId, events); if (update.replica === current.canonicalAgent) return {}; accepted = true; + const records = patchThreadRecord(state.records, threadId, { canonicalAgent: update.replica }); + const reconciled = reconcileCanonicalRuntime(records, state.runningThreadIds, threadId); return { - records: patchThreadRecord(state.records, threadId, { canonicalAgent: update.replica }), + records: reconciled.records, + ...(reconciled.runningThreadIds === state.runningThreadIds ? {} : { runningThreadIds: reconciled.runningThreadIds }), }; }); if ( From 6e510d2775af10ff0d6e0f2da0be3103375ee2e0 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Sun, 27 Sep 2026 23:21:31 +0100 Subject: [PATCH 04/12] fix(desktop): reconnect the IPC push relay with backoff The relay opened one socket and never recovered: a dropped connection silently degraded the renderer to WebSocket-only delivery for the rest of the session. Reconnect in the main process with bounded exponential backoff, reset on successful connect, and stop when the window dies or the relay is torn down. The renderer re-suppresses push channels as frames resume, so no reconnect handshake is needed. --- .../connection/__tests__/ipc-relay.test.ts | 93 +++++++++++++++ .../server-runtime/connection/ipc-relay.ts | 110 +++++++++++------- 2 files changed, 164 insertions(+), 39 deletions(-) diff --git a/apps/desktop/src/features/server-runtime/connection/__tests__/ipc-relay.test.ts b/apps/desktop/src/features/server-runtime/connection/__tests__/ipc-relay.test.ts index 26997292b..559ba06e6 100644 --- a/apps/desktop/src/features/server-runtime/connection/__tests__/ipc-relay.test.ts +++ b/apps/desktop/src/features/server-runtime/connection/__tests__/ipc-relay.test.ts @@ -224,4 +224,97 @@ describe("startIpcRelay", () => { expect(win.webContents.send).toHaveBeenCalledWith("ipc-push-message", message); }); }); + + describe("reconnect", () => { + interface TrackedSocket { + on: ReturnType; + destroy: ReturnType; + handlers: Record void>; + } + + function trackedSockets(): TrackedSocket[] { + const sockets: TrackedSocket[] = []; + vi.mocked(NodeNet.connect).mockImplementation(() => { + const socket: TrackedSocket = { + on: vi.fn((event: string, handler: (...args: unknown[]) => void) => { + socket.handlers[event] = handler; + return socket; + }), + destroy: vi.fn(), + handlers: {}, + }; + sockets.push(socket); + return socket as unknown as ReturnType; + }); + return sockets; + } + + it("reconnects with backoff after the socket closes", async () => { + vi.useFakeTimers(); + try { + const sockets = trackedSockets(); + startIpcRelay("/tmp/mcode.sock", makeWindow() as never); + + sockets[0]!.handlers["close"]?.(); + await vi.advanceTimersByTimeAsync(500); + expect(sockets).toHaveLength(2); + + sockets[1]!.handlers["close"]?.(); + await vi.advanceTimersByTimeAsync(999); + expect(sockets).toHaveLength(2); + await vi.advanceTimersByTimeAsync(1); + expect(sockets).toHaveLength(3); + } finally { + vi.useRealTimers(); + } + }); + + it("resets the backoff after a successful connect", async () => { + vi.useFakeTimers(); + try { + const sockets = trackedSockets(); + startIpcRelay("/tmp/mcode.sock", makeWindow() as never); + + sockets[0]!.handlers["close"]?.(); + await vi.advanceTimersByTimeAsync(500); + sockets[1]!.handlers["connect"]?.(); + sockets[1]!.handlers["close"]?.(); + await vi.advanceTimersByTimeAsync(500); + expect(sockets).toHaveLength(3); + } finally { + vi.useRealTimers(); + } + }); + + it("does not reconnect after cleanup", async () => { + vi.useFakeTimers(); + try { + const sockets = trackedSockets(); + const cleanup = startIpcRelay("/tmp/mcode.sock", makeWindow() as never); + + cleanup(); + sockets[0]!.handlers["close"]?.(); + await vi.advanceTimersByTimeAsync(30_000); + expect(sockets).toHaveLength(1); + } finally { + vi.useRealTimers(); + } + }); + + it("stops reconnecting when the window is destroyed", async () => { + vi.useFakeTimers(); + try { + const sockets = trackedSockets(); + const win = makeWindow(false); + startIpcRelay("/tmp/mcode.sock", win as never); + win.isDestroyed.mockReturnValue(true); + + sockets[0]!.handlers["close"]?.(); + await vi.advanceTimersByTimeAsync(30_000); + expect(sockets).toHaveLength(1); + } finally { + vi.useRealTimers(); + } + }); + }); }); diff --git a/apps/desktop/src/features/server-runtime/connection/ipc-relay.ts b/apps/desktop/src/features/server-runtime/connection/ipc-relay.ts index e3511bee9..c58df402f 100644 --- a/apps/desktop/src/features/server-runtime/connection/ipc-relay.ts +++ b/apps/desktop/src/features/server-runtime/connection/ipc-relay.ts @@ -12,6 +12,11 @@ import * as NodeNet from "node:net"; * a corrupt or malicious length prefix; the socket is destroyed immediately. */ const MAX_FRAME_SIZE = 8 * 1024 * 1024; +/** Reconnect backoff bounds. A dropped relay silently degraded the renderer to + * the WebSocket path forever before this existed. */ +const MIN_RECONNECT_MS = 500; +const MAX_RECONNECT_MS = 15_000; + /** Minimal subset of BrowserWindow required by the relay. */ interface RelayWindow { isDestroyed(): boolean; @@ -28,57 +33,84 @@ interface RelayWindow { * Wire format: each frame is a 4-byte big-endian length prefix followed by * the UTF-8 encoded JSON body. * + * On socket close the renderer is told to fall back to WebSocket and the + * relay reconnects with exponential backoff. The renderer re-suppresses + * channels as frames resume, so no reconnect handshake is required. + * * @returns A cleanup function. Call it when the window closes to destroy * the socket and prevent a named-pipe handle leak on Windows. */ export function startIpcRelay(ipcPath: string, window: RelayWindow): () => void { if (!ipcPath) return () => { /* no-op: no socket was opened */ }; - const socket = NodeNet.connect(ipcPath); - const chunks: Buffer[] = []; - let totalLen = 0; + let stopped = false; + let socket: NodeNet.Socket | null = null; + let reconnectTimer: NodeJS.Timeout | null = null; + let reconnectDelay = MIN_RECONNECT_MS; - socket.on("data", (chunk: Buffer) => { - chunks.push(chunk); - totalLen += chunk.length; + const windowAlive = () => !window.isDestroyed() && !window.webContents.isDestroyed(); - // Avoid concat overhead when only one chunk is buffered. - let buffer = chunks.length === 1 ? chunks[0] : Buffer.concat(chunks, totalLen); - chunks.length = 0; - totalLen = 0; + const connect = (): void => { + if (stopped || !windowAlive()) return; + const current = NodeNet.connect(ipcPath); + socket = current; + const chunks: Buffer[] = []; + let totalLen = 0; - while (buffer.length >= 4) { - const frameLen = buffer.readUInt32BE(0); - if (frameLen > MAX_FRAME_SIZE) { - socket.destroy(); - return; - } - if (buffer.length < 4 + frameLen) break; + current.on("connect", () => { + reconnectDelay = MIN_RECONNECT_MS; + }); + + current.on("data", (chunk: Buffer) => { + chunks.push(chunk); + totalLen += chunk.length; - const json = buffer.subarray(4, 4 + frameLen).toString("utf-8"); - buffer = buffer.subarray(4 + frameLen); + // Avoid concat overhead when only one chunk is buffered. + let buffer = chunks.length === 1 ? chunks[0] : Buffer.concat(chunks, totalLen); + chunks.length = 0; + totalLen = 0; - try { - const data = JSON.parse(json) as unknown; - if (!window.isDestroyed() && !window.webContents.isDestroyed()) { - window.webContents.send("ipc-push-message", data); + while (buffer.length >= 4) { + const frameLen = buffer.readUInt32BE(0); + if (frameLen > MAX_FRAME_SIZE) { + current.destroy(); + return; } - } catch { /* malformed frame - skip */ } - } - - // Retain leftover bytes for the next data event. - if (buffer.length > 0) { - chunks.push(buffer); - totalLen = buffer.length; - } - }); - - socket.on("error", () => socket.destroy()); - socket.on("close", () => { - if (!window.isDestroyed() && !window.webContents.isDestroyed()) { + if (buffer.length < 4 + frameLen) break; + + const json = buffer.subarray(4, 4 + frameLen).toString("utf-8"); + buffer = buffer.subarray(4 + frameLen); + + try { + const data = JSON.parse(json) as unknown; + if (windowAlive()) { + window.webContents.send("ipc-push-message", data); + } + } catch { /* malformed frame - skip */ } + } + + // Retain leftover bytes for the next data event. + if (buffer.length > 0) { + chunks.push(buffer); + totalLen = buffer.length; + } + }); + + current.on("error", () => current.destroy()); + current.on("close", () => { + if (socket === current) socket = null; + if (stopped || !windowAlive()) return; window.webContents.send("ipc-push-disconnect"); - } - }); + reconnectTimer = setTimeout(connect, reconnectDelay); + reconnectDelay = Math.min(reconnectDelay * 2, MAX_RECONNECT_MS); + }); + }; + + connect(); - return () => socket.destroy(); + return () => { + stopped = true; + if (reconnectTimer) clearTimeout(reconnectTimer); + socket?.destroy(); + }; } From eeb139c237514e2b267b235359c904ee22a66ec5 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Sun, 27 Sep 2026 23:21:36 +0100 Subject: [PATCH 05/12] refactor(web): drop the unused isThreadRunning selector No caller remains: consumers read runningThreadIds directly. --- apps/web/src/stores/threadStore.ts | 7 ------- 1 file changed, 7 deletions(-) diff --git a/apps/web/src/stores/threadStore.ts b/apps/web/src/stores/threadStore.ts index 19f5bf881..d10ef1e9e 100644 --- a/apps/web/src/stores/threadStore.ts +++ b/apps/web/src/stores/threadStore.ts @@ -186,8 +186,6 @@ interface ThreadState { clearMessages: () => void; /** Deactivate the selected conversation and invalidate any active hydration commit. */ deactivateConversation: () => void; - /** Returns true if an agent is actively executing on the given thread. */ - isThreadRunning: (threadId: string) => boolean; /** Set questions received from the model and show the wizard. */ setPlanQuestions: (threadId: string, questions: PlanQuestion[]) => void; /** Record the user's answer for one question. */ @@ -3370,11 +3368,6 @@ export const useThreadStore = create((zustandSet, get) => { threadHydrator.deactivate(); }, - /** Check whether an agent is currently executing on the given thread. */ - isThreadRunning: (threadId) => { - return get().runningThreadIds.has(threadId); - }, - /** Return per-thread settings, preferring in-memory overrides then DB-persisted values then defaults. */ getThreadSettings: (threadId) => { const stored = get().records.get(threadId); From ae6758103269f36d8a8547a383670892fda496d6 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Mon, 28 Sep 2026 01:57:40 +0100 Subject: [PATCH 06/12] feat(server): stream parent narrative recovery as canonical item events Parent narrative recovery wrote items directly into canonical_agent_items, so subscribed replicas only learned about them through snapshots. Each mutation now commits an item.recorded envelope inside the existing elapsed-bounded batch writer, discards become narrativeRecoveryDiscarded tombstones, and the semantic writer flushes the envelopes through its buffered publication channel. The web store reruns recoverParentNarrative when a canonical batch touches narrative items, letting parent tool calls and segments reach the chat through the canonical stream. --- .../canonical-agent-writer-client.test.ts | 2 +- ...anonical-execution-semantic-writer.test.ts | 5 +- .../canonical/canonical-agent-boundary.ts | 212 +++++++++++------- .../canonical-execution-semantic-writer.ts | 8 +- apps/web/src/stores/threadStore.ts | 20 +- 5 files changed, 151 insertions(+), 96 deletions(-) diff --git a/apps/server/src/features/agents/canonical/__tests__/canonical-agent-writer-client.test.ts b/apps/server/src/features/agents/canonical/__tests__/canonical-agent-writer-client.test.ts index 7956e6a07..592e5c3b2 100644 --- a/apps/server/src/features/agents/canonical/__tests__/canonical-agent-writer-client.test.ts +++ b/apps/server/src/features/agents/canonical/__tests__/canonical-agent-writer-client.test.ts @@ -284,7 +284,7 @@ describe("canonical SQLite writer", () => { const receipt = await writer.transactSemantic(delta, (events) => published.push(...events.map((event) => event.eventId))); expect(receipt).toMatchObject({ kind: "committed", operationId: "narrative-lease:2" }); expect(created).toBe(2); - expect(published).toEqual([]); + expect(published).toEqual([expect.stringContaining(`narrative:${EXECUTION_ID}:toolCall:writer-recovery-tool:`)]); expect(db.prepare("SELECT COUNT(*) AS count FROM canonical_agent_items WHERE id = ?") .get("toolCall:writer-recovery-tool")).toEqual({ count: 1 }); await writer.close(); diff --git a/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts b/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts index 7d007df48..38e00ba2b 100644 --- a/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts +++ b/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts @@ -863,7 +863,7 @@ describe("CanonicalExecutionSemanticWriter through ExecutionWorkerHandler", () = expect((await writer.transact(first)).kind).toBe("committed"); initialDelta.acknowledge(); expect(canonical.loadParentNarrativeRecovery(TURN_ID)).toHaveLength(1); - expect(published).toHaveLength(publicationCount); + expect(published).toHaveLength(publicationCount + 1); const discardedDelta = reducer.prepare([]); if (!discardedDelta) throw new Error("Expected a discarded narrative delta"); const discard = operation(3, { kind: "narrative-delta", input: { @@ -880,7 +880,8 @@ describe("CanonicalExecutionSemanticWriter through ExecutionWorkerHandler", () = const receipt = await writer.transact(latest); expect(receipt).toMatchObject({ kind: "committed", operationId: "lease-1:4" }); latestDelta.acknowledge(); - expect(published).toHaveLength(publicationCount); + // One envelope per persist and one per discard tombstone. + expect(published).toHaveLength(publicationCount + 3); db.close(true); db = openDatabase({ dbPath: path }); diff --git a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts index 1edf36c69..e892960bd 100644 --- a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts +++ b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts @@ -206,10 +206,6 @@ export type CanonicalParentTurnFinishInput = ParentTurnFinishInput; /** Canonical alias for the structured parent recovery durability input. */ export type ParentNarrativeRecoveryCommitInput = ParentNarrativeRecoveryCommit; -type ParentNarrativeRecoveryOperation = - | { kind: "persist"; item: ParentNarrativeRecoveryItem } - | { kind: "discard"; itemId: string }; - export type { CanonicalChildTurnFinishInput, CodexChildDeliveryInput, @@ -1746,17 +1742,20 @@ export class CanonicalAgentBoundary implements ParentTurnDurability, CodexCollab /** * Upsert the current structured parent narrative recovery projection before * its corresponding provider event reaches the renderer. + * + * Items commit as `item.recorded` envelopes instead of silent writes so + * subscribed replicas observe every mutation between snapshots. */ recordParentNarrativeRecovery( input: ParentNarrativeRecoveryCommitInput, ): boolean { - const turn = this.loadTurnByExecution(input.executionId); - if (!turn) return false; - const thread = this.loadThread(turn.threadId); - if (!thread) throw new Error(`Canonical parent thread not found: ${turn.threadId}`); - if (input.items.length === 0 && (input.discardedItemIds?.length ?? 0) === 0) return true; - const now = new Date().toISOString(); - this.persistParentNarrativeRecoveryBatched(input, thread, turn, now); + const prepared = this.prepareParentNarrativeRecovery(input); + if (prepared === undefined) return false; + if (prepared === null) return true; + const envelopes: CanonicalAgentEventEnvelope[] = []; + runBoundedWriteBatchesSync(this.narrativeRecoveryBatches(prepared, envelopes)); + if (envelopes.length > 0) this.publish(envelopes); + this.dropDiscardedNarrativeItems(prepared.discardedItemIds); return true; } @@ -1764,91 +1763,148 @@ export class CanonicalAgentBoundary implements ParentTurnDurability, CodexCollab async recordParentNarrativeRecoveryYielding( input: ParentNarrativeRecoveryCommitInput, ): Promise { - const turn = this.loadTurnByExecution(input.executionId); - if (!turn) return false; - const thread = this.loadThread(turn.threadId); - if (!thread) throw new Error(`Canonical parent thread not found: ${turn.threadId}`); - if (input.items.length === 0 && (input.discardedItemIds?.length ?? 0) === 0) return true; - const batch = this.parentNarrativeRecoveryBatchInput(input, thread, turn, new Date().toISOString()); - if (!batch) return true; + const prepared = this.prepareParentNarrativeRecovery(input); + if (prepared === undefined) return false; + if (prepared === null) return true; + const envelopes: CanonicalAgentEventEnvelope[] = []; await runBoundedWriteBatches({ - ...batch, + ...this.narrativeRecoveryBatches(prepared, envelopes), beginImmediate: true, // A macrotask-only yield can reacquire SQLite before another connection's busy waiter wakes. yieldControl: () => new Promise((resolve) => setTimeout(resolve, 2)), }); + if (envelopes.length > 0) this.publish(envelopes); + this.dropDiscardedNarrativeItems(prepared.discardedItemIds); return true; } + /** Commits one `item.recorded` draft per batch row so the elapsed-bounded batcher keeps its pacing. */ + private narrativeRecoveryBatches( + prepared: { + commit: Omit; + drafts: readonly CanonicalAgentEventDraft[]; + }, + envelopes: CanonicalAgentEventEnvelope[], + ): RunBoundedWriteBatchesInput { + return { + db: this.db, + items: [...prepared.drafts], + limits: ACTIVE_TURN_WRITE_BATCH_LIMITS, + byteLength: (draft) => Buffer.byteLength(JSON.stringify(draft), "utf8"), + write: (draft) => envelopes.push( + ...this.commitInsideTransaction({ ...prepared.commit, events: [draft] }).events, + ), + }; + } + + /** A crash between the tombstone envelope and the delete leaves the tombstone, which is already the correct end state. */ + private dropDiscardedNarrativeItems(itemIds: readonly string[]): void { + for (const itemId of itemIds) { + this.displayMaterializer.discardItem(itemId); + this.orm.delete(canonicalAgentItems).where(eq(canonicalAgentItems.id, itemId)).run(); + } + } + /** Completes the canonical-to-display startup migration before provider recovery begins. */ async materializeConversationDisplay(): Promise { await this.displayMaterializer.runToCompletion(); } - private persistParentNarrativeRecoveryBatched( - input: ParentNarrativeRecoveryCommitInput, - thread: AgentThread, - turn: AgentTurn, - now: string, - ): void { - const batch = this.parentNarrativeRecoveryBatchInput(input, thread, turn, now); - if (batch) runBoundedWriteBatchesSync(batch); - } - - private parentNarrativeRecoveryBatchInput( + /** + * Resolve the recovery input into commit-scoped item drafts. + * Returns undefined when the turn is unknown, null when there is nothing to write. + */ + private prepareParentNarrativeRecovery( input: ParentNarrativeRecoveryCommitInput, - thread: AgentThread, - turn: AgentTurn, - now: string, - ): RunBoundedWriteBatchesInput | null { + ): { commit: Omit; drafts: CanonicalAgentEventDraft[]; discardedItemIds: string[] } | null | undefined { + const turn = this.loadTurnByExecution(input.executionId); + if (!turn) return undefined; + const thread = this.loadThread(turn.threadId); + if (!thread) throw new Error(`Canonical parent thread not found: ${turn.threadId}`); + if (input.items.length === 0 && (input.discardedItemIds?.length ?? 0) === 0) return null; const checkpoint = this.loadCheckpoint(input.executionId); if (!checkpoint) throw new Error(`Canonical parent checkpoint was not found: ${input.executionId}`); if (checkpoint.terminalOutcome) return null; - const operations: ParentNarrativeRecoveryOperation[] = [ - ...input.items.map((item) => ({ kind: "persist" as const, item })), - ...(input.discardedItemIds ?? []).map((itemId) => ({ kind: "discard" as const, itemId })), - ]; - const byteLength = (operation: (typeof operations)[number]) => operation.kind === "persist" - ? Buffer.byteLength(JSON.stringify(operation.item), "utf8") - : Buffer.byteLength(operation.itemId, "utf8"); + const drafts = this.parentNarrativeRecoveryDrafts(input, thread, turn); assertActiveTurnRecoveryRetention( - operations.length, - operations.reduce((total, operation) => total + byteLength(operation), 0), + drafts.length, + drafts.reduce((total, draft) => total + Buffer.byteLength(JSON.stringify(draft), "utf8"), 0), ); return { - db: this.db, - items: operations, - limits: ACTIVE_TURN_WRITE_BATCH_LIMITS, - byteLength, - write: (operation) => { - if (operation.kind === "persist") { - this.persistParentNarrativeRecoveryItem(operation.item, thread, turn, now); - return; - } - this.discardParentNarrativeRecoveryItem(operation.itemId, thread, turn); + commit: { + threadId: thread.id, + turnId: turn.id, + executionId: input.executionId, + phase: checkpoint.phase, }, + drafts, + discardedItemIds: [...(input.discardedItemIds ?? [])], }; } - private persistParentNarrativeRecoveryItem( - item: ParentNarrativeRecoveryItem, + /** Build one item.recorded draft per recovery mutation; discards become tombstones so replicas drop them. */ + private parentNarrativeRecoveryDrafts( + input: ParentNarrativeRecoveryCommitInput, thread: AgentThread, turn: AgentTurn, - now: string, - ): void { - const itemId = this.parentNarrativeRecoveryItemId(item); - this.persistItem({ - id: itemId, - threadId: thread.id, - turnId: turn.id, - ...this.parentNarrativeRecoveryParent(item), - kind: this.parentNarrativeRecoveryItemKind(item), - providerIdentities: turn.providerIdentities, - payload: this.parentNarrativeRecoveryPayload(item, itemId), - createdAt: item.record.started_at, - updatedAt: now, + ): CanonicalAgentEventDraft[] { + const persisted = input.items.map((item): CanonicalAgentEventDraft => { + const itemId = this.parentNarrativeRecoveryItemId(item); + return { + eventId: `narrative:${input.executionId}:${itemId}:${hashCodexKey(JSON.stringify(item))}`, + routing: { threadId: thread.id, turnId: turn.id, executionId: input.executionId, itemId }, + sourceProviderId: thread.providerId, + sourceIdentities: [], + payload: { + type: "item.recorded", + item: { + id: itemId, + threadId: thread.id, + turnId: turn.id, + ...this.parentNarrativeRecoveryParent(item), + kind: this.parentNarrativeRecoveryItemKind(item), + providerIdentities: turn.providerIdentities, + payload: this.parentNarrativeRecoveryPayload(item, itemId), + createdAt: item.record.started_at, + // Deterministic so a retried commit dedupes on eventId instead of conflicting. + updatedAt: item.kind === "toolCall" + ? item.record.completed_at ?? item.record.started_at + : item.record.ended_at ?? item.record.started_at, + }, + }, + }; }); - this.displayMaterializer.materializeItems([itemId]); + const discarded = (input.discardedItemIds ?? []).map((itemId) => + this.parentNarrativeRecoveryDiscardDraft(input, thread, turn, itemId)); + return [...persisted, ...discarded]; + } + + /** A discard re-records the item under a marker projection so replicas drop it without a delete event. */ + private parentNarrativeRecoveryDiscardDraft( + input: ParentNarrativeRecoveryCommitInput, + thread: AgentThread, + turn: AgentTurn, + itemId: string, + ): CanonicalAgentEventDraft { + const existing = this.loadItem(itemId); + const projection = typeof existing?.payload.projection === "string" ? existing.payload.projection : null; + if (!existing || existing.threadId !== thread.id || existing.turnId !== turn.id + || (projection !== "narrativeRecovery" && projection !== "narrativeRecoveryDiscarded")) { + throw new Error(`Canonical narrative recovery item was not found: ${itemId}`); + } + return { + eventId: `narrative-discard:${input.executionId}:${itemId}:${hashCodexKey(JSON.stringify(existing.payload))}`, + routing: { threadId: thread.id, turnId: turn.id, executionId: input.executionId, itemId }, + sourceProviderId: thread.providerId, + sourceIdentities: [], + payload: { + type: "item.recorded", + item: { + ...existing, + payload: { ...existing.payload, projection: "narrativeRecoveryDiscarded" }, + }, + }, + }; } private parentNarrativeRecoveryItemId(item: ParentNarrativeRecoveryItem): string { @@ -1902,24 +1958,6 @@ export class CanonicalAgentBoundary implements ParentTurnDurability, CodexCollab if (typeof value === "string") metadata[key] = value; } - private discardParentNarrativeRecoveryItem( - itemId: string, - thread: AgentThread, - turn: AgentTurn, - ): void { - const existing = this.loadItem(itemId); - if (!existing || existing.threadId !== thread.id || existing.turnId !== turn.id) { - throw new Error(`Canonical narrative recovery item was not found: ${itemId}`); - } - const projection = typeof existing.payload.projection === "string" - ? existing.payload.projection - : null; - if (projection === "narrativeRecovery" || projection === "narrativeRecoveryDiscarded") { - this.displayMaterializer.discardItem(itemId); - this.orm.delete(canonicalAgentItems).where(eq(canonicalAgentItems.id, itemId)).run(); - } - } - /** Load the newest durable structured narrative snapshot for an unfinished parent turn. */ loadParentNarrativeRecovery(turnId: string): ParentNarrativeRecoveryItem[] { const rows = this.orm.select({ payloadJson: canonicalAgentItems.payloadJson }) diff --git a/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts b/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts index 5a4973fbf..ff858f229 100644 --- a/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts +++ b/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts @@ -550,14 +550,14 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter private recordNarrativeDelta(operation: ExecutionSemanticOperation, hash: string): ExecutionWriteReceipt { const mutation = operation.mutation; if (mutation.kind !== "narrative-delta") return conflict(operation); - return this.db.transaction(() => { + return this.withBufferedPublication(() => this.db.transaction(() => { const head = this.requireNextHead(operation); this.requireUnfinishedCheckpoint(operation.execution); this.persistNarrativeDelta(mutation.input); const receipt = committed(operation, head.durableRevision); this.storeHead({ ...head, ordinal: operation.ordinal }); return this.storeReceipt(operation, hash, receipt); - }).immediate(); + }).immediate()); } private applyLiveRecovery(operation: ExecutionSemanticOperation, hash: string): ExecutionWriteReceipt { @@ -569,7 +569,7 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter private recordLiveEvent(operation: ExecutionSemanticOperation, hash: string): ExecutionWriteReceipt { const mutation = operation.mutation; if (mutation.kind !== "live-event") return conflict(operation); - const receipt = this.db.transaction(() => { + const receipt = this.withBufferedPublication(() => this.db.transaction(() => { const head = this.requireNextHead(operation); if (!validPublicationProvider(operation, head.providerId)) throw new SemanticConflict(); this.requireUnfinishedCheckpoint(operation.execution); @@ -580,7 +580,7 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter this.storeHead({ ...head, ordinal: operation.ordinal, ...(mutation.message ? { assignedMessageId: mutation.message.messageId } : {}) }); return this.storeReceipt(operation, hash, committedReceipt, livePublication); - }).immediate(); + }).immediate()); if (mutation.text.kind === "reclassify") { this.assistantText.discardRecoveryJournal(operation.execution.executionId); } diff --git a/apps/web/src/stores/threadStore.ts b/apps/web/src/stores/threadStore.ts index d10ef1e9e..e1c1810d6 100644 --- a/apps/web/src/stores/threadStore.ts +++ b/apps/web/src/stores/threadStore.ts @@ -115,6 +115,14 @@ function deriveRunningThreadIds(records: Map): Set ); } +/** True when a canonical batch mutates a parent narrative recovery item. */ +function batchTouchesParentNarrative(events: readonly CanonicalAgentEventEnvelope[]): boolean { + return events.some((event) => + event.payload.type === "item.recorded" + && (event.payload.item.payload.projection === "narrativeRecovery" + || event.payload.item.payload.projection === "narrativeRecoveryDiscarded")); +} + function preserveRunningThreadIds(previous: Set, next: Set): Set { return previous.size === next.size && [...next].every((id) => previous.has(id)) ? previous : next; } @@ -2760,7 +2768,10 @@ export const useThreadStore = create((zustandSet, get) => { if (update.installedSnapshot) streamingTextByteSizes.delete(recovery.threadId); records = patchThreadRecord(records, recovery.threadId, { canonicalAgent: update.replica, - ...(update.installedSnapshot ? recoverParentNarrative(recovery.threadId, update.replica.state) : {}), + ...((update.installedSnapshot + || (recovery.mode === "delta" && batchTouchesParentNarrative(recovery.events))) + ? recoverParentNarrative(recovery.threadId, update.replica.state) + : {}), }); const reconciled = reconcileCanonicalRuntime(records, runningThreadIds, recovery.threadId); records = reconciled.records; @@ -2786,7 +2797,12 @@ export const useThreadStore = create((zustandSet, get) => { const update = applyCanonicalPushEvents(current.canonicalAgent, threadId, events); if (update.replica === current.canonicalAgent) return {}; accepted = true; - const records = patchThreadRecord(state.records, threadId, { canonicalAgent: update.replica }); + const records = patchThreadRecord(state.records, threadId, { + canonicalAgent: update.replica, + ...(batchTouchesParentNarrative(events) + ? recoverParentNarrative(threadId, update.replica.state) + : {}), + }); const reconciled = reconcileCanonicalRuntime(records, state.runningThreadIds, threadId); return { records: reconciled.records, From 1248e4179c1f10cfadb39a61aa1c3c4e7f14d9c5 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Tue, 29 Sep 2026 00:25:27 +0100 Subject: [PATCH 07/12] feat(server): record renderer publications as canonical events Numbered livePublication intents are the enriched, sanitized events the agent.event broadcast sends today. Recording each one as a canonical publication.recorded envelope inside the same operation transaction makes the canonical stream carry the same renderer-facing payload, so replay and gap recovery deliver identical effects without the legacy channel. The reducer treats the envelope as a passthrough: identity and dedup bookkeeping apply, but no thread or item state mutates. On the client, canonical pushes and reconnect deltas dispatch the embedded event through handleAgentEvent with its publicationId, so the existing stable publication cursor accepts whichever copy arrives first and drops the other. Envelopes committed by unbuffered operation paths (control, stage-terminal, assistant text, terminal batches) queue on the writer and flush after the outer transaction commits, matching the buffered flush used by live operations. --- ...anonical-execution-semantic-writer.test.ts | 4 +- .../canonical/canonical-agent-boundary.ts | 3 +- .../canonical-execution-semantic-writer.ts | 64 +++++++++++++++- .../canonical-publication-dispatch.test.ts | 75 +++++++++++++++++++ apps/web/src/stores/threadStore.ts | 19 +++++ packages/agent-model/src/events.ts | 8 ++ packages/agent-model/src/reducer.ts | 2 + 7 files changed, 171 insertions(+), 4 deletions(-) create mode 100644 apps/web/src/__tests__/canonical-publication-dispatch.test.ts diff --git a/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts b/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts index 38e00ba2b..97f6ff703 100644 --- a/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts +++ b/apps/server/src/features/agents/canonical/__tests__/canonical-execution-semantic-writer.test.ts @@ -950,7 +950,9 @@ describe("CanonicalExecutionSemanticWriter through ExecutionWorkerHandler", () = expect(db.prepare("SELECT last_sequence FROM canonical_writer_live_publication_heads WHERE thread_id = ?") .get(THREAD_ID)).toEqual({ last_sequence: 1 }); expect(new CanonicalAgentBoundary(db, () => {}).loadParentNarrativeRecovery(TURN_ID)).toEqual([]); - expect(published).toHaveLength(publicationCount); + // The committed publication envelope stays durable; the rolled-back reclassification added nothing. + expect(published).toHaveLength(publicationCount + 1); + expect(published).toContain(`publication:${THREAD_ID}:1`); db.run("DROP TRIGGER fail_compound_receipt"); const reclassifiedReceipt = await writer.transact(reclassified); expect(reclassifiedReceipt).toMatchObject({ kind: "committed", livePublication: [ diff --git a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts index e892960bd..901e41574 100644 --- a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts +++ b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts @@ -495,7 +495,8 @@ export class CanonicalAgentBoundary implements ParentTurnDurability, CodexCollab : this.eventStore.commit(this.toEventStoreCommitInput(input)); } - private commitInsideTransaction(input: CanonicalAgentCommitInput): CanonicalAgentCommitResult { + /** Commit inside the caller's transaction. The caller publishes result.events after that commit. */ + commitInsideTransaction(input: CanonicalAgentCommitInput): CanonicalAgentCommitResult { return serverWorkTrace ? serverWorkTrace.measure("canonical-write", input.threadId, input.executionId, () => this.eventStore.applyWithinTransaction(this.toEventStoreCommitInput(input))) diff --git a/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts b/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts index ff858f229..2fd0c3f37 100644 --- a/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts +++ b/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts @@ -12,6 +12,7 @@ import { ProviderIdSchema, ProviderIdentitySchema, TurnOutcomeSchema, + type CanonicalAgentEventEnvelope, type ProviderIdentity, type ParentNarrativeRecoveryItem, type TurnOutcome, @@ -33,7 +34,7 @@ import type { ExecutionPlanQuestionsReceipt, ExecutionTerminalPersistenceReceipt, } from "../execution/execution-worker-handler.js"; -import type { CanonicalAgentEventPublisher } from "./canonical-agent-boundary.js"; +import type { CanonicalAgentEventDraft, CanonicalAgentEventPublisher } from "./canonical-agent-boundary.js"; import { CanonicalAgentBoundary } from "./canonical-agent-boundary.js"; import { APPEND_GROUP_LIMITS, isGroupableAppend } from "./canonical-append-group.js"; import { CanonicalCommittedProviderProjector } from "./canonical-committed-provider-projector.js"; @@ -240,6 +241,8 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter private bufferedPublication: Parameters[0][] | null = null; private bufferedPublicationCount = 0; private groupedPublication: Parameters[0][] | null = null; + /** Canonical envelopes committed by operation paths that own no publication buffer. */ + private pendingCanonicalEvents: CanonicalAgentEventEnvelope[][] = []; constructor(private readonly db: Database, private readonly publish: CanonicalAgentEventPublisher) { this.turns = new CanonicalParentTurnWrite(db, (events) => { @@ -287,8 +290,12 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter return existing; } try { - return await this.applySupported(operation, hash); + const receipt = await this.applySupported(operation, hash); + this.flushCanonicalEvents(); + return receipt; } catch (error) { + // A rolled-back transaction leaves staged envelopes that must never publish. + this.pendingCanonicalEvents = []; if (error instanceof SemanticConflict) return conflict(operation); throw error; } @@ -1041,6 +1048,7 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter publicationOverride?: readonly ExecutionLivePublicationIntent[], ): Extract { const publication = this.numberLivePublication(operation, publicationOverride); + if (publication?.length) this.commitLivePublicationEvents(operation, hash, publication); const storedReceipt = publication ? { ...receipt, livePublication: publication } : receipt; const receiptJson = JSON.stringify({ ...storedReceipt, publicationVersion: 1 }); if (operation.livePublication && Buffer.byteLength(receiptJson, "utf8") > MAX_LIVE_RECEIPT_BYTES) { @@ -1063,6 +1071,51 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter return storedReceipt; } + /** + * Records each numbered live publication as a canonical envelope inside the same operation + * transaction, so a canonical replay or gap recovery carries the same renderer-facing events. + */ + private commitLivePublicationEvents( + operation: ExecutionSemanticOperation, + hash: string, + publication: readonly ExecutionLivePublicationReceipt[], + ): void { + const thread = this.canonical.loadThread(operation.execution.threadId); + const checkpoint = this.canonical.loadCheckpoint(operation.execution.executionId); + if (!thread || !checkpoint) throw new SemanticConflict(); + const result = this.canonical.commitInsideTransaction({ + threadId: operation.execution.threadId, + turnId: operation.execution.turnId, + executionId: operation.execution.executionId, + phase: checkpoint.phase, + events: publication.map((entry): CanonicalAgentEventDraft => ({ + eventId: `publication:${operation.execution.threadId}:${entry.publicationId}`, + routing: { + threadId: operation.execution.threadId, + turnId: operation.execution.turnId, + executionId: operation.execution.executionId, + }, + sourceProviderId: thread.providerId, + sourceIdentities: [], + payload: { + type: "publication.recorded", + publicationId: entry.publicationId, + event: entry.event, + }, + })), + }); + if (result.outcome !== "committed" && result.outcome !== "duplicate" + && result.outcome !== "terminal-outcome-confirmed") { + throw new SemanticConflict(); + } + if (result.events.length === 0) return; + this.storePublicationChunks(operation.execution.executionId, hash, + result.events.map((envelope) => envelope.acceptedSequence)); + // Unbuffered operation paths publish from the pending queue only after their transaction commits. + if (this.bufferedPublication) this.bufferPublication([...result.events]); + else this.pendingCanonicalEvents.push([...result.events]); + } + private numberLivePublication( operation: ExecutionSemanticOperation, publicationOverride?: readonly ExecutionLivePublicationIntent[], @@ -1216,6 +1269,13 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter if (this.groupedPublication) this.groupedPublication.push(events); else this.publish(events); } + + private flushCanonicalEvents(): void { + const pending = this.pendingCanonicalEvents; + if (pending.length === 0) return; + this.pendingCanonicalEvents = []; + for (const events of pending) this.publishCommitted(events); + } } function validOperation(operation: ExecutionSemanticOperation): boolean { diff --git a/apps/web/src/__tests__/canonical-publication-dispatch.test.ts b/apps/web/src/__tests__/canonical-publication-dispatch.test.ts new file mode 100644 index 000000000..aa9f56207 --- /dev/null +++ b/apps/web/src/__tests__/canonical-publication-dispatch.test.ts @@ -0,0 +1,75 @@ +import { describe, expect, it, beforeEach, vi } from "vitest"; +import type { AgentEvent, CanonicalAgentEventEnvelope } from "@mcode/contracts"; +import { useThreadStore } from "@/stores/threadStore"; +import { + resetThreadStoreForTests, + seedThreadRecord, + readThreadField, +} from "@/stores/thread-store-test-utils"; +import { clearRecordCache } from "@/features/conversation/hydration/record-cache"; +import { mockTransport } from "./mocks/transport"; + +vi.mock("@/transport", async () => ({ + ...(await vi.importActual("@/transport")), + getTransport: () => mockTransport, +})); + +const THREAD_ID = "thread-publication"; +const TURN_ID = "turn-1"; +const EXECUTION_ID = "00000000-0000-4000-8000-000000000001"; +const EPOCH = "00000000-0000-4000-8000-000000000002"; +const CURSOR_KEY = `mcode-agent-publication-v2:${THREAD_ID}`; + +function publicationEnvelope(publicationId: string, event: Record): CanonicalAgentEventEnvelope { + return { + eventId: `publication:${THREAD_ID}:${publicationId}`, + routing: { threadId: THREAD_ID, turnId: TURN_ID, executionId: EXECUTION_ID }, + sourceProviderId: "codex", + sourceIdentities: [], + acceptedSequence: Number(publicationId), + durableRevision: Number(publicationId), + serverTimestamps: { acceptedAt: "2026-08-11T12:00:00.000Z", persistedAt: "2026-08-11T12:00:00.000Z" }, + payload: { type: "publication.recorded", publicationId, event }, + }; +} + +const TOOL_USE: Record = { + type: "toolUse", + threadId: THREAD_ID, + turnExecutionId: EXECUTION_ID, + epoch: EPOCH, + sequence: 8, + toolCallId: "call-1", + toolName: "Read", + toolInput: { path: "/original" }, +}; + +function toolCallCount(): number { + return readThreadField(THREAD_ID, (record) => record.toolCalls)?.length ?? -1; +} + +describe("canonical publication dispatch", () => { + beforeEach(() => { + sessionStorage.removeItem(CURSOR_KEY); + clearRecordCache(); + resetThreadStoreForTests({ currentThreadId: THREAD_ID }); + useThreadStore.setState({ + records: seedThreadRecord(THREAD_ID, { runtimePhase: "running", turnExecutionId: EXECUTION_ID }), + runningThreadIds: new Set([THREAD_ID]), + }); + vi.clearAllMocks(); + }); + + it("applies a recorded publication through the agent event path with its publication identity", () => { + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, [publicationEnvelope("1", TOOL_USE)]); + expect(toolCallCount()).toBe(1); + }); + + it("applies a publication once when the agent.event copy also arrives", () => { + const envelope = publicationEnvelope("1", TOOL_USE); + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, [envelope]); + useThreadStore.getState().handleCanonicalAgentEvents(THREAD_ID, [envelope]); + useThreadStore.getState().handleAgentEvent({ ...TOOL_USE, publicationId: "1" } as AgentEvent); + expect(toolCallCount()).toBe(1); + }); +}); diff --git a/apps/web/src/stores/threadStore.ts b/apps/web/src/stores/threadStore.ts index e1c1810d6..6e7121dd5 100644 --- a/apps/web/src/stores/threadStore.ts +++ b/apps/web/src/stores/threadStore.ts @@ -4,6 +4,7 @@ import type { AgentEvent, CanonicalAgentEventEnvelope, CanonicalAgentReconnectRe import type { DevinMode, PermissionRequest, PermissionDecision } from "@mcode/contracts"; import { recoverParentNarrative } from "./parent-narrative-recovery"; import { + AgentEventSchema, PERMISSION_MODES, INTERACTION_MODES, ProviderIdSchema, @@ -123,6 +124,20 @@ function batchTouchesParentNarrative(events: readonly CanonicalAgentEventEnvelop || event.payload.item.payload.projection === "narrativeRecoveryDiscarded")); } +/** + * Canonical envelopes carrying numbered renderer-facing publications replay through the same + * dispatch as the agent.event copy; the publication cursor accepts whichever copy arrives first. + */ +function dispatchCanonicalPublications(events: readonly CanonicalAgentEventEnvelope[]): void { + const handle = useThreadStore.getState().handleAgentEvent; + for (const envelope of events) { + if (envelope.payload.type !== "publication.recorded") continue; + const parsed = AgentEventSchema().safeParse(envelope.payload.event); + if (!parsed.success) continue; + handle({ ...parsed.data, publicationId: envelope.payload.publicationId }); + } +} + function preserveRunningThreadIds(previous: Set, next: Set): Set { return previous.size === next.size && [...next].every((id) => previous.has(id)) ? previous : next; } @@ -2782,6 +2797,9 @@ export const useThreadStore = create((zustandSet, get) => { ? { records, ...(runningThreadIds === state.runningThreadIds ? {} : { runningThreadIds }) } : {}; }); + for (const recovery of recoveries) { + if (recovery.mode === "delta") dispatchCanonicalPublications(recovery.events); + } }, handleCanonicalAgentEvents: (threadId, events) => { @@ -2816,6 +2834,7 @@ export const useThreadStore = create((zustandSet, get) => { ) { void conversationResidency.refreshVisibleConversation(threadId); } + dispatchCanonicalPublications(events); }, cacheToolCallRecords: (key, records) => { diff --git a/packages/agent-model/src/events.ts b/packages/agent-model/src/events.ts index 2eddf9729..addfe6f34 100644 --- a/packages/agent-model/src/events.ts +++ b/packages/agent-model/src/events.ts @@ -75,6 +75,14 @@ export const CanonicalAgentEventSchema = z.discriminatedUnion("type", [ }) .strict(), z.object({ type: z.literal("item.recorded"), item: AgentItemSchema }).strict(), + z + .object({ + type: z.literal("publication.recorded"), + publicationId: z.string().regex(/^[1-9]\d*$/).max(16), + /** Renderer-facing agent event; opaque here because the canonical layer is provider-agnostic. */ + event: z.record(z.unknown()), + }) + .strict(), z .object({ type: z.literal("collaboration-action.recorded"), diff --git a/packages/agent-model/src/reducer.ts b/packages/agent-model/src/reducer.ts index 29880fa9e..7572ceed0 100644 --- a/packages/agent-model/src/reducer.ts +++ b/packages/agent-model/src/reducer.ts @@ -108,6 +108,8 @@ const AGENT_EVENT_REDUCERS: Record = { event as AgentEventFor<"collaboration-action.recorded">, acceptedInputState, ), + "publication.recorded": (state, _event, acceptedInputState) => + reduceVolatileTruncation(state, acceptedInputState), }; /** Create an empty canonical reducer state. */ From 6f6ba20b54fa5e70e9213e0eb1648b6513112d8b Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Tue, 29 Sep 2026 15:56:51 +0100 Subject: [PATCH 08/12] fix(agents): sanitize canonical live publication payloads Live publication intents embed provider events before the public enrichment stage, so canonical copies would persist raw tool inputs that the legacy wire path strips. Sanitize the embedded event the same way AgentEventPublicationService does before recording publication.recorded. --- .../canonical-execution-semantic-writer.ts | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts b/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts index 2fd0c3f37..aafd41cda 100644 --- a/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts +++ b/apps/server/src/features/agents/canonical/canonical-execution-semantic-writer.ts @@ -12,6 +12,7 @@ import { ProviderIdSchema, ProviderIdentitySchema, TurnOutcomeSchema, + type AgentEvent, type CanonicalAgentEventEnvelope, type ProviderIdentity, type ParentNarrativeRecoveryItem, @@ -36,6 +37,7 @@ import type { } from "../execution/execution-worker-handler.js"; import type { CanonicalAgentEventDraft, CanonicalAgentEventPublisher } from "./canonical-agent-boundary.js"; import { CanonicalAgentBoundary } from "./canonical-agent-boundary.js"; +import { sanitizePublicToolInput } from "../tools/input/public-tool-input.js"; import { APPEND_GROUP_LIMITS, isGroupableAppend } from "./canonical-append-group.js"; import { CanonicalCommittedProviderProjector } from "./canonical-committed-provider-projector.js"; import type { ProviderEventProjection } from "../../providers/composition/provider-event-adapter.js"; @@ -1100,7 +1102,9 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter payload: { type: "publication.recorded", publicationId: entry.publicationId, - event: entry.event, + // Canonical copies reach clients on replay, so they carry the same sanitized + // tool input the wire publication applies in AgentEventPublicationService. + event: sanitizePublicationEvent(entry.event), }, })), }); @@ -1278,6 +1282,17 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter } } +/** Strip raw file contents the way the wire publication does so canonical replay cannot leak them. */ +function sanitizePublicationEvent(event: AgentEvent): AgentEvent { + if (event.type === AgentEventType.ToolUse) { + return { ...event, toolInput: sanitizePublicToolInput(event.toolInput, event.toolName) }; + } + if (event.type === AgentEventType.ToolResult && event.toolInput) { + return { ...event, toolInput: sanitizePublicToolInput(event.toolInput) }; + } + return event; +} + function validOperation(operation: ExecutionSemanticOperation): boolean { return operation.operationId === `${operation.lease.leaseId}:${operation.ordinal}` && operation.operationId !== HEAD_ID && operation.operationId.length <= 256 From 32c71b58077c1cea8129b70b6c6a8883264971e3 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Tue, 29 Sep 2026 21:15:07 +0100 Subject: [PATCH 09/12] fix(web): keep the canonical lifecycle gate when canonical owns the phase The runtime reconcile stamps runtimePhase from a correlated canonical turn but never correlated the record's execution identity, so a canonically owned "running" phase made getCanonicalLifecycleTurn treat the provider-owned child turn as a local execution and suppress its projection. A later terminal canonical event could then no longer clear the stamped phase either. reconcileCanonicalRuntime now stamps turnExecutionId from the claiming turn's executionId on live claims, and the lifecycle gate only treats runtimePhase as locally owned when the canonical runtime turn is not the latest turn. --- .../__tests__/turn-lifecycle.test.tsx | 1 + apps/web/src/stores/thread-lifecycle.ts | 12 +++++++++--- apps/web/src/stores/threadStore.ts | 19 ++++++++++++++----- 3 files changed, 24 insertions(+), 8 deletions(-) diff --git a/apps/web/src/features/conversation/messages/__tests__/turn-lifecycle.test.tsx b/apps/web/src/features/conversation/messages/__tests__/turn-lifecycle.test.tsx index 320346d86..c65a87e26 100644 --- a/apps/web/src/features/conversation/messages/__tests__/turn-lifecycle.test.tsx +++ b/apps/web/src/features/conversation/messages/__tests__/turn-lifecycle.test.tsx @@ -25,6 +25,7 @@ function canonicalTurn(status: AgentTurn["status"], trigger: AgentTurn["trigger" return { id: "canonical-turn", threadId: THREAD, status, trigger, permissionMode: "full", approvalReviewMode: "manual", approvalReviewReason: "manual-requested", + executionId: "00000000-0000-4000-8000-000000000042", providerIdentities: [], startedAt: NOW, endedAt: null, createdAt: NOW, updatedAt: NOW, }; } diff --git a/apps/web/src/stores/thread-lifecycle.ts b/apps/web/src/stores/thread-lifecycle.ts index 6d7cc57cc..491a7149e 100644 --- a/apps/web/src/stores/thread-lifecycle.ts +++ b/apps/web/src/stores/thread-lifecycle.ts @@ -11,10 +11,15 @@ function latestCanonicalTurn(threadId: string, record: ThreadRecord): AgentTurn /** Selects the provider-owned child lifecycle when no local execution is active. */ export function getCanonicalLifecycleTurn(threadId: string, record: ThreadRecord): AgentTurn | undefined { - if (record.runtimePhase === "running" || record.runtimePhase === "finalizing") return undefined; const latest = latestCanonicalTurn(threadId, record); if (latest?.trigger.kind !== "child") return undefined; - if (record.runtimePhase !== "idle" && (latest.status === "Pending" || latest.status === "Running")) return undefined; + // A runtime phase stamped by this same canonical turn must not gate it out; + // only a locally-owned phase suppresses the child lifecycle. + const locallyOwnedPhase = getCanonicalRuntimeTurn(threadId, record)?.id === latest.id + ? "idle" + : record.runtimePhase; + if (locallyOwnedPhase === "running" || locallyOwnedPhase === "finalizing") return undefined; + if (locallyOwnedPhase !== "idle" && (latest.status === "Pending" || latest.status === "Running")) return undefined; return latest; } @@ -40,7 +45,8 @@ export function getCanonicalRuntimeTurn(threadId: string, record: ThreadRecord): return latest; } -function phaseForTurnStatus(status: AgentTurnStatus): TurnRuntimePhase { +/** Runtime phase equivalent of one canonical turn status. */ +export function phaseForTurnStatus(status: AgentTurnStatus): TurnRuntimePhase { switch (status) { case "Pending": case "Running": return "running"; diff --git a/apps/web/src/stores/threadStore.ts b/apps/web/src/stores/threadStore.ts index 6e7121dd5..3e9de01c2 100644 --- a/apps/web/src/stores/threadStore.ts +++ b/apps/web/src/stores/threadStore.ts @@ -106,7 +106,7 @@ import { } from "./thread-store/usage"; export { mergeProviderUsageSnapshot } from "./thread-store/usage"; -import { getCanonicalRuntimePhase, isThreadRuntimeActive } from "./thread-lifecycle"; +import { getCanonicalRuntimeTurn, isThreadRuntimeActive, phaseForTurnStatus } from "./thread-lifecycle"; function deriveRunningThreadIds(records: Map): Set { return new Set( @@ -2745,11 +2745,20 @@ export const useThreadStore = create((zustandSet, get) => { threadId: string, ): { records: Map; runningThreadIds: Set } => { const record = getThreadRecord(records, threadId); - const phase = getCanonicalRuntimePhase(threadId, record); - if (phase === null) return { records, runningThreadIds }; + const turn = getCanonicalRuntimeTurn(threadId, record); + if (!turn) return { records, runningThreadIds }; + const phase = phaseForTurnStatus(turn.status); let nextRecords = records; - if (record.runtimePhase !== phase) { - nextRecords = patchThreadRecord(records, threadId, { runtimePhase: phase }); + const patch: Partial = {}; + if (record.runtimePhase !== phase) patch.runtimePhase = phase; + // Stamp the execution identity on live claims so later reconciles and the + // child lifecycle gate can keep correlating this turn to the record. + if (record.turnExecutionId === null && turn.executionId + && (turn.status === "Pending" || turn.status === "Running")) { + patch.turnExecutionId = turn.executionId; + } + if (Object.keys(patch).length > 0) { + nextRecords = patchThreadRecord(records, threadId, patch); } const running = phase === "running"; if (runningThreadIds.has(threadId) === running) { From d9dc3ea3e34e101e16fdda56a15f5b6f7c05b7e3 Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Tue, 29 Sep 2026 22:00:16 +0100 Subject: [PATCH 10/12] fix(agents): bound narrative recovery retention by persisted items assertActiveTurnRecoveryRetention measured the serialized canonical envelope drafts, letting transport overhead trip the cap before the retained narrative itself did. Measure the persisted item payloads and discarded item ids again, matching the pre-envelope accounting. --- .../agents/canonical/canonical-agent-boundary.ts | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts index 901e41574..c2cf2a6ab 100644 --- a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts +++ b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts @@ -1827,10 +1827,14 @@ export class CanonicalAgentBoundary implements ParentTurnDurability, CodexCollab if (!checkpoint) throw new Error(`Canonical parent checkpoint was not found: ${input.executionId}`); if (checkpoint.terminalOutcome) return null; const drafts = this.parentNarrativeRecoveryDrafts(input, thread, turn); - assertActiveTurnRecoveryRetention( - drafts.length, - drafts.reduce((total, draft) => total + Buffer.byteLength(JSON.stringify(draft), "utf8"), 0), - ); + // Retention bounds the retained narrative, not the envelope transport overhead. + const persistedBytes = input.items.reduce((total, item) => ( + total + Buffer.byteLength(JSON.stringify(item), "utf8") + ), 0); + const discardedBytes = (input.discardedItemIds ?? []).reduce((total, itemId) => ( + total + Buffer.byteLength(itemId, "utf8") + ), 0); + assertActiveTurnRecoveryRetention(drafts.length, persistedBytes + discardedBytes); return { commit: { threadId: thread.id, From dda991d33a68d5efb21898a5653b4b6c9731836c Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Tue, 29 Sep 2026 22:08:50 +0100 Subject: [PATCH 11/12] refactor(agents): extract narrative recovery byte accounting The byte-size computation pushed prepareParentNarrativeRecovery over the complexity cap. Hoisting it keeps the retention accounting intact. --- .../canonical/canonical-agent-boundary.ts | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts index c2cf2a6ab..461220cfa 100644 --- a/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts +++ b/apps/server/src/features/agents/canonical/canonical-agent-boundary.ts @@ -324,6 +324,17 @@ function narrativeParentItem(entry: NarrativeEntry): { parentItemId?: string } { return entry.record.parent_tool_call_id ? { parentItemId: `toolCall:${entry.record.parent_tool_call_id}` } : {}; } +/** Byte size of the retained narrative content: persisted items plus the ids retained for discarded ones. */ +function retainedNarrativeRecoveryBytes(input: ParentNarrativeRecoveryCommitInput): number { + const persistedBytes = input.items.reduce((total, item) => ( + total + Buffer.byteLength(JSON.stringify(item), "utf8") + ), 0); + const discardedBytes = (input.discardedItemIds ?? []).reduce((total, itemId) => ( + total + Buffer.byteLength(itemId, "utf8") + ), 0); + return persistedBytes + discardedBytes; +} + function parentTerminalPayload( outcome: TurnOutcome, error: string | undefined, @@ -1828,13 +1839,7 @@ export class CanonicalAgentBoundary implements ParentTurnDurability, CodexCollab if (checkpoint.terminalOutcome) return null; const drafts = this.parentNarrativeRecoveryDrafts(input, thread, turn); // Retention bounds the retained narrative, not the envelope transport overhead. - const persistedBytes = input.items.reduce((total, item) => ( - total + Buffer.byteLength(JSON.stringify(item), "utf8") - ), 0); - const discardedBytes = (input.discardedItemIds ?? []).reduce((total, itemId) => ( - total + Buffer.byteLength(itemId, "utf8") - ), 0); - assertActiveTurnRecoveryRetention(drafts.length, persistedBytes + discardedBytes); + assertActiveTurnRecoveryRetention(drafts.length, retainedNarrativeRecoveryBytes(input)); return { commit: { threadId: thread.id, From e89182df5c49a0bba28478fa41c61fd5c788506e Mon Sep 17 00:00:00 2001 From: Agent Runtime Fixture Date: Tue, 29 Sep 2026 22:25:06 +0100 Subject: [PATCH 12/12] test(agents): expect committed recovery envelopes after finalize The rollback test still encoded the silent-write behavior where recovery produced zero canonical events. recordParentNarrativeRecovery commits item.recorded envelopes independently, so they survive a finalize rollback and remain in the append-only history after the terminal commit retires the item rows. --- .../agents/turns/__tests__/turn-finalizer.test.ts | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/apps/server/src/features/agents/turns/__tests__/turn-finalizer.test.ts b/apps/server/src/features/agents/turns/__tests__/turn-finalizer.test.ts index fc0bf1f33..c19a8c583 100644 --- a/apps/server/src/features/agents/turns/__tests__/turn-finalizer.test.ts +++ b/apps/server/src/features/agents/turns/__tests__/turn-finalizer.test.ts @@ -495,11 +495,13 @@ describe("TurnFinalizer canonical commit recovery", () => { expect(broadcast).not.toHaveBeenCalledWith("turn.persisted", expect.anything()); expect(checkpoints.restore(executionId)).toBe("answer"); expect(sink.loadParentNarrativeRecovery(turnId)).toHaveLength(2); + // Recovery envelopes commit independently of the finalize projection, so a finalize + // rollback must not remove the durable canonical copies. expect(db.prepare(` SELECT COUNT(*) AS count FROM canonical_agent_events WHERE json_extract(envelope_json, '$.payload.item.payload.projection') = 'narrativeRecovery' - `).get()).toEqual({ count: 0 }); + `).get()).toEqual({ count: 2 }); db.exec("DROP TRIGGER reject_canonical_tool"); await finalizer.finalize(THREAD, "completed", Promise.resolve(), executionId); @@ -516,11 +518,13 @@ describe("TurnFinalizer canonical commit recovery", () => { expect(sink.loadItem("toolCall:tool-0")?.payload).toMatchObject({ projection: "toolCall", }); + // Canonical history is append-only: the recorded recovery envelopes stay committed + // while the terminal commit removes the item rows they hydrated. expect(db.prepare(` SELECT COUNT(*) AS count FROM canonical_agent_events WHERE json_extract(envelope_json, '$.payload.item.payload.projection') = 'narrativeRecovery' - `).get()).toEqual({ count: 0 }); + `).get()).toEqual({ count: 2 }); }); it("replays terminal post-commit effects from the canonical projection", async () => {