diff --git a/extensions/codex/src/app-server/attempt-context.ts b/extensions/codex/src/app-server/attempt-context.ts index a7903ac17eb1..296a8b7d547c 100644 --- a/extensions/codex/src/app-server/attempt-context.ts +++ b/extensions/codex/src/app-server/attempt-context.ts @@ -78,13 +78,16 @@ type CodexWorkspaceBootstrapContext = CodexBootstrapContext & { }; /** Reads mirrored Codex session history for harness hooks. */ -export async function readMirroredSessionHistoryMessages( - sessionFile: string, -): Promise { - const messages = await readCodexMirroredSessionHistoryMessages(sessionFile); +export async function readMirroredSessionHistoryMessages(params: { + agentId?: string; + sessionFile: string; + sessionId: string; + sessionKey?: string; +}): Promise { + const messages = await readCodexMirroredSessionHistoryMessages(params); if (!messages) { embeddedAgentLog.warn("failed to read mirrored session history for codex harness hooks", { - sessionFile, + sessionFile: params.sessionFile, }); } return messages; diff --git a/extensions/codex/src/app-server/event-projector.ts b/extensions/codex/src/app-server/event-projector.ts index 8a98f5b3a758..06d36c93e00b 100644 --- a/extensions/codex/src/app-server/event-projector.ts +++ b/extensions/codex/src/app-server/event-projector.ts @@ -1827,7 +1827,14 @@ export class CodexAppServerEventProjector { } private async readMirroredSessionMessages(): Promise { - return (await readCodexMirroredSessionHistoryMessages(this.params.sessionFile)) ?? []; + return ( + (await readCodexMirroredSessionHistoryMessages({ + agentId: this.params.agentId, + sessionFile: this.params.sessionFile, + sessionId: this.params.sessionId, + sessionKey: this.params.sessionKey, + })) ?? [] + ); } private createAssistantMessage(text: string): AssistantMessage { diff --git a/extensions/codex/src/app-server/run-attempt.ts b/extensions/codex/src/app-server/run-attempt.ts index 70344a93e5f8..c69f161707a8 100644 --- a/extensions/codex/src/app-server/run-attempt.ts +++ b/extensions/codex/src/app-server/run-attempt.ts @@ -849,7 +849,16 @@ export async function runCodexAppServerAttempt( }, }); const hadSessionFile = await pathExists(activeSessionFile); - let historyMessages = (await readMirroredSessionHistoryMessages(activeSessionFile)) ?? []; + const activeTranscriptTarget = { + agentId: sessionAgentId, + sessionFile: activeSessionFile, + sessionId: activeSessionId, + sessionKey: contextSessionKey, + }; + let historyMessages = + !activeContextEngine && initialStartupBindingHadInactiveThreadBootstrap + ? [] + : ((await readMirroredSessionHistoryMessages(activeTranscriptTarget)) ?? []); const hookContextWindowFields = { ...(params.contextWindowInfo?.tokens ? { contextTokenBudget: params.contextWindowInfo.tokens } @@ -907,7 +916,7 @@ export async function runCodexAppServerAttempt( warn: (message) => embeddedAgentLog.warn(message), }); historyMessages = - (await readMirroredSessionHistoryMessages(activeSessionFile)) ?? historyMessages; + (await readMirroredSessionHistoryMessages(activeTranscriptTarget)) ?? historyMessages; } const memoryToolNames = getCodexWorkspaceMemoryToolNames(toolBridge.availableSpecs); const workspaceBootstrapContext = await buildCodexWorkspaceBootstrapContext({ @@ -3039,7 +3048,7 @@ export async function runCodexAppServerAttempt( const activeContextEnginePluginIdLocal = resolveContextEngineOwnerPluginId(activeContextEngine); const finalMessages = - (await readMirroredSessionHistoryMessages(activeSessionFile)) ?? + (await readMirroredSessionHistoryMessages(activeTranscriptTarget)) ?? historyMessages.concat(result.messagesSnapshot); await finalizeHarnessContextEngineTurn({ contextEngine: activeContextEngine, diff --git a/extensions/codex/src/app-server/session-history.test.ts b/extensions/codex/src/app-server/session-history.test.ts index 0427acdea08e..b9ac6a751c4f 100644 --- a/extensions/codex/src/app-server/session-history.test.ts +++ b/extensions/codex/src/app-server/session-history.test.ts @@ -51,6 +51,14 @@ function messageEntry(params: { }; } +function mirroredTarget(sessionFile: string) { + return { + sessionFile, + sessionId: "codex-session", + sessionKey: "codex-session", + }; +} + describe("readCodexMirroredSessionHistoryMessages", () => { it("replays only the branch selected by a leaf control", async () => { const sessionFile = await writeSession([ @@ -75,7 +83,9 @@ describe("readCodexMirroredSessionHistoryMessages", () => { }, ]); - await expect(readCodexMirroredSessionHistoryMessages(sessionFile)).resolves.toMatchObject([ + await expect( + readCodexMirroredSessionHistoryMessages(mirroredTarget(sessionFile)), + ).resolves.toMatchObject([ { role: "user", content: "root prompt" }, { role: "assistant", content: "active answer" }, ]); @@ -93,7 +103,9 @@ describe("readCodexMirroredSessionHistoryMessages", () => { }, ]); - await expect(readCodexMirroredSessionHistoryMessages(sessionFile)).resolves.toEqual([]); + await expect( + readCodexMirroredSessionHistoryMessages(mirroredTarget(sessionFile)), + ).resolves.toEqual([]); }); it("keeps visible history when continuation rows use a disjoint append cursor", async () => { @@ -125,7 +137,9 @@ describe("readCodexMirroredSessionHistoryMessages", () => { }), ]); - await expect(readCodexMirroredSessionHistoryMessages(sessionFile)).resolves.toMatchObject([ + await expect( + readCodexMirroredSessionHistoryMessages(mirroredTarget(sessionFile)), + ).resolves.toMatchObject([ { role: "user", content: "visible prompt" }, { role: "assistant", content: "continued answer" }, ]); @@ -154,7 +168,9 @@ describe("readCodexMirroredSessionHistoryMessages", () => { }), ]); - await expect(readCodexMirroredSessionHistoryMessages(sessionFile)).resolves.toMatchObject([ + await expect( + readCodexMirroredSessionHistoryMessages(mirroredTarget(sessionFile)), + ).resolves.toMatchObject([ { role: "user", content: "visible prompt" }, { role: "assistant", content: "continued answer" }, ]); diff --git a/extensions/codex/src/app-server/session-history.ts b/extensions/codex/src/app-server/session-history.ts index ddeda238a99e..c634958102a6 100644 --- a/extensions/codex/src/app-server/session-history.ts +++ b/extensions/codex/src/app-server/session-history.ts @@ -10,40 +10,59 @@ import { migrateSessionEntries, parseSessionEntries, } from "openclaw/plugin-sdk/agent-sessions"; +import { + resolveSessionTranscriptTarget, + type SessionTranscriptTargetParams, +} from "openclaw/plugin-sdk/session-transcript-runtime"; import { sanitizeCodexHistoryImagePayloads } from "./image-payload-sanitizer.js"; -function isMissingFileError(error: unknown): boolean { - return Boolean( - error && - typeof error === "object" && - "code" in error && - (error as { code?: unknown }).code === "ENOENT", - ); -} +export type CodexMirroredSessionHistoryTarget = { + agentId?: string; + sessionFile: string; + sessionId: string; + sessionKey?: string; +}; /** Returns sanitized session-context messages for a Codex mirrored session file. */ export async function readCodexMirroredSessionHistoryMessages( - sessionFile: string, + target: CodexMirroredSessionHistoryTarget, ): Promise { try { - const raw = await fs.readFile(sessionFile, "utf-8"); + await resolveSessionTranscriptTarget(resolveCodexHistoryTranscriptTarget(target)); + const raw = await fs.readFile(target.sessionFile, "utf-8"); const entries = parseSessionEntries(raw); + if (entries.length === 0) { + return []; + } const firstEntry = entries[0] as { type?: unknown; id?: unknown } | undefined; if (firstEntry?.type !== "session" || typeof firstEntry.id !== "string") { return undefined; } - migrateSessionEntries(entries); - const sessionEntries = entries.filter( - (entry): entry is SessionEntry => entry.type !== "session", - ); + migrateSessionEntries(entries as SessionEntry[]); + const sessionEntries = entries.filter((entry): entry is SessionEntry => { + return ( + entry !== null && + typeof entry === "object" && + !Array.isArray(entry) && + (entry as { type?: unknown }).type !== "session" + ); + }); return sanitizeCodexHistoryImagePayloads( buildSessionContext(sessionEntries).messages, "codex mirrored history", ); - } catch (error) { - if (isMissingFileError(error)) { - return []; - } + } catch { return undefined; } } + +function resolveCodexHistoryTranscriptTarget( + target: CodexMirroredSessionHistoryTarget, +): SessionTranscriptTargetParams { + return { + ...(target.agentId ? { agentId: target.agentId } : {}), + sessionFile: target.sessionFile, + sessionId: target.sessionId, + sessionKey: target.sessionKey ?? "", + }; +} diff --git a/extensions/codex/src/app-server/transcript-mirror.test.ts b/extensions/codex/src/app-server/transcript-mirror.test.ts index dce6d49f26d1..dfd2f60dc9bd 100644 --- a/extensions/codex/src/app-server/transcript-mirror.test.ts +++ b/extensions/codex/src/app-server/transcript-mirror.test.ts @@ -21,13 +21,14 @@ import { mirrorCodexAppServerTranscript, } from "./transcript-mirror.js"; -const emitSessionTranscriptUpdateMock = vi.hoisted(() => vi.fn()); +const publishSessionTranscriptUpdateByIdentityMock = vi.hoisted(() => vi.fn()); -vi.mock("openclaw/plugin-sdk/agent-harness-runtime", async (importOriginal) => { - const actual = await importOriginal(); +vi.mock("openclaw/plugin-sdk/session-transcript-runtime", async (importOriginal) => { + const actual = + await importOriginal(); return { ...actual, - emitSessionTranscriptUpdate: emitSessionTranscriptUpdateMock, + publishSessionTranscriptUpdateByIdentity: publishSessionTranscriptUpdateByIdentityMock, }; }); @@ -44,7 +45,7 @@ const tempDirs: string[] = []; afterEach(async () => { resetGlobalHookRunner(); - emitSessionTranscriptUpdateMock.mockReset(); + publishSessionTranscriptUpdateByIdentityMock.mockReset(); for (const dir of tempDirs.splice(0)) { await fs.rm(dir, { recursive: true, force: true }); } @@ -130,6 +131,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userMessage, assistantMessage, toolResultMessage], idempotencyScope: "scope-1", @@ -164,30 +166,32 @@ describe("mirrorCodexAppServerTranscript", () => { const firstMirror = await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "agent:main:main", messages: [userMessage], idempotencyScope: "codex-app-server:thread-1", }); const secondMirror = await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "agent:main:main", messages: [userMessage], idempotencyScope: "codex-app-server:thread-1", }); - const updates = emitSessionTranscriptUpdateMock.mock.calls.map( - ([update]) => update as Record, + const updates = publishSessionTranscriptUpdateByIdentityMock.mock.calls.map( + ([update]) => update as Record & { update?: Record }, ); expect(updates).toHaveLength(1); expect(updates[0]?.sessionFile).toBe(sessionFile); expect(updates[0]?.sessionKey).toBe("agent:main:main"); - expect(updates[0]?.messageId).toEqual(expect.any(String)); - expect(updates[0]?.message).toMatchObject({ + expect(updates[0]?.update?.messageId).toEqual(expect.any(String)); + expect(updates[0]?.update?.message).toMatchObject({ role: "user", content: [{ type: "text", text: "show me live" }], idempotencyKey: "codex-app-server:thread-1:turn-1:prompt", }); - expect(updates[0]?.messageSeq).toBe(1); + expect(updates[0]?.update?.messageSeq).toBe(1); expect(firstMirror.userMessagesPresent).toHaveLength(1); expect(firstMirror.userMessagesPresent[0]).toMatchObject({ role: "user", @@ -207,6 +211,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "agent:main:main", messages: [ attachCodexMirrorIdentity( @@ -227,14 +232,16 @@ describe("mirrorCodexAppServerTranscript", () => { idempotencyScope: "codex-app-server:thread-1", }); - const updates = emitSessionTranscriptUpdateMock.mock.calls.map( - ([update]) => update as Record, + const updates = publishSessionTranscriptUpdateByIdentityMock.mock.calls.map( + ([update]) => update as Record & { update?: Record }, ); - expect(updates.map((update) => update.messageSeq)).toEqual([1, 2]); - expect(updates.map((update) => (update.message as { role?: string }).role)).toEqual([ - "user", - "assistant", - ]); + expect(updates.map((update) => update.update?.messageSeq)).toEqual([1, 2]); + expect( + updates.map((update) => { + const message = update.update?.message as { role?: string } | undefined; + return message?.role; + }), + ).toEqual(["user", "assistant"]); }); it("creates the transcript directory on first mirror", async () => { @@ -243,6 +250,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [ makeAgentAssistantMessage({ @@ -273,12 +281,14 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [...messages], idempotencyScope: "scope-1", }); await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [...messages], idempotencyScope: "scope-1", @@ -312,6 +322,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [sourceMessage], idempotencyScope: "scope-1", @@ -348,12 +359,14 @@ describe("mirrorCodexAppServerTranscript", () => { const first = await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [sourceMessage], idempotencyScope: "scope-1", }); const second = await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [sourceMessage], idempotencyScope: "scope-1", @@ -394,6 +407,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [sourceMessage], idempotencyScope: "scope-1", @@ -419,6 +433,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [ makeAgentAssistantMessage({ @@ -456,6 +471,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [ makeAgentAssistantMessage({ @@ -534,6 +550,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userMessage, assistantMessage], idempotencyScope: "codex-app-server:thread-X", @@ -547,6 +564,7 @@ describe("mirrorCodexAppServerTranscript", () => { ); await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userMessage, reasoningMessage, assistantMessage], idempotencyScope: "codex-app-server:thread-X", @@ -595,12 +613,14 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userTurn1, assistantTurn1], idempotencyScope: "codex-app-server:thread-X", }); await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userTurn2, assistantTurn2], idempotencyScope: "codex-app-server:thread-X", @@ -638,6 +658,7 @@ describe("mirrorCodexAppServerTranscript", () => { ); await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userTurn1, assistantTurn1], idempotencyScope: "codex-app-server:thread-X", @@ -661,6 +682,7 @@ describe("mirrorCodexAppServerTranscript", () => { // turn 1's entries (with their original identities preserved). await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userTurn1, assistantTurn1, userTurn2, assistantTurn2], idempotencyScope: "codex-app-server:thread-X", @@ -691,6 +713,7 @@ describe("mirrorCodexAppServerTranscript", () => { await mirrorCodexAppServerTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userMessage, assistantMessage], idempotencyScope: "scope-1", diff --git a/extensions/codex/src/app-server/transcript-mirror.ts b/extensions/codex/src/app-server/transcript-mirror.ts index 1973687ffa5e..dee1610e2800 100644 --- a/extensions/codex/src/app-server/transcript-mirror.ts +++ b/extensions/codex/src/app-server/transcript-mirror.ts @@ -1,19 +1,19 @@ // Codex plugin module implements transcript mirror behavior. import { createHash } from "node:crypto"; -import fs from "node:fs/promises"; import { - acquireSessionWriteLock, - appendSessionTranscriptMessage, embeddedAgentLog, - emitSessionTranscriptUpdate, formatErrorMessage, - resolveSessionWriteLockOptions, runAgentHarnessBeforeMessageWriteHook, type AgentMessage, type EmbeddedRunAttemptParams, type EmbeddedRunAttemptResult, - type SessionWriteLockAcquireTimeoutConfig, } from "openclaw/plugin-sdk/agent-harness-runtime"; +import { + publishSessionTranscriptUpdateByIdentity, + withSessionTranscriptWriteLock, + type SessionTranscriptTargetParams, + type SessionTranscriptWriteLockParams, +} from "openclaw/plugin-sdk/session-transcript-runtime"; import { normalizeOptionalString } from "openclaw/plugin-sdk/string-coerce-runtime"; type MirroredAgentMessage = Extract; @@ -273,13 +273,13 @@ function buildMirrorDedupeIdentity(message: MirroredAgentMessage): string { export async function mirrorCodexAppServerTranscript(params: { sessionFile: string; - sessionId?: string; + sessionId: string; cwd?: string; sessionKey?: string; agentId?: string; messages: AgentMessage[]; idempotencyScope?: string; - config?: SessionWriteLockAcquireTimeoutConfig; + config?: SessionTranscriptWriteLockParams["config"]; }): Promise { const messages = params.messages.filter( (message): message is MirroredAgentMessage => @@ -289,129 +289,133 @@ export async function mirrorCodexAppServerTranscript(params: { return { userMessagesPresent: [] }; } - const lock = await acquireSessionWriteLock({ - sessionFile: params.sessionFile, - ...resolveSessionWriteLockOptions(params.config), - }); - const appendedUpdates: Array<{ messageId: string; message: AgentMessage; messageSeq: number }> = - []; - const userMessagesPresent: MirroredUserMessage[] = []; - try { - const mirrorState = await readTranscriptMirrorState(params.sessionFile); - let nextMessageSeq = mirrorState.messageCount; - for (const message of messages) { - const dedupeIdentity = buildMirrorDedupeIdentity(message); - const idempotencyKey = params.idempotencyScope - ? `${params.idempotencyScope}:${dedupeIdentity}` - : undefined; - const transcriptMessage = { - ...message, - ...(idempotencyKey ? { idempotencyKey } : {}), - } as AgentMessage; - if (idempotencyKey && mirrorState.idempotencyKeys.has(idempotencyKey)) { - const persistedUserMessage = mirrorState.userMessagesByIdempotencyKey.get(idempotencyKey); - if (persistedUserMessage) { - userMessagesPresent.push(persistedUserMessage); + const transcriptTarget = resolveCodexMirrorTranscriptTarget(params); + const { appendedUpdates, userMessagesPresent } = await withSessionTranscriptWriteLock( + { ...transcriptTarget, config: params.config }, + async (transcript) => { + const nextAppendedUpdates: Array<{ + messageId: string; + message: AgentMessage; + messageSeq: number; + }> = []; + const nextUserMessagesPresent: MirroredUserMessage[] = []; + const mirrorState = readTranscriptMirrorState(await transcript.readEvents()); + let nextMessageSeq = mirrorState.messageCount; + for (const message of messages) { + const dedupeIdentity = buildMirrorDedupeIdentity(message); + const idempotencyKey = params.idempotencyScope + ? `${params.idempotencyScope}:${dedupeIdentity}` + : undefined; + const transcriptMessage = { + ...message, + ...(idempotencyKey ? { idempotencyKey } : {}), + } as AgentMessage; + if (idempotencyKey && mirrorState.idempotencyKeys.has(idempotencyKey)) { + const persistedUserMessage = mirrorState.userMessagesByIdempotencyKey.get(idempotencyKey); + if (persistedUserMessage) { + nextUserMessagesPresent.push(persistedUserMessage); + } + continue; } - continue; - } - const nextMessage = runAgentHarnessBeforeMessageWriteHook({ - message: transcriptMessage, - agentId: params.agentId, - sessionKey: params.sessionKey, - }); - if (!nextMessage) { - continue; - } - const messageToAppend = ( - idempotencyKey - ? { - ...(nextMessage as unknown as Record), - idempotencyKey, - } - : nextMessage - ) as AgentMessage; - const { messageId, message: appendedMessage } = await appendSessionTranscriptMessage({ - transcriptPath: params.sessionFile, - message: messageToAppend, - idempotencyLookup: idempotencyKey ? "caller-checked" : "scan", - sessionId: params.sessionId, - cwd: params.cwd, - config: params.config, - }); - if (appendedMessage.role === "user") { - userMessagesPresent.push(appendedMessage); + const nextMessage = runAgentHarnessBeforeMessageWriteHook({ + message: transcriptMessage, + agentId: params.agentId, + sessionKey: params.sessionKey, + }); + if (!nextMessage) { + continue; + } + const messageToAppend = ( + idempotencyKey + ? { + ...(nextMessage as unknown as Record), + idempotencyKey, + } + : nextMessage + ) as AgentMessage; + const appended = await transcript.appendMessage({ + message: messageToAppend, + idempotencyLookup: idempotencyKey ? "caller-checked" : "scan", + cwd: params.cwd, + }); + if (!appended) { + continue; + } + const { messageId, message: appendedMessage } = appended; + if (appendedMessage.role === "user") { + nextUserMessagesPresent.push(appendedMessage); + if (idempotencyKey) { + mirrorState.userMessagesByIdempotencyKey.set(idempotencyKey, appendedMessage); + } + } + nextMessageSeq += 1; + nextAppendedUpdates.push({ + messageId, + message: appendedMessage, + messageSeq: nextMessageSeq, + }); if (idempotencyKey) { - mirrorState.userMessagesByIdempotencyKey.set(idempotencyKey, appendedMessage); + mirrorState.idempotencyKeys.add(idempotencyKey); } } - nextMessageSeq += 1; - appendedUpdates.push({ messageId, message: appendedMessage, messageSeq: nextMessageSeq }); - if (idempotencyKey) { - mirrorState.idempotencyKeys.add(idempotencyKey); - } - } - } finally { - await lock.release(); - } + return { appendedUpdates: nextAppendedUpdates, userMessagesPresent: nextUserMessagesPresent }; + }, + ); for (const update of appendedUpdates) { - emitSessionTranscriptUpdate({ - sessionFile: params.sessionFile, - ...(params.sessionKey ? { sessionKey: params.sessionKey } : {}), - ...(params.agentId ? { agentId: params.agentId } : {}), - ...(params.sessionId && params.sessionKey && params.agentId - ? { - target: { - agentId: params.agentId, - sessionId: params.sessionId, - sessionKey: params.sessionKey, - }, - } - : {}), - message: update.message, - messageId: update.messageId, - messageSeq: update.messageSeq, + await publishSessionTranscriptUpdateByIdentity({ + ...transcriptTarget, + update: { + ...(params.sessionKey ? { sessionKey: params.sessionKey } : {}), + ...(params.agentId ? { agentId: params.agentId } : {}), + message: update.message, + messageId: update.messageId, + messageSeq: update.messageSeq, + }, }); } return { userMessagesPresent }; } -async function readTranscriptMirrorState(sessionFile: string): Promise<{ +function resolveCodexMirrorTranscriptTarget(params: { + agentId?: string; + sessionFile: string; + sessionId: string; + sessionKey?: string; +}): SessionTranscriptTargetParams { + return { + ...(params.agentId ? { agentId: params.agentId } : {}), + sessionFile: params.sessionFile, + sessionId: params.sessionId, + sessionKey: params.sessionKey ?? "", + }; +} + +function readTranscriptMirrorState(events: unknown[]): { idempotencyKeys: Set; messageCount: number; userMessagesByIdempotencyKey: Map; -}> { +} { const idempotencyKeys = new Set(); const userMessagesByIdempotencyKey = new Map(); let messageCount = 0; - let raw: string; - try { - raw = await fs.readFile(sessionFile, "utf8"); - } catch (error) { - if ((error as NodeJS.ErrnoException).code !== "ENOENT") { - throw error; - } - return { idempotencyKeys, messageCount, userMessagesByIdempotencyKey }; - } - for (const line of raw.split(/\r?\n/)) { - if (!line.trim()) { + for (const event of events) { + if (!event || typeof event !== "object" || Array.isArray(event)) { continue; } - try { - const parsed = JSON.parse(line) as { message?: AgentMessage & { idempotencyKey?: unknown } }; - if ((parsed as { type?: unknown }).type === "message") { - messageCount += 1; + const parsed = event as { + message?: AgentMessage & { idempotencyKey?: unknown }; + type?: unknown; + }; + if (parsed.type === "message") { + messageCount += 1; + } + if (typeof parsed.message?.idempotencyKey === "string") { + idempotencyKeys.add(parsed.message.idempotencyKey); + if (parsed.message.role === "user") { + userMessagesByIdempotencyKey.set(parsed.message.idempotencyKey, parsed.message); } - if (typeof parsed.message?.idempotencyKey === "string") { - idempotencyKeys.add(parsed.message.idempotencyKey); - if (parsed.message.role === "user") { - userMessagesByIdempotencyKey.set(parsed.message.idempotencyKey, parsed.message); - } - } - } catch { - continue; } } return { idempotencyKeys, messageCount, userMessagesByIdempotencyKey }; diff --git a/extensions/copilot/src/attempt.test.ts b/extensions/copilot/src/attempt.test.ts index 5bede437ee05..8a841c93b788 100644 --- a/extensions/copilot/src/attempt.test.ts +++ b/extensions/copilot/src/attempt.test.ts @@ -2461,11 +2461,13 @@ describe("runCopilotAttempt", () => { expect(dualWriteMock.dualWriteCopilotTranscriptBestEffort).toHaveBeenCalledTimes(1); const args = dualWriteMock.dualWriteCopilotTranscriptBestEffort.mock.calls[0]?.[0] as { sessionFile: string; + sessionId: string; messages: Array<{ role: string }>; idempotencyScope?: string; }; expect(args.sessionFile).toBe("session.json"); - expect(args.idempotencyScope).toMatch(/^copilot:/u); + expect(args.sessionId).toBe("session-1"); + expect(args.idempotencyScope).toBe("copilot:sess-1"); expect(args.messages.length).toBeGreaterThan(0); const roles = args.messages.map((m) => m.role); expect(roles).toContain("user"); @@ -2512,10 +2514,9 @@ describe("runCopilotAttempt", () => { } const identity = message["__openclaw"]?.mirrorIdentity ?? ""; // The terminal assistant carries the turn-stable - // `${runId}:assistant:final` identity attached by attempt.ts - // (rubber-duck-validated identity scheme — survives SDK session - // reuse across turns). Caller-passed history without an - // identity falls through to the positional `${scope}:role:idx` + // `${runId}:assistant:final` identity attached by attempt.ts. + // Caller-passed history without an identity falls through to + // the positional `${scope}:role:idx`. // fingerprint that the existing tagging map applies. if (message.role === "assistant" && index === args.messages.length - 1) { expect(identity).toMatch(/:assistant:final$/u); diff --git a/extensions/copilot/src/attempt.ts b/extensions/copilot/src/attempt.ts index 343f7d65bebf..908a2469c55f 100644 --- a/extensions/copilot/src/attempt.ts +++ b/extensions/copilot/src/attempt.ts @@ -1005,8 +1005,9 @@ export async function runCopilotAttempt( // extension. Identity-tagged so re-emits dedupe. Errors are // swallowed so a mirror failure cannot break the attempt. const sessionFileForMirror = readString(input.sessionFile); - const sessionIdForScope = sessionIdUsed ?? readString(input.sessionId); - if (sessionFileForMirror && messagesSnapshot.length > 0) { + const openClawSessionIdForMirror = readString(input.sessionId); + const mirrorScopeSessionId = sessionIdUsed ?? openClawSessionIdForMirror; + if (sessionFileForMirror && openClawSessionIdForMirror && messagesSnapshot.length > 0) { const taggedMessages = messagesSnapshot.map((message, index) => { if ( message.role !== "user" && @@ -1027,16 +1028,16 @@ export async function runCopilotAttempt( if (hasMirrorIdentity(message)) { return message; } - const identityScope = sdkSessionId ?? sessionIdForScope ?? "attempt"; + const identityScope = sdkSessionId ?? mirrorScopeSessionId ?? "attempt"; return attachCopilotMirrorIdentity(message, `${identityScope}:${message.role}:${index}`); }); await dualWriteCopilotTranscriptBestEffort({ sessionFile: sessionFileForMirror, + sessionId: openClawSessionIdForMirror, sessionKey: readString((input as { sessionKey?: unknown }).sessionKey), - sessionId: readString(input.sessionId), agentId: readString(input.agentId), messages: taggedMessages, - idempotencyScope: sessionIdForScope ? `copilot:${sessionIdForScope}` : undefined, + idempotencyScope: mirrorScopeSessionId ? `copilot:${mirrorScopeSessionId}` : undefined, config: (input as { config?: unknown }).config as never, }).catch((mirrorError: unknown) => { // Defense-in-depth: the best-effort wrapper already swallows diff --git a/extensions/copilot/src/dual-write-transcripts.test.ts b/extensions/copilot/src/dual-write-transcripts.test.ts index b75fff56bfcb..fdae67dbfe20 100755 --- a/extensions/copilot/src/dual-write-transcripts.test.ts +++ b/extensions/copilot/src/dual-write-transcripts.test.ts @@ -86,6 +86,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [userMessage, assistantMessage, toolResultMessage], idempotencyScope: "copilot:session-1", @@ -113,6 +114,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [ makeAgentAssistantMessage({ @@ -143,12 +145,14 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [...messages], idempotencyScope: "copilot:session-1", }); await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [...messages], idempotencyScope: "copilot:session-1", @@ -185,6 +189,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [sourceMessage], idempotencyScope: "copilot:session-1", @@ -210,6 +215,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [ makeAgentAssistantMessage({ @@ -228,6 +234,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", sessionKey: "session-1", messages: [], idempotencyScope: "copilot:session-1", @@ -245,6 +252,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", messages: [message], idempotencyScope: "scope-fp", }); @@ -263,6 +271,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", messages: [tagged], idempotencyScope: "copilot:openclaw-session-1", }); @@ -279,6 +288,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", messages: [ makeAgentAssistantMessage({ content: [{ type: "text", text: "no scope" }], @@ -306,6 +316,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", messages: [userMessage, systemLike], idempotencyScope: "scope", }); @@ -326,6 +337,7 @@ describe("mirrorCopilotTranscript", () => { await mirrorCopilotTranscript({ sessionFile, + sessionId: "session-1", messages: [second], idempotencyScope: "scope", }); @@ -342,6 +354,7 @@ describe("dualWriteCopilotTranscriptBestEffort", () => { await expect( dualWriteCopilotTranscriptBestEffort({ sessionFile, + sessionId: "session-1", messages: [ makeAgentAssistantMessage({ content: [{ type: "text", text: "ok" }], @@ -356,22 +369,34 @@ describe("dualWriteCopilotTranscriptBestEffort", () => { }); it("swallows infrastructure failures and never rejects", async () => { - // Pointing sessionFile at a path under a non-existent root with an - // empty-string segment can fail differently on different platforms; - // instead force failure by passing an invalid type and asserting - // that the wrapper itself does not reject. Use any-cast for the - // bad input shape since we are testing the wrapper's catch. - await expect( - dualWriteCopilotTranscriptBestEffort({ - sessionFile: "" as unknown as string, - messages: [ - makeAgentAssistantMessage({ - content: [{ type: "text", text: "should-not-throw" }], - timestamp: Date.now(), - }), - ], - idempotencyScope: "scope", - }), - ).resolves.toBeUndefined(); + const root = await makeRoot("openclaw-copilot-mirror-invalid-"); + const previousStateDir = process.env.OPENCLAW_STATE_DIR; + process.env.OPENCLAW_STATE_DIR = root; + try { + await expect( + dualWriteCopilotTranscriptBestEffort({ + agentId: "main", + sessionFile: "", + sessionId: "session-1", + sessionKey: "agent:main:session-1", + messages: [ + makeAgentAssistantMessage({ + content: [{ type: "text", text: "should-not-throw" }], + timestamp: Date.now(), + }), + ], + idempotencyScope: "scope", + }), + ).resolves.toBeUndefined(); + await expect( + fs.access(path.join(root, "agents", "main", "sessions", "session-1.jsonl")), + ).rejects.toHaveProperty("code", "ENOENT"); + } finally { + if (previousStateDir === undefined) { + delete process.env.OPENCLAW_STATE_DIR; + } else { + process.env.OPENCLAW_STATE_DIR = previousStateDir; + } + } }); }); diff --git a/extensions/copilot/src/dual-write-transcripts.ts b/extensions/copilot/src/dual-write-transcripts.ts index 72c4d6d745ea..b4dee5975c52 100755 --- a/extensions/copilot/src/dual-write-transcripts.ts +++ b/extensions/copilot/src/dual-write-transcripts.ts @@ -29,16 +29,16 @@ */ import { createHash } from "node:crypto"; -import fs from "node:fs/promises"; import { - acquireSessionWriteLock, - appendSessionTranscriptMessage, - emitSessionTranscriptUpdate, - resolveSessionWriteLockAcquireTimeoutMs, runAgentHarnessBeforeMessageWriteHook, type AgentMessage, - type SessionWriteLockAcquireTimeoutConfig, } from "openclaw/plugin-sdk/agent-harness-runtime"; +import { + publishSessionTranscriptUpdateByIdentity, + withSessionTranscriptWriteLock, + type SessionTranscriptTargetParams, + type SessionTranscriptWriteLockParams, +} from "openclaw/plugin-sdk/session-transcript-runtime"; type MirroredAgentMessage = Extract; @@ -95,8 +95,8 @@ function buildMirrorDedupeIdentity(message: MirroredAgentMessage): string { export interface MirrorCopilotTranscriptParams { sessionFile: string; + sessionId: string; sessionKey?: string; - sessionId?: string; agentId?: string; messages: AgentMessage[]; /** @@ -107,7 +107,7 @@ export interface MirrorCopilotTranscriptParams { * entry collide with its existing on-disk key and be a true no-op. */ idempotencyScope?: string; - config?: SessionWriteLockAcquireTimeoutConfig; + config?: SessionTranscriptWriteLockParams["config"]; } export async function mirrorCopilotTranscript( @@ -121,95 +121,91 @@ export async function mirrorCopilotTranscript( return; } - const lock = await acquireSessionWriteLock({ - sessionFile: params.sessionFile, - timeoutMs: resolveSessionWriteLockAcquireTimeoutMs(params.config), - }); - try { - const existingIdempotencyKeys = await readTranscriptIdempotencyKeys(params.sessionFile); - for (const message of messages) { - const dedupeIdentity = buildMirrorDedupeIdentity(message); - const idempotencyKey = params.idempotencyScope - ? `${params.idempotencyScope}:${dedupeIdentity}` - : undefined; - if (idempotencyKey && existingIdempotencyKeys.has(idempotencyKey)) { - continue; + const transcriptTarget = resolveCopilotMirrorTranscriptTarget(params); + const didAppend = await withSessionTranscriptWriteLock( + { ...transcriptTarget, config: params.config }, + async (transcript) => { + let didAppendMessage = false; + const existingIdempotencyKeys = readTranscriptIdempotencyKeys(await transcript.readEvents()); + for (const message of messages) { + const dedupeIdentity = buildMirrorDedupeIdentity(message); + const idempotencyKey = params.idempotencyScope + ? `${params.idempotencyScope}:${dedupeIdentity}` + : undefined; + if (idempotencyKey && existingIdempotencyKeys.has(idempotencyKey)) { + continue; + } + const transcriptMessage = { + ...message, + ...(idempotencyKey ? { idempotencyKey } : {}), + } as AgentMessage; + const nextMessage = runAgentHarnessBeforeMessageWriteHook({ + message: transcriptMessage, + agentId: params.agentId, + sessionKey: params.sessionKey, + }); + if (!nextMessage) { + continue; + } + const messageToAppend = ( + idempotencyKey + ? { + ...(nextMessage as unknown as Record), + idempotencyKey, + } + : nextMessage + ) as AgentMessage; + const appended = await transcript.appendMessage({ + message: messageToAppend, + idempotencyLookup: idempotencyKey ? "caller-checked" : "scan", + }); + if (!appended) { + continue; + } + didAppendMessage = true; + if (idempotencyKey) { + existingIdempotencyKeys.add(idempotencyKey); + } } - const transcriptMessage = { - ...message, - ...(idempotencyKey ? { idempotencyKey } : {}), - } as AgentMessage; - const nextMessage = runAgentHarnessBeforeMessageWriteHook({ - message: transcriptMessage, - agentId: params.agentId, - sessionKey: params.sessionKey, - }); - if (!nextMessage) { - continue; - } - const messageToAppend = ( - idempotencyKey - ? { - ...(nextMessage as unknown as Record), - idempotencyKey, - } - : nextMessage - ) as AgentMessage; - await appendSessionTranscriptMessage({ - transcriptPath: params.sessionFile, - message: messageToAppend, - config: params.config, - }); - if (idempotencyKey) { - existingIdempotencyKeys.add(idempotencyKey); - } - } - } finally { - await lock.release(); - } + return didAppendMessage; + }, + ); - if (params.sessionKey) { - emitSessionTranscriptUpdate({ - sessionFile: params.sessionFile, - sessionKey: params.sessionKey, - ...(params.agentId ? { agentId: params.agentId } : {}), - ...(params.sessionId && params.agentId - ? { - target: { - agentId: params.agentId, - sessionId: params.sessionId, - sessionKey: params.sessionKey, - }, - } - : {}), + if (didAppend) { + await publishSessionTranscriptUpdateByIdentity({ + ...transcriptTarget, + update: params.sessionKey ? { sessionKey: params.sessionKey } : undefined, }); - } else { - emitSessionTranscriptUpdate(params.sessionFile); } } -async function readTranscriptIdempotencyKeys(sessionFile: string): Promise> { - const keys = new Set(); - let raw: string; - try { - raw = await fs.readFile(sessionFile, "utf8"); - } catch (error) { - if ((error as NodeJS.ErrnoException).code !== "ENOENT") { - throw error; - } - return keys; +function resolveCopilotMirrorTranscriptTarget(params: { + agentId?: string; + sessionFile: string; + sessionId: string; + sessionKey?: string; +}): SessionTranscriptTargetParams { + const sessionFile = params.sessionFile.trim(); + if (!sessionFile) { + throw new Error("Copilot transcript mirror requires a sessionFile target"); } - for (const line of raw.split(/\r?\n/)) { - if (!line.trim()) { + return { + ...(params.agentId ? { agentId: params.agentId } : {}), + sessionFile, + sessionId: params.sessionId, + sessionKey: params.sessionKey ?? "", + }; +} + +function readTranscriptIdempotencyKeys(events: unknown[]): Set { + const keys = new Set(); + for (const event of events) { + if (!event || typeof event !== "object" || Array.isArray(event)) { continue; } - try { - const parsed = JSON.parse(line) as { message?: { idempotencyKey?: unknown } }; - if (typeof parsed.message?.idempotencyKey === "string") { - keys.add(parsed.message.idempotencyKey); - } - } catch { - continue; + const parsed = event as { message?: { idempotencyKey?: unknown } }; + if (typeof parsed.message?.idempotencyKey === "string") { + keys.add(parsed.message.idempotencyKey); } } return keys;