diff --git a/src/cli/index.ts b/src/cli/index.ts index a9b411e9d32..88ab11392d7 100755 --- a/src/cli/index.ts +++ b/src/cli/index.ts @@ -67,7 +67,7 @@ import { quarantinePendingTeardown, } from "../config/pending-teardown"; import { collectStatus, deadProxyRoutingAdviceLines, detectMissingCodexCatalogPath, hubStatusLines, missingCodexCatalogLines, remoteHubBannerLine, remoteHubStatusLines, unusedProxyWarningLines } from "./status"; -import { endpointsToProve, everyEndpointProvenDown, sharedTeardownAuthorized, type UninstallObservation } from "./uninstall-plan"; +import { endpointsToProve, everyEndpointProvenDownAsync, sharedTeardownAuthorized, type UninstallObservation } from "./uninstall-plan"; import { takeFlag } from "./runtime-api"; import { parseStartOptions, StartArgsError } from "./start-args"; @@ -85,7 +85,14 @@ import { SpendLedgerOwnerError } from "../lib/spend-ledger-owner"; import { redactUrlForLog } from "../lib/redact"; import { dispatchCommand, decideBusyPreferredPort, decideStartWithLiveOwner } from "./dispatch"; import { AuxiliaryListenerBindError, findAvailablePort, isAddrInUse, PortUnavailableError, shouldPersistSelectedPort, waitForPortAvailable } from "../server/ports"; -import { findLiveProxy, probeHostname, probePortOwner, START_OWNERSHIP_LIVENESS, type LiveProxy } from "../server/proxy-liveness"; +import { + findLiveProxy, + probeEndpointLiveness, + probeHostname, + probePortOwner, + START_OWNERSHIP_LIVENESS, + type LiveProxy, +} from "../server/proxy-liveness"; import { createReadinessGate } from "../server/readiness"; import { isApiAuthRequired } from "../server/auth-cors"; import { runReady, type ReadyArgs } from "./ready"; @@ -1021,8 +1028,7 @@ async function handleStopUnlocked() { // An obligation that cannot name its endpoint cannot be proven discharged. if (!endpoint) return false; try { - const { probeProxyLiveness } = await import("../update/proxy-liveness-probe.mjs"); - return probeProxyLiveness(endpoint.port, endpoint.hostname) === "dead"; + return await probeEndpointLiveness(endpoint) === "dead"; } catch { // A probe that could not run is not evidence of absence. return false; @@ -1436,11 +1442,10 @@ async function handleUninstall() { /** Definitive "nothing is answering" on the endpoint this home would serve. */ const proxyEndpointProvenDown = async (): Promise => { try { - const { probeProxyLiveness } = await import("../update/proxy-liveness-probe.mjs"); // Every candidate, not just the preferred one: a stale runtime record pointing at a // closed port would otherwise "prove" a live proxy on the configured port is gone. const endpoints = endpointsToProve(readRuntimePort(), loadConfig()); - return everyEndpointProvenDown(endpoints, e => probeProxyLiveness(e.port, e.hostname)); + return await everyEndpointProvenDownAsync(endpoints, probeEndpointLiveness); } catch { return false; } diff --git a/src/cli/resolve.ts b/src/cli/resolve.ts index a4dcc851645..87c3fa69174 100644 --- a/src/cli/resolve.ts +++ b/src/cli/resolve.ts @@ -37,9 +37,14 @@ import { readConfigDiagnostics, type ConfigDiagnostics } from "../config"; import { getConfigDir } from "../config/paths"; import { readRuntimePort } from "../config/process-state"; import { packageVersion } from "../lib/package-version"; -import { findLiveProxy, START_OWNERSHIP_LIVENESS, type LiveProxy } from "../server/proxy-liveness"; -import { endpointsToProve, everyEndpointProvenDown, type ProbeEndpoint } from "./uninstall-plan"; -import { probeProxyLiveness } from "../update/proxy-liveness-probe.mjs"; +import { + findLiveProxy, + probeEndpointLiveness, + START_OWNERSHIP_LIVENESS, + type EndpointLiveness, + type LiveProxy, +} from "../server/proxy-liveness"; +import { endpointsToProve, everyEndpointProvenDownAsync, type ProbeEndpoint } from "./uninstall-plan"; /** Wire version of the resolve document. Bump only on an incompatible shape change. */ export const RESOLVE_SCHEMA = "ocx-resolve/1"; @@ -103,8 +108,8 @@ export interface ResolveIo { findLive?: () => Promise; /** Runtime-port record reader; production default is readRuntimePort. */ readRuntime?: () => { port?: number; hostname?: string } | null; - /** Tri-state endpoint probe; production default is the updater's probeProxyLiveness. */ - probeEndpoint?: (endpoint: ProbeEndpoint) => "live" | "dead" | "unknown"; + /** Tri-state endpoint probe; production default runs in-process for compiled standalone binaries. */ + probeEndpoint?: (endpoint: ProbeEndpoint) => EndpointLiveness | Promise; cliVersion?: () => string; stdout?: { log: (s: string) => void }; stderr?: { error: (s: string) => void }; @@ -174,11 +179,7 @@ export async function runResolve(args: ResolveArgs, io: ResolveIo = {}): Promise const readDiagnostics = io.readDiagnostics ?? readConfigDiagnostics; const findLive = io.findLive ?? (() => findLiveProxy(START_OWNERSHIP_LIVENESS)); const readRuntime = io.readRuntime ?? readRuntimePort; - // The updater's tri-state probe takes (port, hostname) and is plain .mjs (untyped); - // adapt it to the endpoint-shaped seam here. Its own return vocabulary is the - // closed "live" | "dead" | "unknown" set. - const probeEndpoint = io.probeEndpoint - ?? ((endpoint: ProbeEndpoint) => probeProxyLiveness(endpoint.port, endpoint.hostname) as "live" | "dead" | "unknown"); + const probeEndpoint = io.probeEndpoint ?? probeEndpointLiveness; const cliVersion = io.cliVersion ?? packageVersion; const configHome = configDir(); let diagnostics: ConfigDiagnostics; @@ -212,7 +213,7 @@ export async function runResolve(args: ResolveArgs, io: ResolveIo = {}): Promise // authorise starting a second runtime. let provenDown = false; try { - provenDown = everyEndpointProvenDown(endpointsToProve(readRuntime(), diagnostics.config), probeEndpoint); + provenDown = await everyEndpointProvenDownAsync(endpointsToProve(readRuntime(), diagnostics.config), probeEndpoint); } catch { // A probe that cannot run is not evidence of absence. provenDown = false; diff --git a/src/cli/status-probes.ts b/src/cli/status-probes.ts index d3848c95366..d4196ba5cff 100644 --- a/src/cli/status-probes.ts +++ b/src/cli/status-probes.ts @@ -1,5 +1,5 @@ import { readPidFileValue, readRuntimePort } from "../config/process-state"; -import { isOpencodexHealthz, probeHostname } from "../server/proxy-liveness"; +import { isConnectionRefused, isOpencodexHealthz, probeHostname } from "../server/proxy-liveness"; import { directLocalHttpFetch } from "../server/direct-local-http"; import { isProcessAlive } from "../lib/process-control"; @@ -26,23 +26,7 @@ export function proxyHealthFailureReason(error: unknown, signal: AbortSignal): " : "unreachable"; } -/** - * "Nothing is listening" is narrower than "the probe failed". `unreachable` covers every - * non-abort failure, including a socket that was ACCEPTED and then reset — which is what - * an in-flight start looks like mid-bind. Only a connect-phase refusal proves the port is - * free, so this reads the underlying errno instead of the display string. - */ -export function isConnectionRefused(error: unknown): boolean { - for (let current: unknown = error, depth = 0; current instanceof Error && depth < 4; depth++) { - const code = (current as { code?: unknown }).code; - if (code === "ECONNREFUSED" || code === "ConnectionRefused") return true; - // Bun surfaces the refusal as a plain message on some platforms; the errno name is - // still the discriminator, not a substring of arbitrary prose. - if (typeof code === "string" && code.endsWith("ECONNREFUSED")) return true; - current = (current as { cause?: unknown }).cause; - } - return false; -} +export { isConnectionRefused } from "../server/proxy-liveness"; /** * A proxy killed by a native trap or SIGKILL never runs the exit cleanup that removes diff --git a/src/cli/uninstall-plan.ts b/src/cli/uninstall-plan.ts index 0e1df1cd2df..4bfb9d4a7a6 100644 --- a/src/cli/uninstall-plan.ts +++ b/src/cli/uninstall-plan.ts @@ -84,3 +84,12 @@ export function everyEndpointProvenDown( if (endpoints.length === 0) return false; return endpoints.every(e => probe(e) === "dead"); } + +export async function everyEndpointProvenDownAsync( + endpoints: readonly ProbeEndpoint[], + probe: (e: ProbeEndpoint) => Promise<"live" | "dead" | "unknown"> | "live" | "dead" | "unknown", +): Promise { + if (endpoints.length === 0) return false; + const results = await Promise.all(endpoints.map(e => probe(e))); + return results.every(result => result === "dead"); +} diff --git a/src/server/proxy-liveness.ts b/src/server/proxy-liveness.ts index 594be867b02..98202e42d86 100644 --- a/src/server/proxy-liveness.ts +++ b/src/server/proxy-liveness.ts @@ -32,6 +32,8 @@ export interface HealthzIdentity { guiPairCapability?: unknown; } +export type EndpointLiveness = "live" | "dead" | "unknown"; + export interface LivenessIo { fetchFn?: typeof fetch; readPidFn?: () => number | null; @@ -85,6 +87,11 @@ export const START_OWNERSHIP_LIVENESS: Pick Promise; + export interface LiveProxy { pid: number | null; port: number; @@ -148,6 +155,74 @@ export function isOpencodexHealthz(body: HealthzIdentity | null): boolean { return body.status === "ok" && typeof body.version === "string" && typeof body.uptime === "number"; } +/** + * "Nothing is listening" is narrower than "the probe failed". Only a connect-phase refusal + * proves the endpoint is free; a timeout, reset, or other transport failure leaves the + * question open. + */ +export function isConnectionRefused(error: unknown): boolean { + const visit = (current: unknown, depth: number): boolean => { + if (depth >= 4) return false; + if (current === null || (typeof current !== "object" && typeof current !== "function")) return false; + const record = current as { code?: unknown; cause?: unknown; errors?: unknown }; + if (record.code === "ECONNREFUSED" || record.code === "ConnectionRefused") return true; + if (typeof record.code === "string" && record.code.endsWith("ECONNREFUSED")) return true; + if (Array.isArray(record.errors) && record.errors.length > 0) { + // One connect attempt fanned out over several addresses reports a single AggregateError. + // Only a unanimous refusal proves the endpoint is free: a bundle that mixes ECONNREFUSED + // with a timeout means one address answered nothing at all, and an address whose state is + // unreadable is unknown, not absence. Collapsing it to "refused" is how a second runtime + // gets started on a port that already has one. + return record.errors.every(error => visit(error, depth + 1)); + } + return visit(record.cause, depth + 1); + }; + return visit(error, 0); +} + +async function classifyHealthz( + url: string, + fetchFn: LivenessFetch, + timeoutMs: number, +): Promise { + try { + const response = await fetchFn(url, { signal: AbortSignal.timeout(timeoutMs) }); + if (response.status !== 200) return "unknown"; + const body = (await response.json().catch(() => undefined)) as HealthzIdentity | null | undefined; + if (body === undefined) return "unknown"; + return isOpencodexHealthz(body) ? "live" : "dead"; + } catch (error) { + return isConnectionRefused(error) ? "dead" : "unknown"; + } +} + +/** + * Tri-state probe of one endpoint, the in-process counterpart of + * `src/update/proxy-liveness-probe.mjs`. Only a connect-phase refusal or a clean 200 that is + * not ours proves "dead"; a timeout, reset, non-200 or unreadable body leaves the question + * open. Loopback endpoints are checked on both IPv4 and IPv6 because a listener may bind only + * one family. Runs in-process because a compiled standalone binary cannot fork `execPath -e`. + */ +export async function probeEndpointLiveness( + endpoint: { port: number; hostname?: string }, + io: Pick = {}, +): Promise { + if (!Number.isFinite(endpoint.port) || endpoint.port <= 0 || endpoint.port > 65535) return "dead"; + const fetchFn = io.fetchFn ?? directLocalHttpFetch; + const timeoutMs = io.timeoutMs ?? 1500; + let sawUnknown = false; + for (const hostname of loopbackProbeHosts(endpoint.hostname)) { + const result = await classifyHealthz( + `http://${hostname}:${endpoint.port}/healthz`, + fetchFn, + timeoutMs, + ); + if (result === "live") return "live"; + if (result === "unknown") sawUnknown = true; + } + return sawUnknown ? "unknown" : "dead"; +} + /** Identity-checked /healthz probe; null when unreachable, non-OK, or not our proxy. */ export async function proxyIdentityAt( port: number, diff --git a/tests/cli/cli-resolve.test.ts b/tests/cli/cli-resolve.test.ts index 69fbac9764b..a7177ff241b 100644 --- a/tests/cli/cli-resolve.test.ts +++ b/tests/cli/cli-resolve.test.ts @@ -127,6 +127,21 @@ describe("runResolve", () => { expect(parsed.port.effective).toBe(RESOLVE_DEFAULT_PORT); }); + test("accepts async dead probes for every candidate endpoint", async () => { + const lines: string[] = []; + const code = await runResolve({ json: true }, { + configDir: () => "/h", + readDiagnostics: () => ({ config: {}, source: "default", error: null } as ConfigDiagnostics), + findLive: async () => null, + readRuntime: () => ({ port: 10110, hostname: "127.0.0.1" }), + probeEndpoint: async () => "dead", + cliVersion: () => "1.2.3", + stdout: { log: value => lines.push(value) }, + }); + expect(code).toBe(0); + expect((JSON.parse(lines[0]!) as { liveness: { status: string } }).liveness.status).toBe("absent-proven"); + }); + test("an undecidable probe is unknown, and unknown is never answered as absent", async () => { // The launch decision keys on this verdict: a timed-out probe or a listener that // withholds /healthz must exit 1 rather than let the caller start a second runtime. @@ -152,8 +167,8 @@ describe("runResolve", () => { test("absence requires every endpoint dead, not just the configured one", async () => { // The runtime record can point at a live port while the configured port refuses; // answering from the configured port alone would shadow-start over the record. - // everyEndpointProvenDown short-circuits on the first non-dead answer: an unknown - // runtime endpoint defeats the proof without the configured one being probed. + // Every candidate is probed: an unknown runtime endpoint defeats the proof even when the + // configured endpoint is dead. const seen: string[] = []; const code = await runResolve({ json: true }, { configDir: () => "/h", diff --git a/tests/cli/uninstall.test.ts b/tests/cli/uninstall.test.ts index ea984f274fd..1b0e6f399b5 100644 --- a/tests/cli/uninstall.test.ts +++ b/tests/cli/uninstall.test.ts @@ -243,7 +243,7 @@ describe("uninstall gates shared teardown on a proven service stop", () => { }); }); test("proof covers every distinct endpoint, not just the preferred one", async () => { - const { endpointsToProve, everyEndpointProvenDown } = await import("../../src/cli/uninstall-plan"); + const { endpointsToProve, everyEndpointProvenDown, everyEndpointProvenDownAsync } = await import("../../src/cli/uninstall-plan"); // A stale runtime record pointing at a closed port, and the live proxy on the // configured one. Probing only the runtime candidate reports "dead" for a port nobody @@ -267,6 +267,8 @@ describe("uninstall gates shared teardown on a proven service stop", () => { expect(endpointsToProve(null, {})).toEqual([{ hostname: "127.0.0.1", port: 10100 }]); // An empty set is not proof of anything. expect(everyEndpointProvenDown([], () => "dead")).toBe(false); + expect(await everyEndpointProvenDownAsync(endpoints, async () => "dead")).toBe(true); + expect(await everyEndpointProvenDownAsync([], async () => "dead")).toBe(false); // A nonsense runtime port is skipped rather than probed. expect(endpointsToProve({ port: 0 }, { port: 10100 })).toEqual([{ hostname: "127.0.0.1", port: 10100 }]); }); @@ -284,7 +286,7 @@ describe("uninstall gates shared teardown on a proven service stop", () => { .toBeLessThan(windowStep.indexOf("observed.respawnWindowVerified = true;")); // And the proof itself asks every candidate. expect(fn).toContain("endpointsToProve(readRuntimePort(), loadConfig())"); - expect(fn).toContain("everyEndpointProvenDown(endpoints, e => probeProxyLiveness(e.port, e.hostname))"); + expect(fn).toContain("everyEndpointProvenDownAsync(endpoints, probeEndpointLiveness)"); }); const safeTeardown: UninstallObservation = { diff --git a/tests/providers/xai/grok-lifecycle.test.ts b/tests/providers/xai/grok-lifecycle.test.ts index 4c930e57f5d..e6e42c80759 100644 --- a/tests/providers/xai/grok-lifecycle.test.ts +++ b/tests/providers/xai/grok-lifecycle.test.ts @@ -299,7 +299,7 @@ describe("Grok fence lifecycle wiring", () => { expect(noPidBranch).toContain("stopFailed = true;"); expect(noPidBranch).toContain("ownershipBlocked = true;"); const gateFn = sliceFn(CLI_SOURCE, "const abandonedTeardownIsSafeToFinish", "let stopFailed = false;"); - expect(gateFn).toContain('probeProxyLiveness(endpoint.port, endpoint.hostname) === "dead"'); + expect(gateFn).toContain('probeEndpointLiveness(endpoint) === "dead"'); expect(gateFn).toContain("return false;"); }); diff --git a/tests/server/proxy-liveness.test.ts b/tests/server/proxy-liveness.test.ts index d675f04c8b8..ec47ab6de97 100644 --- a/tests/server/proxy-liveness.test.ts +++ b/tests/server/proxy-liveness.test.ts @@ -1,4 +1,5 @@ import { describe, expect, test } from "bun:test"; +import { createServer } from "node:net"; import { createReadinessGate, runStartupReadinessSync, @@ -7,7 +8,9 @@ import { DEFAULT_PROBE_TIMEOUT_MS, findLiveProxy, isOpencodexHealthz, + isConnectionRefused, loopbackProbeHosts, + probeEndpointLiveness, probeHostname, probePortOwner, probeReadiness, @@ -90,6 +93,96 @@ describe("probeHostname", () => { }); }); +describe("probeEndpointLiveness", () => { + test("classifies identity, foreign, non-200, refusal, timeout, and invalid ports", async () => { + const endpoint = { port: 10100, hostname: "127.0.0.1" }; + const fakeFetch = (body: unknown, status = 200) => (async () => healthz(body, status)) as typeof fetch; + expect(await probeEndpointLiveness(endpoint, { fetchFn: fakeFetch(OURS) })).toBe("live"); + expect(await probeEndpointLiveness(endpoint, { fetchFn: fakeFetch({ status: "ok" }) })).toBe("dead"); + expect(await probeEndpointLiveness(endpoint, { fetchFn: fakeFetch(OURS, 503) })).toBe("unknown"); + const refusedServer = createServer(); + await new Promise((resolve, reject) => { + refusedServer.once("error", reject); + refusedServer.listen(0, "127.0.0.1", () => resolve()); + }); + const refusedPort = (refusedServer.address() as { port: number }).port; + await new Promise((resolve, reject) => { + refusedServer.close(error => error ? reject(error) : resolve()); + }); + expect(await probeEndpointLiveness({ port: refusedPort, hostname: "127.0.0.1" })).toBe("dead"); + expect(await probeEndpointLiveness(endpoint, { + fetchFn: (async () => { throw new DOMException("aborted", "AbortError"); }) as typeof fetch, + })).toBe("unknown"); + expect(await probeEndpointLiveness(endpoint, { + fetchFn: (async () => { throw new Error("connection reset"); }) as typeof fetch, + })).toBe("unknown"); + expect(await probeEndpointLiveness({ port: 0 }, { fetchFn: fakeFetch(OURS) })).toBe("dead"); + }); + + test("checks both loopback families sequentially", async () => { + const seen: string[] = []; + const fetchFn = (async (url: string) => { + seen.push(url); + if (url.startsWith("http://127.0.0.1:10100")) { + throw Object.assign(new Error("refused"), { code: "ECONNREFUSED" }); + } + return healthz(OURS); + }) as typeof fetch; + expect(await probeEndpointLiveness({ port: 10100, hostname: "::" }, { fetchFn })).toBe("live"); + expect(seen).toEqual([ + "http://127.0.0.1:10100/healthz", + "http://[::1]:10100/healthz", + ]); + + seen.length = 0; + const refused = (async (url: string) => { + seen.push(url); + throw Object.assign(new Error("refused"), { code: "ECONNREFUSED" }); + }) as typeof fetch; + expect(await probeEndpointLiveness({ port: 10100, hostname: "::" }, { fetchFn: refused })).toBe("dead"); + expect(seen).toEqual([ + "http://127.0.0.1:10100/healthz", + "http://[::1]:10100/healthz", + ]); + + seen.length = 0; + const mixed = (async (url: string) => { + seen.push(url); + if (url.startsWith("http://127.0.0.1:10100")) { + throw Object.assign(new Error("refused"), { code: "ECONNREFUSED" }); + } + throw new DOMException("timed out", "TimeoutError"); + }) as typeof fetch; + expect(await probeEndpointLiveness({ port: 10100, hostname: "::" }, { fetchFn: mixed })).toBe("unknown"); + expect(seen).toEqual([ + "http://127.0.0.1:10100/healthz", + "http://[::1]:10100/healthz", + ]); + }); +}); + +describe("isConnectionRefused", () => { + test("recognizes aggregate socket refusals", () => { + const refused = Object.assign(new Error("refused"), { code: "ECONNREFUSED" }); + expect(isConnectionRefused(new AggregateError([refused]))).toBe(true); + expect(isConnectionRefused(new AggregateError([ + Object.assign(new Error("timeout"), { code: "ETIMEDOUT" }), + ]))).toBe(false); + }); + + test("a mixed aggregate is not proof of absence", () => { + // Happy-eyeballs style fan-out puts every address in one error. If one address refused and + // another never answered, the endpoint's state is unknown: the refusal speaks only for the + // address that produced it. + const refused = Object.assign(new Error("refused"), { code: "ECONNREFUSED" }); + const timedOut = Object.assign(new Error("timeout"), { code: "ETIMEDOUT" }); + expect(isConnectionRefused(new AggregateError([refused, timedOut]))).toBe(false); + expect(isConnectionRefused(new AggregateError([timedOut, refused]))).toBe(false); + expect(isConnectionRefused(new AggregateError([refused, refused]))).toBe(true); + expect(isConnectionRefused(new AggregateError([]))).toBe(false); + }); +}); + describe("proxyIdentityAt", () => { test("returns the reported pid for our proxy", async () => { const identity = await proxyIdentityAt(10100, {}, { fetchFn: (async () => healthz(OURS)) as typeof fetch });