From 3f33e4f0e949ca050494cf3339ce6fe137b95348 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 10:01:47 +0900 Subject: [PATCH 1/3] fix(server): self-heal a replaced package tree via drain-and-restart An out-of-band package-manager upgrade (npm/bun/pnpm global install) under a live proxy replaced the package manifest and fenced every /v1/* request behind 503 package_tree_changed until a manual restart. The integrity guard now accepts an onReplaced callback: after a replaced observation persists past a short debounce (and resets if the manifest goes unreadable mid-install or recovers), the server accepts the existing graceful drain-and-restart path, so the new tree comes up on its own. Source checkouts and standalone binaries still never fence, so development and single-file installs are unaffected. --- src/lib/package-tree-integrity.ts | 44 +++++++-- src/server/index.ts | 15 ++- .../package-tree-integrity.test.ts | 96 +++++++++++++++++++ 3 files changed, 148 insertions(+), 7 deletions(-) diff --git a/src/lib/package-tree-integrity.ts b/src/lib/package-tree-integrity.ts index 30372de183e..45fc75e6e6b 100644 --- a/src/lib/package-tree-integrity.ts +++ b/src/lib/package-tree-integrity.ts @@ -16,6 +16,20 @@ export interface PackageTreeIntegrityGuard { status(): PackageTreeIntegrityStatus; } +export interface PackageTreeIntegrityOptions { + /** + * Called once when a replaced package tree persists past `replacedRestartDelayMs` + * of sustained failure. The intended handler is the graceful drain-and-restart + * acceptor: an out-of-band install (npm/bun/pnpm global upgrade under a live + * proxy) then self-heals instead of serving 503s until someone restarts by hand. + * Only `package_tree_replaced` counts — an unreadable manifest resets the timer, + * so an install still mid-write does not trigger a restart on partial state. + */ + onReplaced?: () => void; + /** Sustained-replacement delay before `onReplaced` fires. 0 fires on first detection. */ + replacedRestartDelayMs?: number; +} + type ObservePackageTree = () => PackageTreeObservation | null; type PackageTreeRuntimeInstall = "bun" | "npm" | "pnpm" | "source"; @@ -70,17 +84,34 @@ const PACKAGE_TREE_RECHECK_MS = 1_000; export function createPackageTreeIntegrityGuard( observe: ObservePackageTree = observePackageManifest, now: () => number = Date.now, + options: PackageTreeIntegrityOptions = {}, ): PackageTreeIntegrityGuard { const boot = observe(); let lastOkAt: number | null = null; + let replacedSinceMs: number | null = null; + let notified = false; + const restartDelayMs = options.replacedRestartDelayMs ?? 5_000; return { status(): PackageTreeIntegrityStatus { const at = now(); if (lastOkAt !== null && at - lastOkAt < PACKAGE_TREE_RECHECK_MS) return { ok: true }; const current = observe(); - if (boot === null || current === null) return { ok: false, reason: "package_tree_unreadable" }; - if (!sameObservation(boot, current)) return { ok: false, reason: "package_tree_replaced" }; + if (boot === null || current === null) { + replacedSinceMs = null; + return { ok: false, reason: "package_tree_unreadable" }; + } + if (!sameObservation(boot, current)) { + if (options.onReplaced && !notified) { + if (replacedSinceMs === null) replacedSinceMs = at; + if (at - replacedSinceMs >= restartDelayMs) { + notified = true; + options.onReplaced(); + } + } + return { ok: false, reason: "package_tree_replaced" }; + } lastOkAt = at; + replacedSinceMs = null; return { ok: true }; }, }; @@ -96,7 +127,8 @@ export function createRuntimePackageTreeIntegrityGuard( installer: PackageTreeRuntimeInstall, observe: ObservePackageTree = observePackageManifest, now: () => number = Date.now, -): PackageTreeIntegrityGuard { - if (installer === "source" || isStandaloneBinary()) return { status: () => ({ ok: true }) }; - return createPackageTreeIntegrityGuard(observe, now); -} + options: PackageTreeIntegrityOptions = {}, + ): PackageTreeIntegrityGuard { + if (installer === "source" || isStandaloneBinary()) return { status: () => ({ ok: true }) }; + return createPackageTreeIntegrityGuard(observe, now, options); + } diff --git a/src/server/index.ts b/src/server/index.ts index 61fbcec7d8a..9cc8c2876da 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -178,6 +178,7 @@ import { import { EXTERNAL_CALL_PREFIX, LiveCallBindings } from "./live-call-bindings"; import { contextEndpoint, contextRelayActivated } from "../codex/context-compat"; import { fetchAllModels, handleManagementAPI, VERSION, type ManagementApiDeps } from "./management-api"; +import { acceptSystemRestart } from "./management/system-restart"; import { createManagementSessionControl, initializeManagementAuthState, @@ -538,7 +539,19 @@ function startServerWithSpendLedgerOwner(port: number | undefined, deps: StartSe // no gate is supplied (tests, ad-hoc starts) a fresh pending gate is created. const readinessGate = deps.readinessGate ?? createReadinessGate(); const packageTreeIntegrity = deps.packageTreeIntegrity - ?? createRuntimePackageTreeIntegrityGuard(detectInstall()); + ?? createRuntimePackageTreeIntegrityGuard(detectInstall(), undefined, undefined, { + onReplaced: () => { + // An out-of-band install replaced the package under this live process. Serve + // the 503 for the triggering request, then let the standard drain-and-restart + // path bring the new tree up instead of refusing traffic until a manual + // restart. acceptSystemRestart is idempotent and supervisored-aware. + try { + acceptSystemRestart(); + } catch (error) { + console.warn("Package tree changed; automatic drain-and-restart failed:", error instanceof Error ? error.message : error); + } + }, + }); // Actual bound port, filled in after Bun.serve binds so /readyz reports the // real ephemeral port for startServer(0). /healthz keeps its existing port // field (the requested listenPort) byte-for-byte. diff --git a/tests/ci-workflows/package-tree-integrity.test.ts b/tests/ci-workflows/package-tree-integrity.test.ts index 13da9e02e7c..3912e3c936a 100644 --- a/tests/ci-workflows/package-tree-integrity.test.ts +++ b/tests/ci-workflows/package-tree-integrity.test.ts @@ -153,6 +153,102 @@ describe("package tree integrity", () => { expect(guard.status()).toEqual({ ok: false, reason: "package_tree_unreadable" }); }); + describe("automatic restart on a replaced package tree", () => { + const base: PackageTreeObservation = { + device: 1n, inode: 10n, contentTimeNs: 100n, size: 500n, + }; + + test("fires once after the replacement persists past the delay", () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 5_000 }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(calls).toBe(0); // inside the debounce window + + clock += 5_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(calls).toBe(1); + + // Idempotent: further refusals never re-arm the restart. + clock += 10_000; + guard.status(); + expect(calls).toBe(1); + }); + + test("a zero delay fires on first detection", () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 0 }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(calls).toBe(1); + }); + + test("recovery or an unreadable manifest resets the debounce", () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 5_000 }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + + // Mid-install the manifest briefly disappears: the timer must restart, + // not fire on a partially written tree. + observation = null; + clock += 4_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_unreadable" }); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + clock += 4_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(calls).toBe(0); + clock += 5_000; + guard.status(); + expect(calls).toBe(1); + }); + + test("source checkouts never auto-restart", () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const guard = createRuntimePackageTreeIntegrityGuard( + "source", + () => observation, + () => clock, + { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 0 }, + ); + + observation = { ...base, inode: 11n }; + clock += 60_000; + expect(guard.status()).toEqual({ ok: true }); + expect(calls).toBe(0); + }); + }); + // BUG-R1: a chmod fenced the whole data plane behind 503. // // These three drive the REAL filesystem rather than a hand-built observation, From 4eac35c04587895a49bab370bf0571d4de5e6bd7 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 11:08:09 +0900 Subject: [PATCH 2/3] fix(server): make package-tree recovery timer-driven and retryable Drive recovery from an unref'd stability timer so a replaced package tree restarts without a second request. Require the same readable replacement identity for the full debounce, retry restart admission failures, cover server wiring and baseline recovery, and document the lifecycle contract. --- src/lib/package-tree-integrity.ts | 90 ++++++++- src/server/index.ts | 27 +-- src/server/index/startup-warnings.ts | 13 ++ structure/runtime.md | 9 + .../package-tree-integrity.test.ts | 189 +++++++++++++++--- 5 files changed, 279 insertions(+), 49 deletions(-) diff --git a/src/lib/package-tree-integrity.ts b/src/lib/package-tree-integrity.ts index 45fc75e6e6b..6049404dbb8 100644 --- a/src/lib/package-tree-integrity.ts +++ b/src/lib/package-tree-integrity.ts @@ -28,10 +28,12 @@ export interface PackageTreeIntegrityOptions { onReplaced?: () => void; /** Sustained-replacement delay before `onReplaced` fires. 0 fires on first detection. */ replacedRestartDelayMs?: number; + /** Test seam; production uses an unref'd timer. */ + schedule?: (callback: () => void, delayMs: number) => void; } -type ObservePackageTree = () => PackageTreeObservation | null; -type PackageTreeRuntimeInstall = "bun" | "npm" | "pnpm" | "source"; +export type ObservePackageTree = () => PackageTreeObservation | null; +export type PackageTreeRuntimeInstall = "bun" | "npm" | "pnpm" | "source"; const packageManifestUrl = new URL("../../package.json", import.meta.url); @@ -88,30 +90,96 @@ export function createPackageTreeIntegrityGuard( ): PackageTreeIntegrityGuard { const boot = observe(); let lastOkAt: number | null = null; - let replacedSinceMs: number | null = null; let notified = false; + let timerGeneration = 0; + let timerScheduled = false; + let waitingForReadableTree = false; + let replacementCandidate: PackageTreeObservation | null = null; const restartDelayMs = options.replacedRestartDelayMs ?? 5_000; + const schedule = options.schedule ?? ((callback, delayMs) => { + const timer = setTimeout(callback, delayMs); + timer.unref?.(); + }); + + const resetRestartTimer = (): void => { + timerGeneration += 1; + timerScheduled = false; + }; + + const armRestartTimer = (delayMs = restartDelayMs): void => { + if (!options.onReplaced || notified || timerScheduled) return; + timerScheduled = true; + const generation = timerGeneration; + const verifyAndNotify = () => { + if (generation !== timerGeneration || notified) return; + timerScheduled = false; + const current = observe(); + if (boot === null || current === null) { + // A package manager may replace package.json before the rest of the tree. + // Wait for a readable tree, then require a fresh full debounce interval. + resetRestartTimer(); + waitingForReadableTree = true; + armRestartTimer(PACKAGE_TREE_RECHECK_MS); + return; + } + if (sameObservation(boot, current)) { + resetRestartTimer(); + waitingForReadableTree = false; + replacementCandidate = null; + return; + } + if (waitingForReadableTree) { + waitingForReadableTree = false; + replacementCandidate = current; + resetRestartTimer(); + armRestartTimer(); + return; + } + if (replacementCandidate === null || !sameObservation(replacementCandidate, current)) { + replacementCandidate = current; + resetRestartTimer(); + armRestartTimer(); + return; + } + try { + options.onReplaced?.(); + notified = true; + } catch { + // A failed restart admission must not leave the proxy fenced forever. + // Re-observe after the normal debounce and try again if replacement persists. + armRestartTimer(Math.max(PACKAGE_TREE_RECHECK_MS, restartDelayMs)); + } + }; + if (delayMs === 0) verifyAndNotify(); + else schedule(verifyAndNotify, delayMs); + }; + return { status(): PackageTreeIntegrityStatus { const at = now(); if (lastOkAt !== null && at - lastOkAt < PACKAGE_TREE_RECHECK_MS) return { ok: true }; const current = observe(); if (boot === null || current === null) { - replacedSinceMs = null; + const wasWatchingReplacement = timerScheduled; + resetRestartTimer(); + if (wasWatchingReplacement && boot !== null) { + waitingForReadableTree = true; + armRestartTimer(PACKAGE_TREE_RECHECK_MS); + } return { ok: false, reason: "package_tree_unreadable" }; } if (!sameObservation(boot, current)) { - if (options.onReplaced && !notified) { - if (replacedSinceMs === null) replacedSinceMs = at; - if (at - replacedSinceMs >= restartDelayMs) { - notified = true; - options.onReplaced(); - } + if (replacementCandidate === null || !sameObservation(replacementCandidate, current)) { + replacementCandidate = current; + resetRestartTimer(); } + armRestartTimer(); return { ok: false, reason: "package_tree_replaced" }; } lastOkAt = at; - replacedSinceMs = null; + resetRestartTimer(); + waitingForReadableTree = false; + replacementCandidate = null; return { ok: true }; }, }; diff --git a/src/server/index.ts b/src/server/index.ts index 9cc8c2876da..8205a00c10d 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -538,20 +538,23 @@ function startServerWithSpendLedgerOwner(port: number | undefined, deps: StartSe // passes it in, and transitions it after the post-startup sync settles. When // no gate is supplied (tests, ad-hoc starts) a fresh pending gate is created. const readinessGate = deps.readinessGate ?? createReadinessGate(); + const acceptPackageTreeRestart = deps.acceptSystemRestart ?? acceptSystemRestart; const packageTreeIntegrity = deps.packageTreeIntegrity - ?? createRuntimePackageTreeIntegrityGuard(detectInstall(), undefined, undefined, { - onReplaced: () => { - // An out-of-band install replaced the package under this live process. Serve - // the 503 for the triggering request, then let the standard drain-and-restart - // path bring the new tree up instead of refusing traffic until a manual - // restart. acceptSystemRestart is idempotent and supervisored-aware. - try { - acceptSystemRestart(); - } catch (error) { - console.warn("Package tree changed; automatic drain-and-restart failed:", error instanceof Error ? error.message : error); - } + ?? createRuntimePackageTreeIntegrityGuard( + deps.packageTreeInstaller ?? detectInstall(), + deps.observePackageTree, + undefined, + { + ...deps.packageTreeIntegrityOptions, + onReplaced: () => { + // An out-of-band install replaced the package under this live process. Serve + // the 503 for the triggering request, then let the standard drain-and-restart + // path bring the new tree up instead of refusing traffic until a manual + // restart. acceptSystemRestart is idempotent and supervisor-aware. + acceptPackageTreeRestart(); + }, }, - }); + ); // Actual bound port, filled in after Bun.serve binds so /readyz reports the // real ephemeral port for startServer(0). /healthz keeps its existing port // field (the requested listenPort) byte-for-byte. diff --git a/src/server/index/startup-warnings.ts b/src/server/index/startup-warnings.ts index 76b9119f735..ebc7486d742 100644 --- a/src/server/index/startup-warnings.ts +++ b/src/server/index/startup-warnings.ts @@ -6,6 +6,11 @@ import { type OwnershipInspection, } from "../../integrations/native/ownership-preflight"; import { registerCodexQuotaAutoRefreshWorker } from "../../codex/quota-auto-refresh"; +import type { + ObservePackageTree, + PackageTreeIntegrityOptions, + PackageTreeRuntimeInstall, +} from "../../lib/package-tree-integrity"; import { consumeForInspection, relaySseWithHeartbeat, @@ -141,6 +146,14 @@ export interface StartServerDeps { readinessGate?: ReadinessGate; /** Test-only package-tree observation; production captures package.json identity at boot. */ packageTreeIntegrity?: PackageTreeIntegrityGuard; + /** Test-only default-guard options; production observes the installed package manifest. */ + packageTreeIntegrityOptions?: PackageTreeIntegrityOptions; + /** Test-only installed-package identity; production detects the current install. */ + packageTreeInstaller?: PackageTreeRuntimeInstall; + /** Test-only manifest observer; production stats the installed package.json. */ + observePackageTree?: ObservePackageTree; + /** Test-only restart acceptor; production uses the normal drain-and-restart path. */ + acceptSystemRestart?: typeof import("../management/system-restart").acceptSystemRestart; /** Test-only seam for observing quota-worker registration ownership. */ registerCodexQuotaAutoRefreshWorker?: typeof registerCodexQuotaAutoRefreshWorker; } diff --git a/structure/runtime.md b/structure/runtime.md index 384b6b394e6..84a852e4397 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -186,6 +186,15 @@ until shutdown. Normal shutdown restores native Codex. Service mode sets `OCX_SERVICE=1`, so managed restarts do not repeatedly restore/reinject; explicit service stop and uninstall still restore. +Installed npm, bun, and pnpm packages bind each server process to the package manifest identity +observed at startup. Replacing that manifest under a live process fences `/healthz`, `/readyz`, and +`/v1/*` with `package_tree_changed`. The first observed replacement starts an unref'd five-second +stability timer; if the same new manifest identity remains readable and distinct, the timer enters the existing +drain-and-restart handoff without waiting for another request. A temporarily unreadable manifest +is polled until readable and then receives a fresh full stability interval, while a return to the +startup identity cancels the pending restart. Failed restart admission retries after the same +bounded delay. Source checkouts and standalone binaries remain outside this integrity fence. + A busy preferred port is never resolved by starting somewhere else. Both questions a start asks about an existing proxy — the pre-bind owner check and the port-is-busy check in `src/cli/index.ts` — are identity probes with a retry budget, because a start that answers "nobody is there" on one diff --git a/tests/ci-workflows/package-tree-integrity.test.ts b/tests/ci-workflows/package-tree-integrity.test.ts index 3912e3c936a..e0dfc80592e 100644 --- a/tests/ci-workflows/package-tree-integrity.test.ts +++ b/tests/ci-workflows/package-tree-integrity.test.ts @@ -157,95 +157,191 @@ describe("package tree integrity", () => { const base: PackageTreeObservation = { device: 1n, inode: 10n, contentTimeNs: 100n, size: 500n, }; + const createScheduler = () => { + const pending: Array<() => void> = []; + return { + pending, + schedule: (callback: () => void) => { pending.push(callback); }, + runNext: () => { + const callback = pending.shift(); + if (!callback) throw new Error("expected a scheduled callback"); + callback(); + }, + }; + }; - test("fires once after the replacement persists past the delay", () => { + test("fires once from its timer without waiting for another request", () => { let observation: PackageTreeObservation | null = base; let clock = 0; let calls = 0; + const scheduler = createScheduler(); const guard = createPackageTreeIntegrityGuard( () => observation, () => clock, - { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 5_000 }, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 5_000, + schedule: scheduler.schedule, + }, ); expect(guard.status()).toEqual({ ok: true }); observation = { ...base, inode: 11n, contentTimeNs: 200n }; - clock += 2_000; expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); - expect(calls).toBe(0); // inside the debounce window + expect(calls).toBe(0); - clock += 5_000; - expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + scheduler.runNext(); expect(calls).toBe(1); + expect(scheduler.pending).toHaveLength(0); - // Idempotent: further refusals never re-arm the restart. clock += 10_000; guard.status(); expect(calls).toBe(1); }); - test("a zero delay fires on first detection", () => { + test("retries when restart acceptance throws", () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let attempts = 0; + const scheduler = createScheduler(); + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { + onReplaced: () => { + attempts += 1; + if (attempts === 1) throw new Error("restart unavailable"); + }, + replacedRestartDelayMs: 5_000, + schedule: scheduler.schedule, + }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + scheduler.runNext(); + expect(attempts).toBe(1); + expect(scheduler.pending).toHaveLength(1); + scheduler.runNext(); + expect(attempts).toBe(2); + expect(scheduler.pending).toHaveLength(0); + }); + + test("baseline recovery cancels the old timer and starts a fresh debounce", () => { let observation: PackageTreeObservation | null = base; let clock = 0; let calls = 0; + const scheduler = createScheduler(); const guard = createPackageTreeIntegrityGuard( () => observation, () => clock, - { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 0 }, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 5_000, + schedule: scheduler.schedule, + }, ); expect(guard.status()).toEqual({ ok: true }); - observation = { ...base, inode: 11n, contentTimeNs: 200n }; + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + + observation = base; + clock += 2_000; + expect(guard.status()).toEqual({ ok: true }); + scheduler.runNext(); // stale generation from the first replacement + expect(calls).toBe(0); + + observation = { ...base, inode: 12n }; clock += 2_000; expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(calls).toBe(0); + scheduler.runNext(); expect(calls).toBe(1); }); - test("recovery or an unreadable manifest resets the debounce", () => { + test("an unreadable manifest must become readable before a fresh debounce", () => { let observation: PackageTreeObservation | null = base; let clock = 0; let calls = 0; + const scheduler = createScheduler(); const guard = createPackageTreeIntegrityGuard( () => observation, () => clock, - { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 5_000 }, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 5_000, + schedule: scheduler.schedule, + }, ); expect(guard.status()).toEqual({ ok: true }); - observation = { ...base, inode: 11n, contentTimeNs: 200n }; + observation = { ...base, inode: 11n }; clock += 2_000; expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); - - // Mid-install the manifest briefly disappears: the timer must restart, - // not fire on a partially written tree. observation = null; - clock += 4_000; - expect(guard.status()).toEqual({ ok: false, reason: "package_tree_unreadable" }); - observation = { ...base, inode: 11n, contentTimeNs: 200n }; - clock += 4_000; + scheduler.runNext(); + expect(calls).toBe(0); + expect(scheduler.pending).toHaveLength(1); + + observation = { ...base, inode: 11n }; + scheduler.runNext(); // readability poll; arms a fresh full debounce + expect(calls).toBe(0); + expect(scheduler.pending).toHaveLength(1); + scheduler.runNext(); + expect(calls).toBe(1); + }); + + test("a second replacement identity receives its own full debounce", () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const scheduler = createScheduler(); + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 5_000, + schedule: scheduler.schedule, + }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + observation = { ...base, inode: 12n, contentTimeNs: 300n }; + scheduler.runNext(); expect(calls).toBe(0); - clock += 5_000; - guard.status(); + expect(scheduler.pending).toHaveLength(1); + scheduler.runNext(); expect(calls).toBe(1); }); test("source checkouts never auto-restart", () => { let observation: PackageTreeObservation | null = base; - let clock = 0; let calls = 0; + const scheduler = createScheduler(); const guard = createRuntimePackageTreeIntegrityGuard( "source", () => observation, - () => clock, - { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 0 }, + Date.now, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 0, + schedule: scheduler.schedule, + }, ); observation = { ...base, inode: 11n }; - clock += 60_000; expect(guard.status()).toEqual({ ok: true }); expect(calls).toBe(0); + expect(scheduler.pending).toHaveLength(0); }); }); @@ -368,4 +464,45 @@ describe("package tree integrity", () => { await server.stop(true); } }); + + test("the default server guard accepts a restart after a sustained replacement", async () => { + saveConfig(config()); + const base: PackageTreeObservation = { + device: 1n, inode: 10n, contentTimeNs: 100n, size: 500n, + }; + let observation: PackageTreeObservation = base; + const pending: Array<() => void> = []; + let restartAcceptances = 0; + const server = startServer(0, { + packageTreeInstaller: "npm", + observePackageTree: () => observation, + packageTreeIntegrityOptions: { + replacedRestartDelayMs: 5_000, + schedule: callback => { pending.push(callback); }, + }, + acceptSystemRestart: () => { + restartAcceptances += 1; + return { + accepted: true, + alreadyDraining: false, + activeTurnCount: 0, + drainTimeoutMs: 60_000, + }; + }, + }); + try { + expect((await fetch(new URL("/healthz", server.url))).status).toBe(200); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + await Bun.sleep(1_100); // expire the guard's successful-observation cache + expect((await fetch(new URL("/healthz", server.url))).status).toBe(503); + expect(restartAcceptances).toBe(0); + expect(pending).toHaveLength(1); + + pending.shift()?.(); + expect(restartAcceptances).toBe(1); + expect(pending).toHaveLength(0); + } finally { + await server.stop(true); + } + }); }); From 89ec3186191875f56d737ac4f971cca338afdb1d Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 12:11:28 +0900 Subject: [PATCH 3/3] fix(server): disarm the package-tree restart timer on reset and shutdown --- src/lib/package-tree-integrity.ts | 37 +++++++- src/server/index.ts | 4 + .../package-tree-integrity.test.ts | 90 ++++++++++++++++++- 3 files changed, 125 insertions(+), 6 deletions(-) diff --git a/src/lib/package-tree-integrity.ts b/src/lib/package-tree-integrity.ts index 6049404dbb8..bff1366311e 100644 --- a/src/lib/package-tree-integrity.ts +++ b/src/lib/package-tree-integrity.ts @@ -14,6 +14,13 @@ export type PackageTreeIntegrityStatus = export interface PackageTreeIntegrityGuard { status(): PackageTreeIntegrityStatus; + /** + * Permanently disarms the guard: cancels any pending restart timer and + * invalidates queued callbacks. Called from `server.stop()` so a still-queued + * replacement callback cannot schedule a drain-and-restart after shutdown + * has already begun. + */ + dispose(): void; } export interface PackageTreeIntegrityOptions { @@ -28,8 +35,12 @@ export interface PackageTreeIntegrityOptions { onReplaced?: () => void; /** Sustained-replacement delay before `onReplaced` fires. 0 fires on first detection. */ replacedRestartDelayMs?: number; - /** Test seam; production uses an unref'd timer. */ - schedule?: (callback: () => void, delayMs: number) => void; + /** + * Test seam; production uses an unref'd timer. May return a cancellation + * function; when it does, `resetRestartTimer` cancels the pending callback + * instead of leaving it queued behind a generation check. + */ + schedule?: (callback: () => void, delayMs: number) => (() => void) | void; } export type ObservePackageTree = () => PackageTreeObservation | null; @@ -93,17 +104,22 @@ export function createPackageTreeIntegrityGuard( let notified = false; let timerGeneration = 0; let timerScheduled = false; + let cancelScheduled: (() => void) | null = null; let waitingForReadableTree = false; let replacementCandidate: PackageTreeObservation | null = null; const restartDelayMs = options.replacedRestartDelayMs ?? 5_000; const schedule = options.schedule ?? ((callback, delayMs) => { const timer = setTimeout(callback, delayMs); timer.unref?.(); + return () => clearTimeout(timer); }); const resetRestartTimer = (): void => { timerGeneration += 1; timerScheduled = false; + const cancel = cancelScheduled; + cancelScheduled = null; + cancel?.(); }; const armRestartTimer = (delayMs = restartDelayMs): void => { @@ -151,10 +167,21 @@ export function createPackageTreeIntegrityGuard( } }; if (delayMs === 0) verifyAndNotify(); - else schedule(verifyAndNotify, delayMs); + else { + // The seam may run the callback synchronously; only keep its cancel + // function when this generation is still the live one afterwards. + const cancel = schedule(verifyAndNotify, delayMs); + if (generation === timerGeneration && timerScheduled && typeof cancel === "function") { + cancelScheduled = cancel; + } + } }; return { + dispose(): void { + resetRestartTimer(); + notified = true; + }, status(): PackageTreeIntegrityStatus { const at = now(); if (lastOkAt !== null && at - lastOkAt < PACKAGE_TREE_RECHECK_MS) return { ok: true }; @@ -197,6 +224,8 @@ export function createRuntimePackageTreeIntegrityGuard( now: () => number = Date.now, options: PackageTreeIntegrityOptions = {}, ): PackageTreeIntegrityGuard { - if (installer === "source" || isStandaloneBinary()) return { status: () => ({ ok: true }) }; + if (installer === "source" || isStandaloneBinary()) { + return { status: () => ({ ok: true }), dispose: () => {} }; + } return createPackageTreeIntegrityGuard(observe, now, options); } diff --git a/src/server/index.ts b/src/server/index.ts index 8205a00c10d..539f68c9ec9 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -773,6 +773,10 @@ function startServerWithSpendLedgerOwner(port: number | undefined, deps: StartSe value: async (closeActiveConnections?: boolean): Promise => { remoteWorkspaceStopping = true; liveCallBindings.clear(); + // Disarm the package-tree restart timer before listener teardown: a queued + // replacement callback must not call acceptSystemRestart() after stop() has + // begun, or it would schedule a drain-and-restart on a stopped server. + packageTreeIntegrity.dispose(); // The orchestration lives in `runListenerShutdown` so its two competing properties — // cleanup completes, failure propagates — are testable without a live socket. await runListenerShutdown( diff --git a/tests/ci-workflows/package-tree-integrity.test.ts b/tests/ci-workflows/package-tree-integrity.test.ts index e0dfc80592e..b2f14ff043c 100644 --- a/tests/ci-workflows/package-tree-integrity.test.ts +++ b/tests/ci-workflows/package-tree-integrity.test.ts @@ -161,7 +161,13 @@ describe("package tree integrity", () => { const pending: Array<() => void> = []; return { pending, - schedule: (callback: () => void) => { pending.push(callback); }, + schedule: (callback: () => void) => { + pending.push(callback); + return () => { + const index = pending.indexOf(callback); + if (index >= 0) pending.splice(index, 1); + }; + }, runNext: () => { const callback = pending.shift(); if (!callback) throw new Error("expected a scheduled callback"); @@ -253,7 +259,9 @@ describe("package tree integrity", () => { observation = base; clock += 2_000; expect(guard.status()).toEqual({ ok: true }); - scheduler.runNext(); // stale generation from the first replacement + // Baseline recovery cancels the armed timer through its cancel handle, + // so the stale callback is already gone from the pending queue. + expect(scheduler.pending).toHaveLength(0); expect(calls).toBe(0); observation = { ...base, inode: 12n }; @@ -264,6 +272,37 @@ describe("package tree integrity", () => { expect(calls).toBe(1); }); + test("dispose cancels the pending restart timer and blocks late callbacks", () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const scheduler = createScheduler(); + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 5_000, + schedule: scheduler.schedule, + }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(scheduler.pending).toHaveLength(1); + + guard.dispose(); + expect(scheduler.pending).toHaveLength(0); + + // A further observation cannot re-arm the guard after disposal. + clock += 10_000; + guard.status(); + expect(scheduler.pending).toHaveLength(0); + expect(calls).toBe(0); + }); + test("an unreadable manifest must become readable before a fresh debounce", () => { let observation: PackageTreeObservation | null = base; let clock = 0; @@ -434,6 +473,7 @@ describe("package tree integrity", () => { saveConfig(config()); const packageTreeIntegrity = { status: () => ({ ok: false as const, reason: "package_tree_replaced" as const }), + dispose: () => {}, }; const server = startServer(0, { packageTreeIntegrity }); try { @@ -505,4 +545,50 @@ describe("package tree integrity", () => { await server.stop(true); } }); + + test("server.stop() disarms a pending package-tree restart callback", async () => { + saveConfig(config()); + const base: PackageTreeObservation = { + device: 1n, inode: 10n, contentTimeNs: 100n, size: 500n, + }; + let observation: PackageTreeObservation = base; + const pending: Array<() => void> = []; + let restartAcceptances = 0; + const server = startServer(0, { + packageTreeInstaller: "npm", + observePackageTree: () => observation, + packageTreeIntegrityOptions: { + replacedRestartDelayMs: 5_000, + schedule: callback => { + pending.push(callback); + return () => { + const index = pending.indexOf(callback); + if (index >= 0) pending.splice(index, 1); + }; + }, + }, + acceptSystemRestart: () => { + restartAcceptances += 1; + return { + accepted: true, + alreadyDraining: false, + activeTurnCount: 0, + drainTimeoutMs: 60_000, + }; + }, + }); + try { + expect((await fetch(new URL("/healthz", server.url))).status).toBe(200); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + await Bun.sleep(1_100); // expire the guard's successful-observation cache + expect((await fetch(new URL("/healthz", server.url))).status).toBe(503); + expect(pending).toHaveLength(1); + + await server.stop(true); + expect(pending).toHaveLength(0); + expect(restartAcceptances).toBe(0); + } finally { + await server.stop(true).catch(() => {}); + } + }); });