diff --git a/src/auto-reply/reply/session.init-conflict-retry.test.ts b/src/auto-reply/reply/session.init-conflict-retry.test.ts index 284404ff12f6..081acb4bec18 100644 --- a/src/auto-reply/reply/session.init-conflict-retry.test.ts +++ b/src/auto-reply/reply/session.init-conflict-retry.test.ts @@ -2,6 +2,11 @@ import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import { describe, expect, it, vi } from "vitest"; +import { + loadSessionEntry, + upsertSessionEntryCore, +} from "../../config/sessions/session-accessor.js"; +import { replaceSessionEntrySync } from "../../config/sessions/session-accessor.sqlite-entry.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { finalizeInboundContext } from "./inbound-context.js"; import { @@ -18,6 +23,9 @@ const initSessionState = ( const commitConflictControl = vi.hoisted(() => ({ abortController: undefined as AbortController | undefined, + beforeEntryMutation: undefined as + | ((params: { sessionKey: string; storePath: string }) => Promise | void) + | undefined, commitCalls: 0, remainingFailures: 0, })); @@ -43,7 +51,22 @@ vi.mock("../../config/sessions/session-accessor.js", async (importOriginal) => { revision: `forced-conflict-${commitConflictControl.commitCalls}`, }; } - return await actual.commitReplySessionInitialization(...args); + const [params] = args; + const beforeEntryMutation = commitConflictControl.beforeEntryMutation; + return await actual.commitReplySessionInitialization({ + ...params, + ...(beforeEntryMutation + ? { + beforeEntryMutation: async (context) => { + await params.beforeEntryMutation?.(context); + await beforeEntryMutation({ + sessionKey: params.sessionKey, + storePath: params.storePath, + }); + }, + } + : {}), + }); }, }; }); @@ -203,6 +226,79 @@ describe("runWithSessionInitConflictRetry", () => { }); describe("initSessionState conflict retry wiring", () => { + it("retries a late same-session lifecycle conflict without losing either update", async () => { + const root = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-session-late-conflict-")); + const storePath = path.join(root, "sessions.json"); + let lateWrites = 0; + commitConflictControl.commitCalls = 0; + commitConflictControl.beforeEntryMutation = ({ sessionKey, storePath: targetStorePath }) => { + if (lateWrites >= 2) { + return; + } + const currentEntry = loadSessionEntry({ + readConsistency: "latest", + sessionKey, + storePath: targetStorePath, + }); + if (!currentEntry) { + throw new Error("expected the projected reply session row"); + } + lateWrites += 1; + replaceSessionEntrySync( + { sessionKey, storePath: targetStorePath }, + { + ...currentEntry, + lastHeartbeatSentAt: 100 + lateWrites, + lastHeartbeatText: `concurrent metadata ${lateWrites}`, + }, + ); + if (lateWrites === 2) { + commitConflictControl.beforeEntryMutation = undefined; + } + }; + + try { + await upsertSessionEntryCore( + { sessionKey: SESSION_KEY, storePath }, + { + displayName: "before reply initialization", + sessionId: "existing-session", + updatedAt: Date.now(), + }, + ); + + const result = await initSessionState({ + cfg: { session: { store: storePath } } as OpenClawConfig, + commandAuthorized: true, + ctx: { + Body: "hello", + SessionKey: SESSION_KEY, + ThreadLabel: "reply initialization update", + }, + }); + + expect(commitConflictControl.commitCalls).toBe(3); + expect(lateWrites).toBe(2); + expect(result.sessionEntry).toMatchObject({ + displayName: "reply initialization update", + lastHeartbeatSentAt: 102, + lastHeartbeatText: "concurrent metadata 2", + sessionId: "existing-session", + }); + expect( + loadSessionEntry({ readConsistency: "latest", sessionKey: SESSION_KEY, storePath }), + ).toMatchObject({ + displayName: "reply initialization update", + lastHeartbeatSentAt: 102, + lastHeartbeatText: "concurrent metadata 2", + sessionId: "existing-session", + }); + } finally { + commitConflictControl.beforeEntryMutation = undefined; + await fs.rm(root, { recursive: true, force: true }); + } + }); + it("cancels the production backoff through the initializer signal", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-session-conflict-abort-")); const controller = new AbortController(); diff --git a/src/config/sessions/session-accessor.lifecycle-types.ts b/src/config/sessions/session-accessor.lifecycle-types.ts index 8bb39176ee2d..47f1afbad7cb 100644 --- a/src/config/sessions/session-accessor.lifecycle-types.ts +++ b/src/config/sessions/session-accessor.lifecycle-types.ts @@ -130,6 +130,13 @@ export type SessionEntryLifecycleRemoval = SessionEntryLifecycleRemovalBase & } ); +export class SessionEntryLifecycleUpsertConflictError extends Error { + constructor(readonly sessionKey: string) { + super(`SQLite session entry changed before lifecycle upsert for ${sessionKey}`); + this.name = "SessionEntryLifecycleUpsertConflictError"; + } +} + export type SessionEntryLifecycleUpsert = { sessionKey: string; resetBoundaryReason?: SessionResetBoundaryReason; diff --git a/src/config/sessions/session-accessor.reset.ts b/src/config/sessions/session-accessor.reset.ts index 3fe5b3db19c4..3aad11263c71 100644 --- a/src/config/sessions/session-accessor.reset.ts +++ b/src/config/sessions/session-accessor.reset.ts @@ -13,9 +13,13 @@ import { import { listSessionEntriesCore, listSessionEntriesReadOnly, + loadSessionEntry, resolveSessionEntryFromStore, } from "./session-accessor.entry.js"; -import type { SessionEntryLifecycleUpsert } from "./session-accessor.lifecycle-types.js"; +import { + SessionEntryLifecycleUpsertConflictError, + type SessionEntryLifecycleUpsert, +} from "./session-accessor.lifecycle-types.js"; import { applySessionEntryLifecycleMutation } from "./session-accessor.lifecycle.js"; import type { SessionLifecycleTranscriptInfo, @@ -119,6 +123,18 @@ export function loadReplySessionInitializationSnapshot(params: { }; } +function createStaleReplySessionInitializationResult( + currentEntry: SessionEntry | undefined, + storePath: string, +): ReplySessionInitializationCommitResult { + return { + ok: false, + ...(currentEntry ? { currentEntry } : {}), + reason: "stale-snapshot", + revision: createReplySessionInitializationRevision({ entry: currentEntry, storePath }), + }; +} + /** * Persists one reply-session initialization result and archives the previous * transcript after metadata commits. SQLite adapters map the guarded write to a @@ -165,12 +181,7 @@ export async function commitReplySessionInitialization(params: { storePath, }); if (revision !== params.expectedRevision) { - return { - ok: false, - ...(currentEntry ? { currentEntry } : {}), - reason: "stale-snapshot", - revision, - }; + return createStaleReplySessionInitializationResult(currentEntry, storePath); } const readEntry = (sessionKey: string) => { @@ -190,12 +201,7 @@ export async function commitReplySessionInitialization(params: { sessionEntry: preparedSessionEntry, storePath, }); - let staleCommit: - | { - currentEntry?: SessionEntry; - revision: string; - } - | undefined; + let staleCommit: SessionEntry | null | undefined; let committedSessionEntry = sessionEntry; let beforeEntryMutationDone = false; const upserts: SessionEntryLifecycleUpsert[] = [ @@ -213,10 +219,7 @@ export async function commitReplySessionInitialization(params: { storePath, }); if (commitRevision !== params.expectedRevision) { - staleCommit = { - ...(commitEntry ? { currentEntry: { ...commitEntry } } : {}), - revision: commitRevision, - }; + staleCommit = commitEntry ? { ...commitEntry } : null; return null; } // The identity-only guard allows commits when background activity @@ -245,23 +248,36 @@ export async function commitReplySessionInitialization(params: { const retiredEntry = params.retiredEntry; upserts.push({ sessionKey: retiredEntry.key, - buildEntry: () => (staleCommit ? null : retiredEntry.entry), + buildEntry: () => (staleCommit === undefined ? retiredEntry.entry : null), }); } - await applySessionEntryLifecycleMutation({ - activeSessionKey: params.activeSessionKey, - agentId: params.agentId, - maintenanceOverride: params.maintenanceConfig, - storePath, - upserts, - }); - if (staleCommit) { - return { - ok: false, - ...(staleCommit.currentEntry ? { currentEntry: staleCommit.currentEntry } : {}), - reason: "stale-snapshot", - revision: staleCommit.revision, - }; + try { + await applySessionEntryLifecycleMutation({ + activeSessionKey: params.activeSessionKey, + agentId: params.agentId, + maintenanceOverride: params.maintenanceConfig, + storePath, + upserts, + }); + } catch (error) { + if ( + !(error instanceof SessionEntryLifecycleUpsertConflictError) || + error.sessionKey !== resolved.normalizedKey + ) { + throw error; + } + return createStaleReplySessionInitializationResult( + loadSessionEntry({ + agentId: params.agentId, + readConsistency: "latest", + sessionKey: error.sessionKey, + storePath, + }), + storePath, + ); + } + if (staleCommit !== undefined) { + return createStaleReplySessionInitializationResult(staleCommit ?? undefined, storePath); } store[resolved.normalizedKey] = committedSessionEntry; if (params.retiredEntry) { diff --git a/src/config/sessions/session-accessor.sqlite-projection.ts b/src/config/sessions/session-accessor.sqlite-projection.ts index e4b63f38adcd..d8c10b99e8e8 100644 --- a/src/config/sessions/session-accessor.sqlite-projection.ts +++ b/src/config/sessions/session-accessor.sqlite-projection.ts @@ -10,7 +10,10 @@ import { runOpenClawAgentWriteTransaction, type OpenClawAgentDatabase, } from "../../state/openclaw-agent-db.js"; -import type { SessionArchivedTranscriptCleanupRule } from "./session-accessor.lifecycle-types.js"; +import { + SessionEntryLifecycleUpsertConflictError, + type SessionArchivedTranscriptCleanupRule, +} from "./session-accessor.lifecycle-types.js"; import { prunePublishedSessionArchivesByRetention, publishSessionStateArchives, @@ -357,7 +360,7 @@ export async function applySessionEntryLifecycleMutation(params: { if (sameKeyRemoval) { throw new Error(`SQLite session entry has stale lifecycle state for ${sessionKey}`); } - throw new Error(`SQLite session entry changed before lifecycle upsert for ${sessionKey}`); + throw new SessionEntryLifecycleUpsertConflictError(sessionKey); } if (sameKeyRemoval && !shouldRemoveSessionEntry(currentEntry, sameKeyRemoval.removal)) { throw new Error(`SQLite session entry has stale lifecycle state for ${sessionKey}`);