From fe625db116fac2d7321035060bfa3e9c642ea59e Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 13 Sep 2026 10:53:00 -0700 Subject: [PATCH] fix: serialize reclamation worker shutdown with foreground writes (#147178) --- ...-accessor.sqlite-reclamation-reuse.test.ts | 108 +++++++++++++++++- ...sion-accessor.sqlite-reclamation-worker.ts | 19 ++- 2 files changed, 122 insertions(+), 5 deletions(-) diff --git a/src/config/sessions/session-accessor.sqlite-reclamation-reuse.test.ts b/src/config/sessions/session-accessor.sqlite-reclamation-reuse.test.ts index 0e40f9af3bdb..1283710aa768 100644 --- a/src/config/sessions/session-accessor.sqlite-reclamation-reuse.test.ts +++ b/src/config/sessions/session-accessor.sqlite-reclamation-reuse.test.ts @@ -16,6 +16,7 @@ import { setOpenClawAgentDatabaseValidation, type OpenClawAgentDatabaseValidation, } from "../../state/openclaw-agent-db-validation-cache.js"; +import { withOpenClawAgentDatabaseWrite } from "../../state/openclaw-agent-db-write.js"; import { closeOpenClawAgentDatabaseByPath, closeOpenClawAgentDatabaseByPathAsync, @@ -24,6 +25,7 @@ import { getOpenClawAgentDatabaseIfOpen, openOpenClawAgentDatabase, } from "../../state/openclaw-agent-db.js"; +import { runOpenClawAgentWriteAdmission } from "../../state/openclaw-agent-write-admission.js"; import { closeOpenClawStateDatabaseForTest, openOpenClawStateDatabase, @@ -374,16 +376,118 @@ test("retires the previous database before opening a different agent store", asy const second = createFixture(); invalidateOpenClawAgentDatabaseValidation(first.database.path); invalidateOpenClawAgentDatabaseValidation(second.database.path); - const spawned = observeReclamationWorkers(); + const closeEntered = createDeferredCore(); + let closeRetained: (() => Promise) | undefined; + const withWorker = reclamationWorker.withSqliteReclamationWorker; + vi.spyOn(reclamationWorker, "withSqliteReclamationWorker").mockImplementation( + (options, claim, run, assertRequestCurrent) => + withWorker( + options, + claim, + async (worker) => { + if (options.path === first.database.path) { + const close = worker.close.bind(worker); + closeRetained = close; + vi.spyOn(worker, "close").mockImplementation(() => { + const pending = close(); + closeEntered.resolve(); + return pending; + }); + } + return run(worker); + }, + assertRequestCurrent, + ), + ); + let firstSpawned: Worker | undefined; + let predecessorExitedAtSuccessorSpawn: boolean | undefined; + const spawned = observeReclamationWorkers((worker) => { + if (firstSpawned) { + predecessorExitedAtSuccessorSpawn = firstSpawned.threadId === -1; + } else { + firstSpawned = worker; + } + }); await runSqliteSessionReclamation({ forceInProcess: false, plan: first.plans[0]! }); expect(leasesFor(first)).toHaveLength(2); - await runSqliteSessionReclamation({ forceInProcess: false, plan: second.plans[0]! }); + const previous = spawned[0]; + const survivor = first.scopes[1]; + const retirePrevious = closeRetained; + if (!previous || !survivor || !retirePrevious) { + throw new Error("Expected the first real reclamation Worker and its retained owner"); + } + const order: string[] = []; + const exited = once(previous, "exit").then(() => { + order.push("worker-exit"); + }); + void exited.catch(() => undefined); + const options = { ...first.options, path: first.database.path }; + let followingWrite: Promise | undefined; + const postMessage = previous.postMessage.bind(previous); + vi.spyOn(previous, "postMessage").mockImplementation((message: unknown, transferList) => { + if ( + typeof message === "object" && + message !== null && + "type" in message && + message.type === "close" + ) { + order.push("close-dispatch"); + followingWrite = withOpenClawAgentDatabaseWrite( + options, + () => { + order.push("following-write"); + return appendTranscriptEventSync(survivor, { type: "after-reclamation-close" }); + }, + first.database.db, + ); + void followingWrite.catch(() => undefined); + } + postMessage(message, transferList); + }); + const entered = createDeferredCore(); + const release = createDeferredCore(); + const holding = runOpenClawAgentWriteAdmission(options, async () => { + order.push("foreground-enter"); + entered.resolve(); + await release.promise; + order.push("foreground-release"); + }); + await entered.promise; + const switching = runSqliteSessionReclamation({ forceInProcess: false, plan: second.plans[0]! }); + void switching.catch(() => undefined); + try { + await Promise.race([closeEntered.promise, switching]); + await yieldToEventLoop(); + } finally { + release.resolve(); + await Promise.allSettled([holding, switching]); + await Promise.allSettled([retirePrevious(), exited]); + if (followingWrite) { + await Promise.allSettled([followingWrite]); + } + } + await expect(switching).resolves.toMatchObject({ + kind: "lifecycle-artifacts", + value: { removedEntries: 1 }, + }); + await expect(followingWrite).resolves.toEqual({ ok: true, value: true }); + expect({ order, predecessorExitedAtSuccessorSpawn }).toEqual({ + order: [ + "foreground-enter", + "foreground-release", + "close-dispatch", + "worker-exit", + "following-write", + ], + predecessorExitedAtSuccessorSpawn: true, + }); expect(spawned).toHaveLength(2); expect(fullChecks()).toBe(2); expect(spawned[0]?.threadId).toBe(-1); expect(leasesFor(first)).toHaveLength(1); expect(leasesFor(second)).toHaveLength(2); expect(loadSessionEntryReadOnly(first.scopes[1]!)).toMatchObject({ sessionId: "second" }); + expect(await loadTranscriptEvents(survivor)).toContainEqual({ type: "after-reclamation-close" }); expect(loadSessionEntryReadOnly(second.scopes[1]!)).toMatchObject({ sessionId: "second" }); }); diff --git a/src/config/sessions/session-accessor.sqlite-reclamation-worker.ts b/src/config/sessions/session-accessor.sqlite-reclamation-worker.ts index a4ff34246379..40ef5d368fed 100644 --- a/src/config/sessions/session-accessor.sqlite-reclamation-worker.ts +++ b/src/config/sessions/session-accessor.sqlite-reclamation-worker.ts @@ -11,6 +11,7 @@ import { type OpenClawAgentDatabaseWorkerLeaseReceipt, } from "../../state/openclaw-agent-db-lease.js"; import { registerOpenClawAgentDatabaseAsyncResource } from "../../state/openclaw-agent-db-lifecycle.js"; +import { runOpenClawAgentWorkerWrite } from "../../state/openclaw-agent-write-admission.js"; import { captureOpenClawStateDatabaseReadAdmission, registerOpenClawStateDatabaseAsyncResource, @@ -319,9 +320,21 @@ class SqliteReclamationWorker { } return (this.closing ??= (async () => { await this.active?.catch(() => {}); - this.worker?.ref(); - this.worker?.postMessage({ type: "close" }); - await this.exited; + const worker = this.worker; + if (worker) { + worker.ref(); + await runOpenClawAgentWorkerWrite(this.options, async () => { + try { + worker.postMessage({ type: "close" }, []); + } catch (error) { + await worker.terminate(); + throw error; + } finally { + // Checkpoint and native close retain admission even if dispatch fails. + await this.exited; + } + }); + } // Native exit is joined before exact receipt cleanup; PID-wide cleanup is never safe. if (this.lease) { releaseExitedOpenClawAgentDatabaseWorkerLease(this.lease);