diff --git a/gui/src/api.ts b/gui/src/api.ts index 1e99ed6f73b..420c04d3f39 100644 --- a/gui/src/api.ts +++ b/gui/src/api.ts @@ -9,7 +9,9 @@ import { adminTokenPromptAllowed, standaloneApiTargets, type ApiPlane, type ApiT export const SESSION_UNAVAILABLE_EVENT = "opencodex:session-unavailable"; const LEGACY_TOKEN_KEY = "opencodex-api-token"; -const ADMIN_TOKEN_VALIDATION_PATH = "/api/settings"; +// Any guarded route answers 401 for a bad token; this one is a cheap config read. /api/settings +// also resolves the Codex runtime and startup health, which made the token prompt hang. +const ADMIN_TOKEN_VALIDATION_PATH = "/api/combos"; const SESSION_REBOOTSTRAP_TIMEOUT_MS = 10_000; const RESOLUTION_WATCHDOG_MS = 15_000; const MACHINE_SESSION_HEADER = "X-OpenCodex-Machine-Session"; diff --git a/gui/src/startup-health-ui.ts b/gui/src/startup-health-ui.ts index 105dfad8a02..5fcf5fa6172 100644 --- a/gui/src/startup-health-ui.ts +++ b/gui/src/startup-health-ui.ts @@ -105,12 +105,15 @@ export function probeNeedsFastRetry(probe: StartupHealthProbe | undefined | null /** * Settings may only seed startup health while it is still unknown or a hard error. * After `/api/startup-health` has produced a real status, it stays authoritative. + * A `diagnosticStale` seed never colors the chip: it is a placeholder while the server + * refreshes (on a cold server, a synthetic not-installed reading), so a healthy service + * would flash `at-risk` until the dedicated probe answered. */ export function seedStartupHealthFromSettings( previous: StartupHealthStatus | null, seeded: { status: "native" | "protected" | "at-risk"; diagnosticStale: boolean } | null | undefined, ): StartupHealthStatus | null { - if (!seeded) return previous; + if (!seeded || seeded.diagnosticStale) return previous; // A prior hard "error" (or unknown) may be replaced by a settings seed; a real // status from the dedicated probe must not be overwritten. if (previous !== null && previous !== "error") return previous; diff --git a/src/codex/app-server-processes.ts b/src/codex/app-server-processes.ts index a33a1091fd6..c7fa8b16fa5 100644 --- a/src/codex/app-server-processes.ts +++ b/src/codex/app-server-processes.ts @@ -1095,6 +1095,38 @@ export async function collectCodexAppServerCatalogStateForRequest( return servableStale ?? flight.promise; } +/** + * Dashboard read of the request-path catalog state, bounded by `deadlineMs`. + * + * The request collector keeps the event loop free, but a cold Windows probe still + * takes 4-7s (CIM walk plus GetOwner per candidate), and the Models and Subagents + * pages cannot render their rosters until these routes answer. The reading only + * drives an advisory banner, so a slow probe answers `unknown` (which renders no + * banner) and keeps running; its result is cached for the next poll. + * + * The deadline bounds the Windows path, where the request collector is asynchronous. On + * other platforms the request collector keeps its existing synchronous read (a /proc walk on + * Linux, `ps` on macOS, typically tens of milliseconds), which completes before the deadline + * timer can fire; making those reads asynchronous is outside this change. + */ +export async function collectCodexAppServerCatalogStateWithin( + deadlineMs: number, + io: CodexAppServerProcessIo = {}, +): Promise { + let timer: ReturnType | undefined; + const deadline = new Promise(resolve => { + timer = setTimeout(() => resolve({ state: "unknown", processes: [], catalogMtimeMs: null }), deadlineMs); + }); + try { + return await Promise.race([collectCodexAppServerCatalogStateForRequest(io), deadline]); + } finally { + clearTimeout(timer); + } +} + +/** Deadline for dashboard catalog-state reads; a cached reading answers well within it. */ +export const DASHBOARD_CATALOG_STATE_DEADLINE_MS = 250; + /** * Drop memoized catalog state after a relevant catalog/cache write and before * the post-write state read. Advancing the generation prevents an older diff --git a/src/codex/app-server-restart-service.ts b/src/codex/app-server-restart-service.ts index d919fac86c1..aa1679c5550 100644 --- a/src/codex/app-server-restart-service.ts +++ b/src/codex/app-server-restart-service.ts @@ -20,6 +20,8 @@ */ import { collectCodexAppServerCatalogState, + collectCodexAppServerCatalogStateWithin, + DASHBOARD_CATALOG_STATE_DEADLINE_MS, listCodexAppServerProcesses, readProcessStartMsBatch, resetCodexAppServerCatalogStateCache, @@ -59,6 +61,8 @@ export interface CodexRestartServiceIo { */ listenPort?: () => number | undefined; collectState?: typeof collectCodexAppServerCatalogState; + /** Deadline for the default dashboard state read (ms). */ + stateDeadlineMs?: number; listProcesses?: typeof listCodexAppServerProcesses; restart?: typeof restartCodexAppServers; resetStateCache?: () => void; @@ -88,10 +92,23 @@ class CodexAppServerIdentityChanged extends Error { */ let inFlight: Promise | null = null; -export function readCodexAppServerState( +/** + * Read-only state for the dashboard route. The default classifier is the bounded + * request-path read: on Windows the synchronous one runs PowerShell CIM enumeration + * inline, which measured 4-7s per call and stalled every listener, proxy traffic + * included, each time the Models page opened. A probe slower than the deadline answers + * `unknown` and keeps refreshing the cache. The restart path below keeps the + * synchronous reading because it acts on it. + */ +export async function readCodexAppServerState( io: CodexRestartServiceIo = {}, -): CodexAppServerStateResponse { - const status = (io.collectState ?? collectCodexAppServerCatalogState)(io.processIo ?? {}); +): Promise { + const status = io.collectState + ? io.collectState(io.processIo ?? {}) + : await collectCodexAppServerCatalogStateWithin( + io.stateDeadlineMs ?? DASHBOARD_CATALOG_STATE_DEADLINE_MS, + io.processIo ?? {}, + ); return { state: status.state, runningCount: status.processes.length }; } diff --git a/src/codex/runtime.ts b/src/codex/runtime.ts index db7a739cd4e..4efd5ce480c 100644 --- a/src/codex/runtime.ts +++ b/src/codex/runtime.ts @@ -1,4 +1,4 @@ -import { execFileSync } from "node:child_process"; +import { execFile, execFileSync } from "node:child_process"; import { existsSync, mkdirSync, mkdtempSync, readFileSync, readdirSync, rmSync, statSync, unlinkSync } from "node:fs"; import { homedir, tmpdir } from "node:os"; import { delimiter, join } from "node:path"; @@ -70,11 +70,18 @@ export type RuntimeExecFile = ( }, ) => string; +/** Async twin of RuntimeExecFile, used by resolveCodexRuntimeAsync; resolves with stdout. */ +export type RuntimeExecFileAsync = ( + ...args: Parameters +) => Promise; + export interface ResolveCodexRuntimeDeps { env?: NodeJS.ProcessEnv; platform?: NodeJS.Platform; configDir?: string; execFileSync?: RuntimeExecFile; + /** Probe exec for resolveCodexRuntimeAsync. Injected impls bypass the process memo like execFileSync. */ + execFile?: RuntimeExecFileAsync; existsSync?: (path: string) => boolean; readFileSync?: (path: string, encoding: "utf8") => string; now?: () => number; @@ -397,10 +404,10 @@ export function clearPersistedCodexRuntime(deps: ResolveCodexRuntimeDeps = {}): } } -function probeVersion( - command: string, - deps: ResolveCodexRuntimeDeps, -): { ok: true; version: string | null } | { ok: false; reason: string } { +type VersionProbeResult = { ok: true; version: string | null } | { ok: false; reason: string }; + +/** Checks that need no exec; a non-null answer settles the candidate without running `--version`. */ +function versionProbePrecheck(command: string, deps: ResolveCodexRuntimeDeps): VersionProbeResult | null { const platform = deps.platform ?? process.platform; if (command.includes("/") || command.includes("\\") || /^[A-Za-z]:/.test(command)) { const exists = deps.existsSync ?? existsSync; @@ -414,7 +421,54 @@ function probeVersion( // dashboard probe that cost made Windows report Codex as missing even when // the App install was sitting under LOCALAPPDATA/OpenAI/Codex/bin (issue 4458). if (deps.probeVersion === false) return { ok: true, version: null }; - const execFile = deps.execFileSync ?? (execFileSync as unknown as RuntimeExecFile); + return null; +} + +function versionProbeInvocation(command: string, deps: ResolveCodexRuntimeDeps, probeHome: string) { + const invocation = codexExecInvocation(command, ["--version"], deps.platform ?? process.platform, { + env: deps.env, + exists: deps.existsSync, + }); + return { + file: invocation.file, + args: invocation.args, + options: { + encoding: "utf8" as const, + stdio: ["ignore", "pipe", "ignore"] as ["ignore", "pipe", "ignore"], + timeout: 8_000, + windowsHide: true, + env: { ...(deps.env ?? process.env), CODEX_HOME: probeHome }, + ...invocation.options, + }, + }; +} + +function versionProbeOutput(output: string): VersionProbeResult { + const version = parseCodexVersionOutput(output); + if (!version) return { ok: false, reason: "unrecognized --version output" }; + return { ok: true, version }; +} + +function versionProbeError(error: unknown, probeHome: string | undefined): VersionProbeResult { + if (!probeHome) return { ok: false, reason: "probe sandbox unavailable" }; + if ((error as NodeJS.ErrnoException)?.code === "ENOENT") { + return { ok: false, reason: CODEX_PROGRAM_NOT_FOUND_REASON }; + } + const message = error instanceof Error ? error.message : String(error); + const redacted = redactUserPath(redactSecretString(message)).slice(0, 160); + return { ok: false, reason: `failed --version (${redacted})` }; +} + +function removeProbeHome(probeHome: string | undefined): void { + if (!probeHome) return; + // Nested catch: a transient Windows EBUSY must never mask the probe result. + try { rmSync(probeHome, { recursive: true, force: true }); } catch { /* best effort */ } +} + +function probeVersion(command: string, deps: ResolveCodexRuntimeDeps): VersionProbeResult { + const settled = versionProbePrecheck(command, deps); + if (settled) return settled; + const exec = deps.execFileSync ?? (execFileSync as unknown as RuntimeExecFile); // Sandbox the probe's CODEX_HOME: a real Codex CLI creates state (tmp/, logs) under // CODEX_HOME even for `--version`, and the probe inherits the caller's env — so a // read-only `ocx status` would dirty the user's CODEX_HOME. Redirect it to a @@ -423,36 +477,36 @@ function probeVersion( let probeHome: string | undefined; try { probeHome = mkdtempSync(join(tmpdir(), "ocx-codex-probe-")); - const invocation = codexExecInvocation(command, ["--version"], platform, { - env: deps.env, - exists: deps.existsSync, - }); - const output = execFile(invocation.file, invocation.args, { - encoding: "utf8", - stdio: ["ignore", "pipe", "ignore"], - timeout: 8_000, - windowsHide: true, - env: { ...(deps.env ?? process.env), CODEX_HOME: probeHome }, - ...invocation.options, - }); - const version = parseCodexVersionOutput(output); - if (!version) { - return { ok: false, reason: "unrecognized --version output" }; - } - return { ok: true, version }; + const { file, args, options } = versionProbeInvocation(command, deps, probeHome); + return versionProbeOutput(exec(file, args, options)); } catch (error) { - if (!probeHome) return { ok: false, reason: "probe sandbox unavailable" }; - if ((error as NodeJS.ErrnoException)?.code === "ENOENT") { - return { ok: false, reason: CODEX_PROGRAM_NOT_FOUND_REASON }; - } - const message = error instanceof Error ? error.message : String(error); - const redacted = redactUserPath(redactSecretString(message)).slice(0, 160); - return { ok: false, reason: `failed --version (${redacted})` }; + return versionProbeError(error, probeHome); } finally { - if (probeHome) { - // Nested catch: a transient Windows EBUSY must never mask the probe result. - try { rmSync(probeHome, { recursive: true, force: true }); } catch { /* best effort */ } - } + removeProbeHome(probeHome); + } +} + +const defaultExecFileAsync: RuntimeExecFileAsync = (file, args, options) => + new Promise((resolve, reject) => { + // execFile always pipes stdout/stderr; it takes no stdio option. + const { stdio: _stdio, ...execOptions } = options; + execFile(file, args, execOptions, (error, stdout) => (error ? reject(error) : resolve(String(stdout)))); + }); + +/** probeVersion with the same sandbox and outcomes, but the exec never holds the event loop. */ +async function probeVersionAsync(command: string, deps: ResolveCodexRuntimeDeps): Promise { + const settled = versionProbePrecheck(command, deps); + if (settled) return settled; + const exec = deps.execFile ?? defaultExecFileAsync; + let probeHome: string | undefined; + try { + probeHome = mkdtempSync(join(tmpdir(), "ocx-codex-probe-")); + const { file, args, options } = versionProbeInvocation(command, deps, probeHome); + return versionProbeOutput(await exec(file, args, options)); + } catch (error) { + return versionProbeError(error, probeHome); + } finally { + removeProbeHome(probeHome); } } @@ -566,23 +620,6 @@ interface RankedCandidate { source: CodexRuntimeSource; } -function tryCandidate( - candidate: RankedCandidate, - failures: RuntimeProbeFailure[], - deps: ResolveCodexRuntimeDeps, -): ResolvedCodexRuntime | null { - const probed = probeVersion(candidate.command, deps); - if (!probed.ok) { - failures.push({ command: candidate.command, source: candidate.source, reason: probed.reason }); - return null; - } - return { - command: candidate.command, - version: probed.version, - source: candidate.source, - }; -} - function sameRuntimeCommand(a: string, b: string): boolean { return a.trim().toLowerCase() === b.trim().toLowerCase(); } @@ -704,10 +741,23 @@ export function clearCodexRuntimeResolveCache(): void { clearResolveCache(); } -/** Observe only an unexpired successful process memo; never resolves or probes. */ +/** + * Observe the successful process memo; never resolves or probes. + * + * An expired memo stays observable while an async refresh for the same inputs is in flight + * and nothing has invalidated it since. The sync resolver used to hold the event loop for the + * whole probe, so nothing could read in that window; with the refresh off the loop, dropping + * the memo there would send catalog gather to the persisted runtime (or none) and have + * convergence reject its process-local candidate every expiry. A published result or a clear + * still retires it; the deferred selection never enters this memo (#4458). + */ export function peekCodexRuntimeProcessCache(): CodexRuntimeProcessCachePeek { const memo = resolveCache; - if (!memo || Date.now() - memo.at >= RESOLVE_CACHE_MS) { + const refreshing = memo !== null + && asyncResolveInflight !== null + && asyncResolveInflight.key === memo.key + && asyncResolveInflight.epoch === resolveCacheEpoch; + if (!memo || (Date.now() - memo.at >= RESOLVE_CACHE_MS && !refreshing)) { return Object.freeze({ kind: "unavailable" as const, epoch: resolveCacheEpoch }); } return Object.freeze({ @@ -735,6 +785,7 @@ function resolveCacheKey(deps: ResolveCodexRuntimeDeps): string | null { // Only memoize uninjected process-env resolves (settings/status hot paths). if ( deps.execFileSync + || deps.execFile || deps.existsSync || deps.readFileSync || deps.readdirSync @@ -799,6 +850,56 @@ export function resolveCodexRuntime(deps: ResolveCodexRuntimeDeps = {}): Resolve return cloneAndDeepFreeze(result); } +let asyncResolveInflight: { key: string; epoch: number; promise: Promise } | null = null; + +/** + * resolveCodexRuntime for server request paths: same selection and process memo, but every + * `codex --version` runs through async exec. The sync resolver blocked the whole proxy for + * ~0.6s per memo expiry while the dashboard polled settings, stalling unrelated requests. + */ +export async function resolveCodexRuntimeAsync( + deps: ResolveCodexRuntimeDeps = {}, +): Promise { + // A deferred selection never execs, so the sync path is already nonblocking. + if (deps.probeVersion === false) return resolveCodexRuntime(deps); + const cacheKey = resolveCacheKey(deps); + if (cacheKey && resolveCache && resolveCache.key === cacheKey && Date.now() - resolveCache.at < RESOLVE_CACHE_MS) { + return cloneAndDeepFreeze(resolveCache.value); + } + if (!cacheKey) return cloneAndDeepFreeze(await resolveCodexRuntimeUncachedAsync(deps)); + const startedEpoch = resolveCacheEpoch; + if (asyncResolveInflight?.key === cacheKey && asyncResolveInflight.epoch === startedEpoch) { + return cloneAndDeepFreeze(await asyncResolveInflight.promise); + } + const promise = resolveCodexRuntimeUncachedAsync(deps).then(result => { + // A persist or clear while the probes ran makes this answer the previous selection's; + // publishing it would resurrect authority the write just revoked. + if (resolveCacheEpoch === startedEpoch) publishResolveCache(cacheKey, Date.now(), result); + return result; + }).finally(() => { + if (asyncResolveInflight?.promise === promise) asyncResolveInflight = null; + }); + asyncResolveInflight = { key: cacheKey, epoch: startedEpoch, promise }; + return cloneAndDeepFreeze(await promise); +} + +/** + * Stale-while-revalidate runtime read for hot UI paths; never waits on a probe. + * + * Returns the process memo when fresh, otherwise the last memo for the same inputs, and kicks + * resolveCodexRuntimeAsync to refresh it. With no memo yet (cold start, or right after a runtime + * switch cleared it) the answer is the exec-free deferred selection, so a version and + * newerAvailable appear once the first background refresh lands. + */ +export function getCodexRuntimeSnapshot(): ResolveCodexRuntimeResult { + const cacheKey = resolveCacheKey({}); + const memo = resolveCache?.key === cacheKey ? resolveCache : null; + if (memo && Date.now() - memo.at < RESOLVE_CACHE_MS) return cloneAndDeepFreeze(memo.value); + resolveCodexRuntimeAsync().catch(() => { /* the next read retries */ }); + if (memo) return cloneAndDeepFreeze(memo.value); + return resolveCodexRuntime({ discoverAlternatives: false, probeVersion: false }); +} + /** Test-only: drop the short-lived process resolve cache. */ export function resetCodexRuntimeResolveCacheForTests(): void { clearCodexRuntimeResolveCache(); @@ -814,7 +915,14 @@ export function setCodexRuntimeResolveCacheForTests( publishResolveCache(key, Date.now(), value); } -function resolveCodexRuntimeUncached(deps: ResolveCodexRuntimeDeps = {}): ResolveCodexRuntimeResult { +/** + * Candidate ranking and selection, written once for both drivers: each yield hands a candidate to + * the caller, which sends back its `--version` outcome. The sync driver keeps CLI behavior; the + * async one lets a server path resolve without blocking every other request on the probes. + */ +function* resolveCodexRuntimeSteps( + deps: ResolveCodexRuntimeDeps, +): Generator { const env = deps.env ?? process.env; const failures: RuntimeProbeFailure[] = []; const ordered: RankedCandidate[] = []; @@ -868,9 +976,12 @@ function resolveCodexRuntimeUncached(deps: ResolveCodexRuntimeDeps = {}): Resolv ) { continue; } - const resolved = tryCandidate(candidate, failures, deps); - if (!resolved) continue; - valid.push(resolved); + const probed: VersionProbeResult = yield candidate; + if (!probed.ok) { + failures.push({ command: candidate.command, source: candidate.source, reason: probed.reason }); + continue; + } + valid.push({ command: candidate.command, version: probed.version, source: candidate.source }); } if (valid.length === 0) { @@ -954,6 +1065,20 @@ function resolveCodexRuntimeUncached(deps: ResolveCodexRuntimeDeps = {}): Resolv }; } +function resolveCodexRuntimeUncached(deps: ResolveCodexRuntimeDeps = {}): ResolveCodexRuntimeResult { + const steps = resolveCodexRuntimeSteps(deps); + let step = steps.next(); + while (!step.done) step = steps.next(probeVersion(step.value.command, deps)); + return step.value; +} + +async function resolveCodexRuntimeUncachedAsync(deps: ResolveCodexRuntimeDeps): Promise { + const steps = resolveCodexRuntimeSteps(deps); + let step = steps.next(); + while (!step.done) step = steps.next(await probeVersionAsync(step.value.command, deps)); + return step.value; +} + /** Resolve and persist a successful selection (unless source is ephemeral fallback-only with no path). */ export function resolveAndPersistCodexRuntime( deps: ResolveCodexRuntimeDeps = {}, diff --git a/src/remote-control/workspace-codex-runtime.ts b/src/remote-control/workspace-codex-runtime.ts index b064e3c9b8c..5e6c2472d6f 100644 --- a/src/remote-control/workspace-codex-runtime.ts +++ b/src/remote-control/workspace-codex-runtime.ts @@ -1,7 +1,7 @@ import { chmodSync, linkSync, mkdirSync, mkdtempSync, realpathSync, symlinkSync } from "node:fs"; import { tmpdir } from "node:os"; import { dirname, isAbsolute, join } from "node:path"; -import { resolveCodexRuntime } from "../codex/runtime"; +import { resolveCodexRuntimeAsync } from "../codex/runtime"; import { remoteWorkspaceThreadStartParams } from "./workspace-coordinator"; import { startRemoteWorkspaceToolBridge } from "./workspace-tool-bridge"; import { truncateRemoteWorkspaceUtf8 } from "./workspace-utf8"; @@ -299,7 +299,8 @@ export class CodexRemoteWorkspaceRuntimeFactory implements RemoteWorkspaceRuntim if (this.options.command && this.options.command.length > 0) { return { available: true, version: this.options.version ?? "test" }; } - const resolved = resolveCodexRuntime(); + // Async probes: a Hub availability check must not freeze every other proxy request. + const resolved = await resolveCodexRuntimeAsync(); const compatibility = codexRemotePermissionProfileCompatibility(); if (!compatibility.compatible) return { available: false, reason: compatibility.reason }; return resolved.runtime.version @@ -310,7 +311,7 @@ export class CodexRemoteWorkspaceRuntimeFactory implements RemoteWorkspaceRuntim async start(options: Parameters[0]): Promise { const command = this.options.command ? [...this.options.command] - : [resolveCodexRuntime().runtime.command]; + : [(await resolveCodexRuntimeAsync()).runtime.command]; if (command.length < 1) throw new Error("Codex CLI is unavailable on this Hub"); const executablePath = isAbsolute(command[0]!) ? command[0]! : findExecutableOnPath(command[0]!); if (!executablePath) throw new Error("Codex CLI executable could not be resolved on this Hub"); diff --git a/src/server/management/agent-settings-routes.ts b/src/server/management/agent-settings-routes.ts index 8e6c430cc12..31527919efc 100644 --- a/src/server/management/agent-settings-routes.ts +++ b/src/server/management/agent-settings-routes.ts @@ -753,9 +753,13 @@ export async function handleAgentSettingsRoutes(ctx: ManagementContext): Promise ...[...new Set(chosen)].filter(model => !selectableSet.has(model)), ]; // #857: let CLI/GUI show when a running Codex app-server keeps an older - // in-memory catalog than the one on disk. - const { collectCodexAppServerCatalogState } = await import("../../codex/app-server-processes"); - const catalogState = collectCodexAppServerCatalogState(); + // in-memory catalog than the one on disk. Bounded request-path read: the synchronous + // collector blocked the event loop for the whole Windows CIM walk (4-7s measured). + const { + collectCodexAppServerCatalogStateWithin, + DASHBOARD_CATALOG_STATE_DEADLINE_MS, + } = await import("../../codex/app-server-processes"); + const catalogState = await collectCodexAppServerCatalogStateWithin(DASHBOARD_CATALOG_STATE_DEADLINE_MS); return jsonResponse({ chosen, available, catalogState, pickerAvailable: [...new Set(filterCatalogVisibleModels(models, config).map(catalogModelSlug).filter(slug => slug.includes("/")))], diff --git a/src/server/management/config-routes.ts b/src/server/management/config-routes.ts index 9682b062649..d439abb137c 100644 --- a/src/server/management/config-routes.ts +++ b/src/server/management/config-routes.ts @@ -119,10 +119,10 @@ import type { PersistedUsageAttempt } from "../../usage/log"; import { isAllowedRequestOrigin, jsonResponse, providerManagementConfigError, publicProviderBaseUrl, safeConfigDTO } from "../auth-cors"; import { withProviderCatalogCapabilityDTO } from "./provider-capability-config"; import { applySystemEnvToggle } from "../system-env"; -import { getCachedStartupHealth, invalidateStartupHealthCache } from "../startup-health-cache"; +import { getCachedStartupHealth, getStartupHealthSnapshot, invalidateStartupHealthCache } from "../startup-health-cache"; import { runWindowsTrayAction } from "../windows-tray-control"; import { runStartupInstallAction, type StartupInstallAction } from "../startup-action-control"; -import { displayCodexRuntimePath, effortClampAppliesToRuntime, liveRemovedEfforts, loadLastEffortClamp, resolveCodexRuntime } from "../../codex/runtime"; +import { displayCodexRuntimePath, effortClampAppliesToRuntime, getCodexRuntimeSnapshot, liveRemovedEfforts, loadLastEffortClamp } from "../../codex/runtime"; import { isPlainRecord, parseDebugLogQuery, tokPerSecondResult, unavailableCostReason, costResult, requestLogDto, stripRegistryOnlyStaticHeaders, fetchAllModels } from "./shared"; import type { MetricUnavailableReason, TokPerSecondResult, CostEstimateReason, CostResult, MetricSource } from "./shared"; @@ -291,6 +291,10 @@ function publicVisionSidecarSettings( export async function handleConfigRoutes(ctx: ManagementContext): Promise { const { req, url, config, deps, convergeCodexCatalog, syncClaudeAgentDefsBestEffort } = ctx; const readStartupHealth = deps.getCachedStartupHealth ?? getCachedStartupHealth; + // Settings only seed the dashboard chip; /api/startup-health owns the bounded fresh read. + // Waiting on the Windows service-manager probe here held settings reads and saves open + // for up to 15s, including the admin-token check. An injected reader stays authoritative. + const readStartupHealthSnapshot = deps.getCachedStartupHealth ?? getStartupHealthSnapshot; if (url.pathname === "/api/config" && req.method === "GET") { return jsonResponse(withProviderCatalogCapabilityDTO(safeConfigDTO(config), config)); } @@ -300,10 +304,12 @@ export async function handleConfigRoutes(ctx: ManagementContext): Promise; + let resolved: ReturnType; try { - // Full alternative discovery (memoized) so newerAvailable warnings work. - resolved = resolveCodexRuntime(); + // Full alternative discovery so newerAvailable warnings work, served stale-while- + // revalidate: the sync resolver ran `codex --version` per candidate on this request and + // froze the whole proxy for ~0.6s every memo expiry while the dashboard polled here. + resolved = getCodexRuntimeSnapshot(); } catch { resolved = { runtime: { command: "codex", version: null, source: "fallback" }, @@ -363,7 +369,7 @@ export async function handleConfigRoutes(ctx: ManagementContext): Promise; +} +const storageScanFlights = new Map(); + +export function sharedStorageScan( + codexHome: string, + scan: (codexHome: string) => Promise = scanStorageAsync, +): Promise { + const epoch = storageMutationEpoch(); + const pending = storageScanFlights.get(codexHome); + if (pending && pending.epoch === epoch) return pending.promise; + const flight: StorageScanFlight = { + epoch, + promise: scan(codexHome).finally(() => { + // A newer flight may already own this key; only retire our own entry. + if (storageScanFlights.get(codexHome) === flight) storageScanFlights.delete(codexHome); + }), + }; + storageScanFlights.set(codexHome, flight); + return flight.promise; +} + +/** + * Log Guard results that are refused before anything under CODEX_HOME is touched. Anything + * else (success, a database or config write failure, `busy` which compaction can report + * after partial progress, a post-compaction integrity failure) may have changed storage. + */ +const REFUSED_BEFORE_CHANGE: ReadonlySet = new Set([ + "unsupported_schema", + "codex_running", + "process_enumeration_failed", + "unsafe_path", + "trigger_collision", + "auto_vacuum_not_incremental", +]); + +export function logGuardResultMayHaveChangedStorage( + result: CodexLogGuardMutationResult | CodexLogGuardCompactionResult, +): boolean { + if (result.ok) return true; + if (result.error === "integrity_check_failed") return result.phase === "after"; + return !REFUSED_BEFORE_CHANGE.has(result.error); +} + +/** Invalidate in-flight storage scans only when the Log Guard operation may have changed storage. */ +function afterLogGuardMutation( + result: CodexLogGuardMutationResult | CodexLogGuardCompactionResult, + response: Response, +): Response { + if (logGuardResultMayHaveChangedStorage(result)) noteStorageMutationCompleted(); + return response; +} + /** Codex Log Guard diagnostics plus explicit protection and maintenance mutations. */ export async function handleStorageLogGuardRoutes(ctx: ManagementContext): Promise { const { req, url, config, deps } = ctx; @@ -129,22 +191,26 @@ export async function handleStorageLogGuardRoutes(ctx: ManagementContext): Promi if (req.method !== "POST") return null; const mode = await readProtectMode(ctx); if (mode instanceof Response) return mode; - return mutationResponse(protectCodexLogs(mode, protectionDeps), ctx); + const result = protectCodexLogs(mode, protectionDeps); + return afterLogGuardMutation(result, mutationResponse(result, ctx)); } if (url.pathname === "/api/storage/codex-logs/unprotect") { if (req.method !== "POST") return null; - return mutationResponse(unprotectCodexLogs(protectionDeps), ctx); + const result = unprotectCodexLogs(protectionDeps); + return afterLogGuardMutation(result, mutationResponse(result, ctx)); } if (url.pathname === "/api/storage/codex-logs/repair") { if (req.method !== "POST") return null; - return mutationResponse(repairCodexLogGuardProtection(protectionDeps), ctx); + const result = repairCodexLogGuardProtection(protectionDeps); + return afterLogGuardMutation(result, mutationResponse(result, ctx)); } if (url.pathname === "/api/storage/codex-logs/compact") { if (req.method !== "POST") return null; - return compactResponse(compactCodexLogs(deps.codexLogGuardMaintenanceDeps), ctx); + const result = compactCodexLogs(deps.codexLogGuardMaintenanceDeps); + return afterLogGuardMutation(result, compactResponse(result, ctx)); } if (url.pathname !== "/api/storage" || req.method !== "GET") return null; @@ -154,7 +220,7 @@ export async function handleStorageLogGuardRoutes(ctx: ManagementContext): Promi // silently folded into CODEX_HOME totals. let storage; try { - storage = scanStorage(); + storage = await sharedStorageScan(resolveCodexHomeDir()); } catch { const fallback = { codexHome: resolveCodexHomeDir(), diff --git a/src/server/management/system-routes.ts b/src/server/management/system-routes.ts index 301baa52231..849cde12540 100644 --- a/src/server/management/system-routes.ts +++ b/src/server/management/system-routes.ts @@ -217,7 +217,7 @@ export async function handleSystemRoutes(ctx: ManagementContext): Promise | null = null; let generation = 0; @@ -64,7 +72,11 @@ export function getStartupHealthSnapshot( const now = deps.now ?? Date.now; if (cached && now() - cached.timestamp < CACHE_TTL_MS) return cached.value; refreshInBackground(config, deps); - return cached ? markStartupHealthDiagnosticStale(cached.value) : conservativeFallback(config); + return staleOrFallback(config); +} + +function staleOrFallback(config: Pick): StartupHealth { + return lastReading ? markStartupHealthDiagnosticStale(lastReading) : conservativeFallback(config); } export function markStartupHealthDiagnosticStale(value: StartupHealth): StartupHealth { @@ -140,7 +152,7 @@ function runProbe(config: Pick): Promise { if (startedGeneration === generation) { cached = { timestamp: (deps.now ?? Date.now)(), value }; + if (!value.diagnosticStale) lastReading = value; } return value; }) - .catch(() => cached ? markStartupHealthDiagnosticStale(cached.value) : conservativeFallback(config)) + .catch(() => staleOrFallback(config)) .finally(() => { // An invalidated probe must never clear the newer generation's flight. if (inflight === probe) inflight = null; @@ -189,7 +202,7 @@ export async function getCachedStartupHealth( ])); if (settled) return settled; } - return cached ? markStartupHealthDiagnosticStale(cached.value) : conservativeFallback(config); + return staleOrFallback(config); } export function invalidateStartupHealthCache(): void { @@ -197,3 +210,9 @@ export function invalidateStartupHealthCache(): void { cached = null; inflight = null; } + +/** Test-only: also forget the last reading, returning to the never-probed state. */ +export function resetStartupHealthCacheForTests(): void { + invalidateStartupHealthCache(); + lastReading = null; +} diff --git a/src/storage/scanner.ts b/src/storage/scanner.ts index b3f081cb6ba..4d4af984aca 100644 --- a/src/storage/scanner.ts +++ b/src/storage/scanner.ts @@ -1,4 +1,5 @@ import { readdirSync, statSync } from "node:fs"; +import { readdir, stat as statAsync } from "node:fs/promises"; import { join } from "node:path"; import { pathToFileURL } from "node:url"; import { Database, constants } from "bun:sqlite"; @@ -110,6 +111,77 @@ function walkFiles(dir: string, relPrefix: string, out: FileEntry[]): void { } } +/** Upper bound on concurrent readdir/stat calls within one async scan. */ +const SCAN_FS_CONCURRENCY = 64; + +type FsLimit = (op: () => Promise) => Promise; + +/** + * FIFO limiter for filesystem calls. A finishing call hands its slot straight to the next + * waiter, so the bound holds exactly. Only leaf fs calls take a slot: a directory walk + * never holds one while awaiting its children, so recursion cannot deadlock the pool. + */ +function createFsLimit(max: number): FsLimit { + let active = 0; + const waiting: Array<() => void> = []; + return async (op: () => Promise): Promise => { + if (active < max) active += 1; + else await new Promise(resolve => waiting.push(resolve)); + try { + return await op(); + } finally { + const next = waiting.shift(); + if (next) next(); + else active -= 1; + } + }; +} + +/** Entries started together at one directory level of an async scan. */ +const SCAN_BATCH_SIZE = 256; + +/** + * Map `items` through `fn` a batch at a time, keeping input order. A directory with tens of + * thousands of entries then has at most one batch of pending tasks at its level, instead of + * one promise per entry queued behind the fs limiter. + */ +async function mapInBatches(items: readonly T[], fn: (item: T) => Promise): Promise { + const results: R[] = []; + for (let start = 0; start < items.length; start += SCAN_BATCH_SIZE) { + const batch = await Promise.all(items.slice(start, start + SCAN_BATCH_SIZE).map(fn)); + for (const result of batch) results.push(result); + } + return results; +} + +/** + * Async twin of {@link walkFiles}: same skip rules, and entries come back in readdir + * order so the report (including `largest` tie order) matches the synchronous scan. + */ +async function walkFilesAsync(dir: string, relPrefix: string, limit: FsLimit): Promise { + let entries; + try { + entries = await limit(() => readdir(dir, { withFileTypes: true })); + } catch { + return []; + } + const parts = await mapInBatches(entries, async (entry): Promise => { + const full = join(dir, entry.name); + const relPath = relPrefix ? `${relPrefix}/${entry.name}` : entry.name; + try { + if (entry.isDirectory()) return await walkFilesAsync(full, relPath, limit); + if (entry.isFile()) { + const stat = await limit(() => statAsync(full)); + return [{ relPath, bytes: stat.size, mtimeMs: stat.mtimeMs }]; + } + } catch { + /* entry vanished mid-scan — diagnostics tolerate racy trees */ + } + return []; + }); + return parts.flat(); +} + function buildBucket(key: StorageBucketKey, files: FileEntry[]): StorageBucket { const bucket: StorageBucket = { key, @@ -169,8 +241,8 @@ function newestVersionedDb(names: string[], pattern: RegExp): string | null { return best; } -export function scanStorage(codexHome: string = resolveCodexHomeDir()): StorageReport { - const files: Record = { +function emptyFiles(): Record { + return { sessions: [], archived_sessions: [], logs_db: [], @@ -179,36 +251,27 @@ export function scanStorage(codexHome: string = resolveCodexHomeDir()): StorageR deletion_manifests: [], other: [], }; +} - let rootNames: string[] = []; - try { - rootNames = readdirSync(codexHome); - } catch (error) { - // A missing home is a normal fresh-machine state — report zeros. Anything else - // (e.g. ENOTDIR: CODEX_HOME points at a file) is a broken setup the caller - // must surface as a scan failure, not silently render as an empty home. - const code = (error as NodeJS.ErrnoException).code; - if (code !== "ENOENT") throw error; - } +/** A missing home is a normal fresh-machine state; any other readdir failure is rethrown. */ +function rootReadFailure(error: unknown): string[] { + // A missing home is a normal fresh-machine state — report zeros. Anything else + // (e.g. ENOTDIR: CODEX_HOME points at a file) is a broken setup the caller + // must surface as a scan failure, not silently render as an empty home. + const code = (error as NodeJS.ErrnoException).code; + if (code !== "ENOENT") throw error; + return []; +} - for (const name of rootNames) { - const full = join(codexHome, name); - let stat; - try { - stat = statSync(full); - } catch { - continue; - } - if (stat.isDirectory()) { - // Quarantine trash (Phase 2) must not inflate "other" or totals. - if (name === TRASH_DIR) continue; - walkFiles(full, name, files[DIR_BUCKETS[name] ?? "other"]); - } else if (stat.isFile()) { - const key: StorageBucketKey = STATE_DB_FILE.test(name) ? "state_db" : LOGS_DB_FILE.test(name) ? "logs_db" : "other"; - files[key].push({ relPath: name, bytes: stat.size, mtimeMs: stat.mtimeMs }); - } - } +function rootFileBucket(name: string): StorageBucketKey { + return STATE_DB_FILE.test(name) ? "state_db" : LOGS_DB_FILE.test(name) ? "logs_db" : "other"; +} +function finishReport( + codexHome: string, + rootNames: string[], + files: Record, +): StorageReport { const buckets = (Object.keys(files) as StorageBucketKey[]).map(key => buildBucket(key, files[key])); const stateDbName = newestVersionedDb(rootNames, STATE_DB_FILE); @@ -236,3 +299,84 @@ export function scanStorage(codexHome: string = resolveCodexHomeDir()): StorageR buckets, }; } + +export function scanStorage(codexHome: string = resolveCodexHomeDir()): StorageReport { + const files = emptyFiles(); + + let rootNames: string[] = []; + try { + rootNames = readdirSync(codexHome); + } catch (error) { + rootNames = rootReadFailure(error); + } + + for (const name of rootNames) { + const full = join(codexHome, name); + let stat; + try { + stat = statSync(full); + } catch { + continue; + } + if (stat.isDirectory()) { + // Quarantine trash (Phase 2) must not inflate "other" or totals. + if (name === TRASH_DIR) continue; + walkFiles(full, name, files[DIR_BUCKETS[name] ?? "other"]); + } else if (stat.isFile()) { + files[rootFileBucket(name)].push({ relPath: name, bytes: stat.size, mtimeMs: stat.mtimeMs }); + } + } + + return finishReport(codexHome, rootNames, files); +} + +/** + * Same report as {@link scanStorage}, but the tree walk uses async fs calls so the server's + * event loop keeps serving while it runs. The synchronous walk measured ~3.7s for a + * 40k-file / 7.5GB CODEX_HOME on Windows, during which every listener (proxy traffic + * included) stalled each time the Storage page loaded. The two sqlite row counts stay + * synchronous: immutable readonly opens measured ~0.1s on the same home. + */ +export async function scanStorageAsync(codexHome: string = resolveCodexHomeDir()): Promise { + const files = emptyFiles(); + const limit = createFsLimit(SCAN_FS_CONCURRENCY); + + let rootNames: string[] = []; + try { + rootNames = await readdir(codexHome); + } catch (error) { + rootNames = rootReadFailure(error); + } + + const roots = await mapInBatches(rootNames, async name => { + const full = join(codexHome, name); + try { + return { name, full, stat: await limit(() => statAsync(full)) }; + } catch { + return null; + } + }); + const walks = await mapInBatches(roots, async root => { + if (!root) return; + if (root.stat.isDirectory()) { + // Quarantine trash (Phase 2) must not inflate "other" or totals. + if (root.name === TRASH_DIR) return; + return { key: DIR_BUCKETS[root.name] ?? "other", entries: await walkFilesAsync(root.full, root.name, limit) }; + } + if (root.stat.isFile()) { + return { + key: rootFileBucket(root.name), + entries: [{ relPath: root.name, bytes: root.stat.size, mtimeMs: root.stat.mtimeMs }], + }; + } + }); + // Appended in root order so bucket contents match the synchronous scan. One entry at a + // time: spreading a large bucket into push() can exceed the engine's argument limit. + for (const walk of walks) { + if (!walk) continue; + const bucket = files[walk.key]; + for (const entry of walk.entries) bucket.push(entry); + } + + return finishReport(codexHome, rootNames, files); +} diff --git a/src/storage/storage-mutation-coordinator.ts b/src/storage/storage-mutation-coordinator.ts index c365420c6aa..bee4b30b2c5 100644 --- a/src/storage/storage-mutation-coordinator.ts +++ b/src/storage/storage-mutation-coordinator.ts @@ -34,8 +34,24 @@ export const MAX_ACTIVE_STORAGE_HOME_SLOTS = 32; const slots = new Map(); const slotGate = createAdmissionGate("storage_home_slots", MAX_ACTIVE_STORAGE_HOME_SLOTS); let releaseMisses = 0; +let mutationEpoch = 0; let testHooks: StorageMutationCoordinatorTestHooks | null = null; +/** + * Advances whenever a CODEX_HOME storage mutation finishes (coordinated cleanup, restore + * and policy runs on lease release; Log Guard maintenance through + * {@link noteStorageMutationCompleted}). An in-flight storage scan that started under an + * older epoch describes the tree before that mutation and must not answer a later read. + */ +export function storageMutationEpoch(): number { + return mutationEpoch; +} + +/** Record a CODEX_HOME mutation that does not go through a coordinator slot. */ +export function noteStorageMutationCompleted(): void { + mutationEpoch += 1; +} + function slotKey(codexHome?: string): string { return resolve(codexHome ?? resolveCodexHomeDir()); } @@ -79,6 +95,7 @@ export function tryBeginStorageMutation( release() { if (!active) return; active = false; + mutationEpoch += 1; const owner = slots.get(key); if (owner?.lease === ownerLease) slots.delete(key); lease.release(); diff --git a/tests/codex-integration/codex-app-server-restart-service.test.ts b/tests/codex-integration/codex-app-server-restart-service.test.ts index 64c1d1c3bb1..e0e9bb23946 100644 --- a/tests/codex-integration/codex-app-server-restart-service.test.ts +++ b/tests/codex-integration/codex-app-server-restart-service.test.ts @@ -14,6 +14,7 @@ import { resetCodexRestartInFlightForTests, } from "../../src/codex/app-server-restart-service"; import type { CodexRestartServiceIo } from "../../src/codex/app-server-restart-service"; +import { resetCodexAppServerCatalogStateCache } from "../../src/codex/app-server-processes"; import type { CodexAppServerProcess } from "../../src/codex/app-server-processes"; import { isCodexRestartResponse } from "../../src/lib/codex-restart-contract"; @@ -275,8 +276,8 @@ describe("performCodexRestart", () => { }); describe("readCodexAppServerState", () => { - test("reports the classifier verdict and a running count", () => { - const state = readCodexAppServerState({ + test("reports the classifier verdict and a running count", async () => { + const state = await readCodexAppServerState({ collectState: () => ({ state: "stale", processes: [{ pid: 1, startedAtMs: 1 }, { pid: 2, startedAtMs: 2 }], @@ -287,13 +288,68 @@ describe("readCodexAppServerState", () => { expect(state).toEqual({ state: "stale", runningCount: 2 }); }); - test("passes unknown through instead of guessing not_running", () => { - const state = readCodexAppServerState({ + test("passes unknown through instead of guessing not_running", async () => { + const state = await readCodexAppServerState({ collectState: () => ({ state: "unknown", processes: [], catalogMtimeMs: null }), }); expect(state).toEqual({ state: "unknown", runningCount: 0 }); }); + + test("the default classifier yields to the event loop while Windows enumeration is slow", async () => { + // The dashboard route calls this with no collectState. The synchronous classifier + // parked the event loop for the whole CIM walk (4-7s on Windows), stalling proxy + // traffic every time the Models page opened. + resetCodexAppServerCatalogStateCache(); + let releaseSnapshots: ((snapshots: CodexAppServerProcess[]) => void) | undefined; + const snapshots = new Promise(resolve => { + releaseSnapshots = resolve; + }); + const reading = readCodexAppServerState({ + stateDeadlineMs: 60_000, + processIo: { + platform: "win32", + listSnapshotsAsync: () => snapshots, + readStartMsBatchAsync: async pids => new Map(pids.map(pid => [pid, 500])), + catalogMtimeMs: () => 1_000, + }, + }); + + const first = await Promise.race([ + reading.then(() => "reading"), + new Promise<"timer">(resolve => setTimeout(() => resolve("timer"), 10)), + ]); + expect(first).toBe("timer"); + + releaseSnapshots?.([proc(42, "/usr/local/bin/codex app-server")]); + await expect(reading).resolves.toEqual({ state: "stale", runningCount: 1 }); + resetCodexAppServerCatalogStateCache(); + }); + + test("a probe slower than the deadline answers unknown, then serves the finished reading", async () => { + resetCodexAppServerCatalogStateCache(); + let releaseSnapshots: ((snapshots: CodexAppServerProcess[]) => void) | undefined; + const snapshots = new Promise(resolve => { + releaseSnapshots = resolve; + }); + const processIo = { + platform: "win32" as const, + listSnapshotsAsync: () => snapshots, + readStartMsBatchAsync: async (pids: readonly number[]) => new Map(pids.map(pid => [pid, 500])), + catalogMtimeMs: () => 1_000, + }; + + await expect(readCodexAppServerState({ stateDeadlineMs: 20, processIo })) + .resolves.toEqual({ state: "unknown", runningCount: 0 }); + + // The probe kept running behind the deadline; once it lands, the next read is served + // from the cache it wrote instead of starting another walk. + releaseSnapshots?.([proc(42, "/usr/local/bin/codex app-server")]); + await new Promise(resolve => setTimeout(resolve, 0)); + await expect(readCodexAppServerState({ stateDeadlineMs: 20, processIo })) + .resolves.toEqual({ state: "stale", runningCount: 1 }); + resetCodexAppServerCatalogStateCache(); + }); }); describe("identity and concurrency protection", () => { diff --git a/tests/codex-integration/codex-restart-route.test.ts b/tests/codex-integration/codex-restart-route.test.ts index f917c83a331..949597c3c50 100644 --- a/tests/codex-integration/codex-restart-route.test.ts +++ b/tests/codex-integration/codex-restart-route.test.ts @@ -50,11 +50,11 @@ function contextFor( } function stubService(overrides: Partial<{ - readState: () => CodexAppServerStateResponse; + readState: () => Promise; performRestart: () => Promise; }> = {}) { return { - readState: overrides.readState ?? (() => STALE_STATE), + readState: overrides.readState ?? (async () => STALE_STATE), performRestart: overrides.performRestart ?? (async () => STOPPED), } as NonNullable; } diff --git a/tests/config/settings-startup-health-seam.test.ts b/tests/config/settings-startup-health-seam.test.ts index f17df8c63e8..b2831e414f6 100644 --- a/tests/config/settings-startup-health-seam.test.ts +++ b/tests/config/settings-startup-health-seam.test.ts @@ -1,6 +1,7 @@ -import { expect, test } from "bun:test"; +import { afterEach, expect, test } from "bun:test"; import { handleManagementAPI, type ManagementApiDeps } from "../../src/server/management-api"; import type { OcxConfig } from "../../src/types"; +import { getCachedStartupHealth, invalidateStartupHealthCache } from "../../src/server/startup-health-cache"; import { ManagementRequest as Request } from "../helpers/management-auth"; import { startupHealthFixture } from "../helpers/startup-health"; @@ -44,3 +45,30 @@ test("settings PUT uses the injected startup-health reader", async () => { startupHealth: { diagnosticStale: true, status: "native" }, }); }); + +afterEach(() => invalidateStartupHealthCache()); + +/** Leave a probe in flight that never settles until released, as a slow Windows service query does. */ +async function holdStartupHealthProbe(): Promise<() => void> { + invalidateStartupHealthCache(); + let release!: () => void; + const pending = new Promise>(resolve => { + release = () => resolve(startupHealthFixture()); + }); + await getCachedStartupHealth({}, { probe: () => pending, waitForProbe: async () => null }); + return release; +} + +test("settings GET answers from the snapshot instead of waiting on a pending probe", async () => { + const release = await holdStartupHealthProbe(); + try { + const req = new Request("http://127.0.0.1:10100/api/settings"); + const response = await handleManagementAPI(req, new URL(req.url), baseConfig(), {}); + + expect(response?.status).toBe(200); + const body = await response!.json() as { startupHealth: { diagnosticStale: boolean } }; + expect(body.startupHealth.diagnosticStale).toBe(true); + } finally { + release(); + } +}); diff --git a/tests/config/settings-stream-mode.test.ts b/tests/config/settings-stream-mode.test.ts index 9385ea78b43..d1977a87ca3 100644 --- a/tests/config/settings-stream-mode.test.ts +++ b/tests/config/settings-stream-mode.test.ts @@ -7,7 +7,7 @@ * backup-and-defaults repair path), and settable alone via PUT (legacy * codexAutoStart-only PUTs keep working). */ -import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; import { Database } from "bun:sqlite"; import { spawnSync } from "node:child_process"; import { existsSync, mkdirSync, mkdtempSync, readFileSync, writeFileSync } from "node:fs"; @@ -233,6 +233,7 @@ describe("GET /api/settings", () => { const { persistEffortClamp, resetCodexRuntimeResolveCacheForTests, + resolveCodexRuntimeAsync, } = await import("../../src/codex/runtime"); resetCodexRuntimeResolveCacheForTests(); @@ -258,6 +259,9 @@ describe("GET /api/settings", () => { try { process.env.CODEX_CLI_PATH = fakeCodex; process.env.PATH = ""; + // Settings serve the runtime stale-while-revalidate; land the probe first so this + // asserts the validated projection rather than the cold deferred answer. + await resolveCodexRuntimeAsync(); const body = await (await getSettings(baseConfig()))!.json() as { codexRuntime?: { path?: string; @@ -295,6 +299,138 @@ describe("GET /api/settings", () => { }); }); +describe("settings codexRuntime snapshot", () => { + /** A launcher whose `--version` takes ~2s, as a real Codex probe can under load. */ + function slowFakeCodex(version: string): string { + mkdirSync(join(TEST_DIR, "slow-bin"), { recursive: true }); + if (process.platform === "win32") { + const path = join(TEST_DIR, "slow-bin", "codex.cmd"); + writeFileSync( + path, + `@echo off\r\n"%SystemRoot%\\System32\\ping.exe" -n 3 127.0.0.1 >nul\r\necho codex-cli ${version}\r\n`, + "utf8", + ); + return path; + } + const path = join(TEST_DIR, "slow-bin", "codex"); + writeFileSync(path, `#!/bin/sh\nsleep 2\necho 'codex-cli ${version}'\n`, { encoding: "utf8", mode: 0o755 }); + return path; + } + + async function withRuntimeEnv(command: string, run: () => Promise): Promise { + const keys = ["CODEX_CLI_PATH", "PATH", "LOCALAPPDATA", "HOME"] as const; + const previous = Object.fromEntries(keys.map(key => [key, process.env[key]])); + try { + process.env.CODEX_CLI_PATH = command; + // No other codex on PATH or in the install roots: only the launcher above is probed. + process.env.PATH = process.platform === "win32" ? "" : "/usr/bin:/bin"; + process.env.LOCALAPPDATA = join(TEST_DIR, "no-codex-app"); + process.env.HOME = join(TEST_DIR, "no-codex-home"); + await run(); + } finally { + for (const key of keys) { + if (previous[key] === undefined) delete process.env[key]; + else process.env[key] = previous[key]; + } + } + } + + test("GET answers without waiting on the runtime probe and serves it once it lands", async () => { + const { resetCodexRuntimeResolveCacheForTests, resolveCodexRuntimeAsync } = await import("../../src/codex/runtime"); + resetCodexRuntimeResolveCacheForTests(); + const launcher = slowFakeCodex("0.200.0"); + try { + await withRuntimeEnv(launcher, async () => { + type Body = { codexRuntime: { version: string | null; source: string } }; + const coldStarted = performance.now(); + const cold = await (await getSettings(baseConfig()))!.json() as Body; + // The sync resolver made this request take the whole ~2s probe. + expect(performance.now() - coldStarted).toBeLessThan(1_000); + expect(cold.codexRuntime).toMatchObject({ version: null, source: "environment" }); + + // The refresh the GET started runs on async exec: timers keep firing meanwhile. + const refresh = resolveCodexRuntimeAsync(); + const tickStarted = performance.now(); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(performance.now() - tickStarted).toBeLessThan(250); + expect((await refresh).runtime.version).toBe("0.200.0"); + + const warmStarted = performance.now(); + const warm = await (await getSettings(baseConfig()))!.json() as Body; + expect(performance.now() - warmStarted).toBeLessThan(1_000); + expect(warm.codexRuntime).toMatchObject({ version: "0.200.0", source: "environment" }); + }); + } finally { + resetCodexRuntimeResolveCacheForTests(); + } + }, 30_000); + + test("an expired memo stays observable while its refresh runs, then gives way to the result", async () => { + // Catalog gather and convergence read the memo through peek. With the refresh off the + // event loop they can now read during it; an expired memo reported as unavailable there + // sent gather to the persisted runtime and got convergence's candidate rejected. + const { + peekCodexRuntimeProcessCache, + resetCodexRuntimeResolveCacheForTests, + resolveCodexRuntimeAsync, + } = await import("../../src/codex/runtime"); + resetCodexRuntimeResolveCacheForTests(); + const launcher = slowFakeCodex("0.200.0"); + const realNow = Date.now.bind(Date); + let offset = 0; + const clock = spyOn(Date, "now").mockImplementation(() => realNow() + offset); + try { + await withRuntimeEnv(launcher, async () => { + await resolveCodexRuntimeAsync(); + const first = peekCodexRuntimeProcessCache(); + expect(first.kind).toBe("available"); + + offset = 20_000; + expect(peekCodexRuntimeProcessCache().kind).toBe("unavailable"); + + const refresh = resolveCodexRuntimeAsync(); + const during = peekCodexRuntimeProcessCache(); + expect(during.kind).toBe("available"); + if (during.kind === "available" && first.kind === "available") { + expect(during.valueIdentity).toBe(first.valueIdentity); + } + + await refresh; + const after = peekCodexRuntimeProcessCache(); + expect(after.kind).toBe("available"); + if (after.kind === "available" && first.kind === "available") { + expect(after.valueIdentity).not.toBe(first.valueIdentity); + } + }); + } finally { + clock.mockRestore(); + resetCodexRuntimeResolveCacheForTests(); + } + }, 30_000); + + test("a runtime switch during the background probe keeps its result out of the memo", async () => { + const { + clearCodexRuntimeResolveCache, + peekCodexRuntimeProcessCache, + resetCodexRuntimeResolveCacheForTests, + resolveCodexRuntimeAsync, + } = await import("../../src/codex/runtime"); + resetCodexRuntimeResolveCacheForTests(); + const launcher = slowFakeCodex("0.200.0"); + try { + await withRuntimeEnv(launcher, async () => { + const refresh = resolveCodexRuntimeAsync(); + // persistCodexRuntime and clearPersistedCodexRuntime invalidate through this. + clearCodexRuntimeResolveCache(); + expect((await refresh).runtime.version).toBe("0.200.0"); + expect(peekCodexRuntimeProcessCache().kind).toBe("unavailable"); + }); + } finally { + resetCodexRuntimeResolveCacheForTests(); + } + }, 30_000); +}); + describe("usage summary retained-store accounting", () => { test("accounts cached summaries and centralized oldest eviction exactly", async () => { for (const range of ["30d", "7d"]) { diff --git a/tests/gui/startup-health-ui.test.ts b/tests/gui/startup-health-ui.test.ts index b66a907f3a8..2a050300c47 100644 --- a/tests/gui/startup-health-ui.test.ts +++ b/tests/gui/startup-health-ui.test.ts @@ -38,7 +38,9 @@ describe("startup health UI decisions", () => { test("settings may seed while unknown or hard-error, but not overwrite a real status", () => { expect(seedStartupHealthFromSettings(null, { status: "protected", diagnosticStale: false })).toBe("protected"); expect(seedStartupHealthFromSettings("error", { status: "protected", diagnosticStale: false })).toBe("protected"); - expect(seedStartupHealthFromSettings(null, { status: "at-risk", diagnosticStale: true })).toBe("at-risk"); + // A stale seed is a server-side placeholder and must not color the chip. + expect(seedStartupHealthFromSettings(null, { status: "at-risk", diagnosticStale: true })).toBeNull(); + expect(seedStartupHealthFromSettings("error", { status: "at-risk", diagnosticStale: true })).toBe("error"); expect(seedStartupHealthFromSettings("at-risk", { status: "protected", diagnosticStale: false })).toBe("at-risk"); }); }); diff --git a/tests/helpers/server-auth-fixture.ts b/tests/helpers/server-auth-fixture.ts index 473eff61a4c..4c9581447e1 100644 --- a/tests/helpers/server-auth-fixture.ts +++ b/tests/helpers/server-auth-fixture.ts @@ -88,7 +88,7 @@ export async function startManagementServerFixture( // response deadline. Keep the production startup and the ordinary 5s test limit. // Neither management endpoint exercises native Codex synchronization or the // developer's installed service. Keep those external owners outside this fixture. - const runtime = spyOn(codexRuntime, "resolveCodexRuntime").mockReturnValue({ + const runtime = spyOn(codexRuntime, "getCodexRuntimeSnapshot").mockReturnValue({ runtime: { command: "codex-fixture", version: null, source: "fallback" }, failures: [], }); try { diff --git a/tests/service/autostart-health.test.ts b/tests/service/autostart-health.test.ts index b681fa7a775..5c96349515b 100644 --- a/tests/service/autostart-health.test.ts +++ b/tests/service/autostart-health.test.ts @@ -5,7 +5,7 @@ import { classifyCodexRouting, hasInjectedCodexRouting } from "../../src/codex/i import { isCodexClientProcess, listCodexClientProcesses } from "../../src/codex/native-profile-processes"; import { collectRoutingAdoption, deriveRoutingAdoption } from "../../src/codex/routing-adoption"; import { handleManagementAPI } from "../../src/server/management-api"; -import { getCachedStartupHealth, getStartupHealthSnapshot, invalidateStartupHealthCache, markStartupHealthDiagnosticStale } from "../../src/server/startup-health-cache"; +import { getCachedStartupHealth, getStartupHealthSnapshot, invalidateStartupHealthCache, markStartupHealthDiagnosticStale, resetStartupHealthCacheForTests } from "../../src/server/startup-health-cache"; import type { OcxConfig } from "../../src/types"; const base = { @@ -280,6 +280,40 @@ describe("Codex startup health", () => { invalidateStartupHealthCache(); }); + test("invalidation keeps the last reading for the snapshot instead of a not-installed fallback", async () => { + // A settings PUT invalidates, then reads the snapshot. Answering with the synthetic + // fallback (serviceInstalled: false, at-risk) flashed a healthy service as at risk + // until the dedicated probe finished. + resetStartupHealthCacheForTests(); + const healthy = { + ...base, + routingKind: "native" as const, + serviceInstalled: true, + serviceViable: true, + serviceEnabled: true, + serviceRunning: true, + }; + const reading = await getCachedStartupHealth( + { codexAutoStart: true }, + { probe: async () => deriveStartupHealth(healthy) }, + ); + expect(reading.diagnosticStale).toBe(false); + expect(reading.serviceInstalled).toBe(true); + + invalidateStartupHealthCache(); + let releaseProbe!: (value: ReturnType) => void; + const pendingProbe = new Promise>(resolve => { + releaseProbe = resolve; + }); + const snapshot = getStartupHealthSnapshot({ codexAutoStart: true }, { probe: async () => pendingProbe }); + expect(snapshot).toEqual(markStartupHealthDiagnosticStale(reading)); + expect(snapshot.serviceInstalled).toBe(true); + + releaseProbe(deriveStartupHealth(healthy)); + await pendingProbe; + resetStartupHealthCacheForTests(); + }); + test("settings snapshot starts a probe without waiting for it", async () => { invalidateStartupHealthCache(); let releaseProbe!: (value: ReturnType) => void; diff --git a/tests/storage/storage-scanner.test.ts b/tests/storage/storage-scanner.test.ts index 67f21d6ec03..ef5cfbb3582 100644 --- a/tests/storage/storage-scanner.test.ts +++ b/tests/storage/storage-scanner.test.ts @@ -3,8 +3,14 @@ import { Database } from "bun:sqlite"; import { mkdirSync, mkdtempSync, readdirSync, statSync, unlinkSync, utimesSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { scanStorage, type StorageBucket, type StorageReport } from "../../src/storage/scanner"; +import { scanStorage, scanStorageAsync, type StorageBucket, type StorageReport } from "../../src/storage/scanner"; import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { logGuardResultMayHaveChangedStorage, sharedStorageScan } from "../../src/server/management/storage-log-guard-routes"; +import { + noteStorageMutationCompleted, + storageMutationEpoch, + tryBeginStorageMutation, +} from "../../src/storage/storage-mutation-coordinator"; const OLD_MTIME = new Date("2026-01-02T03:04:05Z"); const MID_MTIME = new Date("2026-03-04T05:06:07Z"); @@ -269,3 +275,115 @@ describe("scanStorage", () => { expect(otherPaths.some(p => p.includes(".trash"))).toBe(false); }, 15_000); }); + +describe("scanStorageAsync", () => { + const withoutTimestamp = (report: StorageReport) => ({ ...report, generatedAt: 0 }); + + test("produces the same report as the synchronous scan", async () => { + fixtureHome = buildFixtureHome(); + mkdirSync(join(fixtureHome, ".trash", "123"), { recursive: true }); + writeFileSync(join(fixtureHome, ".trash", "123", "rollout-quarantined.jsonl"), "q".repeat(5000)); + + const asyncReport = await scanStorageAsync(fixtureHome); + expect(withoutTimestamp(asyncReport)).toEqual(withoutTimestamp(scanStorage(fixtureHome))); + }, 15_000); + + test("reports zeros for a missing home and rejects a home that is a file", async () => { + fixtureHome = buildFixtureHome(); + const missing = await scanStorageAsync(join(fixtureHome, "does-not-exist")); + expect(missing.total).toEqual({ bytes: 0, fileCount: 0 }); + + const filePath = join(fixtureHome, "not-a-dir"); + writeFileSync(filePath, "x"); + await expect(scanStorageAsync(filePath)).rejects.toThrow(); + }, 15_000); + + test("yields to the event loop while walking the tree", async () => { + // GET /api/storage runs this on the server's event loop; the synchronous walk parked + // it for the whole CODEX_HOME tree (~3.7s on a 40k-file home). A timer queued before + // the scan must get to run before the scan finishes. + fixtureHome = buildFixtureHome(); + let ticked = false; + setTimeout(() => { ticked = true; }, 0); + let tickedBeforeDone = false; + await scanStorageAsync(fixtureHome).then(() => { tickedBeforeDone = ticked; }); + expect(tickedBeforeDone).toBe(true); + }, 15_000); + + test("performs zero writes under CODEX_HOME (read-only invariant)", async () => { + fixtureHome = buildFixtureHome(); + const before = snapshotTree(fixtureHome); + + await scanStorageAsync(fixtureHome); + + const after = snapshotTree(fixtureHome); + expect([...after.keys()].sort()).toEqual([...before.keys()].sort()); + for (const [path, stat] of before) { + expect(after.get(path)).toEqual(stat); + } + }, 15_000); +}); + +describe("shared storage scan flights", () => { + const report = (bytes: number): StorageReport => ({ + codexHome: "home", + generatedAt: 1, + total: { bytes, fileCount: 1 }, + buckets: [], + }); + + test("a read after a completed storage mutation starts a new scan instead of joining the old one", async () => { + const home = join(tmpdir(), `ocx-scan-flight-${process.pid}-${Date.now()}`); + const releases: Array<(value: StorageReport) => void> = []; + let calls = 0; + const scan = () => { + calls += 1; + return new Promise(resolve => releases.push(resolve)); + }; + + const first = sharedStorageScan(home, scan); + const joined = sharedStorageScan(home, scan); + expect(calls).toBe(1); + + // e.g. /api/storage/codex-logs/compact finishing while the first walk is still running + noteStorageMutationCompleted(); + const afterMutation = sharedStorageScan(home, scan); + expect(calls).toBe(2); + + // The older walk settling must not retire the newer flight. + releases[0]!(report(100)); + expect(await first).toEqual(report(100)); + expect(await joined).toEqual(report(100)); + const laterRead = sharedStorageScan(home, scan); + expect(calls).toBe(2); + + releases[1]!(report(40)); + expect(await afterMutation).toEqual(report(40)); + expect(await laterRead).toEqual(report(40)); + }); + + test("only Log Guard results that may have changed storage invalidate scans", () => { + // A refused compact (e.g. unsupported_schema) exits before touching the database; + // invalidating on it let repeated refused requests start overlapping scans. + for (const error of ["unsupported_schema", "codex_running", "process_enumeration_failed", "unsafe_path", "auto_vacuum_not_incremental"] as const) { + expect(logGuardResultMayHaveChangedStorage({ ok: false, error })).toBe(false); + } + expect(logGuardResultMayHaveChangedStorage({ ok: false, error: "trigger_collision" })).toBe(false); + expect(logGuardResultMayHaveChangedStorage({ ok: false, error: "integrity_check_failed", phase: "before" })).toBe(false); + + expect(logGuardResultMayHaveChangedStorage({ ok: false, error: "integrity_check_failed", phase: "after" })).toBe(true); + expect(logGuardResultMayHaveChangedStorage({ ok: false, error: "busy" })).toBe(true); + expect(logGuardResultMayHaveChangedStorage({ ok: false, error: "database_error" })).toBe(true); + expect(logGuardResultMayHaveChangedStorage({ ok: false, error: "config_write_failed" })).toBe(true); + }); + + test("releasing a coordinated mutation lease advances the storage mutation epoch", () => { + const home = join(tmpdir(), `ocx-scan-epoch-${process.pid}-${Date.now()}`); + const before = storageMutationEpoch(); + const gate = tryBeginStorageMutation("cleanup", home); + expect(gate.acquired).toBe(true); + expect(storageMutationEpoch()).toBe(before); + if (gate.acquired) gate.lease.release(); + expect(storageMutationEpoch()).toBe(before + 1); + }); +});