From 9512294e8f629cc7e2ab0df01f593c6cc662c8dd Mon Sep 17 00:00:00 2001 From: Josh Lehman Date: Tue, 23 Jun 2026 16:14:53 -0700 Subject: [PATCH] fix: bridge ACP metadata to session accessors (#96195) * fix: bridge ACP metadata to session accessors * fix: simplify ACP accessor key ownership * fix: bind ACP metadata after session canonicalization --- scripts/check-session-accessor-boundary.mjs | 5 + src/acp/runtime/session-meta.test.ts | 117 ++++++++++++ src/acp/runtime/session-meta.ts | 175 ++++++++++-------- src/agents/acp-spawn.test.ts | 65 +++++++ src/agents/acp-spawn.ts | 56 +++--- src/config/sessions/session-accessor.ts | 42 +++++ src/config/sessions/store.ts | 56 ++++-- .../check-session-accessor-boundary.test.ts | 3 + 8 files changed, 402 insertions(+), 117 deletions(-) diff --git a/scripts/check-session-accessor-boundary.mjs b/scripts/check-session-accessor-boundary.mjs index d427c8b04a6b..ba6c5ddf442e 100644 --- a/scripts/check-session-accessor-boundary.mjs +++ b/scripts/check-session-accessor-boundary.mjs @@ -74,6 +74,8 @@ export const allowedSessionStoreRuntimeFileBackedCompatExports = new Set([ export const migratedSessionAccessorFiles = new Set([ "packages/memory-host-sdk/src/host/session-files.ts", + "src/acp/runtime/session-meta.ts", + "src/agents/acp-spawn.ts", "src/agents/embedded-agent-runner/compaction-successor-transcript.ts", "src/agents/embedded-agent-runner/run/attempt.ts", "src/agents/embedded-agent-runner/tool-result-truncation.ts", @@ -124,6 +126,7 @@ export const migratedBundledPluginSessionAccessorFiles = new Set([ ]); export const migratedSessionAccessorWriteFiles = new Set([ + "src/acp/runtime/session-meta.ts", "src/agents/command/attempt-execution.shared.ts", "src/agents/command/session-store.ts", "src/agents/embedded-agent-runner/run.ts", @@ -535,6 +538,7 @@ export async function main() { "extensions/discord/src/monitor", "extensions/memory-core/src", "extensions/telegram/src", + "src/acp", "src/agents", "src/auto-reply", "src/commands", @@ -546,6 +550,7 @@ export async function main() { "src/tui", ]); const writeSourceRoots = resolveSourceRoots(repoRoot, [ + "src/acp", "src/agents", "src/auto-reply", "src/commands", diff --git a/src/acp/runtime/session-meta.test.ts b/src/acp/runtime/session-meta.test.ts index 31ad51851253..cdf1e6166a2f 100644 --- a/src/acp/runtime/session-meta.test.ts +++ b/src/acp/runtime/session-meta.test.ts @@ -152,6 +152,24 @@ describe("ACP session metadata SQLite store", () => { })?.acp?.runtimeSessionName, ).toBe("codex-normalized"); expect(loadSessionStore(storePath)[storeSessionKey]?.acp).toBeUndefined(); + const legacyEmbeddedEntry = loadSessionStore(storePath)[storeSessionKey]; + expect(legacyEmbeddedEntry).toBeDefined(); + if (!legacyEmbeddedEntry) { + throw new Error("expected normalized ACP session entry"); + } + await writeSessionStoreForTestAsync(storePath, { + [storeSessionKey]: { + ...legacyEmbeddedEntry, + acp: { + backend: "acpx", + agent: "codex", + runtimeSessionName: "legacy-embedded", + mode: "persistent", + state: "idle", + lastActivityAt: 120, + }, + }, + }); await upsertAcpSessionMeta({ cfg, @@ -170,6 +188,105 @@ describe("ACP session metadata SQLite store", () => { sessionKey: storeSessionKey, })?.acp, ).toBeUndefined(); + expect(loadSessionStore(storePath)[storeSessionKey]?.acp).toBeUndefined(); + }); + }); + + it("keeps SQLite ACP metadata visible when legacy store keys are canonicalized", async () => { + await withTempDir({ prefix: "openclaw-acp-meta-" }, async (dir) => { + const storePath = path.join(dir, "sessions.json"); + const databasePath = path.join(dir, "state", "openclaw.sqlite"); + const cfg = { session: { store: storePath } } as OpenClawConfig; + const legacyStoreSessionKey = "agent:CODEX:acp:legacy-runtime"; + const canonicalSessionKey = "agent:codex:acp:legacy-runtime"; + await fs.writeFile( + storePath, + JSON.stringify({ + [legacyStoreSessionKey]: { + sessionId: "sess-acp", + updatedAt: 100, + }, + }), + "utf8", + ); + + await upsertAcpSessionMeta({ + cfg, + databasePath, + sessionKey: canonicalSessionKey, + now: () => 200, + mutate: () => ({ + backend: "acpx", + agent: "codex", + runtimeSessionName: "codex-canonicalized", + mode: "persistent", + state: "idle", + lastActivityAt: 123, + }), + }); + + const store = loadSessionStore(storePath); + expect(store[legacyStoreSessionKey]).toBeUndefined(); + expect(store[canonicalSessionKey]?.sessionId).toBe("sess-acp"); + expect( + readAcpSessionEntry({ + cfg, + databasePath, + sessionKey: canonicalSessionKey, + })?.acp?.runtimeSessionName, + ).toBe("codex-canonicalized"); + expect(await listAcpSessionEntries({ cfg, databasePath })).toHaveLength(1); + }); + }); + + it("binds ACP metadata to the final accessor-selected entry for alias writes", async () => { + await withTempDir({ prefix: "openclaw-acp-meta-" }, async (dir) => { + const storePath = path.join(dir, "sessions.json"); + const databasePath = path.join(dir, "state", "openclaw.sqlite"); + const cfg = { session: { store: storePath } } as OpenClawConfig; + const canonicalSessionKey = "agent:codex:acp:alias-runtime"; + const legacyStoreSessionKey = "agent:CODEX:acp:alias-runtime"; + await fs.writeFile( + storePath, + JSON.stringify({ + [canonicalSessionKey]: { + sessionId: "sess-canonical", + updatedAt: 100, + }, + [legacyStoreSessionKey]: { + sessionId: "sess-legacy", + updatedAt: 150, + }, + }), + "utf8", + ); + + await upsertAcpSessionMeta({ + cfg, + databasePath, + sessionKey: legacyStoreSessionKey, + now: () => 200, + mutate: () => ({ + backend: "acpx", + agent: "codex", + runtimeSessionName: "codex-alias", + mode: "persistent", + state: "idle", + lastActivityAt: 123, + }), + }); + + const store = loadSessionStore(storePath); + expect(store[legacyStoreSessionKey]).toBeUndefined(); + expect(store[canonicalSessionKey]?.sessionId).toBe("sess-legacy"); + expect( + readAcpSessionEntry({ + cfg, + databasePath, + sessionKey: canonicalSessionKey, + })?.acp?.runtimeSessionName, + ).toBe("codex-alias"); + expect(await listAcpSessionEntries({ cfg, databasePath })).toHaveLength(1); }); }); diff --git a/src/acp/runtime/session-meta.ts b/src/acp/runtime/session-meta.ts index 8d79675c404a..f28b0c61634d 100644 --- a/src/acp/runtime/session-meta.ts +++ b/src/acp/runtime/session-meta.ts @@ -4,7 +4,11 @@ import { normalizeLowercaseStringOrEmpty } from "@openclaw/normalization-core/st import type { Insertable, Selectable } from "kysely"; import { getRuntimeConfig } from "../../config/config.js"; import { resolveStorePath } from "../../config/sessions/paths.js"; -import { loadSessionStore } from "../../config/sessions/store-load.js"; +import { + listSessionEntries, + patchSessionEntryWithKey, + type SessionEntrySummary, +} from "../../config/sessions/session-accessor.js"; import { mergeSessionEntry, type AcpSessionRuntimeOptions, @@ -43,30 +47,24 @@ type AcpSessionsTable = OpenClawStateKyselyDatabase["acp_sessions"]; type AcpSessionMetaDatabase = Pick; type AcpSessionRow = Selectable; -let sessionStoreRuntimePromise: - | Promise - | undefined; - -function loadSessionStoreRuntime() { - sessionStoreRuntimePromise ??= import("../../config/sessions/store.runtime.js"); - return sessionStoreRuntimePromise; -} - -function resolveStoreSessionKey(store: Record, sessionKey: string): string { +function resolveStoreSessionKey( + entries: readonly SessionEntrySummary[], + sessionKey: string, +): string { const normalized = sessionKey.trim(); if (!normalized) { return ""; } - if (store[normalized]) { + if (entries.some((entry) => entry.sessionKey === normalized)) { return normalized; } const lower = normalizeLowercaseStringOrEmpty(normalized); - if (store[lower]) { + if (entries.some((entry) => entry.sessionKey === lower)) { return lower; } - for (const key of Object.keys(store)) { - if (normalizeLowercaseStringOrEmpty(key) === lower) { - return key; + for (const entry of entries) { + if (normalizeLowercaseStringOrEmpty(entry.sessionKey) === lower) { + return entry.sessionKey; } } return lower; @@ -371,12 +369,13 @@ function readSessionEntryFromStore(params: { env: params.env, }); try { - const store = loadSessionStore( + const entries = listSessionEntries({ storePath, - params.clone === false ? { clone: false } : undefined, - ); - const storeSessionKey = resolveStoreSessionKey(store, params.sessionKey); - return { cfg, storePath, storeSessionKey, entry: store[storeSessionKey] }; + ...(params.clone === false ? { clone: false } : {}), + }); + const storeSessionKey = resolveStoreSessionKey(entries, params.sessionKey); + const entry = entries.find((candidate) => candidate.sessionKey === storeSessionKey)?.entry; + return { cfg, storePath, storeSessionKey, entry }; } catch { return { cfg, @@ -437,14 +436,19 @@ export async function listAcpSessionEntries(params: { cfg, env: params.env, }); - let store: Record; + let sessionEntries: SessionEntrySummary[]; try { - store = loadSessionStore(storePath, params.clone === false ? { clone: false } : undefined); + sessionEntries = listSessionEntries({ + storePath, + ...(params.clone === false ? { clone: false } : {}), + }); } catch { continue; } - const storeSessionKey = resolveStoreSessionKey(store, sessionKey); - const entry = store[storeSessionKey]; + const storeSessionKey = resolveStoreSessionKey(sessionEntries, sessionKey); + const entry = sessionEntries.find( + (candidate) => candidate.sessionKey === storeSessionKey, + )?.entry; if (!entry || !acpSessionRowMatchesEntry(row, entry)) { continue; } @@ -518,63 +522,86 @@ export async function upsertAcpSessionMeta(params: { current, current ? mergeAcpForReturn(preparedEntry, current) : entry, ); - if (nextMeta === null) { + }, + { env: params.env, path: params.databasePath }, + ); + const metaToPersist = nextMeta; + if (metaToPersist === undefined) { + return current ? mergeAcpForReturn(entry, current) : (entry ?? null); + } + if (metaToPersist === null) { + const patched = entry + ? await patchSessionEntryWithKey( + { storePath: storeEntry.storePath, sessionKey: storageSessionKey }, + (currentEntry) => { + const next = { ...currentEntry }; + delete next.acp; + return next; + }, + { + ...sessionStoreUpdateOptions({ ...params, sessionKey: storageSessionKey }), + replaceEntry: true, + }, + ) + : null; + runOpenClawStateWriteTransaction( + (database) => { + const sessionKeysToDelete = new Set([storageSessionKey]); + if (patched?.sessionKey) { + sessionKeysToDelete.add(patched.sessionKey); + } + for (const key of sessionKeysToDelete) { + executeSqliteQuerySync( + database.db, + getAcpSessionKysely(database.db) + .deleteFrom("acp_sessions") + .where("session_key", "=", key), + ); + } + }, + { env: params.env, path: params.databasePath }, + ); + return patched?.entry ?? null; + } + const persisted = await patchSessionEntryWithKey( + { storePath: storeEntry.storePath, sessionKey: storageSessionKey }, + (currentEntry) => { + const next = mergeSessionEntry(currentEntry, { + updatedAt, + }); + delete next.acp; + return next; + }, + { + ...sessionStoreUpdateOptions({ ...params, sessionKey: storageSessionKey }), + fallbackEntry: preparedEntry, + replaceEntry: true, + }, + ); + if (!persisted) { + return null; + } + runOpenClawStateWriteTransaction( + (database) => { + upsertAcpSessionMetaRow( + database.db, + bindAcpSessionMeta({ + sessionKey: persisted.sessionKey, + sessionId: persisted.entry.sessionId, + meta: metaToPersist, + updatedAt, + }), + ); + if (persisted.sessionKey !== storageSessionKey) { executeSqliteQuerySync( database.db, getAcpSessionKysely(database.db) .deleteFrom("acp_sessions") .where("session_key", "=", storageSessionKey), ); - return; - } - if (nextMeta !== undefined) { - upsertAcpSessionMetaRow( - database.db, - bindAcpSessionMeta({ - sessionKey: storageSessionKey, - sessionId: preparedEntry.sessionId, - meta: nextMeta, - updatedAt, - }), - ); } }, { env: params.env, path: params.databasePath }, ); - if (nextMeta === undefined) { - return current ? mergeAcpForReturn(entry, current) : (entry ?? null); - } - if (nextMeta === null) { - if (!entry) { - return null; - } - const { updateSessionStore } = await loadSessionStoreRuntime(); - return await updateSessionStore( - storeEntry.storePath, - (store) => { - const storeSessionKey = resolveStoreSessionKey(store, storageSessionKey); - const next = { ...(store[storeSessionKey] ?? entry) }; - delete next.acp; - store[storeSessionKey] = next; - return next; - }, - sessionStoreUpdateOptions({ ...params, sessionKey: storageSessionKey }), - ); - } - const { updateSessionStore } = await loadSessionStoreRuntime(); - const persisted = await updateSessionStore( - storeEntry.storePath, - (store) => { - const storeSessionKey = resolveStoreSessionKey(store, storageSessionKey); - const next = mergeSessionEntry(store[storeSessionKey], { - sessionId: preparedEntry?.sessionId, - updatedAt, - }); - delete next.acp; - store[storeSessionKey] = next; - return next; - }, - sessionStoreUpdateOptions({ ...params, sessionKey: storageSessionKey }), - ); - return mergeAcpForReturn(persisted, nextMeta); + return mergeAcpForReturn(persisted.entry, metaToPersist); } diff --git a/src/agents/acp-spawn.test.ts b/src/agents/acp-spawn.test.ts index 6ad77a4dfd37..4452854ac873 100644 --- a/src/agents/acp-spawn.test.ts +++ b/src/agents/acp-spawn.test.ts @@ -71,6 +71,68 @@ const hoisted = vi.hoisted(() => { const countActiveRunsForSessionMock = vi.fn(); const getSubagentRunByChildSessionKeyMock = vi.fn(); const listTasksForOwnerKeyMock = vi.fn(); + const createSessionAccessorMock = () => { + const resolveMockStorePath = (scope: { + agentId?: string; + env?: NodeJS.ProcessEnv; + storePath?: string; + }): string => + scope.storePath ?? + resolveStorePathMock(undefined, { + agentId: scope.agentId, + env: scope.env, + }); + const loadMockEntry = (scope: { + agentId?: string; + env?: NodeJS.ProcessEnv; + sessionKey: string; + storePath?: string; + }): SessionEntry | undefined => { + const store = loadSessionStoreMock(resolveMockStorePath(scope)) as Record< + string, + SessionEntry + >; + return store[scope.sessionKey]; + }; + return { + listSessionEntries: ( + scope: { + agentId?: string; + env?: NodeJS.ProcessEnv; + storePath?: string; + } = {}, + ) => { + const store = loadSessionStoreMock(resolveMockStorePath(scope)) as Record< + string, + SessionEntry + >; + return Object.entries(store).map(([sessionKey, entry]) => ({ sessionKey, entry })); + }, + loadSessionEntry: loadMockEntry, + resolveSessionTranscriptRuntimeTarget: async (scope: { + agentId: string; + sessionId: string; + sessionKey: string; + storePath?: string; + threadId?: string | number; + }) => { + const store = scope.storePath + ? (loadSessionStoreMock(scope.storePath) as Record) + : undefined; + const resolved = await resolveSessionTranscriptFileMock({ + ...scope, + ...(store ? { sessionStore: store } : {}), + sessionEntry: loadMockEntry(scope), + }); + return { + agentId: scope.agentId, + sessionFile: resolved.sessionFile, + sessionId: scope.sessionId, + sessionKey: scope.sessionKey, + }; + }, + }; + }; const state = { cfg: createDefaultSpawnConfig(), }; @@ -98,6 +160,7 @@ const hoisted = vi.hoisted(() => { countActiveRunsForSessionMock, getSubagentRunByChildSessionKeyMock, listTasksForOwnerKeyMock, + createSessionAccessorMock, state, }; }); @@ -134,6 +197,8 @@ vi.mock("../config/sessions/store.js", () => ({ loadSessionStore: hoisted.loadSessionStoreMock, })); +vi.mock("../config/sessions/session-accessor.js", () => hoisted.createSessionAccessorMock()); + vi.mock("../config/sessions.js", () => ({ loadSessionStore: hoisted.loadSessionStoreMock, resolveStorePath: hoisted.resolveStorePathMock, diff --git a/src/agents/acp-spawn.ts b/src/agents/acp-spawn.ts index 3bf8a3162e87..e5089d723fce 100644 --- a/src/agents/acp-spawn.ts +++ b/src/agents/acp-spawn.ts @@ -47,8 +47,11 @@ import { } from "../config/agent-limits.js"; import { getRuntimeConfig } from "../config/config.js"; import { resolveStorePath } from "../config/sessions/paths.js"; -import { loadSessionStore } from "../config/sessions/store.js"; -import { resolveSessionTranscriptFile } from "../config/sessions/transcript.js"; +import { + listSessionEntries, + loadSessionEntry, + resolveSessionTranscriptRuntimeTarget, +} from "../config/sessions/session-accessor.js"; import type { SessionAcpMeta, SessionEntry } from "../config/sessions/types.js"; import type { OpenClawConfig } from "../config/types.openclaw.js"; import { callGateway } from "../gateway/call.js"; @@ -268,7 +271,6 @@ type AcpSpawnInitializedRuntime = { runtimeCloseHandle: AcpSpawnRuntimeCloseHandle; sessionId?: string; sessionEntry: SessionEntry | undefined; - sessionStore: Record; storePath: string; }; @@ -445,8 +447,11 @@ function hasSessionLocalHeartbeatRelayRoute(params: { const storePath = resolveStorePath(params.cfg.session?.store, { agentId: params.requesterAgentId, }); - const sessionStore = loadSessionStore(storePath); - const parentEntry = sessionStore[params.parentSessionKey]; + const parentEntry = loadSessionEntry({ + storePath, + sessionKey: params.parentSessionKey, + clone: false, + }); const parentDeliveryContext = deliveryContextFromSession(parentEntry); return Boolean(parentDeliveryContext?.channel && parentDeliveryContext.to); } @@ -595,23 +600,26 @@ async function persistAcpSpawnSessionFileBestEffort(params: { sessionId: string; sessionKey: string; sessionEntry: SessionEntry | undefined; - sessionStore: Record; storePath: string; agentId: string; threadId?: string | number; stage: "spawn" | "thread-bind"; }): Promise { try { - const resolvedSessionFile = await resolveSessionTranscriptFile({ + const resolvedSessionFile = await resolveSessionTranscriptRuntimeTarget({ sessionId: params.sessionId, sessionKey: params.sessionKey, - sessionEntry: params.sessionEntry, - sessionStore: params.sessionStore, storePath: params.storePath, agentId: params.agentId, threadId: params.threadId, }); - return resolvedSessionFile.sessionEntry; + return ( + loadSessionEntry({ + storePath: params.storePath, + sessionKey: resolvedSessionFile.sessionKey, + clone: false, + }) ?? params.sessionEntry + ); } catch (error) { log.warn( `ACP session-file persistence failed during ${params.stage} for ${params.sessionKey}: ${summarizeError(error)}`, @@ -964,9 +972,8 @@ function validateAcpResumeSessionOwnership(params: { } const storePath = resolveStorePath(params.cfg.session?.store, { agentId: params.targetAgentId }); - const sessionStore = loadSessionStore(storePath); - for (const [sessionKey, entry] of Object.entries(sessionStore)) { - const acp = readAcpSessionMeta({ sessionKey }); + for (const { sessionKey, entry } of listSessionEntries({ storePath, clone: false })) { + const acp = readAcpSessionMeta({ sessionKey, cfg: params.cfg }); if (!sessionEntryMatchesAcpResumeSessionId(acp, resumeSessionId)) { continue; } @@ -1064,14 +1071,16 @@ async function initializeAcpSpawnRuntime(params: { cwd?: string; }): Promise { const storePath = resolveStorePath(params.cfg.session?.store, { agentId: params.targetAgentId }); - const sessionStore = loadSessionStore(storePath); - let sessionEntry: SessionEntry | undefined = sessionStore[params.sessionKey]; + let sessionEntry = loadSessionEntry({ + storePath, + sessionKey: params.sessionKey, + clone: false, + }); const sessionId = sessionEntry?.sessionId; if (sessionId) { sessionEntry = await persistAcpSpawnSessionFileBestEffort({ sessionId, sessionKey: params.sessionKey, - sessionStore, storePath, sessionEntry, agentId: params.targetAgentId, @@ -1098,7 +1107,6 @@ async function initializeAcpSpawnRuntime(params: { }, sessionId, sessionEntry, - sessionStore, storePath, }; } @@ -1170,7 +1178,6 @@ async function bindPreparedAcpThread(params: { sessionEntry = await persistAcpSpawnSessionFileBestEffort({ sessionId: params.initializedRuntime.sessionId, sessionKey: params.sessionKey, - sessionStore: params.initializedRuntime.sessionStore, storePath: params.initializedRuntime.storePath, sessionEntry, agentId: params.targetAgentId, @@ -1532,16 +1539,19 @@ export async function spawnAcpDirect( childSessionKey: sessionKey, }) : undefined; + const parentAgentId = parentSessionKey + ? resolveAgentIdFromSessionKey(parentSessionKey) + : undefined; // Resolve parent session delivery context so system events route to the // correct thread/topic instead of falling back to the main DM. const parentDeliveryCtx = effectiveStreamToParent && parentSessionKey ? deliveryContextFromSession( - loadSessionStore( - resolveStorePath(cfg.session?.store, { - agentId: resolveAgentIdFromSessionKey(parentSessionKey), - }), - )[parentSessionKey], + loadSessionEntry({ + sessionKey: parentSessionKey, + ...(parentAgentId ? { agentId: parentAgentId } : {}), + clone: false, + }), ) : undefined; diff --git a/src/config/sessions/session-accessor.ts b/src/config/sessions/session-accessor.ts index bacda11f6655..a664d8204be5 100644 --- a/src/config/sessions/session-accessor.ts +++ b/src/config/sessions/session-accessor.ts @@ -49,6 +49,7 @@ import { loadSessionStore, applySessionEntryPatchProjection as applyFileSessionEntryPatchProjection, patchSessionEntry as patchFileSessionEntry, + patchSessionEntryWithKey as patchFileSessionEntryWithKey, purgeDeletedAgentSessionEntries as purgeFileDeletedAgentSessionEntries, readSessionUpdatedAt as readFileSessionUpdatedAt, resolveSessionStoreEntry, @@ -394,8 +395,14 @@ export type SessionEntryPatchOptions = { maintenanceConfig?: ResolvedSessionMaintenanceConfig; /** Keep the previous updatedAt value when the patch should not count as activity. */ preserveActivity?: boolean; + /** Throw when best-effort store recovery cannot confirm the requested write. */ + requireWriteSuccess?: boolean; /** Replace the whole entry instead of merging the returned patch. */ replaceEntry?: boolean; + /** Skip prune/cap/rotation maintenance for specialized internal updates. */ + skipMaintenance?: boolean; + /** Let the writer cache retain the updated object without cloning. */ + takeCacheOwnership?: boolean; }; export type SessionEntryPatchContext = { @@ -403,6 +410,13 @@ export type SessionEntryPatchContext = { existingEntry?: SessionEntry; }; +export type SessionEntryPatchResult = { + /** Exact persisted key for the patched entry after alias normalization. */ + sessionKey: string; + /** Persisted entry returned by the backing store. */ + entry: SessionEntry; +}; + export type RestartRecoveryLifecycleEntry = { /** Exact persisted key for the restart recovery candidate row. */ sessionKey: string; @@ -754,7 +768,35 @@ export async function patchSessionEntry( fallbackEntry: options.fallbackEntry, maintenanceConfig: options.maintenanceConfig, preserveActivity: options.preserveActivity, + requireWriteSuccess: options.requireWriteSuccess, replaceEntry: options.replaceEntry, + skipMaintenance: options.skipMaintenance, + takeCacheOwnership: options.takeCacheOwnership, + update, + }); +} + +/** + * Applies an atomic patch and returns the persisted key selected by the backing + * store. Use when a caller must keep sidecar state keyed to the final row. + */ +export async function patchSessionEntryWithKey( + scope: SessionAccessScope, + update: ( + entry: SessionEntry, + context: SessionEntryPatchContext, + ) => Promise | null> | Partial | null, + options: SessionEntryPatchOptions = {}, +): Promise { + return await patchFileSessionEntryWithKey({ + ...scope, + fallbackEntry: options.fallbackEntry, + maintenanceConfig: options.maintenanceConfig, + preserveActivity: options.preserveActivity, + requireWriteSuccess: options.requireWriteSuccess, + replaceEntry: options.replaceEntry, + skipMaintenance: options.skipMaintenance, + takeCacheOwnership: options.takeCacheOwnership, update, }); } diff --git a/src/config/sessions/store.ts b/src/config/sessions/store.ts index 74f4f5739d20..8ec5459c0d71 100644 --- a/src/config/sessions/store.ts +++ b/src/config/sessions/store.ts @@ -1773,18 +1773,29 @@ export async function applySessionStoreEntryPatch(params: { }); } +type SessionEntryPatchParams = SessionEntryWorkflowOptions & { + sessionKey: string; + fallbackEntry?: SessionEntry; + preserveActivity?: boolean; + requireWriteSuccess?: boolean; + replaceEntry?: boolean; + skipMaintenance?: boolean; + takeCacheOwnership?: boolean; + update: ( + entry: SessionEntry, + context: { existingEntry?: SessionEntry }, + ) => Promise | null> | Partial | null; +}; + export async function patchSessionEntry( - params: SessionEntryWorkflowOptions & { - sessionKey: string; - fallbackEntry?: SessionEntry; - preserveActivity?: boolean; - replaceEntry?: boolean; - update: ( - entry: SessionEntry, - context: { existingEntry?: SessionEntry }, - ) => Promise | null> | Partial | null; - }, + params: SessionEntryPatchParams, ): Promise { + return (await patchSessionEntryWithKey(params))?.entry ?? null; +} + +export async function patchSessionEntryWithKey( + params: SessionEntryPatchParams, +): Promise<{ sessionKey: string; entry: SessionEntry } | null> { const storePath = resolveSessionWorkflowStorePath(params); return await runExclusiveSessionStoreWrite(storePath, async () => { const store = loadMutableSessionStoreForWriter(storePath); @@ -1797,22 +1808,27 @@ export async function patchSessionEntry( existingEntry: resolved.existing ? cloneSessionEntry(resolved.existing) : undefined, }); if (!patch) { - return existing; + return { sessionKey: resolved.normalizedKey, entry: existing }; } const next = params.replaceEntry ? cloneSessionEntry(patch as SessionEntry) : params.preserveActivity ? mergeSessionEntryPreserveActivity(existing, patch) : mergeSessionEntry(existing, patch); - return await persistResolvedSessionEntry({ - storePath, - store, - resolved, - next, - maintenanceConfig: params.maintenanceConfig, - takeCacheOwnership: true, - returnDetached: true, - }); + return { + sessionKey: resolved.normalizedKey, + entry: await persistResolvedSessionEntry({ + storePath, + store, + resolved, + next, + maintenanceConfig: params.maintenanceConfig, + requireWriteSuccess: params.requireWriteSuccess, + skipMaintenance: params.skipMaintenance, + takeCacheOwnership: params.takeCacheOwnership ?? true, + returnDetached: params.takeCacheOwnership !== true, + }), + }; }); } diff --git a/test/scripts/check-session-accessor-boundary.test.ts b/test/scripts/check-session-accessor-boundary.test.ts index 8d790c47641c..de2ad2d84afd 100644 --- a/test/scripts/check-session-accessor-boundary.test.ts +++ b/test/scripts/check-session-accessor-boundary.test.ts @@ -24,6 +24,8 @@ describe("session accessor boundary guard", () => { expect(migratedSessionAccessorFiles).toEqual( new Set([ "packages/memory-host-sdk/src/host/session-files.ts", + "src/acp/runtime/session-meta.ts", + "src/agents/acp-spawn.ts", "src/agents/embedded-agent-runner/compaction-successor-transcript.ts", "src/agents/embedded-agent-runner/run/attempt.ts", "src/agents/embedded-agent-runner/tool-result-truncation.ts", @@ -82,6 +84,7 @@ describe("session accessor boundary guard", () => { it("ratchets only files migrated to session accessor writes", () => { expect(migratedSessionAccessorWriteFiles).toEqual( new Set([ + "src/acp/runtime/session-meta.ts", "src/agents/command/attempt-execution.shared.ts", "src/agents/command/session-store.ts", "src/agents/embedded-agent-runner/run.ts",