diff --git a/docs/cli/doctor/checks.md b/docs/cli/doctor/checks.md index ed16a88ce9de..e4be654b2815 100644 --- a/docs/cli/doctor/checks.md +++ b/docs/cli/doctor/checks.md @@ -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. diff --git a/src/cli/update-cli/update-command-doctor-authority-callers.test.ts b/src/cli/update-cli/update-command-doctor-authority-callers.test.ts index 5b5eafd412a4..dc458efd1b22 100644 --- a/src/cli/update-cli/update-command-doctor-authority-callers.test.ts +++ b/src/cli/update-cli/update-command-doctor-authority-callers.test.ts @@ -358,6 +358,7 @@ describe("unproved Doctor authority callers", () => { signal: new AbortController().signal, run: (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 () => { diff --git a/src/cli/update-cli/update-command-lifecycle-repair.test-support.ts b/src/cli/update-cli/update-command-lifecycle-repair.test-support.ts index d0e271504c3a..ca9804048033 100644 --- a/src/cli/update-cli/update-command-lifecycle-repair.test-support.ts +++ b/src/cli/update-cli/update-command-lifecycle-repair.test-support.ts @@ -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 diff --git a/src/cli/update-cli/update-command-lifecycle.test.ts b/src/cli/update-cli/update-command-lifecycle.test.ts index eeda7b989ad0..4c25937b373c 100644 --- a/src/cli/update-cli/update-command-lifecycle.test.ts +++ b/src/cli/update-cli/update-command-lifecycle.test.ts @@ -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: (operation: () => T): T => operation(), repairSqliteNoCow: async () => {}, + cleanupRetainedRuntimes: async () => {}, releaseState: async () => { record("release-state"); }, @@ -1004,6 +1006,7 @@ describe("update plugin lifecycle lease boundaries", () => { run: (operation: () => T): T => operation(), release: async () => {}, repairSqliteNoCow: async () => {}, + cleanupRetainedRuntimes: async () => {}, releaseState: async () => {}, finish: async () => { warnings.push(warning); diff --git a/src/commands/doctor-config-preflight.state-migration.test-harness.ts b/src/commands/doctor-config-preflight.state-migration.test-harness.ts index 7e2ec28ce47f..ca7b43077bfa 100644 --- a/src/commands/doctor-config-preflight.state-migration.test-harness.ts +++ b/src/commands/doctor-config-preflight.state-migration.test-harness.ts @@ -72,6 +72,7 @@ const beginDoctorMaintenance = vi.hoisted(() => run: (operation: () => T): T => operation(), releaseState: vi.fn(async () => {}), repairSqliteNoCow: vi.fn(async () => {}), + cleanupRetainedRuntimes: vi.fn(async () => {}), release: doctorMaintenanceRelease, finish: vi.fn(async () => {}), })), diff --git a/src/commands/doctor-maintenance-state.ts b/src/commands/doctor-maintenance-state.ts index be3435606d1c..5fc111c21c04 100644 --- a/src/commands/doctor-maintenance-state.ts +++ b/src/commands/doctor-maintenance-state.ts @@ -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(); diff --git a/src/commands/doctor-maintenance-types.ts b/src/commands/doctor-maintenance-types.ts index 458067273f79..8665f1c5d299 100644 --- a/src/commands/doctor-maintenance-types.ts +++ b/src/commands/doctor-maintenance-types.ts @@ -25,6 +25,7 @@ export type DoctorMaintenance = { signal: AbortSignal; releaseState(): Promise; repairSqliteNoCow(paths: readonly string[]): Promise; + cleanupRetainedRuntimes(): Promise; release(): Promise; finish( cfg: OpenClawConfig | undefined, diff --git a/src/commands/doctor-maintenance.ts b/src/commands/doctor-maintenance.ts index 202141e8b4e2..658f00dca475 100644 --- a/src/commands/doctor-maintenance.ts +++ b/src/commands/doctor-maintenance.ts @@ -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."); diff --git a/src/commands/doctor-retained-runtime.ts b/src/commands/doctor-retained-runtime.ts index ed5a7307eacb..02e16b6c88c7 100644 --- a/src/commands/doctor-retained-runtime.ts +++ b/src/commands/doctor-retained-runtime.ts @@ -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 { + 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 => { + 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"); + } + }; } diff --git a/src/commands/doctor/shared/temporary-directories.ts b/src/commands/doctor/shared/temporary-directories.ts index 8602078393f4..239eaea386c0 100644 --- a/src/commands/doctor/shared/temporary-directories.ts +++ b/src/commands/doctor/shared/temporary-directories.ts @@ -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, { diff --git a/src/flows/doctor-health-contribution-runners.state.ts b/src/flows/doctor-health-contribution-runners.state.ts index e28647cce5c2..15c674922ea4 100644 --- a/src/flows/doctor-health-contribution-runners.state.ts +++ b/src/flows/doctor-health-contribution-runners.state.ts @@ -44,8 +44,13 @@ export async function runLegacyPluginSourceCapturesHealth( } export async function runRetainedUpdateRuntimesHealth(ctx: DoctorHealthFlowContext): Promise { - 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( diff --git a/src/flows/doctor-health.migration-refusal.test.ts b/src/flows/doctor-health.migration-refusal.test.ts index 053fb33737a0..b9fd81796b61 100644 --- a/src/flows/doctor-health.migration-refusal.test.ts +++ b/src/flows/doctor-health.migration-refusal.test.ts @@ -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( diff --git a/src/flows/doctor-health.ts b/src/flows/doctor-health.ts index 4685af1e97e9..eefcd1e05958 100644 --- a/src/flows/doctor-health.ts +++ b/src/flows/doctor-health.ts @@ -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; diff --git a/src/flows/doctor-retained-runtime.settlement.test.ts b/src/flows/doctor-retained-runtime.settlement.test.ts new file mode 100644 index 000000000000..aa640817c291 --- /dev/null +++ b/src/flows/doctor-retained-runtime.settlement.test.ts @@ -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>(), + activeNativeWork: false, + actualNativeSource: false, +})); + +vi.mock("../../packages/terminal-core/src/note.js", async (importOriginal) => ({ + ...(await importOriginal()), + 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(); + type NativeSource = ReturnType; + const sources = new WeakMap(); + return { + ...actual, + captureRetainedNativeWorkerSource: ( + ...args: Parameters + ) => { + 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}`); + } + }); + }, +); diff --git a/src/infra/temp-artifact-cleanup.ts b/src/infra/temp-artifact-cleanup.ts index dde84796c258..29f691f2535a 100644 --- a/src/infra/temp-artifact-cleanup.ts +++ b/src/infra/temp-artifact-cleanup.ts @@ -88,6 +88,7 @@ export async function maintainRetainedUpdateRuntimes(params: { temporaryDirectories?: readonly string[]; repair: boolean; assertCurrent: () => void; + assertResourcesSettled?: () => void; }): Promise { 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) => { diff --git a/src/infra/worker-native-lifecycle.custody.test-support.ts b/src/infra/worker-native-lifecycle.custody.test-support.ts new file mode 100644 index 000000000000..64baeb826156 --- /dev/null +++ b/src/infra/worker-native-lifecycle.custody.test-support.ts @@ -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(); + 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); + }, + }; +} diff --git a/src/infra/worker-native-lifecycle.runtime.test-support.ts b/src/infra/worker-native-lifecycle.runtime.test-support.ts index 56f534bfa9e5..0aed0cc8a1fe 100644 --- a/src/infra/worker-native-lifecycle.runtime.test-support.ts +++ b/src/infra/worker-native-lifecycle.runtime.test-support.ts @@ -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" ? { diff --git a/src/infra/worker-native-lifecycle.test-support.ts b/src/infra/worker-native-lifecycle.test-support.ts index 8872de48417d..e05d79899943 100644 --- a/src/infra/worker-native-lifecycle.test-support.ts +++ b/src/infra/worker-native-lifecycle.test-support.ts @@ -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") { diff --git a/src/infra/worker-native-lifecycle.test.ts b/src/infra/worker-native-lifecycle.test.ts index e517f74f7459..64d45de5bef6 100644 --- a/src/infra/worker-native-lifecycle.test.ts +++ b/src/infra/worker-native-lifecycle.test.ts @@ -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, diff --git a/src/infra/worker-native-lifecycle.ts b/src/infra/worker-native-lifecycle.ts index d2f3e3815881..eb867ebbfa10 100644 --- a/src/infra/worker-native-lifecycle.ts +++ b/src/infra/worker-native-lifecycle.ts @@ -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; + /** Join an idle broker without closing registered domains or their independent work. */ + retireIdleBroker(): Promise; }; type NativeSource = RetainedNativeWorkerSource & { @@ -59,6 +62,7 @@ type NativeSource = RetainedNativeWorkerSource & { runtimeGeneration?: RuntimeWorkerGeneration; runtime?: NativeRuntime; broker?: SpawnBrokerHost; + retiringBrokers: Set>; automaticBrokerClose?: Promise; brokerModuleUrl: URL; closing: boolean; @@ -91,12 +95,29 @@ function forgetNativeSource(source: NativeSource): void { } function closeNativeBroker(source: NativeSource): Promise { + let closing: Promise; 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[]): Promise { + 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()); diff --git a/src/state/openclaw-state-read-source.test.ts b/src/state/openclaw-state-read-source.test.ts index fd3195690323..0b8fe4525545 100644 --- a/src/state/openclaw-state-read-source.test.ts +++ b/src/state/openclaw-state-read-source.test.ts @@ -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); diff --git a/src/state/openclaw-state-read-worker.ts b/src/state/openclaw-state-read-worker.ts index 952283f4cefb..62518e792081 100644 --- a/src/state/openclaw-state-read-worker.ts +++ b/src/state/openclaw-state-read-worker.ts @@ -125,6 +125,21 @@ function readRuntimes() { }); } +/** Retire cached workers only after every independently admitted reader has joined. */ +export async function retireIdleOpenClawStateReadWorkers( + nativeSource: RetainedNativeWorkerSource, +): Promise { + 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) {