Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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: [
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -506,7 +506,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)))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ import {
ProviderIdSchema,
ProviderIdentitySchema,
TurnOutcomeSchema,
type AgentEvent,
type CanonicalAgentEventEnvelope,
type ProviderIdentity,
type ParentNarrativeRecoveryItem,
type TurnOutcome,
Expand All @@ -33,8 +35,9 @@ 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 { 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";
Expand Down Expand Up @@ -240,6 +243,8 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter
private bufferedPublication: Parameters<CanonicalAgentEventPublisher>[0][] | null = null;
private bufferedPublicationCount = 0;
private groupedPublication: Parameters<CanonicalAgentEventPublisher>[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) => {
Expand Down Expand Up @@ -287,8 +292,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;
}
Expand Down Expand Up @@ -1041,6 +1050,7 @@ export class CanonicalExecutionSemanticWriter implements ExecutionSemanticWriter
publicationOverride?: readonly ExecutionLivePublicationIntent[],
): Extract<ExecutionWriteReceipt, { kind: "committed" }> {
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) {
Expand All @@ -1063,6 +1073,53 @@ 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,
// Canonical copies reach clients on replay, so they carry the same sanitized
// tool input the wire publication applies in AgentEventPublicationService.
event: sanitizePublicationEvent(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[],
Expand Down Expand Up @@ -1216,6 +1273,24 @@ 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);
}
}

/** 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 {
Expand Down
75 changes: 75 additions & 0 deletions apps/web/src/__tests__/canonical-publication-dispatch.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown>): 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<string, unknown> = {
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);
});
});
19 changes: 19 additions & 0 deletions apps/web/src/stores/threadStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<string>, next: Set<string>): Set<string> {
return previous.size === next.size && [...next].every((id) => previous.has(id)) ? previous : next;
}
Expand Down Expand Up @@ -2791,6 +2806,9 @@ export const useThreadStore = create<ThreadState>((zustandSet, get) => {
? { records, ...(runningThreadIds === state.runningThreadIds ? {} : { runningThreadIds }) }
: {};
});
for (const recovery of recoveries) {
if (recovery.mode === "delta") dispatchCanonicalPublications(recovery.events);
}
},

handleCanonicalAgentEvents: (threadId, events) => {
Expand Down Expand Up @@ -2825,6 +2843,7 @@ export const useThreadStore = create<ThreadState>((zustandSet, get) => {
) {
void conversationResidency.refreshVisibleConversation(threadId);
}
dispatchCanonicalPublications(events);
},

cacheToolCallRecords: (key, records) => {
Expand Down
8 changes: 8 additions & 0 deletions packages/agent-model/src/events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Expand Down
2 changes: 2 additions & 0 deletions packages/agent-model/src/reducer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,8 @@ const AGENT_EVENT_REDUCERS: Record<AgentEventType, AgentEventReducer> = {
event as AgentEventFor<"collaboration-action.recorded">,
acceptedInputState,
),
"publication.recorded": (state, _event, acceptedInputState) =>
reduceVolatileTruncation(state, acceptedInputState),
};

/** Create an empty canonical reducer state. */
Expand Down
Loading