fix: Doctor leaves abandoned updater runtimes behind (#162886)

Doctor's own inspection workers and cached shared-state readers kept the conservative holder census from reclaiming abandoned updater runtimes.

Settle those resources through their existing owners while retaining maintenance authority, then run the unchanged census. Preserve independent holders, runtime-only service boundaries, and the refusal to restart after resource settlement fails.

Validated 53 focused/regression cases on Node and fork Bun, published 2026.9.7 upgrades and positive/negative custody on both runtimes, scoped-clean P2 review, and green exact-head CI/security checks.
This commit is contained in:
Peter Steinberger 2026-10-01 12:24:19 -07:00 • committed by GitHub
parent 6ec9fc1f02
commit 4f575823aa
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
22 changed files with 516 additions and 48 deletions

View file

@ -21,6 +21,8 @@ postures and maintenance modes documented on the other pages.
## Config writes and backups
After its checks finish, `doctor --fix` settles its own inspection workers while retaining maintenance ownership, then checks whether abandoned updater runtimes can be removed. Independent OpenClaw processes, Worker threads, and shared-broker work still prevent removal. Doctor reports the holder PIDs and asks you to let their work finish before rerunning `openclaw doctor --fix`.
- On npm global installs, Doctor reports retained `.openclaw.package-backup-*.databases` directories (and `.openclaw-package-backup-*.databases`, the name a failed cleanup retires them under) beside the installed package, with their total regular-file size in bytes and human-readable units and a quoted removal command for each directory. The scan is bounded; incomplete sizes are lower bounds. If inspection is incomplete before any snapshot is found, Doctor warns and asks you to list the npm global root manually, including hidden entries. A missing global root produces no warning. This is warning-only, including with `--fix`: confirm no update is in progress and no recovery needs the snapshots before removing them manually. Updater-driven Doctor passes defer this check so they do not report the active update's snapshots; run standalone Doctor after the update settles.
- Any config write (including a `--fix` repair) rotates a backup to `~/.openclaw/openclaw.json.bak` (with a numbered `.bak.1`..`.bak.4` ring). `--fix` also drops unknown config keys reported by schema validation, listing each removal; it skips this while an update is in progress so partially written upgrade state is not stripped before its migration finishes.
- If `openclaw.json` cannot be parsed and no last-known-good config can be recovered, `doctor --fix` leaves the file unchanged and exits with an error instead of writing a partial replacement. The error points to `openclaw config validate` for the exact parse position and explains how to edit or regenerate the config.

View file

@ -358,6 +358,7 @@ describe("unproved Doctor authority callers", () => {
signal: new AbortController().signal,
run: <T>(operation: () => T) => operation(),
repairSqliteNoCow: async () => {},
cleanupRetainedRuntimes: async () => {},
releaseState: vi.fn(async () => {}),
finish: vi.fn(async () => {}),
release: vi.fn(async () => {}),
@ -471,6 +472,7 @@ describe("unproved Doctor authority callers", () => {
signal: new AbortController().signal,
run: (operation) => operation(),
repairSqliteNoCow: async () => {},
cleanupRetainedRuntimes: async () => {},
releaseState: async () => {},
release: async () => {},
finish: async () => {

View file

@ -132,6 +132,7 @@ export function registerRepairCustodyTests(mocks: {
release,
releaseState,
repairSqliteNoCow: vi.fn(async () => {}),
cleanupRetainedRuntimes: vi.fn(async () => {}),
});
vi.spyOn(updateCheck, "resolveUpdateInstallKind").mockResolvedValue("package");
// Observe reconciliation of the selected old run without inventing a live

View file

@ -441,6 +441,7 @@ describe("update plugin lifecycle lease boundaries", () => {
finish: vi.fn(async () => {}),
release: vi.fn(async () => {}),
repairSqliteNoCow: async () => {},
cleanupRetainedRuntimes: async () => {},
releaseState: vi.fn(async () => {}),
};
mocks.maintenance.mockResolvedValueOnce(maintenance);
@ -824,6 +825,7 @@ describe("update plugin lifecycle lease boundaries", () => {
signal: new AbortController().signal,
run: <T>(operation: () => T): T => operation(),
repairSqliteNoCow: async () => {},
cleanupRetainedRuntimes: async () => {},
releaseState: async () => {
record("release-state");
},
@ -1004,6 +1006,7 @@ describe("update plugin lifecycle lease boundaries", () => {
run: <T>(operation: () => T): T => operation(),
release: async () => {},
repairSqliteNoCow: async () => {},
cleanupRetainedRuntimes: async () => {},
releaseState: async () => {},
finish: async () => {
warnings.push(warning);

View file

@ -72,6 +72,7 @@ const beginDoctorMaintenance = vi.hoisted(() =>
run: <T>(operation: () => T): T => operation(),
releaseState: vi.fn(async () => {}),
repairSqliteNoCow: vi.fn(async () => {}),
cleanupRetainedRuntimes: vi.fn(async () => {}),
release: doctorMaintenanceRelease,
finish: vi.fn(async () => {}),
})),

View file

@ -1,8 +1,12 @@
import path from "node:path";
import { resolveStateDir } from "../config/paths.js";
import { formatErrorMessage } from "../infra/errors.js";
import { resolveGatewayStateOwnerPath } from "../infra/gateway-state-owner.js";
import { createSqliteReadOnlyWorkerScope } from "../infra/sqlite-readonly-worker.js";
import { createUpdateDoctorDatabaseWriteCapture } from "../infra/update-doctor-result.js";
import {
createUpdateDoctorDatabaseWriteCapture,
DoctorMaintenanceRefusalError,
} from "../infra/update-doctor-result.js";
import {
createOpenClawDatabaseMaintenanceScope,
type OpenClawDatabaseMaintenanceScope,
@ -213,6 +217,54 @@ export function createDoctorMaintenanceState(options: {
await enterResources(owner!);
}
},
async cleanupRetainedRuntimes(inspectService: boolean) {
const { captureRetainedNativeWorkerSource } =
await import("../infra/worker-native-lifecycle.js");
const { retireIdleOpenClawStateReadWorkers } =
await import("../state/openclaw-state-read-worker.js");
const { prepareRetainedUpdateRuntimeCleanup } = await import("./doctor-retained-runtime.js");
options.assertCurrent?.();
owner!.assertCurrent();
const cleanup = await state.run(() =>
prepareRetainedUpdateRuntimeCleanup(selectedEnv, { inspectService }),
);
const nativeSource = captureRetainedNativeWorkerSource();
// This phase runs after the tracked Doctor callback has settled.
let readersRetired: boolean;
let brokerRetired: boolean;
try {
await closeResources();
readersRetired = await retireIdleOpenClawStateReadWorkers(nativeSource);
brokerRetired = await nativeSource.retireIdleBroker();
} catch (cause) {
throw new DoctorMaintenanceRefusalError(
`Doctor inspection resource cleanup failed: ${formatErrorMessage(cause)}. Resolve this cleanup failure before restarting the Gateway or rerunning openclaw doctor --fix.`,
{ kind: "data-at-risk", reason: "active-mutation" },
{ cause },
);
}
try {
await owner!.run(() =>
cleanup(true, {
assertCurrent() {
options.assertCurrent?.();
owner!.assertCurrent();
owner!.assertDatabaseAccess(resolveOpenClawStateSqlitePath(selectedEnv));
},
assertResourcesSettled() {
// Worker threads share this PID and are invisible to the process census.
if (!readersRetired || !brokerRetired || nativeSource.hasActiveWorkers) {
throw new Error(
`independent native work in this process (PID: ${process.pid}); let these holders finish, then rerun openclaw doctor --fix`,
);
}
},
}),
);
} finally {
await enterResources(owner!);
}
},
async release() {
await closeResources();
await settleCapture();

View file

@ -25,6 +25,7 @@ export type DoctorMaintenance = {
signal: AbortSignal;
releaseState(): Promise<void>;
repairSqliteNoCow(paths: readonly string[]): Promise<void>;
cleanupRetainedRuntimes(): Promise<void>;
release(): Promise<void>;
finish(
cfg: OpenClawConfig | undefined,

View file

@ -607,6 +607,12 @@ export async function beginDoctorMaintenance(
warn(message);
}
},
async cleanupRetainedRuntimes() {
if (this !== maintenance || custody !== "held") {
throw new Error("Updater runtime cleanup requires its original live maintenance owner.");
}
await settle(() => state.cleanupRetainedRuntimes(serviceUpdateVerdict !== undefined));
},
async release() {
if (this !== maintenance) {
throw new Error("Gateway restoration requires its original live maintenance owner.");

View file

@ -4,29 +4,38 @@ import { maintainRetainedUpdateRuntimes } from "../infra/temp-artifact-cleanup.j
import { getOpenClawDatabaseMaintenanceScope } from "../state/openclaw-state-db-async-lifecycle.js";
import { inspectDoctorTemporaryDirectories } from "./doctor/shared/temporary-directories.js";
export async function noteRetainedUpdateRuntimes(
export async function prepareRetainedUpdateRuntimeCleanup(
env: NodeJS.ProcessEnv,
shouldRepair: boolean,
): Promise<void> {
options?: { inspectService?: boolean },
) {
const maintenance = getOpenClawDatabaseMaintenanceScope();
const { directories, warnings } = await inspectDoctorTemporaryDirectories(env);
const lines = await maintainRetainedUpdateRuntimes({
packageRoots: resolveOpenClawPackageRootsSync({
moduleUrl: import.meta.url,
argv1: process.argv[1],
cwd: process.cwd(),
}),
temporaryDirectories: directories,
repair: shouldRepair,
assertCurrent() {
if (!maintenance?.ownsSchemaMaintenance) {
throw new Error("Doctor does not hold Gateway maintenance");
}
maintenance.assertAdmission();
},
const { directories, warnings } = await inspectDoctorTemporaryDirectories(env, options);
const packageRoots = resolveOpenClawPackageRootsSync({
moduleUrl: import.meta.url,
argv1: process.argv[1],
cwd: process.cwd(),
});
lines.push(...warnings);
if (lines.length) {
note(lines.join("\n"), "Updater runtimes");
}
return async (
shouldRepair: boolean,
authority?: { assertCurrent: () => void; assertResourcesSettled: () => void },
): Promise<void> => {
const lines = await maintainRetainedUpdateRuntimes({
packageRoots,
temporaryDirectories: directories,
repair: shouldRepair,
assertCurrent:
authority?.assertCurrent ??
(() => {
if (!maintenance?.ownsSchemaMaintenance) {
throw new Error("Doctor does not hold Gateway maintenance");
}
maintenance.assertAdmission();
}),
assertResourcesSettled: authority?.assertResourcesSettled,
});
lines.push(...warnings);
if (lines.length) {
note(lines.join("\n"), "Updater runtimes");
}
};
}

View file

@ -4,7 +4,10 @@ import { resolveRequiredHomeDir } from "../../../infra/home-dir.js";
import { resolveEnvironmentValue } from "../../../infra/process-env.js";
/** Service scratch can outlive the shell environment that originally selected it. */
export async function inspectDoctorTemporaryDirectories(env: NodeJS.ProcessEnv): Promise<{
export async function inspectDoctorTemporaryDirectories(
env: NodeJS.ProcessEnv,
options?: { inspectService?: boolean },
): Promise<{
directories: string[];
warnings: string[];
}> {
@ -19,6 +22,9 @@ export async function inspectDoctorTemporaryDirectories(env: NodeJS.ProcessEnv):
),
);
const warnings: string[] = [];
if (options?.inspectService === false) {
return { directories, warnings };
}
try {
const { resolveGatewayService } = await import("../../../daemon/service.js");
const command = await resolveGatewayService().readCommand(env, {

View file

@ -44,8 +44,13 @@ export async function runLegacyPluginSourceCapturesHealth(
}
export async function runRetainedUpdateRuntimesHealth(ctx: DoctorHealthFlowContext): Promise<void> {
const { noteRetainedUpdateRuntimes } = await import("../commands/doctor-retained-runtime.js");
await noteRetainedUpdateRuntimes(ctx.env ?? process.env, ctx.prompter.shouldRepair);
if (ctx.gatewayMaintenanceActive && ctx.prompter.shouldRepair) {
return;
}
const { prepareRetainedUpdateRuntimeCleanup } =
await import("../commands/doctor-retained-runtime.js");
const cleanup = await prepareRetainedUpdateRuntimeCleanup(ctx.env ?? process.env);
await cleanup(ctx.prompter.shouldRepair);
}
export async function runReleaseConfiguredPluginInstallsHealth(

View file

@ -60,6 +60,7 @@ const maintenance = vi.hoisted(() => ({
finish: vi.fn(),
releaseState: vi.fn(),
repairSqliteNoCow: vi.fn(),
cleanupRetainedRuntimes: vi.fn(),
release: vi.fn(),
}));
const resultWriter = await vi.importActual<typeof import("../infra/update-doctor-result.js")>(

View file

@ -490,6 +490,9 @@ async function runDoctorHealthFlowWithResult(
await maintenance.repairSqliteNoCow(sqliteNoCowPaths);
}
}
if (ctx && maintenance && ctx.prompter.shouldRepair) {
await maintenance.cleanupRetainedRuntimes();
}
} catch (error) {
failure = error;
throw error;

View file

@ -0,0 +1,204 @@
import "./doctor-health.test-support.js";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import { afterEach, expect, it, vi } from "vitest";
import {
resolveGatewayStateOwnerPath,
tryAcquireGatewayStateOwner,
} from "../infra/gateway-state-owner.js";
import * as packageRoots from "../infra/openclaw-root.js";
import {
readOnlyWorkerScope,
type SqliteReadOnlyWorkerScope,
} from "../infra/sqlite-readonly-worker-context.js";
import { getOpenClawDatabaseMaintenanceScope } from "../state/openclaw-state-db-async-lifecycle.js";
import { resolveOpenClawStateSqlitePath } from "../state/openclaw-state-db.paths.js";
import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js";
import { runRetainedUpdateRuntimesHealth } from "./doctor-health-contribution-runners.state.js";
import { useDoctorHealthFixture } from "./doctor-health.fixture.test-support.js";
import { runDoctorHealthFlow } from "./doctor-health.js";
const { mocks } = await import("./doctor-health.test-support.js");
const custody = vi.hoisted(() => ({
note: vi.fn(),
census: vi.fn(),
retireBroker: vi.fn<() => Promise<boolean>>(),
activeNativeWork: false,
actualNativeSource: false,
}));
vi.mock("../../packages/terminal-core/src/note.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../../packages/terminal-core/src/note.js")>()),
note: custody.note,
}));
vi.mock("../infra/openclaw-process-census.js", () => ({
inspectOtherOpenClawProcesses: custody.census,
}));
vi.mock("../infra/worker-native-lifecycle.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("../infra/worker-native-lifecycle.js")>();
type NativeSource = ReturnType<typeof actual.captureRetainedNativeWorkerSource>;
const sources = new WeakMap<NativeSource, NativeSource>();
return {
...actual,
captureRetainedNativeWorkerSource: (
...args: Parameters<typeof actual.captureRetainedNativeWorkerSource>
) => {
const source = actual.captureRetainedNativeWorkerSource(...args);
let captured = sources.get(source);
if (!captured) {
captured = {
...source,
retireIdleBroker: () =>
custody.actualNativeSource ? source.retireIdleBroker() : custody.retireBroker(),
get hasActiveWorkers() {
return custody.actualNativeSource ? source.hasActiveWorkers : custody.activeNativeWork;
},
};
sources.set(source, captured);
}
return captured;
},
};
});
const { materializeSharedStateDatabase } = useDoctorHealthFixture();
afterEach(() => {
vi.restoreAllMocks();
custody.note.mockReset();
custody.census.mockReset();
custody.retireBroker.mockReset();
custody.activeNativeWork = false;
custody.actualNativeSource = false;
});
it.each([
"idle",
"idle-read-pool",
"independent-owner",
"independent-thread",
"late-independent-thread",
"settlement-failed",
] as const)(
"reclaims retained runtimes only after Doctor-owned readers settle (%s broker)",
async (broker) => {
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
custody.actualNativeSource = broker === "idle-read-pool";
materializeSharedStateDatabase(state.env);
const root = fs.realpathSync(state.root);
const packageRoot = path.join(root, "checkout");
const artifact = path.join(root, "openclaw-update-runtime-Old001");
const base = path.parse(packageRoot).root;
const marker = path.join(
artifact,
"tree",
Buffer.from(base).toString("hex"),
path.relative(base, packageRoot),
"package.json",
);
fs.mkdirSync(packageRoot);
fs.writeFileSync(path.join(packageRoot, "package.json"), '{"name":"openclaw"}');
fs.mkdirSync(path.dirname(marker), { recursive: true });
fs.writeFileSync(marker, '{"name":"openclaw"}');
vi.spyOn(os, "tmpdir").mockReturnValue(root);
vi.spyOn(packageRoots, "resolveOpenClawPackageRootsSync").mockReturnValue([packageRoot]);
mocks.service.mockReturnValue({
readCommand: async () => null,
readRuntime: async () => ({ status: "stopped", missingUnit: true }),
isLoaded: async () => false,
});
const database = resolveOpenClawStateSqlitePath(state.env);
const ownerPath = resolveGatewayStateOwnerPath(database);
let ownerIdentity: Buffer | undefined;
let doctorReaders: SqliteReadOnlyWorkerScope | undefined;
let brokerLive = false;
let readSource:
| import("../infra/worker-native-lifecycle.js").RetainedNativeWorkerSource
| undefined;
let checksFinished = false;
const assertMaintenanceHeld = () => {
const competingOwner = tryAcquireGatewayStateOwner(database);
competingOwner?.release();
expect(competingOwner).toBeNull();
expect(fs.readFileSync(ownerPath)).toEqual(ownerIdentity);
};
custody.retireBroker.mockImplementation(async () => {
assertMaintenanceHeld();
expect(checksFinished).toBe(true);
expect(doctorReaders?.active).toBe(false);
if (broker === "settlement-failed") {
throw new Error("fixture broker could not settle native cleanup");
}
// Native lifecycle tests cover real shared owners and SQLite locks. Here
// the census reflects their custody while exercising the complete Doctor flow.
if (broker === "independent-owner" || broker === "independent-thread") {
return false;
}
brokerLive = false;
return true;
});
custody.census.mockImplementation(() => {
assertMaintenanceHeld();
const maintenance = getOpenClawDatabaseMaintenanceScope();
expect(maintenance?.ownsSchemaMaintenance).toBe(true);
maintenance!.assertAdmission();
custody.activeNativeWork = broker === "late-independent-thread";
return {
pids: readSource
? readSource.hasActiveWorkers
? [42423]
: []
: [...(doctorReaders?.active ? [42421] : []), ...(brokerLive ? [42422] : [])],
};
});
mocks.runContributions.mockImplementation(async (ctx) => {
ownerIdentity = fs.readFileSync(ownerPath);
doctorReaders = readOnlyWorkerScope.getStore();
expect(doctorReaders?.active).toBe(true);
brokerLive = broker !== "independent-thread";
if (broker === "idle-read-pool") {
const { executeExistingOpenClawStateRead } =
await import("../state/openclaw-state-db-readonly.js");
const { captureRetainedNativeWorkerSource } =
await import("../infra/worker-native-lifecycle.js");
await executeExistingOpenClawStateRead(
{ path: database, env: state.env },
{ type: "fleet.list" },
);
readSource = captureRetainedNativeWorkerSource();
expect(readSource.hasActiveWorkers).toBe(true);
}
await runRetainedUpdateRuntimesHealth(ctx);
expect(fs.existsSync(artifact)).toBe(true);
checksFinished = true;
});
const runtime = { log: vi.fn(), error: vi.fn(), exit: vi.fn() };
const completed = runDoctorHealthFlow(runtime, { repair: true, nonInteractive: true });
if (broker === "settlement-failed") {
await expect(completed).rejects.toThrow("fixture broker could not settle native cleanup");
expect(custody.census).not.toHaveBeenCalled();
expect(fs.existsSync(artifact)).toBe(true);
return;
}
await completed;
const output = custody.note.mock.calls.map(([message]) => String(message)).join("\n");
expect(custody.census).toHaveBeenCalledOnce();
expect(fs.existsSync(artifact)).toBe(broker !== "idle" && broker !== "idle-read-pool");
if (broker === "independent-owner") {
expect(brokerLive).toBe(true);
expect(output).toContain("PIDs: 42422");
expect(output).toContain("let these holders finish, then rerun openclaw doctor --fix");
expect(output).not.toContain("Removed abandoned updater runtime");
} else if (broker === "independent-thread" || broker === "late-independent-thread") {
expect(output).toContain(`independent native work in this process (PID: ${process.pid})`);
expect(output).toContain("let these holders finish, then rerun openclaw doctor --fix");
expect(output).not.toContain("Removed abandoned updater runtime");
} else {
expect(output).toContain(`Removed abandoned updater runtime: ${artifact}`);
}
});
},
);

View file

@ -88,6 +88,7 @@ export async function maintainRetainedUpdateRuntimes(params: {
temporaryDirectories?: readonly string[];
repair: boolean;
assertCurrent: () => void;
assertResourcesSettled?: () => void;
}): Promise<string[]> {
const messages: string[] = [];
const packages = params.packageRoots.map(resolveRealpathOrAbsolute);
@ -146,9 +147,10 @@ export async function maintainRetainedUpdateRuntimes(params: {
}
if (census.pids.length) {
throw new Error(
`other OpenClaw processes are still running (PIDs: ${census.pids.join(", ")})`,
`other OpenClaw processes are still running (PIDs: ${census.pids.join(", ")}); let these holders finish, then rerun openclaw doctor --fix`,
);
}
params.assertResourcesSettled?.();
params.assertCurrent();
let failure: string | undefined;
await removeTemporaryArtifacts(directory, "Updater runtime", (error) => {

View file

@ -0,0 +1,61 @@
import assert from "node:assert/strict";
import type { DatabaseSync } from "node:sqlite";
import { createDeferredCore } from "../shared/deferred.js";
import {
captureRetainedNativeWorkerSource,
createRetainedNativeWorker,
type RetainedNativeWorkerSource,
} from "./worker-native-lifecycle.js";
export function assertNativeResourceCustody(
facts: { constructions: number; childClosed: boolean; childPid: number },
exited: boolean,
observer: DatabaseSync,
): void {
assert.equal(exited, false);
assert.equal(facts.constructions, 1);
assert.equal(facts.childClosed, false);
process.kill(facts.childPid, 0);
assert.throws(() => observer.exec("BEGIN IMMEDIATE"), /locked|busy/i);
}
export function createIdleBrokerRetirementProof(source: RetainedNativeWorkerSource) {
let domainClosed = false;
source.retain({}, async () => {
domainClosed = true;
});
return {
result: {
independentCustodyPreserved: true,
idleBrokerJoined: true,
sourceReusable: true,
},
async assertRefused(brokerPid: number, assertHeld: () => void) {
assert.equal(await source.retireIdleBroker(), false);
assert.equal(domainClosed, false);
process.kill(brokerPid, 0);
assertHeld();
},
async assertRetired(brokerPid: number) {
assert.equal(await source.retireIdleBroker(), true);
assert.throws(() => process.kill(brokerPid, 0), { code: "ESRCH" });
assert.equal(domainClosed, false);
assert.equal(captureRetainedNativeWorkerSource({ runtimeGeneration: undefined }), source);
const next = createRetainedNativeWorker(
'require("node:worker_threads").parentPort.postMessage(42)',
{ eval: true, execArgv: [] },
source,
);
assert.equal(source.hasActiveWorkers, true);
const nextReply = createDeferredCore<unknown>();
next.on("message", nextReply.resolve);
next.on("error", nextReply.reject);
try {
assert.equal(await nextReply.promise, 42);
} finally {
await next.terminate();
}
assert.equal(source.hasActiveWorkers, false);
},
};
}

View file

@ -57,6 +57,7 @@ export async function runNativeResourceLifecycle(
supervisorLoss = false,
closeBeforeLoss = false,
edge?:
| "idle-broker"
| "late-attachment"
| "owner-reply-loss"
| "auto-close-success"
@ -71,6 +72,8 @@ export async function runNativeResourceLifecycle(
]);
const { captureRetainedNativeWorkerSource, createRetainedNativeWorker } =
await import("./worker-native-lifecycle.js");
const { assertNativeResourceCustody, createIdleBrokerRetirementProof } =
await import("./worker-native-lifecycle.custody.test-support.js");
const { SpawnBrokerHost } = await import("../process/spawn-broker/host.js");
const { drainGlobalSingletonLifecycleState } = await import("../shared/global-singleton.js");
const autoCloseEdge =
@ -190,6 +193,8 @@ export async function runNativeResourceLifecycle(
const control = createControl();
const { channel, facts, order, waitFor, permit } = control;
const source = captureRetainedNativeWorkerSource({ runtimeGeneration: undefined });
const idleBrokerProof =
edge === "idle-broker" ? createIdleBrokerRetirementProof(source) : undefined;
const resource = source.captureResource(
resolveRuntimeWorkerUrl(nativeWorkerResourceEntrypoint),
"nativeResource",
@ -320,14 +325,11 @@ export async function runNativeResourceLifecycle(
const observer = new DatabaseSync(databasePath);
database = observer;
observer.exec("PRAGMA busy_timeout=0");
const assertHeld = () => {
assert.equal(exited, false);
assert.equal(facts.constructions, 1);
assert.equal(facts.childClosed, false);
process.kill(facts.childPid, 0);
assert.throws(() => observer.exec("BEGIN IMMEDIATE"), /locked|busy/i);
};
const assertHeld = () => assertNativeResourceCustody(facts, exited, observer);
assertHeld();
if (idleBrokerProof) {
await idleBrokerProof.assertRefused(facts.brokerPid, assertHeld);
}
if (edge === "late-attachment") {
// Cross the original cold broker's 15-second readiness deadline using real elapsed time.
await delay(15_050);
@ -565,6 +567,9 @@ export async function runNativeResourceLifecycle(
assert.equal(facts.childClosed, true);
assert.equal(exited, true);
assert.equal(target.threadId, -1);
if (idleBrokerProof) {
await idleBrokerProof.assertRetired(facts.brokerPid);
}
if (shutdown) {
await supervisorJoined.promise;
await nextTurn();
@ -636,18 +641,15 @@ export async function runNativeResourceLifecycle(
observer.exec("ROLLBACK");
console.log(
JSON.stringify({
ending: autoCloseEdge
ending: edge
? `resource-${edge}`
: edge === "late-attachment"
? "resource-late-attachment"
: edge === "owner-reply-loss"
? "resource-owner-reply-loss"
: closeBeforeLoss
? "resource-close-supervisor-loss"
: supervisorLoss
? "resource-supervisor-loss"
: "native-resource",
: closeBeforeLoss
? "resource-close-supervisor-loss"
: supervisorLoss
? "resource-supervisor-loss"
: "native-resource",
...(closeBeforeLoss ? { originalCloseJoined: true } : {}),
...idleBrokerProof?.result,
...(edge === "late-attachment" ? { lateSameBrokerAttached: true } : {}),
...(edge === "owner-reply-loss"
? {

View file

@ -622,6 +622,7 @@ assert.ok(
ending === "explicit-unbound" ||
ending === "supervisor-loss" ||
ending === "native-resource" ||
ending === "resource-idle-broker" ||
ending === "resource-supervisor-loss" ||
ending === "resource-auto-close-success" ||
ending === "resource-auto-close-failure" ||
@ -643,6 +644,8 @@ if (ending === "generation") {
await runSupervisorLoss();
} else if (ending === "native-resource") {
await runNativeResourceLifecycle(directory, serviceNativeUntil);
} else if (ending === "resource-idle-broker") {
await runNativeResourceLifecycle(directory, serviceNativeUntil, false, false, "idle-broker");
} else if (ending === "resource-cold-supervisor-loss") {
await runNativeColdRecovery(directory, serviceNativeUntil);
} else if (ending === "resource-supervisor-loss") {

View file

@ -15,6 +15,7 @@ async function runFixture(
| "explicit-unbound"
| "supervisor-loss"
| "native-resource"
| "resource-idle-broker"
| "resource-supervisor-loss"
| "resource-auto-close-success"
| "resource-auto-close-failure"
@ -43,6 +44,19 @@ async function runFixture(
}
describe("retained native worker lifecycle", () => {
it("retires an idle broker only after independent resource custody joins", async () => {
expect(await runFixture("resource-idle-broker")).toEqual({
ending: "resource-idle-broker",
independentCustodyPreserved: true,
idleBrokerJoined: true,
sourceReusable: true,
firstCloseRejected: true,
sameOwnerRetried: true,
childClosedBeforeStopped: true,
sqliteReusable: true,
});
}, 20_000);
it("preserves constructor ALS for callbacks serviced from another context", async () => {
const expected = ["message", "error", "exit"].map((event) => ({
event,

View file

@ -39,6 +39,7 @@ type NativeRuntime = NativeWorkerRuntime & {
};
export type RetainedNativeWorkerSource = {
readonly hasActiveWorkers: boolean;
create(
filename: string | URL,
options?: WorkerOptions,
@ -52,6 +53,8 @@ export type RetainedNativeWorkerSource = {
): NativeWorkerResourceDescriptor;
/** The domain joins admitted work and its resources before shared native cleanup. */
retain(owner: object, closeOwner: () => Promise<void>): void;
/** Join an idle broker without closing registered domains or their independent work. */
retireIdleBroker(): Promise<boolean>;
};
type NativeSource = RetainedNativeWorkerSource & {
@ -59,6 +62,7 @@ type NativeSource = RetainedNativeWorkerSource & {
runtimeGeneration?: RuntimeWorkerGeneration;
runtime?: NativeRuntime;
broker?: SpawnBrokerHost;
retiringBrokers: Set<Promise<void>>;
automaticBrokerClose?: Promise<void>;
brokerModuleUrl: URL;
closing: boolean;
@ -91,12 +95,29 @@ function forgetNativeSource(source: NativeSource): void {
}
function closeNativeBroker(source: NativeSource): Promise<void> {
let closing: Promise<void>;
try {
return source.broker?.close() ?? Promise.resolve();
closing = source.broker?.close() ?? Promise.resolve();
} catch (error) {
const failed = createDeferredCore();
failed.reject(error);
return failed.promise;
closing = failed.promise;
}
return source.retiringBrokers.size
? joinNativeBrokerCloses([...source.retiringBrokers, closing])
: closing;
}
async function joinNativeBrokerCloses(attempts: readonly Promise<void>[]): Promise<void> {
const results = await Promise.allSettled(attempts);
const failures = results.flatMap((result) =>
result.status === "rejected" ? [result.reason] : [],
);
if (failures.length === 1) {
throw failures[0];
}
if (failures.length > 1) {
throw new AggregateError(failures, "Native broker retirement failed");
}
}
@ -137,7 +158,7 @@ function nativeRuntime(source: NativeSource): NativeRuntime {
return;
}
source.runtime = undefined;
if (!source.broker) {
if (!source.broker && source.retiringBrokers.size === 0) {
forgetNativeSource(source);
return;
}
@ -313,6 +334,30 @@ export function captureRetainedNativeWorkerSource(options?: {
? captured.runtimeGeneration.resolve(resolveRuntimeProcessEntrypointUrl("spawnBroker"))
: resolveRuntimeProcessEntrypointUrl("spawnBroker"),
closing: false,
retiringBrokers: new Set(),
get hasActiveWorkers() {
return Boolean(source.runtime?.handles.size);
},
async retireIdleBroker() {
// Handles remain until their native execution and resource close receipts join.
if (source.hasActiveWorkers) {
return false;
}
if (source.broker) {
const closing = source.broker.close();
// New work retains this source but must never attach to a closing broker.
source.broker = undefined;
source.retiringBrokers.add(closing);
void closing.then(
() => source.retiringBrokers.delete(closing),
() => {
// Keep the original failed receipt reachable through source finalization.
},
);
}
await joinNativeBrokerCloses([...source.retiringBrokers]);
return true;
},
close() {
return (closingOwners ??= Promise.resolve().then(async () => {
const closing = [...owners.values()].map((owner) => owner.close());

View file

@ -5,10 +5,40 @@ import { AsyncLocalStorage } from "node:async_hooks";
import { expect, it, vi } from "vitest";
import { resolveRuntimeProcessEntrypointUrl } from "../infra/runtime-process-url.js";
import { withRuntimeWorkerGeneration } from "../infra/runtime-worker-generation.js";
import { captureRetainedNativeWorkerSource } from "../infra/worker-native-lifecycle.js";
import { createDeferredCore } from "../shared/deferred.js";
import { captureOpenClawStateReadSource } from "./openclaw-state-read-worker.js";
import { executeExistingOpenClawStateRead } from "./openclaw-state-db-readonly.js";
import {
captureOpenClawStateReadSource,
retireIdleOpenClawStateReadWorkers,
} from "./openclaw-state-read-worker.js";
import { captureOpenClawStateWorkerContext } from "./openclaw-state-worker-context.js";
it("retires cached readers after independent custody joins and permits reuse", async () => {
const { options } = source();
const nativeSource = captureRetainedNativeWorkerSource();
const task = queueTask();
const reading = executeExistingOpenClawStateRead(options, { type: "fleet.list" });
try {
await task.captured;
expect(await retireIdleOpenClawStateReadWorkers(nativeSource)).toBe(false);
expect(mock.closePool).not.toHaveBeenCalled();
} finally {
task.result.resolve(emptyReply);
await reading;
}
expect(mock.closePool).not.toHaveBeenCalled();
expect(await retireIdleOpenClawStateReadWorkers(nativeSource)).toBe(true);
expect(mock.closePool).toHaveBeenCalledOnce();
const next = queueTask();
next.result.resolve(emptyReply);
await expect(executeExistingOpenClawStateRead(options, { type: "fleet.list" })).resolves.toEqual(
emptyReply,
);
expect(mock.create).toHaveBeenCalledTimes(2);
});
it("keeps lazy reads with their captured generation when another generation dispatches them", async () => {
const { options } = source();
const context = captureOpenClawStateWorkerContext(options);

View file

@ -125,6 +125,21 @@ function readRuntimes() {
});
}
/** Retire cached workers only after every independently admitted reader has joined. */
export async function retireIdleOpenClawStateReadWorkers(
nativeSource: RetainedNativeWorkerSource,
): Promise<boolean> {
const state = readRuntimes().get(nativeSource);
if (!state) {
return true;
}
if (state.operations.size > 0) {
return false;
}
await closeReadPool(state);
return true;
}
function readPool(state: ReadRuntime, admitted: boolean): ReadPool {
// Accepted owners may still need a final read while generation admission is sealed.
if ((state.sealed && !admitted) || state.closing) {