mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
perf(sessions): prepare transcript payloads before writer locks (#153818)
Prepare canonical assistant and tool-result payloads once per append and retain their serialized bytes across mutation-conflict retries. Keep authoritative parent, sequence, idempotency, generation, and writer checks inside the transaction. A 200 KiB SessionManager append benchmark reduces median write hold from 7.821 ms to 1.864 ms. Preserve code-mode source decisions, media normalization, public persistence interfaces, and exactly-once behavior. No schema or configuration changes.
This commit is contained in:
parent
a0f9b0b766
commit
6e603b9ba7
10 changed files with 416 additions and 53 deletions
|
|
@ -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",
|
||||
|
|
|
|||
129
scripts/bench-session-manager-write-hold.ts
Normal file
129
scripts/bench-session-manager-write-hold.ts
Normal file
|
|
@ -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<typeof parse>) => {
|
||||
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,
|
||||
),
|
||||
);
|
||||
});
|
||||
|
|
@ -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" });
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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<AgentMessage>,
|
||||
): 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<AgentMessage>,
|
||||
): PersistRecordResult {
|
||||
if (!this.persistenceTarget) {
|
||||
return undefined;
|
||||
}
|
||||
|
|
@ -585,7 +595,7 @@ export class SessionManagerPersistence extends SessionManagerCore {
|
|||
...(options?.appendIntent === "active-branch" ? { appendIntent: options.appendIntent } : {}),
|
||||
} satisfies Parameters<typeof appendTranscriptMessageSnapshotSync>[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,
|
||||
|
|
|
|||
100
src/agents/sessions/session-manager-write-hold.test.ts
Normal file
100
src/agents/sessions/session-manager-write-hold.test.ts
Normal file
|
|
@ -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<typeof stringify>) => {
|
||||
const result = stringify(...args);
|
||||
if (result && result.length > 100_000) {
|
||||
largeJsonHeld.push(db.isTransaction);
|
||||
}
|
||||
return result;
|
||||
}),
|
||||
vi.spyOn(JSON, "parse").mockImplementation((...args: Parameters<typeof parse>) => {
|
||||
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));
|
||||
});
|
||||
},
|
||||
);
|
||||
|
|
@ -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<SessionMessageEntry & { type?: string }>
|
||||
|
|
|
|||
|
|
@ -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<TMessage> = {
|
||||
message: TMessage;
|
||||
messageJson: string;
|
||||
persistedMessage: TMessage;
|
||||
};
|
||||
|
||||
/** SessionManager owns a detached JSON message and retains this preparation across retries. */
|
||||
export function prepareTranscriptMessageAppend<TMessage extends object>(
|
||||
options: Pick<TranscriptMessageAppendOptions<TMessage>, "message" | "config">,
|
||||
): PreparedTranscriptMessageAppend<TMessage> | 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<TMessage>(
|
||||
database: OpenClawAgentDatabase,
|
||||
resolved: ResolvedTranscriptScope,
|
||||
|
|
@ -80,6 +104,7 @@ export function appendTranscriptMessageInTransaction<TMessage>(
|
|||
messageAlreadyRedacted?: boolean;
|
||||
appendMode?: "side";
|
||||
},
|
||||
preparedMessage?: PreparedTranscriptMessageAppend<TMessage>,
|
||||
): TranscriptMessageAppendResult<TMessage> | undefined {
|
||||
const pending = resolveSessionPendingInputAppend(database, resolved, options.message);
|
||||
if (
|
||||
|
|
@ -89,7 +114,10 @@ export function appendTranscriptMessageInTransaction<TMessage>(
|
|||
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<TMessage>(
|
|||
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<TMessage>(
|
|||
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) {
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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<TMessage>(
|
|||
export function appendTranscriptMessageSnapshotSync<TMessage>(
|
||||
scope: SessionTranscriptWriteScope,
|
||||
options: TranscriptMessageAppendOptions<TMessage>,
|
||||
preparedMessage?: PreparedTranscriptMessageAppend<TMessage>,
|
||||
): Result<
|
||||
TranscriptWriteSnapshot<TranscriptMessageAppendResult<TMessage> | undefined>,
|
||||
TranscriptAppendRefusal
|
||||
> {
|
||||
return runTranscriptWriteSnapshotSync(
|
||||
scope,
|
||||
(database, resolved) => appendTranscriptMessageInTransaction(database, resolved, options),
|
||||
(database, resolved) =>
|
||||
appendTranscriptMessageInTransaction(database, resolved, options, preparedMessage),
|
||||
undefined,
|
||||
options.expectedMutationAt,
|
||||
);
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue