diff --git a/src/lib/package-tree-integrity.ts b/src/lib/package-tree-integrity.ts index 30372de183e..e079eeeb446 100644 --- a/src/lib/package-tree-integrity.ts +++ b/src/lib/package-tree-integrity.ts @@ -14,10 +14,37 @@ 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; } -type ObservePackageTree = () => PackageTreeObservation | null; -type PackageTreeRuntimeInstall = "bun" | "npm" | "pnpm" | "source"; +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; + /** + * 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; +export type PackageTreeRuntimeInstall = "bun" | "npm" | "pnpm" | "source"; const packageManifestUrl = new URL("../../package.json", import.meta.url); @@ -70,17 +97,123 @@ 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 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 => { + 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) { + // Defer like the scheduled path: verifyAndNotify can arm the next timer, and a + // synchronous verify inside this frame would re-enter armRestartTimer while this + // arm is still running. + queueMicrotask(verifyAndNotify); + } else { + // The seam may run the callback synchronously; defer the work so + // cancelScheduled ownership is settled before verifyAndNotify can + // re-enter armRestartTimer. + const cancel = schedule(() => { + queueMicrotask(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 }; 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) { + 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 (replacementCandidate === null || !sameObservation(replacementCandidate, current)) { + replacementCandidate = current; + resetRestartTimer(); + } + armRestartTimer(); + return { ok: false, reason: "package_tree_replaced" }; + } lastOkAt = at; + resetRestartTimer(); + waitingForReadableTree = false; + replacementCandidate = null; return { ok: true }; }, }; @@ -96,7 +229,10 @@ 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 }), dispose: () => {} }; + } + return createPackageTreeIntegrityGuard(observe, now, options); + } diff --git a/src/server/index.ts b/src/server/index.ts index 26a68f1a034..12eb31a119c 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -193,13 +193,9 @@ import { createLocalAttestationSecret, } from "../lib/local-management-attestation"; import { createReadinessGate, type ReadinessGate } from "./readiness"; -import { - createRuntimePackageTreeIntegrityGuard, - type PackageTreeIntegrityGuard, -} from "../lib/package-tree-integrity"; -import { detectInstall } from "../update/index"; import { createServeOptions, type ServerIngress } from "./index/serve-options"; import { createClaudeInterceptLifecycle } from "./index/claude-intercept-lifecycle"; +import { createPackageTreeIntegrityGuardForServer } from "./index/package-tree-guard"; import { inspectStartupOwnership, resolveInboundBodyLimitWithWarning, setStartupCacheInvalidationWrite, warnAgentTaskRecoveryStartup, warnPlaintextV2AgentMessagesStartup, type StartServerDeps } from "./index/startup-warnings"; import { acquireSpendLedgerServerLifecycle, recordFailedStartRollback, type SpendLedgerServerLifecycle } from "./index/spend-ledger-lifecycle"; export { waitForFailedStartRollback } from "./index/spend-ledger-lifecycle"; @@ -536,8 +532,7 @@ 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 packageTreeIntegrity = deps.packageTreeIntegrity - ?? createRuntimePackageTreeIntegrityGuard(detectInstall()); + const packageTreeIntegrity = createPackageTreeIntegrityGuardForServer(deps); // 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. @@ -755,6 +750,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/src/server/index/package-tree-guard.ts b/src/server/index/package-tree-guard.ts new file mode 100644 index 00000000000..26c260d8851 --- /dev/null +++ b/src/server/index/package-tree-guard.ts @@ -0,0 +1,34 @@ +import { + createRuntimePackageTreeIntegrityGuard, + type PackageTreeIntegrityGuard, +} from "../../lib/package-tree-integrity"; +import { acceptSystemRestart } from "../management/system-restart"; +import { detectInstall } from "../../update/index"; +import type { StartServerDeps } from "./startup-warnings"; + +// Production guard wiring for the package-tree integrity fence. Tests inject +// deps.packageTreeIntegrity and never reach this path; everything else is the +// same default the inline construction used to build. +export function createPackageTreeIntegrityGuardForServer( + deps: StartServerDeps, +): PackageTreeIntegrityGuard { + if (deps.packageTreeIntegrity) { + return deps.packageTreeIntegrity; + } + const acceptPackageTreeRestart = deps.acceptSystemRestart ?? acceptSystemRestart; + return 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(); + }, + }, + ); +} 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/ops/docs-and-release.md b/structure/ops/docs-and-release.md index dcd15fe1bac..0509bf827bc 100644 --- a/structure/ops/docs-and-release.md +++ b/structure/ops/docs-and-release.md @@ -273,6 +273,17 @@ Invariants: - Public docs (root READMEs + `docs-site` installation pages, all locales) state Node 18+ as the only prerequisite. Do not reintroduce "install Bun first" / "bun must be on PATH" guidance for npm users. +### Package-tree integrity fence + +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. + ## Release workflow Package release is npm-focused. `package.json` exposes `opencodex` and `ocx`, `prepublishOnly` runs diff --git a/structure/runtime.md b/structure/runtime.md index cbcb7199933..d760d81c8af 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -188,6 +188,8 @@ described in [OpenAI quota ownership](providers/openai-tiers.md#public-provider- 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. +The package-tree integrity fence for live package replacement follows the +[update transaction contract](ops/docs-and-release.md#package-tree-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` diff --git a/tests/ci-workflows/package-tree-integrity.test.ts b/tests/ci-workflows/package-tree-integrity.test.ts index 13da9e02e7c..29c119a6ffc 100644 --- a/tests/ci-workflows/package-tree-integrity.test.ts +++ b/tests/ci-workflows/package-tree-integrity.test.ts @@ -7,14 +7,46 @@ import { createRuntimePackageTreeIntegrityGuard, type PackageTreeObservation, } from "../../src/lib/package-tree-integrity"; -import { startServer } from "../../src/server"; +import { startServer, waitForFailedStartRollback } from "../../src/server"; +import { stopServerListener } from "../../src/server/lifecycle"; import type { OcxConfig } from "../../src/types"; import { installIsolatedCodexHome, type IsolatedCodexHome } from "../helpers/isolated-codex-home"; import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { currentServerFixtureConfig, settleServerAuthFixture } from "../helpers/server-auth-fixture"; +import { ownedServiceHomeInspection } from "../helpers/owned-service-home-inspection"; const TEST_DIR = join(import.meta.dir, ".tmp-package-tree-integrity"); const previousOpencodexHome = process.env.OPENCODEX_HOME; let isolatedCodexHome: IsolatedCodexHome | null = null; +let ownedServer: ReturnType | null = null; +let caseAbort = new AbortController(); +let caseWork: Promise | undefined; +let closing = false; + +async function prepareServer(deps: Parameters[1]): Promise { + // Package integrity uses the real listener and guard, not native client sync. + saveConfig(currentServerFixtureConfig({ ...config(), clientIntegrations: { codex: false } })); + try { + ownedServer = startServer(0, { + inspectNativeCodexOwnership: ownedServiceHomeInspection("package integrity sandbox"), + ...deps, + }); + } catch (error) { + await waitForFailedStartRollback(error); + throw error; + } +} + +function runServerCase(work: (server: ReturnType, signal: AbortSignal) => Promise): Promise { + const server = ownedServer; + if (!server) throw new Error("package integrity server was not prepared"); + const signal = caseAbort.signal; + caseWork = Promise.resolve().then(() => work(server, signal)).catch(error => { + if (closing && signal.aborted && error instanceof Error && error.name === "AbortError") return; + throw error; + }); + return caseWork; +} function config(): OcxConfig { return { @@ -32,13 +64,26 @@ function config(): OcxConfig { } beforeEach(() => { + caseAbort = new AbortController(); + caseWork = undefined; + closing = false; if (existsSync(TEST_DIR)) removeTreeWithRetry(TEST_DIR); mkdirSync(TEST_DIR, { recursive: true }); process.env.OPENCODEX_HOME = TEST_DIR; isolatedCodexHome = installIsolatedCodexHome("ocx-package-tree-integrity-"); }); -afterEach(() => { +afterEach(async () => { + closing = true; + caseAbort.abort(); + // A runner timeout does not settle the test body. Start the shared real stop + // before waiting for the body; neither may outlive restoration/removal of its home. + const stopped = Promise.allSettled([ownedServer ? stopServerListener(ownedServer) : Promise.resolve()]); + await Promise.allSettled(caseWork ? [caseWork] : []); + const [result] = await stopped; + if (result.status === "rejected") throw result.reason; + ownedServer = null; + await settleServerAuthFixture(TEST_DIR, isolatedCodexHome?.path); if (previousOpencodexHome === undefined) delete process.env.OPENCODEX_HOME; else process.env.OPENCODEX_HOME = previousOpencodexHome; isolatedCodexHome?.restore(); @@ -153,6 +198,406 @@ 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, + }; + const createScheduler = () => { + const pending: Array<() => void> = []; + return { + pending, + schedule: (callback: () => void) => { + pending.push(callback); + return () => { + const index = pending.indexOf(callback); + if (index >= 0) pending.splice(index, 1); + }; + }, + runNext: async () => { + const callback = pending.shift(); + if (!callback) throw new Error("expected a scheduled callback"); + callback(); + // The scheduled wrapper defers its verify step to a microtask; drain + // it here so callers keep one runNext = one verify semantics. + await Promise.resolve(); + }, + }; + }; + + test("fires once from its timer without waiting for another request", async () => { + 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, contentTimeNs: 200n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(calls).toBe(0); + + await scheduler.runNext(); + expect(calls).toBe(1); + expect(scheduler.pending).toHaveLength(0); + + clock += 10_000; + guard.status(); + expect(calls).toBe(1); + }); + + test("retries when restart acceptance throws", async () => { + 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" }); + await scheduler.runNext(); + expect(attempts).toBe(1); + expect(scheduler.pending).toHaveLength(1); + await scheduler.runNext(); + expect(attempts).toBe(2); + expect(scheduler.pending).toHaveLength(0); + }); + + test("baseline recovery cancels the old timer and starts a fresh debounce", async () => { + 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; + clock += 2_000; + expect(guard.status()).toEqual({ ok: true }); + // 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 }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + expect(calls).toBe(0); + await scheduler.runNext(); + 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("a zero restart delay defers verification past the current frame", async () => { + // The zero-delay branch used to call verifyAndNotify synchronously inside the + // status() call that armed it. The deferred path must keep the same contract: + // the arming call returns before the verify, a verify that re-enters status() + // still fires exactly once, and a dispose or baseline recovery before the + // microtask runs means it never fires. + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { + onReplaced: () => { + calls += 1; + // Re-entering status() from inside the callback must not arm a second + // verification behind this one. + guard.status(); + }, + replacedRestartDelayMs: 0, + }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + // The verify is deferred: the arming status() returned without firing. + expect(calls).toBe(0); + + await Promise.resolve(); + expect(calls).toBe(1); + }); + + test("a zero-delay verify does not fire after dispose or a baseline recovery", async () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const disposed = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { onReplaced: () => { calls += 1; }, replacedRestartDelayMs: 0 }, + ); + + expect(disposed.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(disposed.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + disposed.dispose(); + await Promise.resolve(); + expect(calls).toBe(0); + + // The same hold applies when the tree returns to the baseline before the + // deferred verify runs: the queued microtask observes the recovered state and + // stands down. + observation = base; + clock += 2_000; + expect(disposed.status()).toEqual({ ok: true }); + + let recoveredCalls = 0; + const recovered = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { onReplaced: () => { recoveredCalls += 1; }, replacedRestartDelayMs: 0 }, + ); + expect(recovered.status()).toEqual({ ok: true }); + observation = { ...base, inode: 12n }; + clock += 2_000; + expect(recovered.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + observation = base; + clock += 2_000; + expect(recovered.status()).toEqual({ ok: true }); + await Promise.resolve(); + expect(recoveredCalls).toBe(0); + }); + + test("an unreadable manifest must become readable before a fresh debounce", async () => { + 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 = null; + await scheduler.runNext(); + expect(calls).toBe(0); + expect(scheduler.pending).toHaveLength(1); + + observation = { ...base, inode: 11n }; + await scheduler.runNext(); // readability poll; arms a fresh full debounce + expect(calls).toBe(0); + expect(scheduler.pending).toHaveLength(1); + await scheduler.runNext(); + expect(calls).toBe(1); + }); + + test("a second replacement identity receives its own full debounce", async () => { + 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 }; + await scheduler.runNext(); + expect(calls).toBe(0); + expect(scheduler.pending).toHaveLength(1); + await scheduler.runNext(); + expect(calls).toBe(1); + }); + + // A scheduler that invokes its callback inline used to re-enter + // armRestartTimer before cancelScheduled was assigned, orphaning the + // stale timer. The queueMicrotask wrapper defers the verify step until + // ownership is settled. + const synchronousScheduler = () => ({ + schedule: (callback: () => void) => { + callback(); + return () => {}; + }, + }); + + test("a synchronous schedule implementation still notifies exactly once", async () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 5_000, + schedule: synchronousScheduler().schedule, + }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + // The wrapper ran inline but the verify step is a queued microtask. + expect(calls).toBe(0); + await Promise.resolve(); + expect(calls).toBe(1); + guard.dispose(); + }); + + test("disposing before the queued verify runs prevents notification", async () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let calls = 0; + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 5_000, + schedule: synchronousScheduler().schedule, + }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + guard.dispose(); + await Promise.resolve(); + expect(calls).toBe(0); + }); + + test("a throwing onReplaced still retries under a synchronous schedule", async () => { + let observation: PackageTreeObservation | null = base; + let clock = 0; + let attempts = 0; + const guard = createPackageTreeIntegrityGuard( + () => observation, + () => clock, + { + onReplaced: () => { + attempts += 1; + if (attempts === 1) throw new Error("restart unavailable"); + }, + replacedRestartDelayMs: 5_000, + schedule: synchronousScheduler().schedule, + }, + ); + + expect(guard.status()).toEqual({ ok: true }); + observation = { ...base, inode: 11n }; + clock += 2_000; + expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); + // The first verify throws and re-arms; its retry wrapper queues another + // microtask behind this continuation, so the second attempt needs a + // second flush. + await Promise.resolve(); + expect(attempts).toBe(1); + await Promise.resolve(); + expect(attempts).toBe(2); + guard.dispose(); + }); + + test("source checkouts never auto-restart", () => { + let observation: PackageTreeObservation | null = base; + let calls = 0; + const scheduler = createScheduler(); + const guard = createRuntimePackageTreeIntegrityGuard( + "source", + () => observation, + Date.now, + { + onReplaced: () => { calls += 1; }, + replacedRestartDelayMs: 0, + schedule: scheduler.schedule, + }, + ); + + observation = { ...base, inode: 11n }; + expect(guard.status()).toEqual({ ok: true }); + expect(calls).toBe(0); + expect(scheduler.pending).toHaveLength(0); + }); + }); + // BUG-R1: a chmod fenced the whole data plane behind 503. // // These three drive the REAL filesystem rather than a hand-built observation, @@ -238,14 +683,14 @@ describe("package tree integrity", () => { expect(guard.status()).toEqual({ ok: false, reason: "package_tree_replaced" }); }); - test("degrades health and refuses Responses requests with a restart-required error", async () => { - saveConfig(config()); + describe("restart-required server", () => { const packageTreeIntegrity = { status: () => ({ ok: false as const, reason: "package_tree_replaced" as const }), + dispose: () => {}, }; - const server = startServer(0, { packageTreeIntegrity }); - try { - const health = await fetch(new URL("/healthz", server.url)); + beforeEach(() => prepareServer({ packageTreeIntegrity })); + test("degrades health and refuses Responses requests with a restart-required error", () => runServerCase(async (server, signal) => { + const health = await fetch(new URL("/healthz", server.url), { signal }); expect(health.status).toBe(503); expect(health.headers.get("retry-after")).toBe("5"); expect(await health.json()).toMatchObject({ @@ -256,6 +701,7 @@ describe("package tree integrity", () => { const response = await fetch(new URL("/v1/responses", server.url), { method: "POST", + signal, headers: { "content-type": "application/json" }, body: JSON.stringify({ model: "test/gpt-test", input: "hello" }), }); @@ -268,8 +714,108 @@ describe("package tree integrity", () => { message: expect.stringContaining("restart"), }, }); - } finally { - await server.stop(true); - } + })); + }); + + describe("sustained package replacement", () => { + const base: PackageTreeObservation = { + device: 1n, inode: 10n, contentTimeNs: 100n, size: 500n, + }; + let observation: PackageTreeObservation = base; + const pending: Array<() => void> = []; + let restartAcceptances = 0; + beforeEach(() => { + observation = base; + pending.length = 0; + restartAcceptances = 0; + return prepareServer({ + 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, + }; + }, + }); + }); + test("the default server guard accepts a restart after a sustained replacement", () => runServerCase(async (server, signal) => { + const healthy = await fetch(new URL("/healthz", server.url), { signal }); + expect(healthy.status).toBe(200); + await healthy.text(); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + await Bun.sleep(1_100); // expire the guard's successful-observation cache + signal.throwIfAborted(); + const replaced = await fetch(new URL("/healthz", server.url), { signal }); + expect(replaced.status).toBe(503); + await replaced.text(); + expect(restartAcceptances).toBe(0); + expect(pending).toHaveLength(1); + + pending.shift()?.(); + await Promise.resolve(); // the deferred verify step runs as a microtask + expect(restartAcceptances).toBe(1); + expect(pending).toHaveLength(0); + })); + }); + + describe("pending package restart shutdown", () => { + const base: PackageTreeObservation = { + device: 1n, inode: 10n, contentTimeNs: 100n, size: 500n, + }; + let observation: PackageTreeObservation = base; + const pending: Array<() => void> = []; + let restartAcceptances = 0; + beforeEach(() => { + observation = base; + pending.length = 0; + restartAcceptances = 0; + return prepareServer({ + 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, + }; + }, + }); + }); + test("server.stop() disarms a pending package-tree restart callback", () => runServerCase(async (server, signal) => { + const healthy = await fetch(new URL("/healthz", server.url), { signal }); + expect(healthy.status).toBe(200); + await healthy.text(); + observation = { ...base, inode: 11n, contentTimeNs: 200n }; + await Bun.sleep(1_100); // expire the guard's successful-observation cache + signal.throwIfAborted(); + const replaced = await fetch(new URL("/healthz", server.url), { signal }); + expect(replaced.status).toBe(503); + await replaced.text(); + expect(pending).toHaveLength(1); + + await stopServerListener(server); + expect(pending).toHaveLength(0); + expect(restartAcceptances).toBe(0); + })); }); });