diff --git a/scripts/check-session-accessor-boundary.mjs b/scripts/check-session-accessor-boundary.mjs index 12982f7dd892..a43fa46f9c89 100644 --- a/scripts/check-session-accessor-boundary.mjs +++ b/scripts/check-session-accessor-boundary.mjs @@ -29,6 +29,11 @@ const legacyWriterNames = new Set([ "updateSessionStore", "updateSessionStoreEntry", ]); +const legacyTranscriptWriterNames = new Set([ + "appendSessionTranscriptMessage", + "emitSessionTranscriptUpdate", + "rewriteTranscriptEntriesInSessionFile", +]); export const migratedSessionAccessorFiles = new Set([ "src/agents/embedded-agent-runner/compaction-successor-transcript.ts", @@ -88,6 +93,15 @@ export const migratedSessionAccessorWriteFiles = new Set([ "src/auto-reply/reply/session-usage.ts", ]); +export const migratedTranscriptWriterFiles = new Set([ + "src/agents/command/attempt-execution.ts", + "src/agents/embedded-agent-runner/context-engine-maintenance.ts", + "src/config/sessions/transcript.ts", + "src/gateway/server-methods/chat.ts", + "src/gateway/server-methods/chat-transcript-inject.ts", + "src/sessions/user-turn-transcript.ts", +]); + function normalizeRelativePath(filePath) { return filePath.replaceAll(path.sep, "/"); } @@ -127,7 +141,7 @@ function bindingName(node) { return null; } -function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyKind) { +function findNamedBoundaryViolations(content, fileName, legacyNames, subject) { const sourceFile = ts.createSourceFile(fileName, content, ts.ScriptTarget.Latest, true); const violations = []; @@ -140,7 +154,7 @@ function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyK if (legacyNames.has(importedName)) { violations.push({ line: toLine(sourceFile, specifier), - reason: `imports legacy session store ${legacyKind} "${importedName}"`, + reason: `imports ${subject} "${importedName}"`, }); } } @@ -152,7 +166,7 @@ function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyK if (name && legacyNames.has(name)) { violations.push({ line: toLine(sourceFile, node), - reason: `aliases legacy session store ${legacyKind} "${name}"`, + reason: `aliases ${subject} "${name}"`, }); } } @@ -160,7 +174,7 @@ function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyK if (ts.isPropertyAccessExpression(node) && legacyNames.has(node.name.text)) { violations.push({ line: toLine(sourceFile, node.name), - reason: `references legacy session store ${legacyKind} "${node.name.text}"`, + reason: `references ${subject} "${node.name.text}"`, }); } @@ -171,7 +185,7 @@ function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyK ) { violations.push({ line: toLine(sourceFile, node.argumentExpression), - reason: `references legacy session store ${legacyKind} "${node.argumentExpression.text}"`, + reason: `references ${subject} "${node.argumentExpression.text}"`, }); } @@ -184,7 +198,7 @@ function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyK ) { violations.push({ line: toLine(sourceFile, node.expression), - reason: `calls legacy session store ${legacyKind} "${calleeName}"`, + reason: `calls ${subject} "${calleeName}"`, }); } } @@ -196,6 +210,15 @@ function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyK return violations; } +function findNamedSessionStoreViolations(content, fileName, legacyNames, legacyKind) { + return findNamedBoundaryViolations( + content, + fileName, + legacyNames, + `legacy session store ${legacyKind}`, + ); +} + export function findSessionAccessorBoundaryViolations(content, fileName = "source.ts") { const legacyNames = legacyNamesForFile(fileName); const legacyKind = legacyNames === legacyWholeStoreAccessNames ? "access" : "reader"; @@ -206,6 +229,15 @@ export function findSessionAccessorWriteBoundaryViolations(content, fileName = " return findNamedSessionStoreViolations(content, fileName, legacyWriterNames, "writer"); } +export function findTranscriptWriterBoundaryViolations(content, fileName = "source.ts") { + return findNamedBoundaryViolations( + content, + fileName, + legacyTranscriptWriterNames, + "legacy transcript writer", + ); +} + export async function main() { const repoRoot = resolveRepoRoot(import.meta.url); const readSourceRoots = resolveSourceRoots(repoRoot, [ @@ -220,6 +252,13 @@ export async function main() { "src/infra", ]); const writeSourceRoots = resolveSourceRoots(repoRoot, ["src/agents", "src/auto-reply"]); + const transcriptWriterSourceRoots = resolveSourceRoots(repoRoot, [ + "src/agents/command", + "src/agents/embedded-agent-runner", + "src/config/sessions", + "src/gateway/server-methods", + "src/sessions", + ]); const readViolations = await collectFileViolations({ repoRoot, sourceRoots: readSourceRoots, @@ -241,7 +280,14 @@ export async function main() { ), findViolations: findSessionAccessorWriteBoundaryViolations, }); - const violations = [...readViolations, ...writeViolations]; + const transcriptWriterViolations = await collectFileViolations({ + repoRoot, + sourceRoots: transcriptWriterSourceRoots, + skipFile: (filePath) => + !migratedTranscriptWriterFiles.has(normalizeRelativePath(path.relative(repoRoot, filePath))), + findViolations: findTranscriptWriterBoundaryViolations, + }); + const violations = [...readViolations, ...writeViolations, ...transcriptWriterViolations]; if (violations.length === 0) { console.log("session accessor boundary guard passed."); @@ -253,7 +299,7 @@ export async function main() { console.error(`- ${violation.path}:${violation.line}: ${violation.reason}`); } console.error( - "Use src/config/sessions/session-accessor.ts helpers for migrated read/write paths. Expand this ratchet only after a slice migrates more files.", + "Use src/config/sessions/session-accessor.ts helpers for migrated read/write and transcript-writer paths. Expand this ratchet only after a slice migrates more files.", ); process.exit(1); } diff --git a/src/agents/cli-runner/types.ts b/src/agents/cli-runner/types.ts index 573b62db85ba..acf535ea9c4e 100644 --- a/src/agents/cli-runner/types.ts +++ b/src/agents/cli-runner/types.ts @@ -17,7 +17,7 @@ import type { InputProvenance } from "../../sessions/input-provenance.js"; import type { PersistedUserTurnMessage, UserTurnTranscriptRecorder, -} from "../../sessions/user-turn-transcript.js"; +} from "../../sessions/user-turn-transcript.types.js"; import type { SkillSnapshot } from "../../skills/types.js"; import type { BootstrapContextMode } from "../bootstrap-files.js"; import type { ResolvedCliBackend } from "../cli-backends.js"; diff --git a/src/agents/command/attempt-execution.cli.test.ts b/src/agents/command/attempt-execution.cli.test.ts index d17358d40c1a..e9f010f073b5 100644 --- a/src/agents/command/attempt-execution.cli.test.ts +++ b/src/agents/command/attempt-execution.cli.test.ts @@ -801,7 +801,7 @@ describe("CLI attempt execution", () => { await fs.realpath(sessionFile), ); expect(persisted[sessionKey]?.updatedAt).toBeGreaterThan(sessionEntry.updatedAt); - expect(persisted[sessionKey]?.updatedAt).toBeLessThan(nowCalls.at(-1) ?? 0); + expect(persisted[sessionKey]?.updatedAt).toBeLessThanOrEqual(nowCalls.at(-1) ?? 0); expect(sessionStore[sessionKey]?.updatedAt).toBe(persisted[sessionKey]?.updatedAt); }); @@ -1034,11 +1034,13 @@ describe("CLI attempt execution", () => { throw new Error("Expected CLI transcript session file."); } expect(path.isAbsolute(sessionFile)).toBe(true); - expect( - sessionFile.endsWith( - path.join(".openclaw", "agents", "main", "sessions", `${sessionEntry.sessionId}.jsonl`), - ), - ).toBe(true); + const persistedFirst = JSON.parse(await fs.readFile(storePath, "utf-8")) as Record< + string, + SessionEntry + >; + expect(await fs.realpath(persistedFirst[sessionKey]?.sessionFile ?? "")).toBe( + await fs.realpath(sessionFile), + ); await appendSessionTranscriptMessage({ transcriptPath: sessionFile, diff --git a/src/agents/command/attempt-execution.ts b/src/agents/command/attempt-execution.ts index 6b063091b4ff..e86460718f01 100644 --- a/src/agents/command/attempt-execution.ts +++ b/src/agents/command/attempt-execution.ts @@ -6,11 +6,8 @@ import { sanitizeForLog } from "../../../packages/terminal-core/src/ansi.js"; import { formatAcpErrorChain } from "../../acp/runtime/errors.js"; import { normalizeReplyPayload } from "../../auto-reply/reply/normalize-reply.js"; import type { ThinkLevel, VerboseLevel } from "../../auto-reply/thinking.js"; -import { appendSessionTranscriptMessage } from "../../config/sessions/transcript-append.js"; -import { - readTailAssistantTextFromSessionTranscript, - resolveSessionTranscriptFile, -} from "../../config/sessions/transcript.js"; +import { persistSessionTranscriptTurn } from "../../config/sessions/session-accessor.js"; +import { readTailAssistantTextFromSessionTranscript } from "../../config/sessions/transcript.js"; import type { SessionEntry } from "../../config/sessions/types.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { @@ -24,9 +21,8 @@ import { createSubsystemLogger } from "../../logging/subsystem.js"; import type { PluginMetadataSnapshot } from "../../plugins/plugin-metadata-snapshot.types.js"; import { isSubagentSessionKey } from "../../routing/session-key.js"; import { annotateInterSessionPromptText } from "../../sessions/input-provenance.js"; -import { emitSessionTranscriptUpdate } from "../../sessions/transcript-events.js"; import { - appendUserTurnTranscriptMessage, + preparePersistedUserTurnMessageForTranscriptWrite, type PersistedUserTurnMessage, } from "../../sessions/user-turn-transcript.js"; import { buildWorkspaceSkillSnapshot } from "../../skills/loading/workspace.js"; @@ -47,14 +43,12 @@ import { resolveOpenAIRuntimeProvider } from "../openai-routing.js"; import { resolveAgentRunAbortLifecycleFields } from "../run-termination.js"; import { buildAgentRuntimeAuthPlan } from "../runtime-plan/auth.js"; import type { AgentMessage } from "../runtime/index.js"; -import { acquireSessionWriteLock, resolveSessionWriteLockOptions } from "../session-write-lock.js"; import { buildUsageWithNoCost } from "../stream-message-shared.js"; import { buildClaudeCliFallbackContextPrelude, claudeCliSessionTranscriptHasContent, resolveFallbackRetryPrompt, } from "./attempt-execution.helpers.js"; -import { persistSessionEntry } from "./attempt-execution.shared.js"; import { resolveAgentRunContext } from "./run-context.js"; import { clearCliSessionInStore } from "./session-store.js"; import type { AgentCommandOpts } from "./types.js"; @@ -237,110 +231,73 @@ async function persistTextTurnTranscript( return params.sessionEntry; } - const { sessionFile, sessionEntry } = await resolveSessionTranscriptFile({ - sessionId: params.sessionId, - sessionKey: params.sessionKey, - sessionEntry: params.sessionEntry, - sessionStore: params.sessionStore, - storePath: params.storePath, - agentId: params.sessionAgentId, - threadId: params.threadId, - }); - const lock = await acquireSessionWriteLock({ - sessionFile, - ...resolveSessionWriteLockOptions(params.config), - allowReentrant: true, - }); - let transcriptMarkerUpdatedAt: number | undefined; - try { - let wroteTranscript = false; - const userMessage = params.userMessage; - if (userMessage || promptText) { - await appendUserTurnTranscriptMessage({ - transcriptPath: sessionFile, - sessionId: params.sessionId, - sessionKey: params.sessionKey, - cwd: params.sessionCwd, - config: params.config, - beforeMessageWrite: runAgentHarnessBeforeMessageWriteHook, - ...(userMessage - ? { message: userMessage } - : { - input: { - text: promptText, - timestamp: Date.now(), - }, - }), - updateMode: "none", - }); - wroteTranscript = true; - } + const messages = []; + const userMessage = + params.userMessage ?? + (promptText + ? ({ + role: "user", + content: promptText, + timestamp: Date.now(), + } as PersistedUserTurnMessage) + : undefined); + if (userMessage) { + messages.push({ + message: userMessage, + idempotencyLookup: "scan" as const, + prepareMessageAfterIdempotencyCheck: (message: unknown) => + preparePersistedUserTurnMessageForTranscriptWrite(message as PersistedUserTurnMessage, { + agentId: params.sessionAgentId, + sessionKey: params.sessionKey, + beforeMessageWrite: runAgentHarnessBeforeMessageWriteHook, + }), + }); + } - if (replyText) { - let appendAssistant = true; - if (params.embeddedAssistantGapFill) { + if (replyText) { + messages.push({ + message: { + role: "assistant", + content: [{ type: "text", text: replyText }], + api: params.assistant.api, + provider: params.assistant.provider, + model: params.assistant.model, + usage: resolveTranscriptUsage(params.assistant.usage), + stopReason: "stop", + timestamp: Date.now(), + }, + shouldAppend: async ({ sessionFile }: { sessionFile: string }) => { + if (!params.embeddedAssistantGapFill) { + return true; + } const latest = await readTailAssistantTextFromSessionTranscript(sessionFile); const normalizedReply = normalizeTranscriptMirrorText(replyText); const normalizedLatest = latest?.text ? normalizeTranscriptMirrorText(latest.text) : ""; - if (normalizedLatest && normalizedLatest === normalizedReply) { - appendAssistant = false; - } - } - if (appendAssistant) { - await appendSessionTranscriptMessage({ - transcriptPath: sessionFile, - sessionId: params.sessionId, - cwd: params.sessionCwd, - config: params.config, - message: { - role: "assistant", - content: [{ type: "text", text: replyText }], - api: params.assistant.api, - provider: params.assistant.provider, - model: params.assistant.model, - usage: resolveTranscriptUsage(params.assistant.usage), - stopReason: "stop", - timestamp: Date.now(), - }, - }); - wroteTranscript = true; - } - } - if (wroteTranscript) { - transcriptMarkerUpdatedAt = Date.now(); - } - } finally { - await lock.release(); + return !normalizedLatest || normalizedLatest !== normalizedReply; + }, + }); } - let updatedSessionEntry = sessionEntry; - if (params.sessionStore && params.storePath && transcriptMarkerUpdatedAt !== undefined) { - const currentEntry = params.sessionStore[params.sessionKey] ?? sessionEntry; - if (currentEntry?.sessionId === params.sessionId) { - // Keep updatedAt as the registry marker for transcript writes we own. - // Session reuse checks compare transcript mtime against this marker, not endedAt. - updatedSessionEntry = - (await persistSessionEntry({ - sessionStore: params.sessionStore, - sessionKey: params.sessionKey, - storePath: params.storePath, - entry: { - sessionId: params.sessionId, - sessionFile, - updatedAt: transcriptMarkerUpdatedAt, - }, - preserveTranscriptMarkerUpdatedAt: true, - shouldPersist: (current) => current?.sessionId === params.sessionId, - })) ?? updatedSessionEntry; - } - } - - emitSessionTranscriptUpdate({ - sessionFile, - sessionKey: params.sessionKey, - agentId: params.sessionAgentId, - }); - return updatedSessionEntry; + const turn = await persistSessionTranscriptTurn( + { + sessionId: params.sessionId, + sessionKey: params.sessionKey, + sessionEntry: params.sessionEntry, + sessionStore: params.sessionStore, + storePath: params.storePath, + agentId: params.sessionAgentId, + threadId: params.threadId, + }, + { + config: params.config, + cwd: params.sessionCwd, + messages, + publishWhen: "always", + touchSessionEntry: true, + updateMode: "file-only", + }, + ); + return turn.sessionEntry; } function resolveCliTranscriptReplyText(result: EmbeddedAgentRunResult): string { diff --git a/src/agents/embedded-agent-runner/context-engine-maintenance.test.ts b/src/agents/embedded-agent-runner/context-engine-maintenance.test.ts index dc90800b4bb8..b3f5d401d853 100644 --- a/src/agents/embedded-agent-runner/context-engine-maintenance.test.ts +++ b/src/agents/embedded-agent-runner/context-engine-maintenance.test.ts @@ -27,7 +27,7 @@ const rewriteTranscriptEntriesInSessionManagerMock = vi.fn((_params?: unknown) = bytesFreed: 77, rewrittenEntries: 1, })); -const rewriteTranscriptEntriesInSessionFileMock = vi.fn(async (_params?: unknown) => ({ +const rewriteTranscriptEntriesInRuntimeTranscriptMock = vi.fn(async (_params?: unknown) => ({ changed: true, bytesFreed: 123, rewrittenEntries: 2, @@ -107,8 +107,8 @@ vi.mock("./context-engine-capabilities.js", () => ({ vi.mock("./transcript-rewrite.js", () => ({ rewriteTranscriptEntriesInSessionManager: (params: unknown) => rewriteTranscriptEntriesInSessionManagerMock(params), - rewriteTranscriptEntriesInSessionFile: (params: unknown) => - rewriteTranscriptEntriesInSessionFileMock(params), + rewriteTranscriptEntriesInRuntimeTranscript: (params: unknown) => + rewriteTranscriptEntriesInRuntimeTranscriptMock(params), })); async function loadFreshContextEngineMaintenanceModuleForTest() { @@ -127,13 +127,13 @@ async function loadFreshContextEngineMaintenanceModuleForTest() { describe("buildContextEngineMaintenanceRuntimeContext", () => { beforeEach(async () => { rewriteTranscriptEntriesInSessionManagerMock.mockClear(); - rewriteTranscriptEntriesInSessionFileMock.mockClear(); + rewriteTranscriptEntriesInRuntimeTranscriptMock.mockClear(); resetSystemEventsForTest(); resetTaskRegistryDeliveryRuntimeForTests(); await loadFreshContextEngineMaintenanceModuleForTest(); }); - it("adds a transcript rewrite helper that targets the current session file", async () => { + it("adds a transcript rewrite helper that targets the current runtime session", async () => { const runtimeContext = buildContextEngineMaintenanceRuntimeContext({ sessionId: "session-1", sessionKey: "agent:main:session-1", @@ -157,16 +157,18 @@ describe("buildContextEngineMaintenanceRuntimeContext", () => { bytesFreed: 123, rewrittenEntries: 2, }); - expect(rewriteTranscriptEntriesInSessionFileMock).toHaveBeenCalledWith({ - sessionFile: "/tmp/session.jsonl", - sessionId: "session-1", - sessionKey: "agent:main:session-1", - config: undefined, + expect(rewriteTranscriptEntriesInRuntimeTranscriptMock).toHaveBeenCalledWith({ + scope: { + sessionId: "session-1", + sessionKey: "agent:main:session-1", + sessionFile: "/tmp/session.jsonl", + }, request: { replacements: [ { entryId: "entry-1", message: { role: "user", content: "hi", timestamp: 1 } }, ], }, + config: undefined, }); }); @@ -198,7 +200,7 @@ describe("buildContextEngineMaintenanceRuntimeContext", () => { { entryId: "entry-1", message: { role: "user", content: "hi", timestamp: 1 } }, ], }); - expect(rewriteTranscriptEntriesInSessionFileMock).not.toHaveBeenCalled(); + expect(rewriteTranscriptEntriesInRuntimeTranscriptMock).not.toHaveBeenCalled(); }); it("wraps active session manager rewrites in the supplied lock", async () => { @@ -242,7 +244,7 @@ describe("buildContextEngineMaintenanceRuntimeContext", () => { { entryId: "entry-1", message: { role: "user", content: "hi", timestamp: 1 } }, ], }); - expect(rewriteTranscriptEntriesInSessionFileMock).not.toHaveBeenCalled(); + expect(rewriteTranscriptEntriesInRuntimeTranscriptMock).not.toHaveBeenCalled(); }); it("lets background file rewrites run without the session lane", async () => { @@ -262,7 +264,7 @@ describe("buildContextEngineMaintenanceRuntimeContext", () => { }); await Promise.resolve(); - rewriteTranscriptEntriesInSessionFileMock.mockImplementationOnce( + rewriteTranscriptEntriesInRuntimeTranscriptMock.mockImplementationOnce( async (_params?: unknown) => { events.push("rewrite"); return { @@ -359,7 +361,7 @@ describe("runContextEngineMaintenance", () => { beforeEach(async () => { vi.useRealTimers(); rewriteTranscriptEntriesInSessionManagerMock.mockClear(); - rewriteTranscriptEntriesInSessionFileMock.mockClear(); + rewriteTranscriptEntriesInRuntimeTranscriptMock.mockClear(); await loadFreshContextEngineMaintenanceModuleForTest(); }); @@ -461,11 +463,12 @@ describe("runContextEngineMaintenance", () => { }); expect(rewriteTranscriptEntriesInSessionManagerMock).not.toHaveBeenCalled(); - expect(rewriteTranscriptEntriesInSessionFileMock).toHaveBeenCalledWith({ - sessionFile: "/tmp/session-background-file-rewrite.jsonl", - sessionId: "session-background-file-rewrite", - sessionKey: "agent:main:session-background-file-rewrite", - config: { session: { writeLock: { acquireTimeoutMs: 75_000 } } }, + expect(rewriteTranscriptEntriesInRuntimeTranscriptMock).toHaveBeenCalledWith({ + scope: { + sessionId: "session-background-file-rewrite", + sessionKey: "agent:main:session-background-file-rewrite", + sessionFile: "/tmp/session-background-file-rewrite.jsonl", + }, request: { replacements: [ { @@ -478,6 +481,7 @@ describe("runContextEngineMaintenance", () => { }, ], }, + config: { session: { writeLock: { acquireTimeoutMs: 75_000 } } }, }); }); @@ -541,7 +545,7 @@ describe("runContextEngineMaintenance", () => { { entryId: "entry-1", message: { role: "user", content: "hi", timestamp: 1 } }, ], }); - expect(rewriteTranscriptEntriesInSessionFileMock).not.toHaveBeenCalled(); + expect(rewriteTranscriptEntriesInRuntimeTranscriptMock).not.toHaveBeenCalled(); }); it("defers turn maintenance to a hidden background task when enabled", async () => { @@ -616,11 +620,12 @@ describe("runContextEngineMaintenance", () => { expect(result).toBeUndefined(); await waitForAssertion(() => expect(maintain).toHaveBeenCalledTimes(1)); await waitForAssertion(() => - expect(rewriteTranscriptEntriesInSessionFileMock).toHaveBeenCalledWith({ - sessionFile: "/tmp/session.jsonl", - sessionId: "session-1", - sessionKey, - config: { session: { writeLock: { acquireTimeoutMs: 91_000 } } }, + expect(rewriteTranscriptEntriesInRuntimeTranscriptMock).toHaveBeenCalledWith({ + scope: { + sessionId: "session-1", + sessionKey, + sessionFile: "/tmp/session.jsonl", + }, request: { replacements: [ { @@ -633,6 +638,7 @@ describe("runContextEngineMaintenance", () => { }, ], }, + config: { session: { writeLock: { acquireTimeoutMs: 91_000 } } }, }), ); @@ -1378,7 +1384,7 @@ describe("runContextEngineMaintenance", () => { }; }); - rewriteTranscriptEntriesInSessionFileMock.mockImplementationOnce( + rewriteTranscriptEntriesInRuntimeTranscriptMock.mockImplementationOnce( async (_params?: unknown) => { events.push("rewrite"); return { diff --git a/src/agents/embedded-agent-runner/context-engine-maintenance.ts b/src/agents/embedded-agent-runner/context-engine-maintenance.ts index 3180c568af99..0de8eb5a5920 100644 --- a/src/agents/embedded-agent-runner/context-engine-maintenance.ts +++ b/src/agents/embedded-agent-runner/context-engine-maintenance.ts @@ -33,7 +33,7 @@ import { findActiveSessionTask } from "../session-async-task-status.js"; import { resolveContextEngineCapabilities } from "./context-engine-capabilities.js"; import { log } from "./logger.js"; import { - rewriteTranscriptEntriesInSessionFile, + rewriteTranscriptEntriesInRuntimeTranscript, rewriteTranscriptEntriesInSessionManager, } from "./transcript-rewrite.js"; @@ -336,16 +336,18 @@ export function buildContextEngineMaintenanceRuntimeContext(params: { ? await params.withSessionManagerRewriteLock(rewriteSessionManagerEntries) : rewriteSessionManagerEntries(); } - const rewriteTranscriptEntriesInFile = async () => - await rewriteTranscriptEntriesInSessionFile({ - sessionFile: params.sessionFile, - sessionId: params.sessionId, - sessionKey: params.sessionKey, - agentId: params.agentId, - config: params.config, + const rewriteRuntimeTranscriptEntries = async () => + await rewriteTranscriptEntriesInRuntimeTranscript({ + scope: { + sessionId: params.sessionId, + sessionKey: params.sessionKey ?? params.sessionId, + sessionFile: params.sessionFile, + ...(params.agentId ? { agentId: params.agentId } : {}), + }, request, + config: params.config, }); - return await rewriteTranscriptEntriesInFile(); + return await rewriteRuntimeTranscriptEntries(); }, }; } diff --git a/src/agents/embedded-agent-runner/run/params.ts b/src/agents/embedded-agent-runner/run/params.ts index 2923c1f949d6..d20ac1988bbc 100644 --- a/src/agents/embedded-agent-runner/run/params.ts +++ b/src/agents/embedded-agent-runner/run/params.ts @@ -15,7 +15,7 @@ import type { ImageContent } from "../../../llm/types.js"; import type { PromptImageOrderEntry } from "../../../media/prompt-image-order.js"; import type { CommandQueueEnqueueFn } from "../../../process/command-queue.types.js"; import type { InputProvenance } from "../../../sessions/input-provenance.js"; -import type { UserTurnTranscriptRecorder } from "../../../sessions/user-turn-transcript.js"; +import type { UserTurnTranscriptRecorder } from "../../../sessions/user-turn-transcript.types.js"; import type { SkillSnapshot } from "../../../skills/types.js"; import type { ExecElevatedDefaults, ExecToolDefaults } from "../../bash-tools.exec-types.js"; import type { AgentStreamParams, ClientToolDefinition } from "../../command/shared-types.js"; diff --git a/src/agents/embedded-agent-runner/transcript-rewrite.test.ts b/src/agents/embedded-agent-runner/transcript-rewrite.test.ts index b4af68ce30ee..fe789fb0a810 100644 --- a/src/agents/embedded-agent-runner/transcript-rewrite.test.ts +++ b/src/agents/embedded-agent-runner/transcript-rewrite.test.ts @@ -329,6 +329,95 @@ describe("rewriteTranscriptEntriesInSessionFile", () => { expect(await fs.readFile(storePath, "utf8")).toBe("{}\n"); }); + it("rewrites runtime transcripts through scoped session identity", async () => { + const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-transcript-rewrite-runtime-")); + const storePath = path.join(dir, "sessions.json"); + const sessionManager = SessionManager.create(dir, dir); + const entryIds = appendSessionMessages(sessionManager, [ + asAppendMessage({ + role: "user", + content: "run tool", + timestamp: 1, + }), + asAppendMessage({ + role: "toolResult", + toolCallId: "call_1", + toolName: "exec", + content: createTextContent("before rewrite"), + isError: false, + timestamp: 2, + }), + asAppendMessage({ + role: "assistant", + content: createTextContent("summarized"), + timestamp: 3, + }), + ]); + const sessionFile = requireString(sessionManager.getSessionFile(), "persisted session file"); + const resolvedSessionFile = await fs.realpath(sessionFile); + const sessionId = path.basename(sessionFile, ".jsonl"); + await fs.writeFile( + storePath, + JSON.stringify({ + "agent:main:test": { + sessionFile, + sessionId, + updatedAt: 10, + }, + }), + "utf8", + ); + const toolResultEntryId = entryIds[1]; + const listener = vi.fn(); + const cleanup = onSessionTranscriptUpdate(listener); + + try { + const result = await rewriteTranscriptEntriesInRuntimeTranscript({ + scope: { + agentId: "main", + sessionId, + sessionKey: "agent:main:test", + storePath, + }, + request: { + replacements: [ + { + entryId: toolResultEntryId, + message: createToolResultReplacement("exec", "[runtime rewrite]", 2), + }, + ], + }, + }); + + expect(result.changed).toBe(true); + expect(acquireSessionWriteLockMock).toHaveBeenCalledWith({ + sessionFile: resolvedSessionFile, + staleMs: 1_800_000, + timeoutMs: 60_000, + maxHoldMs: 300_000, + }); + expect(acquireSessionWriteLockReleaseMock).toHaveBeenCalledTimes(1); + expect(listener).toHaveBeenCalledWith({ + agentId: "main", + sessionFile: resolvedSessionFile, + sessionKey: "agent:main:test", + }); + + const rewrittenSession = SessionManager.open(sessionFile); + const branchMessages = getBranchMessages(rewrittenSession); + expect(branchMessages.map((message) => message.role)).toEqual([ + "user", + "toolResult", + "assistant", + ]); + expect((branchMessages[1] as Extract).content).toEqual([ + { type: "text", text: "[runtime rewrite]" }, + ]); + } finally { + cleanup(); + } + }); + it("aborts under the write lock when the active suffix contains an unexpected entry", async () => { const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-transcript-rewrite-guard-")); const sessionManager = SessionManager.create(dir, dir); diff --git a/src/agents/sessions/session-manager.ts b/src/agents/sessions/session-manager.ts index 0f34623bef07..e25ec6a41081 100644 --- a/src/agents/sessions/session-manager.ts +++ b/src/agents/sessions/session-manager.ts @@ -2113,6 +2113,14 @@ export class SessionManager { this.sessionFileSnapshot = rememberedAppend.snapshot; if (rememberedAppend.ownedAppendVerified && publishSnapshot) { publishRememberedSessionFileSnapshot(this.sessionFile, rememberedAppend.snapshot); + } else if (cacheOwnedAppend) { + this.setLoadedSessionFile( + this.sessionFile, + revalidateLoadedSessionFile(this.sessionFile, { + entries: this.fileEntries, + snapshot: beforeAppendSnapshot, + }), + ); } } } diff --git a/src/auto-reply/get-reply-options.types.ts b/src/auto-reply/get-reply-options.types.ts index 83ed607e6ce1..b2ff7146bac7 100644 --- a/src/auto-reply/get-reply-options.types.ts +++ b/src/auto-reply/get-reply-options.types.ts @@ -1,7 +1,7 @@ /** Public option types for reply generation callbacks, streaming, and delivery policy. */ import type { ImageContent } from "../llm/types.js"; import type { PromptImageOrderEntry } from "../media/prompt-image-order.js"; -import type { UserTurnTranscriptRecorder } from "../sessions/user-turn-transcript.js"; +import type { UserTurnTranscriptRecorder } from "../sessions/user-turn-transcript.types.js"; import type { ReplyPayload } from "./reply-payload.js"; import type { TypingController } from "./reply/typing.js"; diff --git a/src/auto-reply/reply/queue/types.ts b/src/auto-reply/reply/queue/types.ts index d3cd32b9b933..9cb3b2f73cd7 100644 --- a/src/auto-reply/reply/queue/types.ts +++ b/src/auto-reply/reply/queue/types.ts @@ -10,7 +10,7 @@ import type { ReplyToMode } from "../../../config/types.base.js"; import type { OpenClawConfig } from "../../../config/types.openclaw.js"; import type { PromptImageOrderEntry } from "../../../media/prompt-image-order.js"; import type { InputProvenance } from "../../../sessions/input-provenance.js"; -import type { UserTurnTranscriptRecorder } from "../../../sessions/user-turn-transcript.js"; +import type { UserTurnTranscriptRecorder } from "../../../sessions/user-turn-transcript.types.js"; import type { SkillSnapshot } from "../../../skills/types.js"; import type { QueuedReplyDeliveryCorrelation, diff --git a/src/config/sessions/session-accessor.test.ts b/src/config/sessions/session-accessor.test.ts index 831e2fd4d80f..436e185f7b67 100644 --- a/src/config/sessions/session-accessor.test.ts +++ b/src/config/sessions/session-accessor.test.ts @@ -12,6 +12,7 @@ import { loadSessionEntry, loadTranscriptEvents, patchSessionEntry, + persistSessionTranscriptTurn, publishTranscriptUpdate, readSessionUpdatedAt, replaceSessionEntry, @@ -20,7 +21,8 @@ import { updateSessionEntry, upsertSessionEntry, } from "./session-accessor.js"; -import { loadSessionStore } from "./store.js"; +import { loadSessionStore, updateSessionStoreEntry } from "./store.js"; +import { withOwnedSessionTranscriptWrites } from "./transcript-write-context.js"; import type { SessionEntry } from "./types.js"; describe("session accessor file-backed seam", () => { @@ -533,6 +535,281 @@ describe("session accessor file-backed seam", () => { ]); }); + it("persists a transcript turn, touches metadata, and publishes after the write", async () => { + const scope = { + agentId: "main", + sessionId: "session-lock-order", + sessionKey: "agent:main:lock-order", + storePath, + }; + await upsertSessionEntry(scope, { + sessionId: scope.sessionId, + updatedAt: 10, + }); + const updates: Array<{ + lineCount: number; + sessionFile: string | undefined; + updatedAt: number | undefined; + }> = []; + const unsubscribe = onSessionTranscriptUpdate((update) => { + const lines = fs.readFileSync(update.sessionFile, "utf8").trim().split("\n"); + updates.push({ + lineCount: lines.length, + sessionFile: loadSessionEntry(scope)?.sessionFile, + updatedAt: loadSessionEntry(scope)?.updatedAt, + }); + }); + + const result = await persistSessionTranscriptTurn(scope, { + cwd: tempDir, + messages: [ + { + message: { + role: "user", + content: "hello", + timestamp: 100, + }, + }, + { + message: { + role: "assistant", + content: "hi there", + timestamp: 200, + }, + }, + ], + publishWhen: "always", + touchSessionEntry: true, + updateMode: "file-only", + }); + unsubscribe(); + + expect(result.appendedCount).toBe(2); + expect(loadSessionEntry(scope)).toMatchObject({ + sessionFile: result.sessionFile, + sessionId: scope.sessionId, + updatedAt: expect.any(Number), + }); + expect(loadSessionEntry(scope)?.updatedAt).toBeGreaterThanOrEqual(10); + const events = await loadTranscriptEvents({ ...scope, sessionFile: result.sessionFile }); + expect(events).toEqual([ + expect.objectContaining({ type: "session" }), + expect.objectContaining({ + id: result.messages[0]?.messageId, + message: expect.objectContaining({ role: "user", content: "hello" }), + parentId: null, + type: "message", + }), + expect.objectContaining({ + id: result.messages[1]?.messageId, + message: expect.objectContaining({ role: "assistant", content: "hi there" }), + parentId: result.messages[0]?.messageId, + type: "message", + }), + ]); + expect(updates).toEqual([ + { + lineCount: 3, + sessionFile: result.sessionFile, + updatedAt: expect.any(Number), + }, + ]); + }); + + it("queues transcript turn appends before taking the file write lock", async () => { + const scope = { + agentId: "main", + sessionId: "session-1", + sessionKey: "agent:main:main", + storePath, + }; + await upsertSessionEntry(scope, { + sessionId: scope.sessionId, + updatedAt: 10, + }); + let markShouldAppendEntered!: () => void; + const shouldAppendEntered = new Promise((resolve) => { + markShouldAppendEntered = resolve; + }); + let resumeShouldAppend!: () => void; + const shouldAppendReleased = new Promise((resolve) => { + resumeShouldAppend = () => resolve(true); + }); + + const turnPromise = persistSessionTranscriptTurn(scope, { + cwd: tempDir, + messages: [ + { + message: { + role: "assistant", + content: "batch reply", + timestamp: 100, + }, + shouldAppend: async () => { + markShouldAppendEntered(); + return await shouldAppendReleased; + }, + }, + ], + publishWhen: "always", + touchSessionEntry: true, + updateMode: "file-only", + }); + + await shouldAppendEntered; + const queuedAppendPromise = appendTranscriptMessage(scope, { + cwd: tempDir, + message: { + role: "user", + content: "queued prompt", + timestamp: 200, + }, + }); + resumeShouldAppend(); + + const results = Promise.all([turnPromise, queuedAppendPromise]); + const completed = await Promise.race([ + results.then(() => true), + new Promise((resolve) => { + setTimeout(() => resolve(false), 1_000); + }), + ]); + expect(completed).toBe(true); + const [turnResult] = await results; + + const events = await loadTranscriptEvents({ ...scope, sessionFile: turnResult.sessionFile }); + expect( + events + .filter( + (event): event is { message?: { content?: unknown }; type?: unknown } => + typeof event === "object" && + event !== null && + (event as { type?: unknown }).type === "message", + ) + .map((event) => event.message?.content), + ).toEqual(["batch reply", "queued prompt"]); + }); + + it("rejects expected-session transcript turns after a queued session rebind", async () => { + const scope = { + agentId: "main", + sessionId: "session-original", + sessionKey: "agent:main:main", + storePath, + }; + await upsertSessionEntry(scope, { + sessionId: scope.sessionId, + updatedAt: 10, + }); + let releaseReset = () => {}; + const resetGate = new Promise((resolve) => { + releaseReset = resolve; + }); + let markResetStarted = () => {}; + const resetStarted = new Promise((resolve) => { + markResetStarted = resolve; + }); + const replacementSessionFile = path.join(tempDir, "session-replacement.jsonl"); + const reset = updateSessionStoreEntry({ + storePath, + sessionKey: scope.sessionKey, + update: async () => { + markResetStarted(); + await resetGate; + return { + sessionFile: replacementSessionFile, + sessionId: "session-replacement", + }; + }, + }); + await resetStarted; + + const turn = persistSessionTranscriptTurn(scope, { + expectedSessionId: scope.sessionId, + messages: [ + { + message: { + role: "assistant", + content: "late reply", + timestamp: 100, + }, + }, + ], + publishWhen: "always", + touchSessionEntry: true, + updateMode: "file-only", + }); + await new Promise((resolve) => { + setImmediate(resolve); + }); + releaseReset(); + + await reset; + const result = await turn; + + expect(result).toMatchObject({ + appendedCount: 0, + rejectedReason: "session-rebound", + }); + expect(fs.existsSync(path.join(tempDir, "session-original.jsonl"))).toBe(false); + expect(fs.existsSync(replacementSessionFile)).toBe(false); + }); + + it("publishes transcript turn appends through an active owned write lock", async () => { + const scope = { + agentId: "main", + sessionFile: transcriptPath, + sessionId: "session-owned-publish", + sessionKey: "agent:main:owned-publish", + storePath, + }; + const publishOptions: Array = []; + const publishedEntryBatches: unknown[][] = []; + + await withOwnedSessionTranscriptWrites( + { + sessionFile: transcriptPath, + sessionKey: scope.sessionKey, + withSessionWriteLock: async (run, options) => { + publishOptions.push(options?.publishOwnedWrite); + const result = await run(); + publishedEntryBatches.push([...(options?.resolvePublishedEntries?.(result) ?? [])]); + return result; + }, + }, + async () => + await persistSessionTranscriptTurn(scope, { + cwd: tempDir, + messages: [ + { + message: { + role: "assistant", + content: "owned batch", + timestamp: 100, + }, + }, + ], + publishWhen: "always", + touchSessionEntry: true, + updateMode: "file-only", + }), + ); + + expect(publishOptions).toEqual([true]); + expect(publishedEntryBatches).toHaveLength(1); + expect(publishedEntryBatches[0]).toEqual([ + expect.objectContaining({ kind: "header" }), + expect.objectContaining({ kind: "id" }), + ]); + await expect(loadTranscriptEvents(scope)).resolves.toEqual([ + expect.objectContaining({ type: "session" }), + expect.objectContaining({ + message: expect.objectContaining({ content: "owned batch" }), + type: "message", + }), + ]); + }); + it("honors thread fallback paths when resolving transcript scope from the store", async () => { const scope = { agentId: "main", @@ -585,6 +862,29 @@ describe("session accessor file-backed seam", () => { expect(loadSessionEntry(scope)?.sessionFile).toBe(target.sessionFile); }); + it("preserves an explicitly resolved runtime transcript file target", async () => { + const explicitSessionFile = path.join(tempDir, "explicit-session.jsonl"); + const scope = { + agentId: "main", + sessionFile: explicitSessionFile, + sessionId: "session-1", + sessionKey: "agent:main:main", + storePath, + }; + + await upsertSessionEntry(scope, { + sessionId: scope.sessionId, + updatedAt: 10, + }); + + const readTarget = await resolveSessionTranscriptRuntimeReadTarget(scope); + const writeTarget = await resolveSessionTranscriptRuntimeTarget(scope); + + expect(readTarget.sessionFile).toBe(explicitSessionFile); + expect(writeTarget.sessionFile).toBe(explicitSessionFile); + expect(loadSessionEntry(scope)?.sessionFile).toBeUndefined(); + }); + it("keeps read and write runtime targets aligned for new topic sessions", async () => { const scope = { agentId: "main", diff --git a/src/config/sessions/session-accessor.ts b/src/config/sessions/session-accessor.ts index 5085eb67b2c5..8b12f37e05d6 100644 --- a/src/config/sessions/session-accessor.ts +++ b/src/config/sessions/session-accessor.ts @@ -1,5 +1,9 @@ import { randomUUID } from "node:crypto"; import path from "node:path"; +import { + acquireSessionWriteLock, + resolveSessionWriteLockOptions, +} from "../../agents/session-write-lock.js"; import { resolveAgentIdFromSessionKey } from "../../routing/session-key.js"; import { emitSessionTranscriptUpdate } from "../../sessions/transcript-events.js"; import type { SessionTranscriptUpdate } from "../../sessions/transcript-events.js"; @@ -7,6 +11,7 @@ import { getRuntimeConfig } from "../io.js"; import type { OpenClawConfig } from "../types.openclaw.js"; import { resolveSessionFilePath, + resolveSessionFilePathOptions, resolveSessionTranscriptPath, resolveSessionTranscriptPathInDir, resolveStorePath, @@ -27,11 +32,20 @@ import { } from "./store.js"; import { parseSessionThreadInfo } from "./thread-info.js"; import { + type AppendSessionTranscriptMessageParams, + type AppendSessionTranscriptMessageResult, appendSessionTranscriptEvent, appendSessionTranscriptMessage, + appendSessionTranscriptMessageWithOwnedWriteLock, + withSessionTranscriptAppendQueue, } from "./transcript-append.js"; +import { resolveSessionTranscriptFile } from "./transcript-file-resolve.js"; import { streamSessionTranscriptLines } from "./transcript-stream.js"; -import { resolveSessionTranscriptFile } from "./transcript.js"; +import { + type OwnedSessionTranscriptPublishedEntry, + resolveOwnedSessionTranscriptWriteLockRunner, + withOwnedSessionTranscriptWrites, +} from "./transcript-write-context.js"; import type { SessionEntry } from "./types.js"; /** @@ -87,6 +101,8 @@ export type SessionTranscriptAccessScope = SessionTranscriptReadScope & { }; export type SessionTranscriptRuntimeScope = SessionAccessScope & { + /** Resolved file-backed artifact for the current runtime target. */ + sessionFile?: string; sessionId: string; threadId?: string | number; }; @@ -141,6 +157,62 @@ export type TranscriptMessageAppendResult = { /** Transcript update fields supplied by callers; sessionFile is resolved here. */ export type TranscriptUpdatePayload = Omit; +export type SessionTranscriptTurnUpdateMode = "inline" | "file-only" | "none"; + +export type SessionTranscriptTurnMessageAppend = TranscriptMessageAppendOptions & { + /** + * Runs inside the file-backed write lock before this message is appended. + * SQLite implementation note: duplicate/skip decisions should be evaluated + * inside the same write transaction as the transcript row append. + */ + shouldAppend?: (context: SessionTranscriptTurnWriteContext) => Promise | boolean; +}; + +export type SessionTranscriptTurnWriteContext = { + agentId?: string; + sessionFile: string; + sessionId?: string; + sessionKey?: string; +}; + +export type SessionTranscriptTurnPersistOptions = { + /** Runtime config used for lock settings, redaction, and header metadata. */ + config?: OpenClawConfig; + /** Working directory recorded in a newly created transcript header. */ + cwd?: string; + /** + * Rejects the turn when the persisted session key no longer points at this + * runtime session id. SQLite implementations must evaluate this guard inside + * the same write transaction as the transcript append and metadata touch. + */ + expectedSessionId?: string; + /** Message rows to append under one transcript write lock. */ + messages: readonly SessionTranscriptTurnMessageAppend[]; + /** Controls whether the update event includes the last appended message. */ + updateMode?: SessionTranscriptTurnUpdateMode; + /** Emit file-only updates even when every candidate message was skipped. */ + publishWhen?: "always" | "when-appended"; + /** + * Touch updatedAt/sessionFile metadata after appending. + * SQLite implementation note: transcript row append(s) plus this session + * metadata touch should be one SQLite write transaction; publish happens + * after that transaction commits. + */ + touchSessionEntry?: boolean; +}; + +export type SessionTranscriptTurnPersistResult = { + appendedCount: number; + messages: TranscriptMessageAppendResult[]; + rejectedReason?: "session-rebound"; + sessionEntry: SessionEntry | undefined; + sessionFile: string; +}; + +type SessionTranscriptTurnAppendRunner = ( + params: AppendSessionTranscriptMessageParams, +) => Promise | undefined>; + export type SessionTranscriptRuntimeTarget = { agentId: string; sessionFile: string; @@ -391,6 +463,241 @@ export async function publishTranscriptUpdate( }); } +/** + * Persists one logical transcript turn through the current file-backed writer. + * The file implementation resolves/rebinds the transcript file, holds one + * session write lock across all message appends, optionally touches session + * metadata, then publishes after the write has completed. + * + * SQLite implementation note: the transcript row append(s), sessionFile marker, + * and requested updatedAt touch become one SQLite write transaction; transcript + * update delivery must run only after commit. + */ +export async function persistSessionTranscriptTurn( + scope: SessionTranscriptWriteScope & { + sessionEntry?: SessionEntry; + sessionStore?: Record; + }, + options: SessionTranscriptTurnPersistOptions, +): Promise { + const expectedSessionId = options.expectedSessionId; + if (expectedSessionId) { + return await persistExpectedSessionTranscriptTurn(scope, { ...options, expectedSessionId }); + } + const target = await resolveTranscriptTurnTarget(scope); + const appendedMessages = await appendTranscriptTurnMessages(target, options); + const appendedCount = countAppendedTranscriptMessages(appendedMessages); + const sessionEntry = await touchTranscriptTurnSessionEntry({ + scope, + target, + shouldTouch: options.touchSessionEntry === true && appendedCount > 0, + }); + await publishTranscriptTurnUpdate({ + target, + updateMode: options.updateMode ?? "inline", + publishWhen: options.publishWhen ?? "when-appended", + appendedMessages, + }); + + return { + appendedCount, + messages: appendedMessages, + sessionEntry, + sessionFile: target.sessionFile, + }; +} + +async function appendTranscriptTurnMessages( + target: SessionTranscriptTurnWriteContext, + options: SessionTranscriptTurnPersistOptions, +): Promise[]> { + const appendedMessages: TranscriptMessageAppendResult[] = []; + const publishedEntries: OwnedSessionTranscriptPublishedEntry[] = []; + const appendMessages = async (appendMessage: SessionTranscriptTurnAppendRunner) => { + for (const append of options.messages) { + const shouldAppend = append.shouldAppend + ? await append.shouldAppend({ + ...(target.agentId ? { agentId: target.agentId } : {}), + sessionFile: target.sessionFile, + ...(target.sessionId ? { sessionId: target.sessionId } : {}), + ...(target.sessionKey ? { sessionKey: target.sessionKey } : {}), + }) + : true; + if (!shouldAppend) { + continue; + } + const result = await appendMessage({ + transcriptPath: target.sessionFile, + message: append.message, + ...(target.sessionId ? { sessionId: target.sessionId } : {}), + ...((append.cwd ?? options.cwd) ? { cwd: append.cwd ?? options.cwd } : {}), + ...((append.config ?? options.config) ? { config: append.config ?? options.config } : {}), + ...(append.idempotencyLookup ? { idempotencyLookup: append.idempotencyLookup } : {}), + ...(append.now !== undefined ? { now: append.now } : {}), + ...(append.prepareMessageAfterIdempotencyCheck + ? { prepareMessageAfterIdempotencyCheck: append.prepareMessageAfterIdempotencyCheck } + : {}), + onHeaderCreated: (header) => { + publishedEntries.push({ kind: "header", serialized: header }); + }, + ...(append.useRawWhenLinear !== undefined + ? { useRawWhenLinear: append.useRawWhenLinear } + : {}), + }); + if (result) { + appendedMessages.push(result); + if (result.appended) { + publishedEntries.push({ kind: "id", id: result.messageId }); + } + } + } + }; + const activeLockRunner = resolveOwnedSessionTranscriptWriteLockRunner({ + sessionFile: target.sessionFile, + sessionKey: target.sessionKey, + }); + const runBatchWithOwnedLock = async () => + await withOwnedSessionTranscriptWrites( + { + sessionFile: target.sessionFile, + sessionKey: target.sessionKey, + withSessionWriteLock: async (run) => await run(), + }, + async () => await appendMessages(appendSessionTranscriptMessageWithOwnedWriteLock), + ); + if (activeLockRunner) { + await activeLockRunner( + () => withSessionTranscriptAppendQueue(target.sessionFile, runBatchWithOwnedLock), + { + publishOwnedWrite: true, + resolvePublishedEntries: () => publishedEntries, + resolvePublishedEntriesAfterFailure: () => publishedEntries, + }, + ); + } else { + await withSessionTranscriptAppendQueue(target.sessionFile, async () => { + const lock = await acquireSessionWriteLock({ + sessionFile: target.sessionFile, + ...resolveSessionWriteLockOptions(options.config), + allowReentrant: true, + }); + try { + await runBatchWithOwnedLock(); + } finally { + await lock.release(); + } + }); + } + return appendedMessages; +} + +function countAppendedTranscriptMessages( + messages: readonly TranscriptMessageAppendResult[], +): number { + return messages.filter((message) => message.appended).length; +} + +async function persistExpectedSessionTranscriptTurn( + scope: SessionTranscriptWriteScope & { + sessionEntry?: SessionEntry; + sessionStore?: Record; + }, + options: SessionTranscriptTurnPersistOptions & { expectedSessionId: string }, +): Promise { + const sessionKey = scope.sessionKey?.trim(); + if (!scope.storePath || !sessionKey) { + throw new Error("Cannot guard a transcript turn without a session store and key"); + } + const expectedSessionId = options.expectedSessionId; + const agentId = scope.agentId ?? resolveAgentIdFromSessionKey(sessionKey); + if (!agentId) { + throw new Error(`Cannot resolve transcript turn without an agent id: ${sessionKey}`); + } + const store = + scope.sessionStore ?? loadSessionStore(scope.storePath, { skipCache: true, clone: false }); + const resolved = resolveSessionStoreEntry({ store, sessionKey }); + let appendedMessages: TranscriptMessageAppendResult[] = []; + let target: SessionTranscriptTurnWriteContext = { + agentId, + sessionFile: + scope.sessionFile ?? + resolveSessionTranscriptPathInDir(expectedSessionId, path.dirname(scope.storePath)), + sessionId: expectedSessionId, + sessionKey: resolved.normalizedKey, + }; + let rejectedEntry: SessionEntry | undefined; + let touchUpdatedAt: number | undefined; + + const updated = await updateSessionEntry( + { + sessionKey: resolved.normalizedKey, + storePath: scope.storePath, + }, + async (currentEntry) => { + if (currentEntry.sessionId !== expectedSessionId) { + rejectedEntry = currentEntry; + return null; + } + const sessionFile = + scope.sessionFile ?? + resolveSessionFilePath( + currentEntry.sessionId, + currentEntry, + resolveSessionFilePathOptions({ + agentId, + storePath: scope.storePath, + }), + ); + target = { + agentId, + sessionFile, + sessionId: currentEntry.sessionId, + sessionKey: resolved.normalizedKey, + }; + appendedMessages = await appendTranscriptTurnMessages(target, options); + const appendedCount = countAppendedTranscriptMessages(appendedMessages); + if (options.touchSessionEntry === true && appendedCount > 0) { + touchUpdatedAt = Date.now(); + } + const patch = { + ...(currentEntry.sessionFile === sessionFile ? {} : { sessionFile }), + ...(touchUpdatedAt !== undefined + ? { updatedAt: Math.max(currentEntry.updatedAt ?? 0, touchUpdatedAt) } + : {}), + }; + return Object.keys(patch).length > 0 ? patch : null; + }, + { skipMaintenance: true }, + ); + + if (rejectedEntry || updated?.sessionId !== expectedSessionId) { + return { + appendedCount: 0, + messages: [], + rejectedReason: "session-rebound", + sessionEntry: rejectedEntry ?? updated ?? undefined, + sessionFile: target.sessionFile, + }; + } + + await publishTranscriptTurnUpdate({ + target, + updateMode: options.updateMode ?? "inline", + publishWhen: options.publishWhen ?? "when-appended", + appendedMessages, + }); + + if (updated && scope.sessionStore) { + scope.sessionStore[resolved.normalizedKey] = updated; + } + return { + appendedCount: countAppendedTranscriptMessages(appendedMessages), + messages: appendedMessages, + sessionEntry: updated ?? scope.sessionEntry, + sessionFile: target.sessionFile, + }; +} + /** * Resolves the current file-backed target for a storage-neutral runtime * transcript scope. Callers use the scope as identity; sessionFile is returned @@ -411,6 +718,14 @@ export async function resolveSessionTranscriptRuntimeTarget( : undefined; const sessionEntry = resolvedStoreEntry?.existing ?? loadSessionEntry(scope); const sessionKey = resolvedStoreEntry?.normalizedKey ?? scope.sessionKey; + if (scope.sessionFile?.trim()) { + return { + agentId, + sessionFile: path.resolve(scope.sessionFile), + sessionId: scope.sessionId, + sessionKey, + }; + } if (sessionStore && scope.storePath) { const sessionsDir = path.dirname(path.resolve(scope.storePath)); const threadId = scope.threadId ?? parseSessionThreadInfo(scope.sessionKey).threadId; @@ -473,6 +788,14 @@ export async function resolveSessionTranscriptRuntimeReadTarget( : undefined; const sessionEntry = resolvedStoreEntry?.existing ?? loadSessionEntry(scope); const sessionKey = resolvedStoreEntry?.normalizedKey ?? scope.sessionKey; + if (scope.sessionFile?.trim()) { + return { + agentId, + sessionFile: path.resolve(scope.sessionFile), + sessionId: scope.sessionId, + sessionKey, + }; + } const matchingSessionEntry = sessionEntry?.sessionId === scope.sessionId ? sessionEntry : undefined; if (scope.storePath) { @@ -569,3 +892,121 @@ async function resolveTranscriptAccess(scope: SessionTranscriptWriteScope): Prom sessionKey: scopeSessionKey, }); } + +async function resolveTranscriptTurnTarget( + scope: SessionTranscriptWriteScope & { + sessionEntry?: SessionEntry; + sessionStore?: Record; + }, +): Promise< + SessionTranscriptTurnWriteContext & { + sessionEntry: SessionEntry | undefined; + } +> { + if (scope.sessionFile?.trim()) { + return { + ...(scope.agentId ? { agentId: scope.agentId } : {}), + sessionFile: scope.sessionFile, + ...(scope.sessionId ? { sessionId: scope.sessionId } : {}), + ...(scope.sessionKey ? { sessionKey: scope.sessionKey } : {}), + sessionEntry: scope.sessionEntry, + }; + } + const sessionKey = scope.sessionKey?.trim(); + if (!sessionKey || !scope.sessionId) { + throw new Error( + "Cannot persist a transcript turn without a session key and session id or explicit session file", + ); + } + const agentId = scope.agentId ?? resolveAgentIdFromSessionKey(sessionKey); + if (!agentId) { + throw new Error(`Cannot resolve transcript turn without an agent id: ${sessionKey}`); + } + const store = + scope.sessionStore ?? + (scope.storePath ? loadSessionStore(scope.storePath, { skipCache: true }) : undefined); + const resolved = store ? resolveSessionStoreEntry({ store, sessionKey }) : undefined; + const sessionEntry = + resolved?.existing ?? scope.sessionEntry ?? loadSessionEntry({ ...scope, sessionKey }); + const resolvedFile = await resolveSessionTranscriptFile({ + agentId, + sessionEntry, + sessionId: scope.sessionId, + sessionKey, + ...(store ? { sessionStore: store } : {}), + ...(scope.storePath ? { storePath: scope.storePath } : {}), + ...(scope.threadId !== undefined ? { threadId: scope.threadId } : {}), + }); + return { + agentId, + sessionFile: resolvedFile.sessionFile, + sessionId: scope.sessionId, + sessionKey: resolved?.normalizedKey ?? sessionKey, + sessionEntry: resolvedFile.sessionEntry, + }; +} + +async function touchTranscriptTurnSessionEntry(params: { + scope: SessionTranscriptWriteScope & { + sessionEntry?: SessionEntry; + sessionStore?: Record; + }; + target: SessionTranscriptTurnWriteContext & { + sessionEntry: SessionEntry | undefined; + }; + shouldTouch: boolean; +}): Promise { + if ( + !params.shouldTouch || + !params.scope.storePath || + !params.target.sessionKey || + !params.target.sessionId + ) { + return params.target.sessionEntry; + } + const markerUpdatedAt = Date.now(); + const updated = await updateSessionEntry( + { + sessionKey: params.target.sessionKey, + storePath: params.scope.storePath, + }, + (current) => + current.sessionId === params.target.sessionId + ? { + sessionFile: params.target.sessionFile, + updatedAt: Math.max(current.updatedAt ?? 0, markerUpdatedAt), + } + : null, + { skipMaintenance: true }, + ); + if (updated && params.scope.sessionStore) { + params.scope.sessionStore[params.target.sessionKey] = updated; + } + return updated ?? params.target.sessionEntry; +} + +async function publishTranscriptTurnUpdate(params: { + target: SessionTranscriptTurnWriteContext; + updateMode: SessionTranscriptTurnUpdateMode; + publishWhen: "always" | "when-appended"; + appendedMessages: TranscriptMessageAppendResult[]; +}): Promise { + if (params.updateMode === "none") { + return; + } + const lastAppended = params.appendedMessages.findLast((message) => message.appended); + if (params.publishWhen === "when-appended" && !lastAppended) { + return; + } + emitSessionTranscriptUpdate({ + ...(params.target.sessionKey ? { sessionKey: params.target.sessionKey } : {}), + ...(params.target.agentId ? { agentId: params.target.agentId } : {}), + ...(params.updateMode === "inline" && lastAppended + ? { + message: lastAppended.message, + messageId: lastAppended.messageId, + } + : {}), + sessionFile: params.target.sessionFile, + }); +} diff --git a/src/config/sessions/transcript-append.ts b/src/config/sessions/transcript-append.ts index b2f9e3f811f6..72eaaa877c7c 100644 --- a/src/config/sessions/transcript-append.ts +++ b/src/config/sessions/transcript-append.ts @@ -338,7 +338,7 @@ async function resolveTranscriptAppendQueueKey(transcriptPath: string): Promise< } } -async function withTranscriptAppendQueue( +export async function withSessionTranscriptAppendQueue( transcriptPath: string, fn: () => Promise, ): Promise { @@ -363,7 +363,7 @@ async function withTranscriptAppendQueue( } } -type AppendSessionTranscriptMessageParams = { +export type AppendSessionTranscriptMessageParams = { transcriptPath: string; message: TMessage; now?: number; @@ -375,9 +375,11 @@ type AppendSessionTranscriptMessageParams = { /** Runs under the transcript write lock after idempotency replay checks and before append. */ prepareMessageAfterIdempotencyCheck?: (message: TMessage) => TMessage | undefined; config?: OpenClawConfig; + /** Internal owned-batch hook for publishing a newly created transcript header. */ + onHeaderCreated?: (serializedHeader: string) => void; }; -type AppendSessionTranscriptMessageResult = { +export type AppendSessionTranscriptMessageResult = { messageId: string; message: TMessage; appended: boolean; @@ -413,7 +415,7 @@ export async function appendSessionTranscriptMessage( let publishedHeader: string | undefined; return await activeLockRunner( () => - withTranscriptAppendQueue(params.transcriptPath, () => + withSessionTranscriptAppendQueue(params.transcriptPath, () => appendSessionTranscriptMessageLocked({ ...params, onHeaderCreated: (header) => { @@ -432,11 +434,36 @@ export async function appendSessionTranscriptMessage( }, ); } - return await withTranscriptAppendQueue(params.transcriptPath, () => + return await withSessionTranscriptAppendQueue(params.transcriptPath, () => withSessionTranscriptWriteLock(params, () => appendSessionTranscriptMessageLocked(params)), ); } +/** + * Appends a message while the caller already owns the transcript write lock and + * append FIFO. Batch writers use this to keep queue-before-lock ordering while + * reusing the same file lock for multiple transcript rows. + */ +export async function appendSessionTranscriptMessageWithOwnedWriteLock( + params: AppendSessionTranscriptMessageParams & { + prepareMessageAfterIdempotencyCheck: (message: TMessage) => TMessage | undefined; + }, +): Promise | undefined>; +export async function appendSessionTranscriptMessageWithOwnedWriteLock( + params: AppendSessionTranscriptMessageParams, +): Promise>; +export async function appendSessionTranscriptMessageWithOwnedWriteLock( + params: AppendSessionTranscriptMessageParams, +): Promise | undefined> { + const activeLockRunner = resolveOwnedSessionTranscriptWriteLockRunner({ + sessionFile: params.transcriptPath, + }); + if (!activeLockRunner) { + throw new Error("Owned transcript write lock is required for batch transcript append"); + } + return await activeLockRunner(() => appendSessionTranscriptMessageLocked(params)); +} + export type AppendSessionTranscriptEventParams = { config?: OpenClawConfig; event: unknown; @@ -453,7 +480,7 @@ export async function appendSessionTranscriptEvent( if (activeLockRunner) { await activeLockRunner( () => - withTranscriptAppendQueue(params.transcriptPath, () => + withSessionTranscriptAppendQueue(params.transcriptPath, () => appendSessionTranscriptEventLocked(params), ), { @@ -465,7 +492,7 @@ export async function appendSessionTranscriptEvent( ); return; } - await withTranscriptAppendQueue(params.transcriptPath, () => + await withSessionTranscriptAppendQueue(params.transcriptPath, () => withSessionTranscriptWriteLock(params, () => appendSessionTranscriptEventLocked(params)), ); } @@ -496,9 +523,7 @@ async function appendSessionTranscriptEventLocked( } async function appendSessionTranscriptMessageLocked( - params: AppendSessionTranscriptMessageParams & { - onHeaderCreated?: (serializedHeader: string) => void; - }, + params: AppendSessionTranscriptMessageParams, ): Promise | undefined> { const now = params.now ?? Date.now(); const serializedHeader = await ensureTranscriptHeader(params.transcriptPath, { diff --git a/src/config/sessions/transcript-file-resolve.ts b/src/config/sessions/transcript-file-resolve.ts new file mode 100644 index 000000000000..00c53bceec17 --- /dev/null +++ b/src/config/sessions/transcript-file-resolve.ts @@ -0,0 +1,59 @@ +// Resolves transcript file targets without depending on transcript read/write facades. +import { + resolveSessionFilePath, + resolveSessionFilePathOptions, + resolveSessionTranscriptPath, +} from "./paths.js"; +import { resolveAndPersistSessionFile } from "./session-file.js"; +import { parseSessionThreadInfo } from "./thread-info.js"; +import type { SessionEntry } from "./types.js"; + +/** + * Resolves the transcript file for a session and persists the resolved target + * when the caller supplies the owning session store. + */ +export async function resolveSessionTranscriptFile(params: { + sessionId: string; + sessionKey: string; + sessionEntry: SessionEntry | undefined; + sessionStore?: Record; + storePath?: string; + agentId: string; + threadId?: string | number; +}): Promise<{ sessionFile: string; sessionEntry: SessionEntry | undefined }> { + const sessionPathOpts = resolveSessionFilePathOptions({ + agentId: params.agentId, + storePath: params.storePath, + }); + let sessionFile = resolveSessionFilePath(params.sessionId, params.sessionEntry, sessionPathOpts); + let sessionEntry = params.sessionEntry; + + if (params.sessionStore && params.storePath) { + // Persisting the resolved transcript path keeps later tail reads and exports on the same file. + const threadIdFromSessionKey = parseSessionThreadInfo(params.sessionKey).threadId; + const fallbackSessionFile = !sessionEntry?.sessionFile + ? resolveSessionTranscriptPath( + params.sessionId, + params.agentId, + params.threadId ?? threadIdFromSessionKey, + ) + : undefined; + const resolvedSessionFile = await resolveAndPersistSessionFile({ + sessionId: params.sessionId, + sessionKey: params.sessionKey, + sessionStore: params.sessionStore, + storePath: params.storePath, + sessionEntry, + agentId: sessionPathOpts?.agentId, + sessionsDir: sessionPathOpts?.sessionsDir, + fallbackSessionFile, + }); + sessionFile = resolvedSessionFile.sessionFile; + sessionEntry = resolvedSessionFile.sessionEntry; + } + + return { + sessionFile, + sessionEntry, + }; +} diff --git a/src/config/sessions/transcript.test.ts b/src/config/sessions/transcript.test.ts index 4554e513c050..6b41f3b00e41 100644 --- a/src/config/sessions/transcript.test.ts +++ b/src/config/sessions/transcript.test.ts @@ -286,7 +286,7 @@ describe("appendAssistantMessageToSessionTranscript", () => { ); expect(result.ok).toBe(true); - expect(events).toEqual(["lock", "lock"]); + expect(events).toEqual(["lock"]); }); it("keeps matching owned transcript appends locked from bound callbacks", async () => { diff --git a/src/config/sessions/transcript.ts b/src/config/sessions/transcript.ts index 9bea21d39514..77e437f3bab2 100644 --- a/src/config/sessions/transcript.ts +++ b/src/config/sessions/transcript.ts @@ -4,24 +4,15 @@ import type { AgentMessage } from "../../agents/runtime/index.js"; import type { SessionManager } from "../../agents/sessions/session-manager.js"; import { redactTranscriptMessage } from "../../agents/transcript-redact.js"; import { formatErrorMessage } from "../../infra/errors.js"; -import { emitSessionTranscriptUpdate } from "../../sessions/transcript-events.js"; import { extractAssistantVisibleText } from "../../shared/chat-message-content.js"; import { isTranscriptOnlyOpenClawAssistantModel } from "../../shared/transcript-only-openclaw-assistant.js"; import type { OpenClawConfig } from "../types.openclaw.js"; -import { - resolveDefaultSessionStorePath, - resolveSessionFilePath, - resolveSessionFilePathOptions, - resolveSessionTranscriptPath, -} from "./paths.js"; +import { resolveDefaultSessionStorePath } from "./paths.js"; +import { persistSessionTranscriptTurn } from "./session-accessor.js"; import { resolveAndPersistSessionFile } from "./session-file.js"; -import { loadSessionStore, resolveSessionStoreEntry, updateSessionStoreEntry } from "./store.js"; -import { parseSessionThreadInfo } from "./thread-info.js"; -import { appendSessionTranscriptMessage } from "./transcript-append.js"; +import { loadSessionStore, resolveSessionStoreEntry } from "./store.js"; import { resolveMirroredTranscriptText } from "./transcript-mirror.js"; import { streamSessionTranscriptLinesReverse } from "./transcript-stream.js"; -import { runWithOwnedSessionTranscriptWriteLock } from "./transcript-write-context.js"; -import type { SessionEntry } from "./types.js"; export type SessionTranscriptAppendResult = | { ok: true; sessionFile: string; messageId: string } @@ -80,6 +71,8 @@ type AssistantTranscriptText = { export type LatestAssistantTranscriptText = AssistantTranscriptText; export type TailAssistantTranscriptText = AssistantTranscriptText; +export { resolveSessionTranscriptFile } from "./transcript-file-resolve.js"; + function parseAssistantTranscriptText( line: string, options?: { excludeTranscriptOnlyOpenClawAssistant?: boolean }, @@ -120,52 +113,6 @@ function isTranscriptOnlyOpenClawAssistantMessage(message: { return isTranscriptOnlyOpenClawAssistantModel(message.provider, message.model); } -export async function resolveSessionTranscriptFile(params: { - sessionId: string; - sessionKey: string; - sessionEntry: SessionEntry | undefined; - sessionStore?: Record; - storePath?: string; - agentId: string; - threadId?: string | number; -}): Promise<{ sessionFile: string; sessionEntry: SessionEntry | undefined }> { - const sessionPathOpts = resolveSessionFilePathOptions({ - agentId: params.agentId, - storePath: params.storePath, - }); - let sessionFile = resolveSessionFilePath(params.sessionId, params.sessionEntry, sessionPathOpts); - let sessionEntry = params.sessionEntry; - - if (params.sessionStore && params.storePath) { - // Persisting the resolved transcript path keeps later tail reads and exports on the same file. - const threadIdFromSessionKey = parseSessionThreadInfo(params.sessionKey).threadId; - const fallbackSessionFile = !sessionEntry?.sessionFile - ? resolveSessionTranscriptPath( - params.sessionId, - params.agentId, - params.threadId ?? threadIdFromSessionKey, - ) - : undefined; - const resolvedSessionFile = await resolveAndPersistSessionFile({ - sessionId: params.sessionId, - sessionKey: params.sessionKey, - sessionStore: params.sessionStore, - storePath: params.storePath, - sessionEntry, - agentId: sessionPathOpts?.agentId, - sessionsDir: sessionPathOpts?.sessionsDir, - fallbackSessionFile, - }); - sessionFile = resolvedSessionFile.sessionFile; - sessionEntry = resolvedSessionFile.sessionEntry; - } - - return { - sessionFile, - sessionEntry, - }; -} - export async function readLatestAssistantTextFromSessionTranscript( sessionFile: string | undefined, ): Promise { @@ -312,140 +259,108 @@ export async function appendExactAssistantMessageToSessionTranscript(params: { return { ok: false, reason: `unknown sessionKey: ${sessionKey}` }; } - let transcriptMarkerUpdatedAt: number | undefined; - let appendedSessionId = entry.sessionId; const appendToSessionFile = async ( - currentEntry: SessionEntry, - sessionFile: string, - ): Promise => - await runWithOwnedSessionTranscriptWriteLock( - { sessionFile, sessionKey: resolved.normalizedKey }, - async (): Promise => { - const explicitIdempotencyKey = - params.idempotencyKey ?? - ((params.message as { idempotencyKey?: unknown }).idempotencyKey as string | undefined); - const message = { - ...params.message, - ...(explicitIdempotencyKey ? { idempotencyKey: explicitIdempotencyKey } : {}), - } as Parameters[0]; - const preparedUnkeyedMessage = - !explicitIdempotencyKey && params.beforeMessageWrite - ? applyBeforeMessageWriteToAssistant({ - message, - beforeMessageWrite: params.beforeMessageWrite, - agentId: params.agentId, - sessionKey: resolved.normalizedKey, - }) - : message; - if (!preparedUnkeyedMessage) { - return { - ok: false, - code: "blocked", - reason: "blocked by before_message_write", - }; - } - const identifiedChannelFinal = - Boolean(explicitIdempotencyKey) && isChannelFinalDeliveryMirror(params.message); - const latestEquivalentAssistantId = - isRedundantDeliveryMirror(params.message) && !identifiedChannelFinal - ? await findLatestEquivalentAssistantMessageId( - sessionFile, - preparedUnkeyedMessage as SessionTranscriptAssistantMessage, - params.config, - ) - : undefined; - // Unidentified delivery mirrors dedupe by latest text. Identified channel finals use their - // idempotency key so repeated replies on separate user turns remain distinct. - if (latestEquivalentAssistantId) { - return { ok: true, sessionFile, messageId: latestEquivalentAssistantId }; - } - const appendedResult = await appendSessionTranscriptMessage({ - transcriptPath: sessionFile, - sessionId: currentEntry.sessionId, - cwd: currentEntry.spawnedCwd, - message: preparedUnkeyedMessage, - ...(explicitIdempotencyKey ? { idempotencyLookup: "scan" } : {}), - ...(explicitIdempotencyKey && params.beforeMessageWrite - ? { - prepareMessageAfterIdempotencyCheck: ( - candidate: Parameters[0], - ) => - applyBeforeMessageWriteToAssistant({ - message: candidate, - beforeMessageWrite: params.beforeMessageWrite, - explicitIdempotencyKey, - agentId: params.agentId, - sessionKey: resolved.normalizedKey, - }), - } - : {}), - config: params.config, - }); - if (!appendedResult) { - return { - ok: false, - code: "blocked", - reason: "blocked by before_message_write", - }; - } - const { messageId, message: appendedMessage, appended } = appendedResult; - if (!appended) { - return { ok: true, sessionFile, messageId }; - } - transcriptMarkerUpdatedAt = Date.now(); - - switch (params.updateMode ?? "inline") { - case "inline": - emitSessionTranscriptUpdate({ - sessionFile, - sessionKey, - ...(params.agentId ? { agentId: params.agentId } : {}), - message: appendedMessage, - messageId, - }); - break; - case "file-only": - emitSessionTranscriptUpdate({ - sessionFile, - sessionKey, - ...(params.agentId ? { agentId: params.agentId } : {}), - }); - break; - case "none": - break; - } - return { ok: true, sessionFile, messageId }; + currentEntry: NonNullable, + sessionFile?: string, + ): Promise => { + const explicitIdempotencyKey = + params.idempotencyKey ?? + ((params.message as { idempotencyKey?: unknown }).idempotencyKey as string | undefined); + const message = { + ...params.message, + ...(explicitIdempotencyKey ? { idempotencyKey: explicitIdempotencyKey } : {}), + } as Parameters[0]; + const preparedUnkeyedMessage = + !explicitIdempotencyKey && params.beforeMessageWrite + ? applyBeforeMessageWriteToAssistant({ + message, + beforeMessageWrite: params.beforeMessageWrite, + agentId: params.agentId, + sessionKey: resolved.normalizedKey, + }) + : message; + if (!preparedUnkeyedMessage) { + return { + ok: false, + code: "blocked", + reason: "blocked by before_message_write", + }; + } + const identifiedChannelFinal = + Boolean(explicitIdempotencyKey) && isChannelFinalDeliveryMirror(params.message); + let latestEquivalentAssistantId: string | undefined; + // Unidentified delivery mirrors dedupe by latest text. Identified channel finals use their + // idempotency key so repeated replies on separate user turns remain distinct. + const turn = await persistSessionTranscriptTurn( + { + sessionId: currentEntry.sessionId, + sessionKey: resolved.normalizedKey, + storePath, + ...(sessionFile ? { sessionFile } : {}), + ...(params.agentId ? { agentId: params.agentId } : {}), + }, + { + cwd: currentEntry.spawnedCwd, + ...(params.expectedSessionId ? { expectedSessionId: params.expectedSessionId } : {}), + ...(params.config ? { config: params.config } : {}), + updateMode: params.updateMode ?? "inline", + touchSessionEntry: true, + messages: [ + { + message: preparedUnkeyedMessage, + ...(explicitIdempotencyKey ? { idempotencyLookup: "scan" } : {}), + ...(explicitIdempotencyKey && params.beforeMessageWrite + ? { + prepareMessageAfterIdempotencyCheck: (candidate: unknown) => + applyBeforeMessageWriteToAssistant({ + message: candidate as Parameters[0], + beforeMessageWrite: params.beforeMessageWrite, + explicitIdempotencyKey, + agentId: params.agentId, + sessionKey: resolved.normalizedKey, + }), + } + : {}), + shouldAppend: async (target) => { + latestEquivalentAssistantId = + isRedundantDeliveryMirror(params.message) && !identifiedChannelFinal + ? await findLatestEquivalentAssistantMessageId( + target.sessionFile, + preparedUnkeyedMessage as SessionTranscriptAssistantMessage, + params.config, + ) + : undefined; + return !latestEquivalentAssistantId; + }, + }, + ], }, ); + if (turn.rejectedReason === "session-rebound") { + return { + ok: false, + code: "session-rebound", + reason: `session rebound for sessionKey: ${sessionKey}`, + }; + } + if (latestEquivalentAssistantId) { + return { ok: true, sessionFile: turn.sessionFile, messageId: latestEquivalentAssistantId }; + } + const appendedResult = turn.messages[0]; + if (!appendedResult) { + return { + ok: false, + code: "blocked", + reason: "blocked by before_message_write", + }; + } + const { messageId } = appendedResult; + return { ok: true, sessionFile: turn.sessionFile, messageId }; + }; let result: SessionTranscriptAppendResult; if (params.expectedSessionId) { - result = { - ok: false, - code: "session-rebound", - reason: `session rebound for sessionKey: ${sessionKey}`, - }; - await updateSessionStoreEntry({ - storePath, - sessionKey: resolved.normalizedKey, - update: async (currentEntry) => { - if (currentEntry.sessionId !== params.expectedSessionId) { - return null; - } - const sessionFile = resolveSessionFilePath( - currentEntry.sessionId, - currentEntry, - resolveSessionFilePathOptions({ - agentId: params.agentId, - storePath, - }), - ); - appendedSessionId = currentEntry.sessionId; - result = await appendToSessionFile(currentEntry, sessionFile); - return currentEntry.sessionFile === sessionFile ? null : { sessionFile }; - }, - skipMaintenance: true, - }); + result = await appendToSessionFile(entry); } else { let sessionFile: string; try { @@ -467,14 +382,6 @@ export async function appendExactAssistantMessageToSessionTranscript(params: { } result = await appendToSessionFile(entry, sessionFile); } - if (result.ok && transcriptMarkerUpdatedAt !== undefined) { - await updateSessionStoreEntry({ - storePath, - sessionKey: resolved.normalizedKey, - update: (current) => - current.sessionId === appendedSessionId ? { updatedAt: transcriptMarkerUpdatedAt } : null, - }); - } return result; } diff --git a/src/gateway/server-methods/chat-transcript-inject.ts b/src/gateway/server-methods/chat-transcript-inject.ts index a3f2c6d5c450..ff05a852de93 100644 --- a/src/gateway/server-methods/chat-transcript-inject.ts +++ b/src/gateway/server-methods/chat-transcript-inject.ts @@ -1,10 +1,9 @@ // Chat transcript injection appends gateway-authored assistant rows while // preserving agent-session parent links and transcript update notifications. import type { SessionManager } from "../../agents/sessions/session-manager.js"; -import { appendSessionTranscriptMessage } from "../../config/sessions/transcript-append.js"; +import { persistSessionTranscriptTurn } from "../../config/sessions/session-accessor.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { formatErrorMessage } from "../../infra/errors.js"; -import { emitSessionTranscriptUpdate } from "../../sessions/transcript-events.js"; type AppendMessageArg = Parameters[0]; @@ -119,21 +118,33 @@ export async function appendInjectedAssistantMessageToTranscript(params: { }; try { - const { messageId, message: appendedMessage } = await appendSessionTranscriptMessage({ - transcriptPath: params.transcriptPath, - message: messageBody, - now, - useRawWhenLinear: true, - config: params.config, - }); - emitSessionTranscriptUpdate({ - sessionFile: params.transcriptPath, - ...(params.sessionKey ? { sessionKey: params.sessionKey } : {}), - ...(params.agentId ? { agentId: params.agentId } : {}), - message: appendedMessage, - messageId, - }); - return { ok: true, messageId, message: appendedMessage as unknown as Record }; + const turn = await persistSessionTranscriptTurn( + { + sessionFile: params.transcriptPath, + sessionKey: params.sessionKey ?? "", + ...(params.agentId ? { agentId: params.agentId } : {}), + }, + { + updateMode: "inline", + ...(params.config ? { config: params.config } : {}), + messages: [ + { + message: messageBody, + now, + useRawWhenLinear: true, + }, + ], + }, + ); + const appended = turn.messages[0]; + if (!appended) { + return { ok: false, error: "gateway-injected assistant message was not appended" }; + } + return { + ok: true, + messageId: appended.messageId, + message: appended.message as Record, + }; } catch (err) { return { ok: false, error: formatErrorMessage(err) }; } diff --git a/src/gateway/server-methods/chat.ts b/src/gateway/server-methods/chat.ts index 2ffc06779e38..db5b17f3c6c4 100644 --- a/src/gateway/server-methods/chat.ts +++ b/src/gateway/server-methods/chat.ts @@ -38,7 +38,7 @@ import { resolveAgentWorkspaceDir, resolveSessionAgentId, } from "../../agents/agent-scope.js"; -import { rewriteTranscriptEntriesInSessionFile } from "../../agents/embedded-agent-runner/transcript-rewrite.js"; +import { rewriteTranscriptEntriesInRuntimeTranscript } from "../../agents/embedded-agent-runner/transcript-rewrite.js"; import { runAgentHarnessBeforeMessageWriteHook } from "../../agents/harness/hook-helpers.js"; import { modelCatalogBrowseRequiresFullDiscovery } from "../../agents/model-catalog-browse.js"; import type { ModelCatalogEntry } from "../../agents/model-catalog.types.js"; @@ -3700,10 +3700,11 @@ export const chatHandlers: GatewayRequestHandlers = { input: baseUserTurnInput, resolveInput: () => userTurnInputPromise, target: () => { - const { storePath: latestStorePath, entry: latestEntry } = loadSessionEntry( - sessionKey, - sessionLoadOptions, - ); + const { + storePath: latestStorePath, + store: latestStore, + entry: latestEntry, + } = loadSessionEntry(sessionKey, sessionLoadOptions); const resolvedSessionId = latestEntry?.sessionId ?? backingSessionId; if (!resolvedSessionId) { return undefined; @@ -3712,6 +3713,7 @@ export const chatHandlers: GatewayRequestHandlers = { sessionId: resolvedSessionId, sessionKey, sessionEntry: latestEntry ?? entry, + sessionStore: latestStore, storePath: latestStorePath, agentId, config: cfg, @@ -4761,11 +4763,14 @@ export const chatHandlers: GatewayRequestHandlers = { allowedSourceReplyMirrorIds.has(entryLocal.id), ) === true; if (canRewriteSourceReplyMirrors) { - const result = await rewriteTranscriptEntriesInSessionFile({ - sessionFile: resolvedTranscriptPath, - sessionKey, - agentId, - config: cfg, + const result = await rewriteTranscriptEntriesInRuntimeTranscript({ + scope: { + sessionId, + sessionKey, + sessionFile: resolvedTranscriptPath, + agentId, + ...(latestStorePath ? { storePath: latestStorePath } : {}), + }, request: { allowedRewriteSuffixEntryIds: [...allowedSourceReplyMirrorIds], replacements: rewriteTargets.map((target) => ({ @@ -4777,6 +4782,7 @@ export const chatHandlers: GatewayRequestHandlers = { } as unknown as AgentMessage, })), }, + config: cfg, }); if (result.changed) { await advanceSessionTranscriptMarker({ diff --git a/src/gateway/server-methods/server-methods.test.ts b/src/gateway/server-methods/server-methods.test.ts index 4fb1870af3dd..555495a764a6 100644 --- a/src/gateway/server-methods/server-methods.test.ts +++ b/src/gateway/server-methods/server-methods.test.ts @@ -2337,7 +2337,7 @@ describe("gateway chat transcript writes (guardrail)", () => { expect(chatSrc.includes("fs.appendFileSync(transcriptPath")).toBe(false); expect(chatSrc).toContain("appendInjectedAssistantMessageToTranscript("); - expect(helperSrc).toContain("appendSessionTranscriptMessage({"); + expect(helperSrc).toContain("persistSessionTranscriptTurn("); expect(helperSrc).toContain("useRawWhenLinear: true"); expect(helperSrc).not.toContain("SessionManager.open(params.transcriptPath)"); }); diff --git a/src/gateway/server.chat.gateway-server-chat-b.test.ts b/src/gateway/server.chat.gateway-server-chat-b.test.ts index 741d3283f0c5..50c13910dc3f 100644 --- a/src/gateway/server.chat.gateway-server-chat-b.test.ts +++ b/src/gateway/server.chat.gateway-server-chat-b.test.ts @@ -97,14 +97,26 @@ async function withGatewayChatHarness( } } -async function writeMainSessionStore() { +function testSessionFilePath(sessionDir: string, sessionId: string): string { + return path.join(sessionDir, `${sessionId}.jsonl`); +} + +async function writeMainSessionStore(sessionDir?: string, sessionId = "sess-main") { await writeSessionStore({ entries: { - main: { sessionId: "sess-main", updatedAt: Date.now() }, + main: { + sessionId, + updatedAt: futureFixtureUpdatedAt(), + ...(sessionDir ? { sessionFile: testSessionFilePath(sessionDir, sessionId) } : {}), + }, }, }); } +function futureFixtureUpdatedAt(): number { + return Date.now() + 60_000; +} + async function writeGatewayConfig(config: Record) { const configPath = process.env.OPENCLAW_CONFIG_PATH; if (!configPath) { @@ -115,8 +127,12 @@ async function writeGatewayConfig(config: Record) { clearConfigCache(); } -async function writeMainSessionTranscript(sessionDir: string, lines: string[]) { - await fs.writeFile(path.join(sessionDir, "sess-main.jsonl"), `${lines.join("\n")}\n`, "utf-8"); +async function writeMainSessionTranscript( + sessionDir: string, + lines: string[], + sessionId = "sess-main", +) { + await fs.writeFile(testSessionFilePath(sessionDir, sessionId), `${lines.join("\n")}\n`, "utf-8"); } async function removeTempDir(dir: string): Promise { @@ -197,13 +213,14 @@ async function prepareMainHistoryHarness(params: { ws: GatewaySocket; createSessionDir: () => Promise; historyMaxBytes?: number; + sessionId?: string; }) { if (params.historyMaxBytes !== undefined) { setMaxChatHistoryMessagesBytesForTest(params.historyMaxBytes); } await connectOk(params.ws); const sessionDir = await params.createSessionDir(); - await writeMainSessionStore(); + await writeMainSessionStore(sessionDir, params.sessionId); return sessionDir; } @@ -1985,6 +2002,7 @@ describe("gateway server chat", () => { await withGatewayChatHarness(async ({ ws, createSessionDir }) => { await connectOk(ws); const sessionDir = await createSessionDir(); + const sessionId = "sess-claude-cli-backfill"; const originalHome = process.env.HOME; const homeDir = path.join(sessionDir, "home"); const cliSessionId = "5b8b202c-f6bb-4046-9475-d2f15fd07530"; @@ -2028,8 +2046,9 @@ describe("gateway server chat", () => { await writeSessionStore({ entries: { main: { - sessionId: "sess-main", - updatedAt: Date.now(), + sessionId, + sessionFile: testSessionFilePath(sessionDir, sessionId), + updatedAt: futureFixtureUpdatedAt(), modelProvider: "claude-cli", model: "claude-sonnet-4-6", cliSessionBindings: { @@ -2737,11 +2756,12 @@ describe("gateway server chat", () => { test("chat.message.get returns archive-backed rows surfaced by history", async () => { await withGatewayChatHarness(async ({ ws, createSessionDir }) => { - const sessionDir = await prepareMainHistoryHarness({ ws, createSessionDir }); + const sessionId = "sess-archive-backed"; + const sessionDir = await prepareMainHistoryHarness({ ws, createSessionDir, sessionId }); await fs.writeFile( - path.join(sessionDir, "sess-main.jsonl.reset.2026-02-16T22-26-34.000Z"), + `${testSessionFilePath(sessionDir, sessionId)}.reset.2026-02-16T22-26-34.000Z`, [ - JSON.stringify({ type: "session", version: 1, id: "sess-main" }), + JSON.stringify({ type: "session", version: 1, id: sessionId }), JSON.stringify({ id: "msg-archive-full-assistant", message: { diff --git a/src/gateway/server.chat.gateway-server-chat.test.ts b/src/gateway/server.chat.gateway-server-chat.test.ts index 50ab504bea30..18a85bc71449 100644 --- a/src/gateway/server.chat.gateway-server-chat.test.ts +++ b/src/gateway/server.chat.gateway-server-chat.test.ts @@ -99,14 +99,19 @@ describe("gateway server chat", () => { }); }; - const withMainSessionStore = async (run: (dir: string) => Promise): Promise => { + const withMainSessionStore = async ( + run: (dir: string) => Promise, + options?: { sessionId?: string }, + ): Promise => { const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-gw-")); try { + const sessionId = options?.sessionId ?? "sess-main"; testState.sessionStorePath = path.join(dir, "sessions.json"); await writeSessionStore({ entries: { main: { - sessionId: "sess-main", + sessionId, + sessionFile: path.join(dir, `${sessionId}.jsonl`), updatedAt: Date.now(), }, }, @@ -1429,93 +1434,100 @@ describe("gateway server chat", () => { }); test("chat.history persists assistant image data URLs as managed image blocks", async () => { - await withMainSessionStore(async (dir) => { - const previousStateDir = process.env.OPENCLAW_STATE_DIR; - process.env.OPENCLAW_STATE_DIR = dir; - const pngB64 = - "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR4nGNgYAAAAAMAASsJTYQAAAAASUVORK5CYII="; - dispatchInboundMessageMock.mockImplementationOnce(async (...args: unknown[]) => { - const [params] = args as [ - { - dispatcher: { - sendFinalReply: (payload: { text?: string; mediaUrls?: string[] }) => boolean; - markComplete: () => void; - waitForIdle: () => Promise; - getQueuedCounts: () => { final: number; block: number; tool: number }; - }; - }, - ]; - params.dispatcher.sendFinalReply({ - mediaUrls: [`data:image/png;base64,${pngB64}`], - }); - params.dispatcher.markComplete(); - await params.dispatcher.waitForIdle(); - return { - queuedFinal: true, - counts: params.dispatcher.getQueuedCounts(), - }; - }); - - try { - const finalPromise = onceMessage( - ws, - (o) => - o.type === "event" && - o.event === "chat" && - o.payload?.state === "final" && - o.payload?.runId === "idem-managed-image-history", - 8000, - ); - const res = await rpcReq(ws, "chat.send", { - sessionKey: "main", - message: "show me an image", - idempotencyKey: "idem-managed-image-history", - }); - - expect(res.ok).toBe(true); - expect(res.payload?.runId).toBe("idem-managed-image-history"); - await finalPromise; - - let assistantMessage: Record | undefined; - await vi.waitFor(async () => { - const historyRes = await rpcReq<{ messages?: unknown[] }>(ws, "chat.history", { - sessionKey: "main", + await withMainSessionStore( + async (dir) => { + const previousStateDir = process.env.OPENCLAW_STATE_DIR; + process.env.OPENCLAW_STATE_DIR = dir; + const pngB64 = + "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR4nGNgYAAAAAMAASsJTYQAAAAASUVORK5CYII="; + dispatchInboundMessageMock.mockImplementationOnce(async (...args: unknown[]) => { + const [params] = args as [ + { + dispatcher: { + sendFinalReply: (payload: { text?: string; mediaUrls?: string[] }) => boolean; + markComplete: () => void; + waitForIdle: () => Promise; + getQueuedCounts: () => { final: number; block: number; tool: number }; + }; + }, + ]; + params.dispatcher.sendFinalReply({ + text: "Image reply", + mediaUrls: [`data:image/png;base64,${pngB64}`], }); - expect(historyRes.ok).toBe(true); - const messages = historyRes.payload?.messages ?? []; - assistantMessage = messages.find( - (message): message is Record => - typeof message === "object" && - message !== null && - (message as { role?: unknown }).role === "assistant", + params.dispatcher.markComplete(); + await params.dispatcher.waitForIdle(); + return { + queuedFinal: true, + counts: params.dispatcher.getQueuedCounts(), + }; + }); + + try { + const finalPromise = onceMessage( + ws, + (o) => + o.type === "event" && + o.event === "chat" && + o.payload?.state === "final" && + o.payload?.runId === "idem-managed-image-history", + 8000, ); - if (!assistantMessage) { - throw new Error("Expected assistant history message"); + const res = await rpcReq(ws, "chat.send", { + sessionKey: "main", + message: "show me an image", + idempotencyKey: "idem-managed-image-history", + }); + + expect(res.ok).toBe(true); + expect(res.payload?.runId).toBe("idem-managed-image-history"); + await finalPromise; + + let assistantMessage: Record | undefined; + await vi.waitFor( + async () => { + const historyRes = await rpcReq<{ messages?: unknown[] }>(ws, "chat.history", { + sessionKey: "main", + }); + expect(historyRes.ok).toBe(true); + const messages = historyRes.payload?.messages ?? []; + assistantMessage = messages.find( + (message): message is Record => + typeof message === "object" && + message !== null && + (message as { role?: unknown }).role === "assistant", + ); + if (!assistantMessage) { + throw new Error("Expected assistant history message"); + } + }, + { timeout: CHAT_RESPONSE_TIMEOUT_MS }, + ); + const assistantContent = (assistantMessage as { content?: unknown[] }).content ?? []; + expect(assistantContent).toHaveLength(2); + expect(assistantContent[0]).toEqual({ type: "text", text: "Image reply" }); + const imageBlock = expectRecordFields(assistantContent[1], { + type: "image", + alt: "Generated image 1", + mimeType: "image/png", + width: 1, + height: 1, + }); + expect(String(imageBlock.url)).toContain("/api/chat/media/outgoing/"); + expect(String(imageBlock.openUrl)).toContain("/full"); + const serializedAssistant = JSON.stringify(assistantMessage); + expect(serializedAssistant).not.toContain("data:image/png;base64"); + expect(serializedAssistant).not.toContain(pngB64); + } finally { + if (previousStateDir == null) { + delete process.env.OPENCLAW_STATE_DIR; + } else { + process.env.OPENCLAW_STATE_DIR = previousStateDir; } - }); - const assistantContent = (assistantMessage as { content?: unknown[] }).content ?? []; - expect(assistantContent).toHaveLength(2); - expect(assistantContent[0]).toEqual({ type: "text", text: "Image reply" }); - const imageBlock = expectRecordFields(assistantContent[1], { - type: "image", - alt: "Generated image 1", - mimeType: "image/png", - width: 1, - height: 1, - }); - expect(String(imageBlock.url)).toContain("/api/chat/media/outgoing/"); - expect(String(imageBlock.openUrl)).toContain("/full"); - const serializedAssistant = JSON.stringify(assistantMessage); - expect(serializedAssistant).not.toContain("data:image/png;base64"); - expect(serializedAssistant).not.toContain(pngB64); - } finally { - if (previousStateDir == null) { - delete process.env.OPENCLAW_STATE_DIR; - } else { - process.env.OPENCLAW_STATE_DIR = previousStateDir; } - } - }); + }, + { sessionId: "sess-managed-image-history" }, + ); }); test("chat.history hides assistant NO_REPLY-only entries and keeps mixed-content assistant entries", async () => { diff --git a/src/gateway/test-helpers.server.ts b/src/gateway/test-helpers.server.ts index 00daf0d9dcb6..ff7d4e1be539 100644 --- a/src/gateway/test-helpers.server.ts +++ b/src/gateway/test-helpers.server.ts @@ -166,7 +166,7 @@ async function persistTestSessionConfig(): Promise { } const nextStoreValue = typeof testState.sessionStorePath === "string" - ? preservedTemplateStore || testState.sessionStorePath + ? testState.sessionStorePath : preservedTemplateStore; for (const configPath of configPaths) { const config = { ...parsedConfigs.get(configPath) }; @@ -223,8 +223,9 @@ export async function writeSessionStore(params: { // file directly; clear the in-process cache so handlers reload the seeded state. clearSessionStoreCacheForTest(); await persistTestSessionConfig(); + const serializedStore = JSON.stringify(store, null, 2); await fs.mkdir(path.dirname(storePath), { recursive: true }); - await fs.writeFile(storePath, JSON.stringify(store, null, 2), "utf-8"); + await fs.writeFile(storePath, serializedStore, "utf-8"); clearSessionStoreCacheForTest(); } diff --git a/src/sessions/input-provenance.ts b/src/sessions/input-provenance.ts index 0da77c0bf760..c9680156e5be 100644 --- a/src/sessions/input-provenance.ts +++ b/src/sessions/input-provenance.ts @@ -1,6 +1,6 @@ // Input provenance helpers normalize source metadata for session messages. import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; -import type { AgentMessage } from "../agents/runtime/index.js"; +import type { AgentMessage } from "../../packages/agent-core/src/types.js"; // Input provenance marks whether a user-role message actually came from an // external user, another session, or an internal system/tool handoff. diff --git a/src/sessions/user-turn-transcript.ts b/src/sessions/user-turn-transcript.ts index 5e4e9c8298c3..bc4992184bbd 100644 --- a/src/sessions/user-turn-transcript.ts +++ b/src/sessions/user-turn-transcript.ts @@ -1,32 +1,29 @@ // User turn transcript helpers extract user-turn text from session transcripts. import path from "node:path"; import { mimeTypeFromFilePath } from "@openclaw/media-core/mime"; -import type { AgentMessage } from "../agents/runtime/index.js"; -import { appendSessionTranscriptMessage } from "../config/sessions/transcript-append.js"; -import { - applyInputProvenanceToUserMessage, - type InputProvenance, - normalizeInputProvenance, -} from "./input-provenance.js"; -import { emitSessionTranscriptUpdate } from "./transcript-events.js"; +import type { AgentMessage } from "../../packages/agent-core/src/types.js"; +import { persistSessionTranscriptTurn } from "../config/sessions/session-accessor.js"; +import type { OpenClawConfig } from "../config/types.openclaw.js"; +import { applyInputProvenanceToUserMessage, normalizeInputProvenance } from "./input-provenance.js"; +import type { + PersistedUserTurnMediaInput, + PersistedUserTurnMessage, + UserTurnBeforeMessageWrite, + UserTurnInput, + UserTurnSessionEntry, + UserTurnTranscriptFileTarget, + UserTurnTranscriptPersistResult, + UserTurnTranscriptRecorder, + UserTurnTranscriptTarget, + UserTurnTranscriptTargetResolver, + UserTurnTranscriptUpdateMode, +} from "./user-turn-transcript.types.js"; -// User-turn transcript helpers persist the selected prompt/media as a user -// message before or during runtime execution, preserving provenance/idempotency. -type TranscriptAppendConfig = Parameters[0]["config"]; - -type UserTurnSessionEntry = { - sessionId: string; - updatedAt: number; - sessionFile?: string; - threadId?: string | number; -} & Record; - -type PersistedUserTurnMediaInput = { - path?: string | null; - url?: string | null; - contentType?: string | null; - kind?: string | null; -}; +export type { + PersistedUserTurnMessage, + UserTurnInput, + UserTurnTranscriptRecorder, +} from "./user-turn-transcript.types.js"; type PersistedUserTurnMediaFields = { MediaPath?: string; @@ -35,25 +32,6 @@ type PersistedUserTurnMediaFields = { MediaTypes?: string[]; }; -export type PersistedUserTurnMessage = Extract; - -export type UserTurnInput = { - text?: string | null; - media?: readonly PersistedUserTurnMediaInput[] | null; - timestamp?: number; - idempotencyKey?: string; - provenance?: InputProvenance; - mediaOnlyText?: string; -}; - -type UserTurnTranscriptUpdateMode = "inline" | "none"; - -export type UserTurnBeforeMessageWrite = (params: { - message: PersistedUserTurnMessage; - agentId?: string; - sessionKey?: string; -}) => AgentMessage | null; - type AppendUserTurnTranscriptMessageParams = { transcriptPath: string; input?: UserTurnInput; @@ -62,7 +40,7 @@ type AppendUserTurnTranscriptMessageParams = { agentId?: string; sessionKey?: string; cwd?: string; - config?: TranscriptAppendConfig; + config?: OpenClawConfig; updateMode?: UserTurnTranscriptUpdateMode; beforeMessageWrite?: UserTurnBeforeMessageWrite; }; @@ -78,60 +56,13 @@ type PersistUserTurnTranscriptParams = { agentId: string; threadId?: string | number; cwd?: string; - config?: TranscriptAppendConfig; + config?: unknown; updateMode?: UserTurnTranscriptUpdateMode; beforeMessageWrite?: UserTurnBeforeMessageWrite; }; -type UserTurnTranscriptPersistenceTarget = Omit< - PersistUserTurnTranscriptParams, - "input" | "message" | "updateMode" ->; - -type UserTurnTranscriptFileTarget = { - transcriptPath: string; - sessionId?: string; - agentId?: string; - sessionKey?: string; - cwd?: string; - config?: TranscriptAppendConfig; -}; - -type UserTurnTranscriptTarget = UserTurnTranscriptPersistenceTarget | UserTurnTranscriptFileTarget; - -type UserTurnTranscriptPersistResult = { - sessionFile: string; - sessionEntry: UserTurnSessionEntry | undefined; - messageId: string; - message: PersistedUserTurnMessage; -}; - -type UserTurnTranscriptTargetResolver = - | UserTurnTranscriptTarget - | (() => UserTurnTranscriptTarget | undefined | Promise); - type UserTurnInputResolver = () => UserTurnInput | undefined | Promise; -export type UserTurnTranscriptRecorder = { - readonly message: PersistedUserTurnMessage | undefined; - resolveMessage: () => Promise; - markRuntimePersistencePending: (pending: Promise) => void; - markRuntimePersisted: (message?: PersistedUserTurnMessage) => void; - markBlocked: () => void; - hasPersisted: () => boolean; - isBlocked: () => boolean; - hasRuntimePersistencePending: () => boolean; - waitForRuntimePersistence: () => Promise; - persistApproved: (params?: { - target?: UserTurnTranscriptTargetResolver; - updateMode?: UserTurnTranscriptUpdateMode; - }) => Promise; - persistFallback: (params?: { - target?: UserTurnTranscriptTargetResolver; - updateMode?: UserTurnTranscriptUpdateMode; - }) => Promise; -}; - type CreateUserTurnTranscriptRecorderParams = { input?: UserTurnInput; message?: PersistedUserTurnMessage; @@ -356,7 +287,8 @@ export function mergePreparedUserTurnMessageForRuntime(params: { } as unknown as AgentMessage; } -function applyBeforeMessageWriteToUserTurn( +/** Applies before-message hooks while preserving user-turn transcript metadata. */ +export function preparePersistedUserTurnMessageForTranscriptWrite( message: PersistedUserTurnMessage, params: Pick< AppendUserTurnTranscriptMessageParams, @@ -406,36 +338,40 @@ export async function appendUserTurnTranscriptMessage( return undefined; } - const appended = await appendSessionTranscriptMessage({ - transcriptPath: params.transcriptPath, - ...(params.sessionId ? { sessionId: params.sessionId } : {}), - ...(params.cwd ? { cwd: params.cwd } : {}), - ...(params.config ? { config: params.config } : {}), - message: resolvedMessage, - idempotencyLookup: "scan", - prepareMessageAfterIdempotencyCheck: (message) => - applyBeforeMessageWriteToUserTurn(message, params), - }); + const turn = await persistSessionTranscriptTurn( + { + sessionFile: params.transcriptPath, + sessionKey: params.sessionKey ?? "", + ...(params.agentId ? { agentId: params.agentId } : {}), + ...(params.sessionId ? { sessionId: params.sessionId } : {}), + }, + { + ...(params.cwd ? { cwd: params.cwd } : {}), + ...(params.config ? { config: params.config } : {}), + updateMode: params.updateMode ?? "inline", + messages: [ + { + message: resolvedMessage, + idempotencyLookup: "scan", + prepareMessageAfterIdempotencyCheck: (message) => + preparePersistedUserTurnMessageForTranscriptWrite( + message as PersistedUserTurnMessage, + params, + ), + }, + ], + }, + ); + const appended = turn.messages[0] as + | { + messageId: string; + message: PersistedUserTurnMessage; + } + | undefined; if (!appended) { return undefined; } - switch (params.updateMode ?? "inline") { - case "inline": - if (appended.appended) { - emitSessionTranscriptUpdate({ - sessionFile: params.transcriptPath, - ...(params.sessionKey ? { sessionKey: params.sessionKey } : {}), - ...(params.agentId ? { agentId: params.agentId } : {}), - message: appended.message, - messageId: appended.messageId, - }); - } - break; - case "none": - break; - } - return { sessionFile: params.transcriptPath, messageId: appended.messageId, @@ -453,38 +389,72 @@ export async function persistUserTurnTranscript( return undefined; } - const { resolveSessionTranscriptFile } = await import("../config/sessions/transcript.js"); - const { sessionFile, sessionEntry } = await resolveSessionTranscriptFile({ - sessionId: params.sessionId, - sessionKey: params.sessionKey, - sessionEntry: params.sessionEntry, - ...(params.sessionStore ? { sessionStore: params.sessionStore } : {}), - ...(params.storePath ? { storePath: params.storePath } : {}), - agentId: params.agentId, - ...(params.threadId !== undefined ? { threadId: params.threadId } : {}), - }); - - const appended = await appendUserTurnTranscriptMessage({ - transcriptPath: sessionFile, - message, - sessionId: params.sessionId, - agentId: params.agentId, - sessionKey: params.sessionKey, - ...(params.cwd ? { cwd: params.cwd } : {}), - ...(params.config ? { config: params.config } : {}), - ...(params.updateMode ? { updateMode: params.updateMode } : {}), - ...(params.beforeMessageWrite ? { beforeMessageWrite: params.beforeMessageWrite } : {}), - }); + const turn = await persistSessionTranscriptTurn( + { + sessionId: params.sessionId, + sessionKey: params.sessionKey, + sessionEntry: params.sessionEntry, + ...(params.sessionStore ? { sessionStore: params.sessionStore } : {}), + ...(params.storePath ? { storePath: params.storePath } : {}), + agentId: params.agentId, + ...(params.threadId !== undefined ? { threadId: params.threadId } : {}), + }, + { + ...(params.cwd ? { cwd: params.cwd } : {}), + ...(params.config ? { config: params.config as OpenClawConfig } : {}), + updateMode: params.updateMode ?? "inline", + messages: [ + { + message, + idempotencyLookup: "scan", + prepareMessageAfterIdempotencyCheck: (candidate) => + preparePersistedUserTurnMessageForTranscriptWrite( + candidate as PersistedUserTurnMessage, + params, + ), + }, + ], + }, + ); + const appended = turn.messages[0] as + | { + messageId: string; + message: PersistedUserTurnMessage; + } + | undefined; if (!appended) { return undefined; } return { ...appended, - sessionEntry, + sessionEntry: turn.sessionEntry, + sessionFile: turn.sessionFile, }; } +async function appendFileTargetUserTurnTranscript(params: { + target: UserTurnTranscriptFileTarget; + message: PersistedUserTurnMessage; + updateMode: UserTurnTranscriptUpdateMode; + beforeMessageWrite?: UserTurnBeforeMessageWrite; +}): Promise { + const { config, ...target } = params.target; + const appended = await appendUserTurnTranscriptMessage({ + ...target, + message: params.message, + updateMode: params.updateMode, + ...(config ? { config: config as OpenClawConfig } : {}), + ...(params.beforeMessageWrite ? { beforeMessageWrite: params.beforeMessageWrite } : {}), + }); + return appended + ? { + ...appended, + sessionEntry: undefined, + } + : undefined; +} + async function resolveUserTurnTranscriptTarget( target: UserTurnTranscriptTargetResolver, ): Promise { @@ -596,19 +566,12 @@ export function createUserTurnTranscriptRecorder( } const updateMode = options.updateMode ?? params.updateMode ?? "inline"; const result = isUserTurnTranscriptFileTarget(target) - ? await appendUserTurnTranscriptMessage({ - ...target, + ? await appendFileTargetUserTurnTranscript({ + target, message: resolvedMessage, updateMode, - ...(params.beforeMessageWrite ? { beforeMessageWrite: params.beforeMessageWrite } : {}), - }).then((appended) => - appended - ? { - ...appended, - sessionEntry: undefined, - } - : undefined, - ) + beforeMessageWrite: params.beforeMessageWrite, + }) : await persistUserTurnTranscript({ ...target, message: resolvedMessage, diff --git a/src/sessions/user-turn-transcript.types.ts b/src/sessions/user-turn-transcript.types.ts new file mode 100644 index 000000000000..ec5e580e1c75 --- /dev/null +++ b/src/sessions/user-turn-transcript.types.ts @@ -0,0 +1,93 @@ +// User-turn transcript type contracts shared by runtime and queue option types. +import type { AgentMessage } from "../../packages/agent-core/src/types.js"; +import type { InputProvenance } from "./input-provenance.js"; + +export type UserTurnSessionEntry = { + sessionId: string; + updatedAt: number; + sessionFile?: string; + threadId?: string | number; +} & Record; + +export type PersistedUserTurnMediaInput = { + path?: string | null; + url?: string | null; + contentType?: string | null; + kind?: string | null; +}; + +export type PersistedUserTurnMessage = Extract; + +export type UserTurnInput = { + text?: string | null; + media?: readonly PersistedUserTurnMediaInput[] | null; + timestamp?: number; + idempotencyKey?: string; + provenance?: InputProvenance; + mediaOnlyText?: string; +}; + +export type UserTurnTranscriptUpdateMode = "inline" | "none"; + +export type UserTurnBeforeMessageWrite = (params: { + message: PersistedUserTurnMessage; + agentId?: string; + sessionKey?: string; +}) => AgentMessage | null; + +export type UserTurnTranscriptPersistenceTarget = { + sessionId: string; + sessionKey: string; + sessionEntry: UserTurnSessionEntry | undefined; + sessionStore?: Record; + storePath?: string; + agentId: string; + threadId?: string | number; + cwd?: string; + config?: unknown; + beforeMessageWrite?: UserTurnBeforeMessageWrite; +}; + +export type UserTurnTranscriptFileTarget = { + transcriptPath: string; + sessionId?: string; + agentId?: string; + sessionKey?: string; + cwd?: string; + config?: unknown; +}; + +export type UserTurnTranscriptTarget = + | UserTurnTranscriptPersistenceTarget + | UserTurnTranscriptFileTarget; + +export type UserTurnTranscriptPersistResult = { + sessionFile: string; + sessionEntry: UserTurnSessionEntry | undefined; + messageId: string; + message: PersistedUserTurnMessage; +}; + +export type UserTurnTranscriptTargetResolver = + | UserTurnTranscriptTarget + | (() => UserTurnTranscriptTarget | undefined | Promise); + +export type UserTurnTranscriptRecorder = { + readonly message: PersistedUserTurnMessage | undefined; + resolveMessage: () => Promise; + markRuntimePersistencePending: (pending: Promise) => void; + markRuntimePersisted: (message?: PersistedUserTurnMessage) => void; + markBlocked: () => void; + hasPersisted: () => boolean; + isBlocked: () => boolean; + hasRuntimePersistencePending: () => boolean; + waitForRuntimePersistence: () => Promise; + persistApproved: (params?: { + target?: UserTurnTranscriptTargetResolver; + updateMode?: UserTurnTranscriptUpdateMode; + }) => Promise; + persistFallback: (params?: { + target?: UserTurnTranscriptTargetResolver; + updateMode?: UserTurnTranscriptUpdateMode; + }) => Promise; +}; diff --git a/test/scripts/check-session-accessor-boundary.test.ts b/test/scripts/check-session-accessor-boundary.test.ts index b0e3dbde4522..53bc19749dab 100644 --- a/test/scripts/check-session-accessor-boundary.test.ts +++ b/test/scripts/check-session-accessor-boundary.test.ts @@ -1,10 +1,12 @@ import { describe, expect, it } from "vitest"; import { findSessionAccessorBoundaryViolations, - migratedBundledPluginSessionAccessorFiles, findSessionAccessorWriteBoundaryViolations, + findTranscriptWriterBoundaryViolations, + migratedBundledPluginSessionAccessorFiles, migratedSessionAccessorFiles, migratedSessionAccessorWriteFiles, + migratedTranscriptWriterFiles, } from "../../scripts/check-session-accessor-boundary.mjs"; describe("session accessor boundary guard", () => { @@ -78,6 +80,19 @@ describe("session accessor boundary guard", () => { ); }); + it("ratchets only the files migrated by the transcript writer slice", () => { + expect(migratedTranscriptWriterFiles).toEqual( + new Set([ + "src/agents/command/attempt-execution.ts", + "src/agents/embedded-agent-runner/context-engine-maintenance.ts", + "src/config/sessions/transcript.ts", + "src/gateway/server-methods/chat.ts", + "src/gateway/server-methods/chat-transcript-inject.ts", + "src/sessions/user-turn-transcript.ts", + ]), + ); + }); + it("flags legacy reader imports", () => { expect( findSessionAccessorBoundaryViolations(` @@ -178,6 +193,52 @@ describe("session accessor boundary guard", () => { ).toEqual([]); }); + it("flags legacy transcript writer imports", () => { + expect( + findTranscriptWriterBoundaryViolations(` + import { appendSessionTranscriptMessage } from "../config/sessions/transcript-append.js"; + import { emitSessionTranscriptUpdate as emitUpdate } from "../sessions/transcript-events.js"; + import { rewriteTranscriptEntriesInSessionFile } from "../agents/embedded-agent-runner/transcript-rewrite.js"; + `), + ).toEqual([ + { line: 2, reason: 'imports legacy transcript writer "appendSessionTranscriptMessage"' }, + { line: 3, reason: 'imports legacy transcript writer "emitSessionTranscriptUpdate"' }, + { + line: 4, + reason: 'imports legacy transcript writer "rewriteTranscriptEntriesInSessionFile"', + }, + ]); + }); + + it("flags direct and namespace legacy transcript writer calls", () => { + expect( + findTranscriptWriterBoundaryViolations(` + appendSessionTranscriptMessage({ transcriptPath, message }); + transcriptEvents.emitSessionTranscriptUpdate({ sessionFile }); + transcriptAppend["appendSessionTranscriptMessage"]({ transcriptPath, message }); + transcriptRewrite.rewriteTranscriptEntriesInSessionFile({ sessionFile, request }); + `), + ).toEqual([ + { line: 2, reason: 'calls legacy transcript writer "appendSessionTranscriptMessage"' }, + { line: 3, reason: 'references legacy transcript writer "emitSessionTranscriptUpdate"' }, + { line: 4, reason: 'references legacy transcript writer "appendSessionTranscriptMessage"' }, + { + line: 5, + reason: 'references legacy transcript writer "rewriteTranscriptEntriesInSessionFile"', + }, + ]); + }); + + it("allows migrated transcript writer helpers", () => { + expect( + findTranscriptWriterBoundaryViolations(` + import { appendTranscriptMessage, publishTranscriptUpdate } from "../config/sessions/session-accessor.js"; + appendTranscriptMessage(scope, { message }); + publishTranscriptUpdate(scope, { messageId }); + `), + ).toEqual([]); + }); + it("ignores comments and strings that describe legacy readers", () => { expect( findSessionAccessorBoundaryViolations(`