diff --git a/package.json b/package.json index 09a4bd278c6b..6bbcf8caf9b6 100644 --- a/package.json +++ b/package.json @@ -2131,6 +2131,7 @@ "test:sessions:history:bench": "node --import ./scripts/tsx.mjs scripts/bench-session-history.ts", "test:sessions:hydration-memory": "node --import ./scripts/tsx.mjs scripts/bench-session-manager-hydration-memory.ts", "test:sessions:paths:bench": "node --import ./scripts/tsx.mjs scripts/bench-session-paths.ts", + "test:sessions:write-hold": "node --import ./scripts/tsx.mjs scripts/bench-session-manager-write-hold.ts", "test:startup:memory": "node --import ./scripts/tsx.mjs scripts/ensure-cli-startup-build.mts && node scripts/check-cli-startup-memory.mjs", "test:ui": "pnpm lint:ui:no-raw-window-open && node --import ./scripts/tsx.mjs scripts/ensure-playwright-chromium.mts && pnpm --dir ui test", "test:ui:e2e": "node --import ./scripts/tsx.mjs scripts/ensure-playwright-chromium.mts && node scripts/run-vitest.mjs run --config test/vitest/vitest.ui-e2e.config.ts --configLoader runner", diff --git a/scripts/bench-session-manager-write-hold.ts b/scripts/bench-session-manager-write-hold.ts new file mode 100644 index 000000000000..9d0f59392afc --- /dev/null +++ b/scripts/bench-session-manager-write-hold.ts @@ -0,0 +1,129 @@ +import assert from "node:assert/strict"; +import path from "node:path"; +import { SessionManager } from "../src/agents/sessions/session-manager.js"; +import { createZeroUsageFixture } from "../src/agents/test-helpers/usage-fixtures.js"; +import { + loadTranscriptEventsSync, + resolveSessionTranscriptDatabasePath, + upsertSessionEntryCore, +} from "../src/config/sessions/session-accessor.js"; +import { openOpenClawAgentDatabase } from "../src/state/openclaw-agent-db.js"; +import { withOpenClawTestState } from "../src/test-utils/openclaw-test-state.js"; + +// node --import ./scripts/tsx.mjs scripts/bench-session-manager-write-hold.ts +await withOpenClawTestState({ label: "session-write-hold" }, async (state) => { + const scope = { + agentId: "main", + sessionId: "write-hold", + sessionKey: "agent:main:write-hold", + storePath: path.join(state.sessionsDir(), "sessions.json"), + }; + await upsertSessionEntryCore(scope, { sessionId: scope.sessionId, updatedAt: 1 }); + const manager = SessionManager.open(scope, state.workspaceDir); + manager.appendMessage({ role: "user", content: "Explain this code", timestamp: 1 }); + const { db } = openOpenClawAgentDatabase({ + agentId: scope.agentId, + path: resolveSessionTranscriptDatabasePath(scope), + }); + const text = "```ts\nconst example = { label: 'synthetic', count: 42 };\n```\n" + .repeat(4000) + .slice(0, 200 * 1024); + const samples: Array<{ + totalMs: number; + holdMs: number; + commitMs: number; + jsonInHoldMs: number; + }> = []; + const exec = db.exec.bind(db); + const stringify = JSON.stringify; + const parse = JSON.parse; + let holdStart: number | undefined; + let holdMs = 0; + let commitMs = 0; + let jsonInHoldMs = 0; + db.exec = (sql) => { + const finishing = (sql === "COMMIT" || sql === "ROLLBACK") && holdStart !== undefined; + if (finishing && holdStart !== undefined) { + holdMs += performance.now() - holdStart; + holdStart = undefined; + } + const start = performance.now(); + exec(sql); + if (sql === "BEGIN IMMEDIATE") { + holdStart = performance.now(); + } else if (finishing) { + commitMs += performance.now() - start; + } + }; + JSON.stringify = new Proxy(stringify, { + apply(target, receiver, args) { + const start = performance.now(); + const held = holdStart !== undefined; + try { + return Reflect.apply(target, receiver, args); + } finally { + if (held) { + jsonInHoldMs += performance.now() - start; + } + } + }, + }); + JSON.parse = (...args: Parameters) => { + const start = performance.now(); + const held = holdStart !== undefined; + try { + return parse(...args); + } finally { + if (held) { + jsonInHoldMs += performance.now() - start; + } + } + }; + try { + for (let index = 0; index < 60; index += 1) { + holdMs = 0; + commitMs = 0; + jsonInHoldMs = 0; + const start = performance.now(); + manager.appendMessage({ + role: "assistant", + content: [{ type: "text", text }], + api: "messages", + provider: "anthropic", + model: "sonnet-4.6", + usage: createZeroUsageFixture(), + stopReason: "stop", + timestamp: index + 2, + }); + const totalMs = performance.now() - start; + if (index >= 10) { + samples.push({ totalMs, holdMs, commitMs, jsonInHoldMs }); + } + } + } finally { + db.exec = exec; + JSON.stringify = stringify; + JSON.parse = parse; + } + assert.equal(loadTranscriptEventsSync(scope).length, 62); + const summary = (key: keyof (typeof samples)[number]) => { + const values = samples.map((sample) => sample[key]).toSorted((a, b) => a - b); + return { median: values[25], p95: values[47], max: values.at(-1) }; + }; + console.log( + JSON.stringify( + { + node: process.version, + payloadBytes: Buffer.byteLength(text), + samples: samples.length, + totalMs: summary("totalMs"), + holdMs: summary("holdMs"), + commitMs: summary("commitMs"), + jsonInHoldMs: summary("jsonInHoldMs"), + maxRssKiB: process.resourceUsage().maxRSS, + }, + null, + 2, + ), + ); +}); diff --git a/src/agents/session-tool-result-guard.transcript-events.test.ts b/src/agents/session-tool-result-guard.transcript-events.test.ts index 770e46aa33df..69fb754ac317 100644 --- a/src/agents/session-tool-result-guard.transcript-events.test.ts +++ b/src/agents/session-tool-result-guard.transcript-events.test.ts @@ -12,6 +12,7 @@ import { makeTextToolResult } from "../../test/helpers/text-tool-result.js"; import { makeUserMessage } from "../../test/helpers/user-message.js"; import { appendTranscriptMessage, + appendTranscriptMessageSync, loadSessionEntry, listSessionPendingInputs, persistCompactionBoundaryWithSessionEntrySync, @@ -34,6 +35,7 @@ import { createUserTurnTranscriptRecorder, type UserTurnTranscriptRecorder, } from "../sessions/user-turn-transcript.js"; +import { openOpenClawAgentDatabase } from "../state/openclaw-agent-db.js"; import { createAssistantErrorTranscript } from "./assistant-error-transcript.js"; import { normalizeAssistantReplayContent } from "./embedded-agent-runner/replay-history.js"; import { runAgentHarnessBeforeMessageWriteHook } from "./harness/hook-helpers.js"; @@ -77,6 +79,72 @@ afterEach(() => { }); describe("guardSessionManager transcript updates", () => { + it("preserves prepared source and redaction when a concurrent append forces a retry", async () => { + const { sessionManager: manager, target } = await openPersistedSessionManager(); + const baseId = manager.appendMessage(makeUserMessage("Compute a value", 1)); + installSessionToolResultGuard(manager, { + config: { logging: { redactPatterns: [String.raw`/opaque\(([^)]+)\)/g`] } }, + }); + const code = "const API_TOKEN = computeToken(); return API_TOKEN;"; + const toolCall = { + type: "toolCall" as const, + id: "retry-source", + name: "exec", + arguments: { code }, + }; + const message = makeAgentAssistantMessage({ + content: [{ type: "text", text: "opaque(abcdefghijklmnopqrst)" }, toolCall], + stopReason: "toolUse", + }); + const stream = createAssistantMessageEventStream(); + stream.push({ type: "done", reason: "toolUse", message }); + const response = await wrapStreamFnCodeModeSource(() => stream, new Set(["exec"]))( + makeProviderModelFixture({ + id: "test-model", + api: "openai-responses", + provider: "openai", + baseUrl: "https://example.invalid", + }), + { messages: [] }, + ); + const emitted = await response.result(); + const { db } = openOpenClawAgentDatabase({ agentId: target.agentId, path: target.storePath }); + const exec = db.exec.bind(db); + let injected = false; + const execSpy = vi.spyOn(db, "exec").mockImplementation((statement) => { + if (statement === "BEGIN IMMEDIATE" && !injected) { + injected = true; + // Commit after validation but before the writer acquires its snapshot. + const concurrent = appendTranscriptMessageSync(target, { + eventId: "concurrent-assistant", + message: makeAgentAssistantMessage({ content: [{ type: "text", text: "Concurrent" }] }), + }); + expect(concurrent.ok).toBe(true); + } + return exec(statement); + }); + let entryId: string; + try { + entryId = manager.appendMessage( + emitted, + prepareCodeModeSourceAppend({}, emitted, takeCodeModeResponseSource(emitted)), + ); + expect(execSpy).toHaveBeenCalledWith("ROLLBACK"); + } finally { + execSpy.mockRestore(); + } + closeOpenClawAgentDatabasesForTest(); + const entries = SessionManager.open(target).getBranch(); + expect(entries.map(({ id, parentId }) => ({ id, parentId }))).toEqual([ + { id: baseId, parentId: null }, + { id: "concurrent-assistant", parentId: baseId }, + { id: entryId, parentId: "concurrent-assistant" }, + ]); + expect(entries.at(-1)).toMatchObject({ + message: { content: [{ type: "text", text: "opaque(abcdef…qrst)" }, toolCall] }, + }); + }); + it("refreshes the deferred error owner when a session manager serves a new run", async () => { const { sessionManager, target } = await openPersistedSessionManager(); const first = createAssistantErrorTranscript({ runId: "run-first" }); diff --git a/src/agents/sessions/session-manager-entries.ts b/src/agents/sessions/session-manager-entries.ts index 70ba9f271d72..1748e0a51fce 100644 --- a/src/agents/sessions/session-manager-entries.ts +++ b/src/agents/sessions/session-manager-entries.ts @@ -7,6 +7,7 @@ import { validatePreparedAssistantAppendSync, type TranscriptEntryAnchor, } from "../../config/sessions/session-accessor.js"; +import { prepareTranscriptMessageAppend } from "../../config/sessions/session-accessor.sqlite-transcript-message-append.js"; import { resolveSessionTranscriptReadFence } from "../../config/sessions/session-transcript-read-fence.js"; import { applyAssistantDeliveryDirectives } from "../../config/sessions/transcript-assistant-delivery.js"; import { isSessionTranscriptSideAppendEntry } from "../../config/sessions/transcript-tree.js"; @@ -135,9 +136,20 @@ export class SessionManagerEntries extends SessionManagerPersistence { expectedMutationAt: validatedMutationAt, }); } + // Keep preparation local to this append: retries must not redact the payload again or + // consume its code-mode source token against a different message object. + const preparedMessage = + this.persistenceTarget && canonicalEntry.type === "message" + ? prepareTranscriptMessageAppend( + copyCodeModeSourceAppendOptions(options, { + message: canonicalEntry.message, + config: options?.config, + }), + ) + : undefined; let persistenceResult; try { - persistenceResult = this.persist(canonicalEntry, attemptOptions); + persistenceResult = this.persistRecord(canonicalEntry, attemptOptions, preparedMessage); } catch (error) { const deliberateBranchAppend = this.pendingDeliberateAppend; const sideBranchAppend = @@ -182,7 +194,7 @@ export class SessionManagerEntries extends SessionManagerPersistence { ...persistenceOptions, expectedMutationAt: this.readPersistedTranscriptMutationAt(), }); - persistenceResult = this.persist(canonicalEntry, retryOptions); + persistenceResult = this.persistRecord(canonicalEntry, retryOptions, preparedMessage); } if (persistenceResult?.adoptedMessageId) { this.reloadPersistedTranscript(); diff --git a/src/agents/sessions/session-manager-persistence.ts b/src/agents/sessions/session-manager-persistence.ts index e95d19aaa459..b8b4f11644ea 100644 --- a/src/agents/sessions/session-manager-persistence.ts +++ b/src/agents/sessions/session-manager-persistence.ts @@ -13,6 +13,7 @@ import { toDatabaseOptions, } from "../../config/sessions/session-accessor.sqlite-scope.js"; import { requireTranscriptEventAppendSnapshot } from "../../config/sessions/session-accessor.sqlite-transcript-append-result.js"; +import type { PreparedTranscriptMessageAppend } from "../../config/sessions/session-accessor.sqlite-transcript-message-append.js"; import { appendTranscriptEventSnapshotSync, appendTranscriptMessageSnapshotSync, @@ -28,6 +29,7 @@ import { type InitialSessionTranscriptWriter, } from "../../config/sessions/transcript-write-context.js"; import { openOpenClawAgentDatabase } from "../../state/openclaw-agent-db.js"; +import type { AgentMessage } from "../runtime/index.js"; import { copyCodeModeSourceAppendOptions } from "../transcript-code-mode-source.js"; import { getSessionCompactionPersistence } from "./session-compaction-persistence.js"; import { isIndexedSessionEntry, parseOpaqueLeafEntry } from "./session-manager-codec.js"; @@ -430,9 +432,13 @@ export class SessionManagerPersistence extends SessionManagerCore { return removedEntries.length; } - protected persistRecord(entry: unknown, options?: PersistRecordOptions): PersistRecordResult { + protected persistRecord( + entry: unknown, + options?: PersistRecordOptions, + preparedMessage?: PreparedTranscriptMessageAppend, + ): PersistRecordResult { if (this.persistenceTarget) { - return this.persistSqliteRecord(entry, options); + return this.persistSqliteRecord(entry, options, preparedMessage); } if (getSessionCompactionPersistence(this)) { throw new Error("Compaction boundary validation failed"); @@ -444,7 +450,11 @@ export class SessionManagerPersistence extends SessionManagerCore { return this.persistRecord(entry, options); } - private persistSqliteRecord(entry: unknown, options?: PersistRecordOptions): PersistRecordResult { + private persistSqliteRecord( + entry: unknown, + options?: PersistRecordOptions, + preparedMessage?: PreparedTranscriptMessageAppend, + ): PersistRecordResult { if (!this.persistenceTarget) { return undefined; } @@ -585,7 +595,7 @@ export class SessionManagerPersistence extends SessionManagerCore { ...(options?.appendIntent === "active-branch" ? { appendIntent: options.appendIntent } : {}), } satisfies Parameters[1]); const loadedVersion = this.transcriptVersion; - const outcome = appendTranscriptMessageSnapshotSync(scope, appendOptions); + const outcome = appendTranscriptMessageSnapshotSync(scope, appendOptions, preparedMessage); if (!outcome.ok) { throw new Error(`Session transcript message was not persisted: ${entry.id}`, { cause: outcome.error, diff --git a/src/agents/sessions/session-manager-write-hold.test.ts b/src/agents/sessions/session-manager-write-hold.test.ts new file mode 100644 index 000000000000..850140c8a9e1 --- /dev/null +++ b/src/agents/sessions/session-manager-write-hold.test.ts @@ -0,0 +1,100 @@ +import path from "node:path"; +import { expect, it, vi } from "vitest"; +import { + loadTranscriptEventsSync, + resolveSessionTranscriptDatabasePath, + upsertSessionEntryCore, +} from "../../config/sessions/session-accessor.js"; +import { openOpenClawAgentDatabase } from "../../state/openclaw-agent-db.js"; +import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js"; +import type { AgentMessage } from "../runtime/index.js"; +import { createZeroUsageFixture } from "../test-helpers/usage-fixtures.js"; +import * as transcriptRedact from "../transcript-redact.js"; +import { SessionManager } from "./session-manager.js"; + +it.each(["assistant", "toolResult"] as const)( + "prepares large %s payloads before taking the SQLite writer lock", + async (role) => { + await withOpenClawTestState({ label: "session-write-hold" }, async (state) => { + const scope = { + agentId: "main", + sessionId: "write-hold", + sessionKey: "agent:main:write-hold", + storePath: path.join(state.sessionsDir(), "sessions.json"), + }; + await upsertSessionEntryCore(scope, { sessionId: scope.sessionId, updatedAt: 1 }); + const manager = SessionManager.open(scope, state.workspaceDir); + const parentId = manager.appendMessage({ role: "user", content: "first", timestamp: 1 }); + const { db } = openOpenClawAgentDatabase({ + agentId: scope.agentId, + path: resolveSessionTranscriptDatabasePath(scope), + }); + const text = "```ts\nconst value = 42;\n```\n".repeat(8000); + const content = [{ type: "text" as const, text }]; + const message: AgentMessage = + role === "assistant" + ? { + role, + content, + api: "messages", + provider: "anthropic", + model: "sonnet-4.6", + usage: createZeroUsageFixture(), + stopReason: "stop", + timestamp: 2, + } + : { role, content, toolCallId: "call-1", toolName: "read", isError: false, timestamp: 2 }; + Object.assign(message, { MediaPaths: ["/media/a.png"], MediaTypes: ["image/png"] }); + const redactionHeld: boolean[] = []; + const largeJsonHeld: boolean[] = []; + const redact = transcriptRedact.redactTranscriptMessage; + const stringify = JSON.stringify; + const parse = JSON.parse; + const spies = [ + vi.spyOn(transcriptRedact, "redactTranscriptMessage").mockImplementation((...args) => { + redactionHeld.push(db.isTransaction); + return redact(...args); + }), + vi.spyOn(JSON, "stringify").mockImplementation((...args: Parameters) => { + const result = stringify(...args); + if (result && result.length > 100_000) { + largeJsonHeld.push(db.isTransaction); + } + return result; + }), + vi.spyOn(JSON, "parse").mockImplementation((...args: Parameters) => { + if (args[0].length > 100_000) { + largeJsonHeld.push(db.isTransaction); + } + return parse(...args); + }), + ]; + let entryId: string; + try { + entryId = manager.appendMessage(message); + } finally { + for (const spy of spies) { + spy.mockRestore(); + } + } + expect(redactionHeld).toEqual([false]); + expect(largeJsonHeld.length).toBeGreaterThan(0); + expect(largeJsonHeld).not.toContain(true); + const entry = manager.getEntry(entryId); + expect(entry).toMatchObject({ + parentId, + message: { + role, + content, + __openclaw: { media: [{ path: "/media/a.png", contentType: "image/png" }] }, + }, + }); + if (entry?.type !== "message") { + throw new Error("Missing committed message"); + } + expect(entry.message).not.toHaveProperty("MediaPaths"); + expect(entry.message).not.toHaveProperty("MediaTypes"); + expect(loadTranscriptEventsSync(scope).at(-1)).toEqual(manager.getEntry(entryId)); + }); + }, +); diff --git a/src/agents/sessions/session-manager.fork-rebase.test.ts b/src/agents/sessions/session-manager.fork-rebase.test.ts index 11cd499ecff6..e8f06dc6a066 100644 --- a/src/agents/sessions/session-manager.fork-rebase.test.ts +++ b/src/agents/sessions/session-manager.fork-rebase.test.ts @@ -7,12 +7,13 @@ import { appendTranscriptMessageSync, loadTranscriptEvents, replaceTranscriptEventsSync, + resolveSessionTranscriptDatabasePath, upsertSessionEntryCore, } from "../../config/sessions/session-accessor.js"; import { createNestedToolActivity } from "../../sessions/nested-tool-activity.js"; +import { openOpenClawAgentDatabase } from "../../state/openclaw-agent-db.js"; import { cleanupSessionStateForTest } from "../../test-utils/session-state-cleanup.js"; import { createZeroUsageFixture } from "../test-helpers/usage-fixtures.js"; -import type { AppendPersistenceOptions } from "./session-manager-types.js"; import { SessionManager, type SessionEntry, type SessionMessageEntry } from "./session-manager.js"; const tempDirs = useAutoCleanupTempDirTracker((cleanup) => @@ -362,44 +363,41 @@ describe("SessionManager stale-parent rebase", () => { now: 2, }); const branchBeforeRetry = manager.getBranch().map((entry) => entry.id); - const persistencePrototype = SessionManager.prototype as unknown as { - persist( - entry: SessionEntry, - options?: AppendPersistenceOptions & { expectedMutationAt?: number | null }, - ): unknown; - }; - // Preserve the unbound implementation so the spy can forward each concurrent manager receiver. - // oxlint-disable-next-line typescript/unbound-method - const originalPersist = persistencePrototype.persist; - let persistCalls = 0; - vi.spyOn(persistencePrototype, "persist").mockImplementation( - function (this: SessionManager, entry, options) { - persistCalls += 1; - if (persistCalls === 1) { - const concurrent = appendTranscriptMessageSync(target, { - appendIntent: "active-branch", - eventId: "new-user", - message: { role: "user", content: "new", timestamp: 3 }, - now: 3, - }); - expect(concurrent.ok).toBe(true); - } - return originalPersist.call(this, entry, options); - }, - ); - - expect(() => - manager.appendMessage({ - role: "assistant", - content: [{ type: "text", text: "stale reply" }], - api: "openai-responses", - provider: "openai", - model: "gpt-5.5", - usage: createZeroUsageFixture(), - stopReason: "stop", - timestamp: 4, - }), - ).toThrow("SQLite transcript changed while preparing rewrite"); + const { db } = openOpenClawAgentDatabase({ + agentId: target.agentId, + path: resolveSessionTranscriptDatabasePath(target), + }); + const exec = db.exec.bind(db); + let injected = false; + const execSpy = vi.spyOn(db, "exec").mockImplementation((statement) => { + if (statement === "BEGIN IMMEDIATE" && !injected) { + injected = true; + const concurrent = appendTranscriptMessageSync(target, { + appendIntent: "active-branch", + eventId: "new-user", + message: { role: "user", content: "new", timestamp: 3 }, + now: 3, + }); + expect(concurrent.ok).toBe(true); + } + return exec(statement); + }); + try { + expect(() => + manager.appendMessage({ + role: "assistant", + content: [{ type: "text", text: "stale reply" }], + api: "openai-responses", + provider: "openai", + model: "gpt-5.5", + usage: createZeroUsageFixture(), + stopReason: "stop", + timestamp: 4, + }), + ).toThrow("SQLite transcript changed while preparing rewrite"); + } finally { + execSpy.mockRestore(); + } expect(manager.getBranch().map((entry) => entry.id)).toEqual(branchBeforeRetry); const messages = ( (await loadTranscriptEvents(target)) as Array diff --git a/src/config/sessions/session-accessor.sqlite-transcript-message-append.ts b/src/config/sessions/session-accessor.sqlite-transcript-message-append.ts index 722c7fca2117..d31ef9164b00 100644 --- a/src/config/sessions/session-accessor.sqlite-transcript-message-append.ts +++ b/src/config/sessions/session-accessor.sqlite-transcript-message-append.ts @@ -2,6 +2,7 @@ import { randomUUID } from "node:crypto"; import { isDeepStrictEqual } from "node:util"; import { resolveTimestampMsToIsoString } from "@openclaw/normalization-core/number-coercion"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; +import { canonicalizePersistedUserMessageMedia } from "../../media/media-facts.js"; import { isOpenClawDeliveryMirrorAssistantMessage, OPENCLAW_TRANSCRIPT_ARTIFACT_API, @@ -73,6 +74,29 @@ function messagesMatchForIdempotentReplay(stored: unknown, candidate: unknown): return isDeepStrictEqual(serializedShape(stored), serializedShape(candidate, legacyMediaMirror)); } +export type PreparedTranscriptMessageAppend = { + message: TMessage; + messageJson: string; + persistedMessage: TMessage; +}; + +/** SessionManager owns a detached JSON message and retains this preparation across retries. */ +export function prepareTranscriptMessageAppend( + options: Pick, "message" | "config">, +): PreparedTranscriptMessageAppend | undefined { + if ( + !isRecord(options.message) || + (options.message.role !== "assistant" && options.message.role !== "toolResult") + ) { + // Pending user custody retains its transaction-owned preparation. + return undefined; + } + const message = redactTranscriptMessageForStorage(options.message, options); + const messageJson = JSON.stringify(canonicalizePersistedUserMessageMedia(message).message); + // SAFETY: Decode the detached canonical message from its own JSON storage bytes. + return { message, messageJson, persistedMessage: JSON.parse(messageJson) as TMessage }; +} + export function appendTranscriptMessageInTransaction( database: OpenClawAgentDatabase, resolved: ResolvedTranscriptScope, @@ -80,6 +104,7 @@ export function appendTranscriptMessageInTransaction( messageAlreadyRedacted?: boolean; appendMode?: "side"; }, + preparedMessage?: PreparedTranscriptMessageAppend, ): TranscriptMessageAppendResult | undefined { const pending = resolveSessionPendingInputAppend(database, resolved, options.message); if ( @@ -89,7 +114,10 @@ export function appendTranscriptMessageInTransaction( throw new Error("Pending input session changed before transcript promotion"); } const serializeForStorage = (message: TMessage): TMessage => - options.messageAlreadyRedacted ? message : redactTranscriptMessageForStorage(message, options); + preparedMessage?.message ?? + (options.messageAlreadyRedacted + ? message + : redactTranscriptMessageForStorage(message, options)); const readAnchor = (params: { message: unknown; messageId: string; @@ -182,9 +210,16 @@ export function appendTranscriptMessageInTransaction( parentId: parentId ?? null, ...(options.appendMode ? { appendMode: options.appendMode } : {}), timestamp: resolveTimestampMsToIsoString(now), - message: finalMessage, + message: preparedMessage?.persistedMessage ?? finalMessage, }; + let eventJson: string | undefined; + if (preparedMessage) { + // The parent is authoritative only after BEGIN; serialize just its small envelope here. + const { message: _message, ...envelope } = event; + eventJson = `${JSON.stringify(envelope).slice(0, -1)},"message":${preparedMessage.messageJson}}`; + } const appended = appendTranscriptEventInTransaction(database, resolved, event, { + eventJson, idempotencyKeyMode: options.idempotencyLookup === "caller-checked" ? "relocate-owner" @@ -224,8 +259,10 @@ export function appendTranscriptMessageInTransaction( if (!appended) { throw new Error(`SQLite transcript append did not insert message ${messageId}.`); } - // SAFETY: Receipt custody comes from this event's exact committed JSON after storage normalization. - const persistedMessage = (JSON.parse(appended) as typeof event).message; + const persistedMessage = + preparedMessage?.persistedMessage ?? + // SAFETY: Receipt custody comes from this event's exact committed JSON after storage normalization. + (JSON.parse(appended) as typeof event).message; const anchor = readAnchor({ message: persistedMessage, messageId }); if (pending) { if (pending.stageRelocation) { diff --git a/src/config/sessions/session-accessor.sqlite-transcript-store.ts b/src/config/sessions/session-accessor.sqlite-transcript-store.ts index 5a8b2d8e40ef..48fecaf6b42c 100644 --- a/src/config/sessions/session-accessor.sqlite-transcript-store.ts +++ b/src/config/sessions/session-accessor.sqlite-transcript-store.ts @@ -56,6 +56,8 @@ import { copyRetainedTranscriptPayload } from "./session-transcript-retained-dat import { createSessionTranscriptHeader } from "./transcript-header.js"; type TranscriptAppendOptions = { + /** Exact event bytes, including canonical media, prepared by the message append owner. */ + eventJson?: string; allowStoredAlias?: boolean; idempotencyKeyMode?: "dedupe" | "preserve-owner" | "relocate-owner"; onProjectionReconcileNeeded?: () => void; @@ -170,7 +172,8 @@ function appendTranscriptEvent( options: TranscriptAppendOptions, cursor: TranscriptAppendCursor = {}, ): string | false { - const persistedEvent = canonicalizeTranscriptEventMedia(event); + const persistedEvent = + options.eventJson === undefined ? canonicalizeTranscriptEventMedia(event) : event; const db = getSessionKysely(database.db); const createdAt = readEventTimestamp(persistedEvent) ?? Date.now(); if (cursor.initialized) { @@ -206,7 +209,7 @@ function appendTranscriptEvent( } const seq = cursor.nextSeq ?? readNextTranscriptSeq(database, scope.sessionId); cursor.insertEvent ??= createTranscriptEventInserter(database, scope.sessionId); - const eventJson = JSON.stringify(persistedEvent); + const eventJson = options.eventJson ?? JSON.stringify(persistedEvent); cursor.insertEvent({ seq, eventJson, createdAt }); cursor.nextSeq = seq + 1; if (options.touchMutation !== false) { diff --git a/src/config/sessions/session-accessor.sqlite-transcript-write.ts b/src/config/sessions/session-accessor.sqlite-transcript-write.ts index 5b8c54b76263..1cb8d680eb8a 100644 --- a/src/config/sessions/session-accessor.sqlite-transcript-write.ts +++ b/src/config/sessions/session-accessor.sqlite-transcript-write.ts @@ -36,7 +36,10 @@ import { toDatabaseOptions, transcriptWriteScopeIsCurrent, } from "./session-accessor.sqlite-scope.js"; -import { appendTranscriptMessageInTransaction } from "./session-accessor.sqlite-transcript-message-append.js"; +import { + appendTranscriptMessageInTransaction, + type PreparedTranscriptMessageAppend, +} from "./session-accessor.sqlite-transcript-message-append.js"; import { readTranscriptMirrorFacts } from "./session-accessor.sqlite-transcript-mirror.js"; import { resolveTranscriptEventAppendParent } from "./session-accessor.sqlite-transcript-parent.js"; import { @@ -478,13 +481,15 @@ export function appendTranscriptMessageSync( export function appendTranscriptMessageSnapshotSync( scope: SessionTranscriptWriteScope, options: TranscriptMessageAppendOptions, + preparedMessage?: PreparedTranscriptMessageAppend, ): Result< TranscriptWriteSnapshot | undefined>, TranscriptAppendRefusal > { return runTranscriptWriteSnapshotSync( scope, - (database, resolved) => appendTranscriptMessageInTransaction(database, resolved, options), + (database, resolved) => + appendTranscriptMessageInTransaction(database, resolved, options, preparedMessage), undefined, options.expectedMutationAt, );