diff --git a/src/client/connect.ts b/src/client/connect.ts index 67ea8d0968b..7aa3e65e005 100644 --- a/src/client/connect.ts +++ b/src/client/connect.ts @@ -104,6 +104,8 @@ export interface LinkClientCredential { } export interface ClientConnectDeps { + /** Abort enrollment network work and refuse subsequent writes; rollback still drains. */ + signal?: AbortSignal; fetchImpl?: typeof fetch; now?: () => Date; lifecycleLockDeps?: ClientLifecycleLockDeps; @@ -539,6 +541,18 @@ export async function connectClient( options: ConnectOptions, deps: ClientConnectDeps = {}, ): Promise { + deps.signal?.throwIfAborted(); + const rawFetch = deps.fetchImpl ?? fetch; + const fetchImpl: typeof fetch = deps.signal ? Object.assign(async (...[input, init = {}]: Parameters) => { + deps.signal!.throwIfAborted(); + const signals = [deps.signal, init.signal, input instanceof Request ? input.signal : undefined] + .filter((signal): signal is AbortSignal => signal != null); + return rawFetch(input, { ...init, signal: AbortSignal.any(signals), redirect: "manual" }); + }, { preconnect: rawFetch.preconnect }) : rawFetch; + const assertActiveConnectingState = (fingerprint?: string) => { + deps.signal?.throwIfAborted(); + assertConnectingState(fingerprint); + }; let serverUrl = ""; let managementUrl = ""; let linkAdmissionToken: string | null = null; @@ -586,7 +600,7 @@ export async function connectClient( throw new Error("link mode requires a valid local config port"); } withClientLifecycleSync(() => withConfigMutationLockSync(() => { - assertConnectingState(); + assertActiveConnectingState(); const externalProvider = currentExternalCodexModelProvider(); if (externalProvider) throw new Error("connect refused: an external Codex provider owns config.toml"); if (linkMode) { @@ -603,7 +617,7 @@ export async function connectClient( } }), deps.lifecycleLockDeps); - const upstreamFetch = deps.fetchImpl ?? fetch; + const upstreamFetch = fetchImpl; const readinessFetch = linkMode ? (async (input, init = {}) => { const url = input instanceof Request ? input.url : String(input); @@ -622,18 +636,18 @@ export async function connectClient( managementUrl, localGuiOrigin(), options.credential.value, - { fetchImpl: deps.fetchImpl }, + { fetchImpl }, ); cleanupCredential = { kind: "gui-session", value: session }; } else if (!linkMode && options.credential.kind === "admin") { cleanupCredential = { kind: "admin", value: options.credential.value }; } if (!linkMode) { - issued = await issueClientKey(managementUrl, cleanupCredential!, clientKeyName(), { fetchImpl: deps.fetchImpl }); + issued = await issueClientKey(managementUrl, cleanupCredential!, clientKeyName(), { fetchImpl }); } const initialFiles = withClientLifecycleSync(() => withConfigMutationLockSync(() => { - assertConnectingState(linkMode ? pendingConnectFingerprint! : undefined); + assertActiveConnectingState(linkMode ? pendingConnectFingerprint! : undefined); const persisted = linkMode ? (() => { const current = readServiceApiTokenState(); @@ -656,7 +670,7 @@ export async function connectClient( if (!admissionToken) throw new Error("client admission credential unavailable"); const apiKeyId = linkMode ? (options.credential as LinkClientCredential).apiKeyId : issued!.id; const catalog = await downloadClientCatalog(serverUrl, admissionToken, { - fetchImpl: deps.fetchImpl, + fetchImpl, timeoutMs: options.catalogTimeoutMs, }); // Fail closed BEFORE the write (#4207). The hub being reachable and the credential working @@ -666,7 +680,7 @@ export async function connectClient( // than writing one and restoring it afterwards. assertClientCatalogCompatible(catalog.body, deps.catalogCompatibility); writtenCatalogFingerprint = withClientLifecycleSync(() => withConfigMutationLockSync(() => { - assertConnectingState(persisted.fingerprint); + assertActiveConnectingState(persisted.fingerprint); atomicWriteFile(DEFAULT_CATALOG_PATH, catalog.body); return sha256(catalog.body); }), deps.lifecycleLockDeps); @@ -679,7 +693,7 @@ export async function connectClient( routingTarget: target, catalogPath: DEFAULT_CATALOG_PATH, journalOwner: { kind: "client", apiKeyId }, - beforeClientWrite: () => assertConnectingState(persisted.fingerprint), + beforeClientWrite: () => assertActiveConnectingState(persisted.fingerprint), }); if (!preflight.success) throw new Error(preflight.message); @@ -688,7 +702,7 @@ export async function connectClient( routingTarget: target, catalogPath: DEFAULT_CATALOG_PATH, journalOwner: { kind: "client", apiKeyId }, - beforeClientWrite: () => assertConnectingState(persisted.fingerprint), + beforeClientWrite: () => assertActiveConnectingState(persisted.fingerprint), }); if (!injected.success || injected.status === "skipped") throw new Error(injected.message); injectionCommitted = true; @@ -715,7 +729,7 @@ export async function connectClient( ...(linkMode ? { transport: "link" as const, link: linkMetadata } : {}), }; withClientLifecycleSync(() => withConfigMutationLockSync(() => { - assertConnectingState(persisted.fingerprint); + assertActiveConnectingState(persisted.fingerprint); clearClientConnectPending(persisted.fingerprint); commitClientConnection(connection); committed = true; diff --git a/src/client/link-join.ts b/src/client/link-join.ts index 5059a4a5515..0527fc21114 100644 --- a/src/client/link-join.ts +++ b/src/client/link-join.ts @@ -1,6 +1,7 @@ import { randomBytes } from "node:crypto"; import { hostname } from "node:os"; import { isPortAvailable } from "../server/ports"; +import { scanListenPidsForAddress, type ListenPidScan } from "../server/port-reclaim"; import { isLinkPort, JOIN_TUNNEL_PORT_MAX, JOIN_TUNNEL_PORT_MIN } from "../link/ports"; import { buildExecArgv, REMOTE_COMMAND_NOT_FOUND, remoteOcxArgv } from "../link/ssh-argv"; import { sshFailureHint, sshRunnerErrorHint, type SshRunner, type SshRunResult } from "../link/ssh-runner"; @@ -22,6 +23,7 @@ import type { OcxConnectedClientId } from "../types"; const JOIN_TUNNEL_READY_TIMEOUT_MS = 15_000; const JOIN_TUNNEL_POLL_MS = 100; +const JOIN_TUNNEL_SPAWN_GRACE_MS = 100; const JOIN_REVOKE_TIMEOUT_MS = 30_000; const JOIN_CONFIRM_TTL_MS = 5 * 60_000; const JOIN_PORT_ATTEMPTS = 32; @@ -74,6 +76,12 @@ export interface ClientLinkJoinDeps { hostname?: () => string; randomBytes?: (size: number) => Uint8Array; fetchImpl?: typeof fetch; + /** + * LISTEN-owner probe for the tunnel port; defaults to the netstat/lsof/ss scan. + * Receives the loopback address the tunnel binds so listeners on unrelated + * addresses do not confuse the readiness check. + */ + scanListenPids?: (port: number, address?: string) => ListenPidScan; spawnTunnel?: (spec: { linkId: string; alias: string; @@ -222,8 +230,10 @@ async function compensateStaleSidecar(deps: ClientLinkJoinDeps): Promise { } } +/** Wait for the live tunnel's authenticated readiness, retaining the deadline after failed ownership rechecks. */ async function waitForReady( deps: ClientLinkJoinDeps, + tunnel: ClientLinkTunnelHandle, port: number, key: string, ): Promise { @@ -231,13 +241,47 @@ async function waitForReady( const now = deps.now ?? Date.now; const sleep = deps.sleep ?? ((ms: number) => new Promise(resolve => setTimeout(resolve, ms))); const deadline = now() + JOIN_TUNNEL_READY_TIMEOUT_MS; + const tunnelExited = tunnel.exited.then(() => { throw new ClientLinkJoinError("join_tunnel_failed"); }); + await Promise.race([ + tunnelExited, + new Promise(resolve => setTimeout(resolve, JOIN_TUNNEL_SPAWN_GRACE_MS)), + ]); + const listenPids = deps.scanListenPids ?? scanListenPidsForAddress; + // The tunnel binds 127.0.0.1; a listener on a different loopback or interface address + // never receives our requests, so ownership is only judged among sockets that serve it. + const tunnelAddress = "127.0.0.1"; for (;;) { try { - const response = await fetchImpl(`http://127.0.0.1:${port}/readyz`, { - headers: { "x-opencodex-api-key": key }, - }); - if (response.status === 200) return; - if (response.status === 401) throw new ClientLinkJoinError("admission_failed"); + // A squatter answering the 401 challenge would otherwise collect the issued key: + // the only listener allowed a keyed request is the ssh process we spawned — it owns + // the port only after a successful bind, and ExitOnForwardFailure makes it exit when + // it cannot take the port. An unverifiable scan stays "not ready", never a pass. + const ownership = listenPids(port, tunnelAddress); + if (ownership.ok && ownership.pids.length === 1 && ownership.pids[0] === tunnel.pid) { + // Never follow redirects: a port occupant must not reroute the challenge, and a + // redirected keyed request would carry the issued key to an unrelated listener. + const probe = await Promise.race([ + tunnelExited, + fetchImpl(`http://127.0.0.1:${port}/readyz`, { redirect: "manual" }), + ]); + if (probe.status === 401) { + // Ownership can flip between the probe and the keyed request (a squatter + // takes the port after the tunnel dies). Re-scan in the same iteration and + // skip only the keyed request on failure, not the deadline check and sleep. + const recheck = listenPids(port, tunnelAddress); + if (recheck.ok && recheck.pids.length === 1 && recheck.pids[0] === tunnel.pid) { + const response = await Promise.race([ + tunnelExited, + fetchImpl(`http://127.0.0.1:${port}/readyz`, { + headers: { "x-opencodex-api-key": key }, + redirect: "manual", + }), + ]); + if (response.status === 200) return; + if (response.status === 401) throw new ClientLinkJoinError("admission_failed"); + } + } + } } catch (error) { if (error instanceof ClientLinkJoinError) throw error; } @@ -256,6 +300,7 @@ function requireConfirmedHost(deps: ClientLinkJoinDeps, alias: string): JoinConf return confirmed; } +/** Issue and enroll a confirmed Home link, compensating failures before committing the connection. */ export async function joinHome(deps: ClientLinkJoinDeps, input: { alias: string }): Promise<{ linkId: string; apiKeyId: string }> { const confirmed = requireConfirmedHost(deps, input.alias); await compensateStaleSidecar(deps); @@ -308,30 +353,52 @@ export async function joinHome(deps: ClientLinkJoinDeps, input: { alias: string configDir: deps.configDir, knownHostsFile: deps.knownHostsFile, }); - await waitForReady(deps, tunnelPort, issued.key); + await waitForReady(deps, tunnel, tunnelPort, issued.key); } catch (error) { const code = error instanceof ClientLinkJoinError ? error.code : "join_tunnel_failed"; await rollback(deps, issued.linkId, tunnel); throw new ClientLinkJoinError(code); } + const enrollmentAbort = new AbortController(); + let enrollment: Promise | undefined; + let enrollmentFinished = false; try { + if (!tunnel) throw new ClientLinkJoinError("join_tunnel_failed"); const connect = deps.connect ?? connectClient; - await connect({ - serverUrl: `http://127.0.0.1:${tunnelPort}`, - managementUrl: `http://127.0.0.1:${tunnelPort}`, - credential: { kind: "link", apiKeyId: issued.apiKeyId, key: issued.key }, - transport: "link", - link: { tunnelPort, linkId: issued.linkId }, - selectedClients: deps.selectedClients ?? ["codex", "claude"], - managementTransport: "direct", - }, { - fetchImpl: deps.fetchImpl, - ...deps.connectDeps, + // Keep watching the tunnel until the connection commits: an exited tunnel + // must not let the issued key ride out to whatever next holds the port. + const tunnelFailure = tunnel.exited.then(() => { + const error = new ClientLinkJoinError("join_tunnel_failed"); + if (!enrollmentFinished) enrollmentAbort.abort(error); + throw error; }); - } catch { + const signal = deps.connectDeps?.signal + ? AbortSignal.any([enrollmentAbort.signal, deps.connectDeps.signal]) : enrollmentAbort.signal; + enrollment = Promise.resolve().then(() => connect({ + serverUrl: `http://127.0.0.1:${tunnelPort}`, + managementUrl: `http://127.0.0.1:${tunnelPort}`, + credential: { kind: "link", apiKeyId: issued.apiKeyId, key: issued.key }, + transport: "link", + link: { tunnelPort, linkId: issued.linkId }, + selectedClients: deps.selectedClients ?? ["codex", "claude"], + managementTransport: "direct", + }, { + fetchImpl: deps.fetchImpl, + ...deps.connectDeps, + signal, + })); + await Promise.race([tunnelFailure, enrollment]); + enrollmentFinished = true; + } catch (error) { + enrollmentAbort.abort(error); + // An observed tunnel exit is not completion of the losing exchange. Drain its + // cancellation and local rollback before stopping/revoking the shared link. + if (enrollment) await Promise.allSettled([enrollment]); await rollback(deps, issued.linkId, tunnel); - throw new ClientLinkJoinError("join_connect_failed"); + throw new ClientLinkJoinError(error instanceof ClientLinkJoinError ? error.code : "join_connect_failed"); + } finally { + enrollmentFinished = true; } await stopTunnel(tunnel); diff --git a/src/server/port-reclaim.ts b/src/server/port-reclaim.ts index 42fdc2305c2..97827cb60ae 100644 --- a/src/server/port-reclaim.ts +++ b/src/server/port-reclaim.ts @@ -18,6 +18,16 @@ export type ListenPidScan = | { ok: true; pids: number[] } | { ok: false; error?: string }; +/** One listening socket with its bound local address (host part only). */ +export interface ListenEntry { + pid: number; + address: string; +} + +export type ListenEntryScan = + | { ok: true; listeners: ListenEntry[] } + | { ok: false; error?: string }; + export type ReclaimListenPortOptions = WaitForPortOptions & { /** * When true AND `onlyKillPids` is a non-empty allowlist, those PIDs may be @@ -60,12 +70,49 @@ export type ReclaimListenPortOptions = WaitForPortOptions & { sleepMs?: (ms: number) => Promise; }; +/** Split `host:port`/`[v6]:port` on a numeric port boundary; returns the host part. */ +function listenHost(token: string): string { + const bracketed = /^(\[[0-9a-fA-F:.]+\]):/.exec(token); + if (bracketed) return bracketed[1].slice(1, -1).toLowerCase(); + // Only a trailing : is a port; a bare "::" or hostname wildcard has none. + const withPort = /^(.*):(\d+)$/.exec(token); + return (withPort ? withPort[1] : token).toLowerCase(); +} + +/** Normalize a listen-address host: strips brackets and the IPv4-mapped prefix. */ +export function normalizeListenAddress(token: string): string { + let host = listenHost(token); + if (host.startsWith("::ffff:")) host = host.slice(7); + return host; +} + +/** Normalize a bare bind address (no port): drops brackets, keeps bare IPv6 whole. */ +function bareListenAddress(address: string): string { + let host = address.replace(/^\[|\]$/g, "").toLowerCase(); + if (host.startsWith("::ffff:")) host = host.slice(7); + return host; +} + +const WILDCARD_LISTEN_HOSTS = new Set(["", "*", "0.0.0.0", "::"]); + /** - * Parse `netstat -ano` (Windows) / `netstat -anlp` listen lines for a port. - * Exported for unit tests. + * Whether a socket bound to `listenerAddress` also serves connections to `bound` — + * exact match, or a wildcard listener, or a wildcard `bound` (the caller listens on + * every address). IPv4-mapped IPv6 forms of the same address are equalized first. */ -export function parseListenPidsFromNetstat(output: string, port: number): number[] { - const pids = new Set(); +export function listenAddressServes(listenerAddress: string, bound: string): boolean { + const listener = normalizeListenAddress(listenerAddress); + const want = bareListenAddress(bound); + return WILDCARD_LISTEN_HOSTS.has(listener) || WILDCARD_LISTEN_HOSTS.has(want) + || listener === want; +} + +/** + * Parse `netstat -ano` (Windows) / `netstat -anlp` listen lines for a port, keeping + * each distinct PID/address pair. Exported for unit tests. + */ +export function parseListenEntriesFromNetstat(output: string, port: number): ListenEntry[] { + const entries = new Map(); const portSuffix = `:${port}`; for (const rawLine of output.split(/\r?\n/)) { const line = rawLine.trim(); @@ -88,9 +135,66 @@ export function parseListenPidsFromNetstat(output: string, port: number): number : unixPid ? Number(unixPid[1]) : NaN; - if (Number.isSafeInteger(pid) && pid > 0) pids.add(pid); + if (Number.isSafeInteger(pid) && pid > 0) { + const address = normalizeListenAddress(parts[localIdx]); + entries.set(`${pid}|${address}`, { pid, address }); + } } - return [...pids]; + return [...entries.values()]; +} + +/** Parse netstat LISTEN owners, deduplicating PIDs after preserving their addresses. */ +export function parseListenPidsFromNetstat(output: string, port: number): number[] { + return [...new Set(parseListenEntriesFromNetstat(output, port).map(entry => entry.pid))]; +} + +/** + * Parse `ss -Hltnp` rows for a port, keeping each distinct PID/address pair. A row + * without a `pid=` attribution (another user's socket) is dropped rather than + * reported unverifiable. Exported for unit tests. + */ +export function parseListenEntriesFromSs(output: string, port: number): ListenEntry[] { + const entries = new Map(); + const portSuffix = `:${port}`; + for (const rawLine of output.split(/\r?\n/)) { + const line = rawLine.trim(); + if (!/^LISTEN\b/i.test(line)) continue; + const parts = line.split(/\s+/); + // LISTEN users:(...) + const localIdx = parts.findIndex(part => part.endsWith(portSuffix) || part.endsWith(`]:${port}`)); + if (localIdx < 0) continue; + const pidMatch = /pid=(\d+)/.exec(line); + const pid = pidMatch ? Number(pidMatch[1]) : NaN; + if (Number.isSafeInteger(pid) && pid > 0) { + const address = normalizeListenAddress(parts[localIdx]); + entries.set(`${pid}|${address}`, { pid, address }); + } + } + return [...entries.values()]; +} + +/** + * Parse `lsof -nP -iTCP: -sTCP:LISTEN` output (without -t), keeping each + * distinct PID/address pair. The NAME column is the last address token, optionally + * followed by `(LISTEN)`; skip the header and nonnumeric PIDs. Exported for tests. + */ +export function parseListenEntriesFromLsof(output: string, port: number): ListenEntry[] { + const entries = new Map(); + const portSuffix = `:${port}`; + for (const rawLine of output.split(/\r?\n/)) { + const line = rawLine.trim(); + if (!line || /^COMMAND\b/.test(line)) continue; + const parts = line.split(/\s+/); + const pid = /^\d+$/.test(parts[1] ?? "") ? Number(parts[1]) : NaN; + if (!Number.isSafeInteger(pid) || pid <= 0) continue; + let addressIdx = parts.length - 1; + if (/^\(.*\)$/.test(parts[addressIdx] ?? "")) addressIdx -= 1; + const address = parts[addressIdx] ?? ""; + if (!address.endsWith(portSuffix) && !address.endsWith(`]:${port}`)) continue; + const normalized = normalizeListenAddress(address); + entries.set(`${pid}|${normalized}`, { pid, address: normalized }); + } + return [...entries.values()]; } function normalizeListenPidScan(result: ListenPidScan | number[]): ListenPidScan { @@ -121,50 +225,84 @@ function readWindowsNetstatAno(): string { } /** - * Scan for PIDs currently LISTENing on `port`. - * Distinguishes probe failure (`ok: false`) from a successful empty result. + * Scan for the sockets currently LISTENing on `port`, with each listener's bound + * local address. Distinguishes probe failure (`ok: false`) from a successful empty + * result. POSIX backends are tried in order — `lsof`, `ss` (iproute2, the only + * scanner on minimal Linux installs), then `netstat` — and a missing scanner falls + * through to the next instead of failing the scan. */ -export function scanListenPids(port: number): ListenPidScan { +export function scanListenEntries(port: number): ListenEntryScan { if (!Number.isFinite(port) || port <= 0 || port > 65535) { return { ok: false, error: "invalid port" }; } + const scanned = Math.trunc(port); try { if (process.platform === "win32") { - return { ok: true, pids: parseListenPidsFromNetstat(readWindowsNetstatAno(), port) }; + return { ok: true, listeners: parseListenEntriesFromNetstat(readWindowsNetstatAno(), scanned) }; } + const errors: string[] = []; try { - const output = execFileSync("lsof", ["-nP", `-iTCP:${port}`, "-sTCP:LISTEN", "-t"], { + const output = execFileSync("lsof", ["-nP", `-iTCP:${scanned}`, "-sTCP:LISTEN"], { encoding: "utf-8", stdio: ["ignore", "pipe", "ignore"], timeout: 3000, }); - return { - ok: true, - pids: output - .split(/\r?\n/) - .map(line => Number(line.trim())) - .filter(pid => Number.isSafeInteger(pid) && pid > 0), - }; - } catch (lsofErr) { - try { - const output = execFileSync("netstat", ["-anlp"], { - encoding: "utf-8", - stdio: ["ignore", "pipe", "ignore"], - timeout: 3000, - }); - return { ok: true, pids: parseListenPidsFromNetstat(output, Math.trunc(port)) }; - } catch (netstatErr) { - return { - ok: false, - error: `lsof/netstat unavailable: ${String(lsofErr)} / ${String(netstatErr)}`, - }; - } + return { ok: true, listeners: parseListenEntriesFromLsof(output, scanned) }; + } catch (error) { + errors.push(`lsof: ${String(error)}`); + } + try { + const output = execFileSync("ss", ["-Hltnp"], { + encoding: "utf-8", + stdio: ["ignore", "pipe", "ignore"], + timeout: 3000, + }); + return { ok: true, listeners: parseListenEntriesFromSs(output, scanned) }; + } catch (error) { + errors.push(`ss: ${String(error)}`); + } + try { + const output = execFileSync("netstat", ["-anlp"], { + encoding: "utf-8", + stdio: ["ignore", "pipe", "ignore"], + timeout: 3000, + }); + return { ok: true, listeners: parseListenEntriesFromNetstat(output, scanned) }; + } catch (error) { + errors.push(`netstat: ${String(error)}`); } + return { ok: false, error: `no listener scanner available (${errors.join(" / ")})` }; } catch (error) { return { ok: false, error: String(error) }; } } +/** + * Scan for PIDs currently LISTENing on `port`. + * Distinguishes probe failure (`ok: false`) from a successful empty result. + */ +export function scanListenPids(port: number): ListenPidScan { + const scan = scanListenEntries(port); + if (!scan.ok) return { ok: false, error: scan.error }; + return { ok: true, pids: [...new Set(scan.listeners.map(entry => entry.pid))] }; +} + +/** + * PIDs LISTENing on `port` that actually serve `address`: listeners bound to that + * exact address plus wildcards (0.0.0.0/::). A listener on a different loopback or + * interface address (e.g. 127.0.0.2 while the tunnel binds 127.0.0.1) never receives + * the connection and must not block or qualify a readiness check. + */ +export function scanListenPidsForAddress(port: number, address = "127.0.0.1"): ListenPidScan { + const scan = scanListenEntries(port); + if (!scan.ok) return { ok: false, error: scan.error }; + const pids = new Set(); + for (const entry of scan.listeners) { + if (listenAddressServes(entry.address, address)) pids.add(entry.pid); + } + return { ok: true, pids: [...pids] }; +} + /** Best-effort PIDs currently LISTENing on `port`. Empty on probe failure. */ export function listListenPids(port: number): number[] { const scan = scanListenPids(port); diff --git a/structure/remote-link.md b/structure/remote-link.md index 81d52007e08..61b8ec4b693 100644 --- a/structure/remote-link.md +++ b/structure/remote-link.md @@ -18,6 +18,10 @@ Every remote `ocx` call goes through `remoteOcxArgv`, which runs `sh -c` with a `src/server/management/link-routes.ts` accepts `POST /api/link/join` with exactly `{ "alias": string }`. The route admits the same dashboard sessions as the Home-side routes (see [Dashboard admission](#dashboard-admission)), so a standalone computer turns itself into a Child from its own dashboard. A Tailscale identity session receives `403 tailscale_session_refused`, any other caller `403 forbidden`, a runtime that is not standalone `409 standalone_required`, and a standalone whose live listener port (`resolveListenPort` in `src/server/management/system-restart.ts`) is not its configured `port`, or cannot be determined, `409 join_port_mismatch`, because the client runtime it restarts into binds exactly the configured port. These gates run before link state is read and before any SSH. The alias must have a confirmed, unexpired host entry in the same route state. Before choosing a port or issuing a new link, a valid stale client sidecar is compensated over SSH unless the machine is already connected to that link; a successful revoke clears the sidecar, while a failed revoke preserves it and returns `join_rollback_failed` with the link id. A corrupt sidecar is left for the next successful write. A successful join issues the Home link through SSH, records the client sidecar, starts the client tunnel and connects the client, then returns `202 { "linkId": string, "alias": string, "restarting": true }`. +During enrollment, `src/client/link-join.ts` watches the spawned SSH tunnel through its 100 ms spawn grace, readiness checks and the connection attempt. Before an unauthenticated `/readyz` probe and again before the keyed request, the only LISTEN owner serving `127.0.0.1:` must be that tunnel PID. Both requests use `redirect: "manual"`; only a 401 challenge permits the keyed request. A failed, empty, foreign or ambiguous ownership recheck withholds the key and reaches the same 15-second deadline check and up-to-100 ms polling delay as any other not-ready iteration. Repeated recheck failures therefore reach rollback instead of bypassing it. The deadline is checked between operations, not an independent per-fetch cancellation timer. An observed tunnel exit winning the readiness or connection race fails the join and runs compensation. + +The enrollment scanner in `src/server/port-reclaim.ts` retains each distinct normalized `(PID, bound address)` pair from Windows `netstat` or the POSIX `lsof`, `ss`, then `netstat` fallback chain. It filters entries for the requested loopback address before deduplicating PIDs, so another socket owned by the same process cannot overwrite the relevant listener. Duplicate rows and IPv4-mapped aliases of the same address still collapse, and the PID-only API continues to return unique PIDs. This enrollment scan does not replace the runtime supervisor's asynchronous ownership check described under [Client link transport](#client-link-transport); the check-to-connect race described there remains. + `src/client/link-state.ts` stores `/link/client-link.json` with mode 0600. The sidecar contains exactly `alias`, `hubHostKeyFingerprint`, `tunnelPort`, `peerListenerPort` and `linkId`; it contains no key. The client tunnel port uses `MIN_LINK_PORT = 1024` through `MAX_LINK_PORT = 65535` and `isLinkPort`; the Home listener port keeps its existing 1–65535 contract. A dashboard join picks a free port at random from `JOIN_TUNNEL_PORT_MIN = 20000` through `JOIN_TUNNEL_PORT_MAX = 29999` (`chooseJoinTunnelPort` in `src/client/link-join.ts`), below the macOS, Windows and Linux ephemeral ranges, so an outgoing connection rarely holds the port when the tunnel comes back after a reboot. The persisted port of an existing link is never rewritten. `src/client/link-tunnel.ts` owns the client `ssh -L 127.0.0.1::127.0.0.1:` process. A client runtime starts that supervisor when link transport and a matching sidecar are present, and it drives the tunnel with `CLIENT_TUNNEL_RETRY_POLICY`, so the Child reconnects by itself after sleep, an outage or a crash. A spawned or respawned tunnel, whether connecting, reconnecting or retrying from failed, is promoted to connected only when a keyed `GET http://127.0.0.1:/readyz` (the link key in `x-opencodex-api-key`, `cache: "no-store"`) proves the link: a 200, or a 503 whose body (read up to 4 KiB) carries `service: "opencodex"`. The Home's link listener answers 401 before it reaches `/readyz`, so that 503 only means the Home's own start-up readiness is pending or failed, which relayed requests do not depend on. Until then the probe backs off from one to five seconds, or up to 30 seconds while the link reads failed; while a request is held it runs on every one-second check instead. One probe runs at a time, detached from the check, and `stop()` or the end of the tunnel it probes aborts it, so a slow Home never delays noticing a disconnect or stopping. While connected the same probe runs every 30 seconds and is display-only: a 401 or 403 reports `probe: "unauthorized"`, the readiness 503 `probe: "home_not_ready"`, and any other failure `probe: "home_unreachable"` in the supervisor status, which the Child's `GET /api/link/status` shows as the child `reason`; it never cuts the tunnel. The key comes from the runtime's cached key source, so no probe reads the token file. The supervisor's one-second check is an unref'd interval that stats the sidecar and `config.json` and parses one again only after it changed. It stops the tunnel and schedules the existing standalone recycle when the connection is no longer connected with link transport, the link id no longer matches, or the sidecar disappears; an unreadable sidecar or connection state acts only after three consecutive checks, so one read during a write never ends a healthy link. Normal shutdown, including recycle, stops the client supervisor before the client listener; it sends TERM, waits at most five seconds, then sends KILL. @@ -56,4 +60,8 @@ Codex keeps the standalone loopback routing: `routingTarget` in `src/client/conn `src/client/link-relay.ts` forwards exactly the `linkRouteAllowed` routes from `src/link/routes.ts` through the tunnel. It drops the caller's `Authorization`, `x-api-key`, `x-opencodex-api-key`, `chatgpt-account-id` and `cookie` and sends the link key as `Authorization: Bearer`, the wire an `env_key` config sent; `GET /v1/usage` takes it as `x-opencodex-api-key`, the only header that route admits. The Home admits the key and serves the Child with its own accounts. The request body is streamed chunk by chunk with the caller's `Content-Length` and a byte-counting cap at the inbound limit (`resolveInboundBodyLimitBytes`, 256 MiB by default); a larger declared or streamed body answers 413. A lone `Transfer-Encoding: chunked` without `Content-Length` is admitted as a standalone admits it, because the listener has already de-chunked the body; any other Transfer-Encoding, or one next to a `Content-Length`, answers 400. The Home's response headers may take up to 300 seconds, and a caller abort ends the wait sooner. SSE passes through chunk by chunk with caller-abort propagation and a 300-second idle limit, other response bodies stream under the same byte cap, and the relay answers 503 with Retry-After while the tunnel is down. The client supervisor is the relay's tunnel gate (`LinkTunnelGate`): only while the tunnel is connecting or reconnecting (including the start of the client runtime) does a relayed request wait, for at most 15 seconds (`LINK_RELAY_HOLD_MS`) from its first wait and with at most 64 requests waiting, before it is forwarded once; a connected tunnel costs one `pending()` call per request, and a failed one answers 503 at once. A forward whose connection was refused sent nothing, so while the tunnel reconnects it may wait again and be sent again inside the same 15 seconds, provided the streamed body was never read or cancelled; any other failure (a reset, a timeout, a failure after the body started) is never replayed. Both the Child's machine listener and the Home's hub-link listener bind with `idleTimeout: 255`, the public listener's limit, so a held or slow turn is not cut by Bun's 10-second default. Like a standalone data route, a relayed request then lifts its own idle timer (`server.timeout(req, 0)` in `src/client/link-ingress.ts`), so a quiet stretch longer than 255 seconds inside a long generation is not cut either; the relay's header deadline, SSE idle limit and caller abort bound the wait instead. Hub transport keeps the 4 MiB management-relay listener bound and its default idle limit. Link mode waits for the configured port without signalling its holder, then binds there or fails; `src/client/runtime.ts` passes the cached link key, tunnel status and tunnel gate through `bindClientListener` to every bind attempt. Link mode turns the management relay off and refuses key rotation and revocation, which belong to the hub. -Regression coverage lives in `tests/clients/link-ssh-argv.test.ts`, `tests/clients/link-ssh-config.test.ts`, `tests/clients/link-tunnel-state.test.ts`, `tests/clients/link-store.test.ts`, `tests/clients/link-boundary.test.ts`, `tests/clients/link-routes.test.ts`, `tests/clients/client-link-connect.test.ts`, `tests/clients/client-link-relay.test.ts`, `tests/clients/client-machine-listener.test.ts`, `tests/clients/client-link-status.test.ts`, `tests/clients/client-link-runtime.test.ts`, `tests/codex-integration/injection-link-websocket.test.ts`, `tests/clients/link-supervisor.test.ts`, `tests/clients/link-status-projection.test.ts`, `tests/clients/link-admission-wait.test.ts`, `tests/clients/link-fingerprint.test.ts`, `tests/cli/cli-link.test.ts`, `tests/server/link-management-routes.test.ts`, `tests/server/link-join-route.test.ts`, `tests/server/link-listener-lifecycle.test.ts`, `tests/clients/client-link-teardown.test.ts` and `gui/tests/remote-link.test.tsx`. +Regression coverage lives in `tests/clients/link-ssh-argv.test.ts`, `tests/clients/link-ssh-config.test.ts`, `tests/clients/link-tunnel-state.test.ts`, `tests/clients/link-store.test.ts`, `tests/clients/link-boundary.test.ts`, `tests/clients/link-routes.test.ts`, `tests/clients/client-link-connect.test.ts`, `tests/clients/client-link-relay.test.ts`, `tests/clients/client-machine-listener.test.ts`, `tests/clients/client-link-status.test.ts`, `tests/clients/client-link-runtime.test.ts`, `tests/codex-integration/injection-link-websocket.test.ts`, `tests/clients/link-supervisor.test.ts`, `tests/clients/link-status-projection.test.ts`, `tests/clients/link-admission-wait.test.ts`, `tests/clients/link-fingerprint.test.ts`, `tests/cli/cli-link.test.ts`, `tests/server/link-management-routes.test.ts`, `tests/server/link-join-route.test.ts`, `tests/server/port-reclaim.test.ts`, `tests/server/link-listener-lifecycle.test.ts`, `tests/clients/client-link-teardown.test.ts` and `gui/tests/remote-link.test.tsx`. + +### Enrollment cancellation + +Tunnel exit aborts the enrollment signal and its physical fetches. Every later enrollment write rechecks that signal; the join waits for the cancelled enrollment and its local rollback before tunnel/key compensation. This does not eliminate the separate listener-observation-to-connect race. diff --git a/tests/clients/client-link-connect.test.ts b/tests/clients/client-link-connect.test.ts index 9f323c9ead2..335ade02440 100644 --- a/tests/clients/client-link-connect.test.ts +++ b/tests/clients/client-link-connect.test.ts @@ -130,6 +130,34 @@ describe("client link connection contracts", () => { }); }); + test("cancellation during catalog download prevents late enrollment writes and drains token rollback", async () => { + await withLinkHome(async home => { + const prior = '{"models":[{"id":"prior"}]}\n'; + writeFileSync(DEFAULT_CATALOG_PATH, prior); + const abort = new AbortController(); + let cancelledFetch = false; + await expect(connectClient(linkOptions(), { + signal: abort.signal, + fetchImpl: async (input, init) => { + if (String(input).endsWith("/readyz")) return Response.json({ + service: "opencodex", version: "0.0.0", uptime: 1, pid: 1, port: 34567, + status: "ready", protocol: 1, minimumClientProtocol: 1, + managementUrl: "http://127.0.0.1:34567", + }); + abort.abort(new Error("fixture enrollment cancelled")); + cancelledFetch = init?.signal?.aborted === true; + // Even a fetch implementation returning after abort cannot authorize a write. + return Response.json({ models: [] }); + }, + lifecycleLockDeps: { lockPath: join(home, "lifecycle.sqlite") }, + })).rejects.toThrow("fixture enrollment cancelled"); + expect(cancelledFetch).toBe(true); + expect(readFileSync(DEFAULT_CATALOG_PATH, "utf8")).toBe(prior); + expect(readServiceApiTokenState()).toEqual({ kind: "absent" }); + expect(readClientConnectionState()).toEqual({ kind: "disconnected" }); + }); + }); + test("catalog failure removes the pending link token and leaves config.client unset", async () => { await withLinkHome(async home => { await expect(connectClient(linkOptions(), { diff --git a/tests/server/link-join-route.test.ts b/tests/server/link-join-route.test.ts index 653e749bf37..8c722c3226b 100644 --- a/tests/server/link-join-route.test.ts +++ b/tests/server/link-join-route.test.ts @@ -1,5 +1,7 @@ import { describe, expect, test, spyOn } from "bun:test"; import { chooseJoinTunnelPort, ClientLinkJoinError, joinHome, type ClientLinkJoinDeps } from "../../src/client/link-join"; +import { spawnClientLinkTunnel } from "../../src/client/link-tunnel"; +import type { ListenPidScan } from "../../src/server/port-reclaim"; import { isLinkPort, JOIN_TUNNEL_PORT_MAX, JOIN_TUNNEL_PORT_MIN } from "../../src/link/ports"; import { handleLinkRoutes, type LinkRouteState } from "../../src/server/management/link-routes"; import type { ManagementContext } from "../../src/server/management/context"; @@ -14,6 +16,7 @@ function isWrappedIssue(argv: readonly string[]): boolean { return (argv.at(-1) ?? "").startsWith(`${quoteRemote(remoteOcxArgv(["link", "issue", "--alias"]))} `); } +/** Select only the wrapped revoke command, not other SSH traffic. */ function revokeCalls(calls: readonly string[][]): string[][] { return calls.filter(argv => argv.at(-1) === REVOKE_COMMAND); } @@ -21,6 +24,7 @@ const API_KEY_ID = "link-key-1"; const KEY = `ocx_data_${"a".repeat(40)}`; const FINGERPRINT = `SHA256:${"a".repeat(32)}`; +/** Record issue and compensation calls without starting an SSH process. */ function runnerFor(calls: string[][], issueResult = true): SshRunner { return { run: async argv => { @@ -41,16 +45,29 @@ function runnerFor(calls: string[][], issueResult = true): SshRunner { }; } +/** Keep the fake tunnel alive until its owner stops it. */ function tunnelFor(order: string[]) { return { pid: 123, - exited: Promise.resolve(0), + exited: new Promise(() => {}), stop: async () => { order.push("stop-tunnel"); }, }; } +/** Return the link-auth challenge before accepting the issued key. */ +function challengedFetch(order?: string[]) { + return async (_input: RequestInfo | URL, init?: RequestInit) => { + const authed = new Headers(init?.headers).get("x-opencodex-api-key") === KEY; + order?.push(authed ? "readyz:key" : "readyz:probe"); + return new Response(null, { status: authed ? 200 : 401 }); + }; +} + +/** Supply isolated join dependencies and attribute the listener to the fake tunnel. */ function joinDeps(overrides: Partial = {}): ClientLinkJoinDeps { - const calls = overrides.runner ? [] : []; + const calls: string[][] = []; + let tunnelPid = 0; + const spawn = overrides.spawnTunnel ?? spawnClientLinkTunnel; return { runner: overrides.runner ?? runnerFor(calls), knownHostsFile: "/tmp/ocx-known-hosts", @@ -61,7 +78,14 @@ function joinDeps(overrides: Partial = {}): ClientLinkJoinDe hostname: () => "client-host", readSidecar: () => null, readConnectionState: () => ({ kind: "disconnected" }), + scheduleRestart: () => {}, ...overrides, + spawnTunnel: (spec, spawnDeps) => { + const handle = spawn(spec, spawnDeps); + tunnelPid = handle.pid; + return handle; + }, + scanListenPids: overrides.scanListenPids ?? (() => ({ ok: true, pids: [tunnelPid] })), }; } @@ -258,10 +282,8 @@ describe("client initiated link join", () => { hostname: () => "client-host", writeState: state => { order.push("write-state"); Object.assign(sidecar, state); }, spawnTunnel: () => { order.push("spawn-tunnel"); return tunnelFor(order); }, - fetchImpl: async (_input, init) => { - order.push(`readyz:${new Headers(init?.headers).get("x-opencodex-api-key") === KEY ? "key" : "missing"}`); - return new Response(null, { status: 200 }); - }, + scanListenPids: () => ({ ok: true, pids: [123] }), + fetchImpl: challengedFetch(order), connect: (async () => { order.push("connect"); }) as typeof import("../../src/client/connect").connectClient, scheduleRestart: () => { order.push("restart"); }, }, { alias: "home" })), @@ -271,7 +293,7 @@ describe("client initiated link join", () => { expect(response?.status).toBe(202); expect(responseBody).toEqual({ linkId: LINK_ID, alias: "home", restarting: true }); expect(sidecar).toMatchObject({ linkId: LINK_ID, tunnelPort: 23456, peerListenerPort: 45678 }); - expect(order).toEqual(["write-state", "spawn-tunnel", "readyz:key", "connect", "stop-tunnel", "restart"]); + expect(order).toEqual(["write-state", "spawn-tunnel", "readyz:probe", "readyz:key", "connect", "stop-tunnel", "restart"]); expect(isWrappedIssue(calls[0] ?? [])).toBe(true); expect(calls[0]?.at(-1)?.endsWith(" '--json'")).toBe(true); }); @@ -297,7 +319,7 @@ describe("client initiated link join", () => { choosePort: undefined, writeState: state => { sidecarPort = state.tunnelPort; }, spawnTunnel: () => tunnelFor([]), - fetchImpl: async () => new Response(null, { status: 200 }), + fetchImpl: challengedFetch(), connect: (async () => {}) as typeof import("../../src/client/connect").connectClient, scheduleRestart: () => {}, }), { alias: "home" }); @@ -330,7 +352,7 @@ describe("client initiated link join", () => { sleep: async () => {}, writeState: () => {}, clearState: () => { cleared += 1; }, - spawnTunnel: () => ({ pid: 1, exited: Promise.resolve(0), stop: async () => { stopped += 1; } }), + spawnTunnel: () => ({ pid: 1, exited: new Promise(() => {}), stop: async () => { stopped += 1; } }), fetchImpl: async () => readiness === "unauthorized" ? new Response(null, { status: 401 }) : new Response(null, { status: 503 }), }); await expect(joinHome(deps, { alias: "home" })).rejects.toMatchObject({ @@ -342,6 +364,230 @@ describe("client initiated link join", () => { } }); + test("never sends the issued key to a listener that skips the link-auth challenge", async () => { + const calls: string[][] = []; + let keyedFetches = 0; + let ticks = 0; + let stopped = 0; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + now: () => (ticks++ === 0 ? 0 : 15_002 * ticks), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => ({ pid: 1, exited: new Promise(() => {}), stop: async () => { stopped += 1; } }), + fetchImpl: async (_input, init) => { + if (new Headers(init?.headers).has("x-opencodex-api-key")) keyedFetches += 1; + return new Response(null, { status: 200 }); + }, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(keyedFetches).toBe(0); + expect(stopped).toBe(1); + expect(revokeCalls(calls)).toHaveLength(1); + }); + + test("a squatter answering the 401 challenge never receives the issued key", async () => { + const calls: string[][] = []; + let keyedFetches = 0; + let ticks = 0; + let stopped = 0; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + now: () => (ticks++ === 0 ? 0 : 15_002 * ticks), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => ({ pid: 123, exited: new Promise(() => {}), stop: async () => { stopped += 1; } }), + // A live SSH process alone does not prove ownership of the listener. + scanListenPids: () => ({ ok: true, pids: [999] }), + fetchImpl: async (_input, init) => { + if (new Headers(init?.headers).has("x-opencodex-api-key")) keyedFetches += 1; + return new Response(null, { status: 401 }); + }, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(keyedFetches).toBe(0); + expect(stopped).toBe(1); + expect(revokeCalls(calls)).toHaveLength(1); + }); + + test("the readiness scan is scoped to the tunnel's loopback address", async () => { + const seenAddresses: Array = []; + const order: string[] = []; + await joinHome(joinDeps({ + runner: runnerFor([]), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => tunnelFor(order), + scanListenPids: (_port, address) => { + seenAddresses.push(address); + return { ok: true, pids: [123] }; + }, + fetchImpl: challengedFetch(order), + connect: (async () => {}) as never, + scheduleRestart: () => {}, + }), { alias: "home" }); + expect(seenAddresses.length).toBeGreaterThan(0); + for (const address of seenAddresses) expect(address).toBe("127.0.0.1"); + }); + + test("a port flip between the probe and the keyed request never receives the key", async () => { + const calls: string[][] = []; + let keyedFetches = 0; + let ticks = 0; + let scans = 0; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + now: () => (ticks++ === 0 ? 0 : 15_002 * ticks), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => ({ pid: 123, exited: new Promise(() => {}), stop: async () => {} }), + scanListenPids: () => ({ ok: true, pids: scans++ === 0 ? [123] : [999] }), + fetchImpl: async (_input, init) => { + if (new Headers(init?.headers).has("x-opencodex-api-key")) keyedFetches += 1; + return new Response(null, { status: 401 }); + }, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(keyedFetches).toBe(0); + expect(revokeCalls(calls)).toHaveLength(1); + }); + + test("a redirect on the readiness probe is never followed with the issued key", async () => { + const calls: string[][] = []; + const redirects: Array = []; + let keyedFetches = 0; + let ticks = 0; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + now: () => (ticks++ === 0 ? 0 : 15_002 * ticks), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => ({ pid: 123, exited: new Promise(() => {}), stop: async () => {} }), + fetchImpl: async (_input, init) => { + redirects.push(init?.redirect); + if (new Headers(init?.headers).has("x-opencodex-api-key")) keyedFetches += 1; + return new Response(null, { status: 302, headers: { location: "http://169.254.1.1/fake-readyz" } }); + }, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(redirects).toEqual(["manual"]); + expect(keyedFetches).toBe(0); + expect(revokeCalls(calls)).toHaveLength(1); + }); + + test("a tunnel that exits during connect cannot commit the connection", async () => { + const calls: string[][] = []; + let releaseExit!: (code: number) => void; + const exited = new Promise(resolve => { releaseExit = resolve; }); + let stopped = 0; + let connectCommitted = false; + let connectDrained = false; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => ({ pid: 1, exited, stop: async () => { stopped += 1; } }), + fetchImpl: challengedFetch(), + connect: (async (_options, deps) => { + releaseExit(255); + try { + await new Promise(resolve => setTimeout(resolve, 10)); + deps?.signal?.throwIfAborted(); + connectCommitted = true; + } finally { connectDrained = true; } + }) as typeof import("../../src/client/connect").connectClient, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(connectDrained).toBe(true); + expect(connectCommitted).toBe(false); + expect(stopped).toBe(1); + expect(revokeCalls(calls)).toHaveLength(1); + }); + + test("does not disclose the issued key when the tunnel exits during its spawn grace", async () => { + const calls: string[][] = []; + let fetches = 0; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => ({ pid: 1, exited: Promise.resolve(255), stop: async () => {} }), + fetchImpl: async () => { + fetches += 1; + return new Response(null, { status: 200 }); + }, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(fetches).toBe(0); + expect(revokeCalls(calls)).toHaveLength(1); + }); + + for (const { name, recheck } of [ + { name: "unavailable", recheck: { ok: false, error: "scanner unavailable" } }, + { name: "empty", recheck: { ok: true, pids: [] } }, + { name: "foreign", recheck: { ok: true, pids: [999] } }, + { name: "ambiguous", recheck: { ok: true, pids: [123, 999] } }, + ] satisfies Array<{ name: string; recheck: ListenPidScan }>) { + test(`repeated ${name} rechecks reach the readiness deadline and revoke the key`, async () => { + const calls: string[][] = []; + const sleeps: number[] = []; + let clock = 1; + let scans = 0; + let keyedFetches = 0; + let stopped = 0; + let cleared = 0; + let connected = 0; + let restarted = 0; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + now: () => clock, + sleep: async ms => { sleeps.push(ms); clock += ms; }, + writeState: () => {}, + clearState: () => { cleared += 1; }, + spawnTunnel: () => ({ pid: 123, exited: new Promise(() => {}), stop: async () => { stopped += 1; } }), + scanListenPids: () => { + // Terminate the broken implementation without hanging the test runner. + // Its bypassed deadline produces the wrong error, so this is not a pass. + if (++scans > 400) throw new ClientLinkJoinError("admission_failed"); + return scans % 2 === 1 ? { ok: true, pids: [123] } : recheck; + }, + fetchImpl: async (_input, init) => { + if (new Headers(init?.headers).has("x-opencodex-api-key")) keyedFetches += 1; + return new Response(null, { status: 401 }); + }, + connect: (async () => { connected += 1; }) as typeof import("../../src/client/connect").connectClient, + scheduleRestart: () => { restarted += 1; }, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(clock).toBe(15_001); + expect(sleeps).toHaveLength(150); + expect(sleeps.every(ms => ms === 100)).toBe(true); + expect(scans).toBe(302); + expect(keyedFetches).toBe(0); + expect(connected).toBe(0); + expect(restarted).toBe(0); + expect(stopped).toBe(1); + expect(cleared).toBe(1); + expect(revokeCalls(calls)).toHaveLength(1); + }); + } + + test("a transient failed recheck polls before retrying and can still join", async () => { + const calls: string[][] = []; + const order: string[] = []; + let clock = 1; + let scans = 0; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + now: () => clock, + sleep: async ms => { clock += ms; order.push(`sleep:${ms}`); }, + writeState: () => {}, + spawnTunnel: () => tunnelFor(order), + scanListenPids: () => ++scans === 2 + ? { ok: false, error: "transient" } + : { ok: true, pids: [123] }, + fetchImpl: challengedFetch(order), + connect: (async () => { order.push("connect"); }) as typeof import("../../src/client/connect").connectClient, + scheduleRestart: () => { order.push("restart"); }, + }), { alias: "home" })).resolves.toEqual({ linkId: LINK_ID, apiKeyId: API_KEY_ID }); + expect(scans).toBe(4); + expect(order).toEqual(["readyz:probe", "sleep:100", "readyz:probe", "readyz:key", "connect", "stop-tunnel", "restart"]); + expect(revokeCalls(calls)).toHaveLength(0); + }); + test("rolls back on connect failure and never exposes the issued key", async () => { const calls: string[][] = []; const logs = spyOn(console, "log").mockImplementation(() => {}); @@ -351,7 +597,7 @@ describe("client initiated link join", () => { writeState: () => {}, clearState: () => {}, spawnTunnel: () => tunnelFor([]), - fetchImpl: async () => new Response(null, { status: 200 }), + fetchImpl: challengedFetch(), connect: (async () => { throw new Error(`connect failed ${KEY}`); }) as typeof import("../../src/client/connect").connectClient, }), { alias: "home" })).rejects.toMatchObject({ code: "join_connect_failed" }); } finally { @@ -392,8 +638,8 @@ describe("client initiated link join", () => { readSidecar: () => sidecarPresent ? sidecar : null, writeState: value => { sidecarPresent = true; Object.assign(sidecar, value); }, clearState: () => { sidecarPresent = false; }, - spawnTunnel: () => ({ pid: 1, exited: Promise.resolve(0), stop: async () => {} }), - fetchImpl: async () => new Response(null, { status: 200 }), + spawnTunnel: () => ({ pid: 1, exited: new Promise(() => {}), stop: async () => {} }), + fetchImpl: challengedFetch(), connect: (async () => { throw new Error("connect failed"); }) as typeof import("../../src/client/connect").connectClient, }); await expect(joinHome(base, { alias: "home" })).rejects.toMatchObject({ code: "join_rollback_failed", linkId: LINK_ID }); @@ -426,7 +672,8 @@ describe("client initiated link join", () => { writeState: state => { sidecar = { ...state }; }, clearState: () => { cleared = true; }, spawnTunnel: () => tunnelFor([]), - fetchImpl: async () => new Response(null, { status: 200 }), + scanListenPids: () => ({ ok: true, pids: [123] }), + fetchImpl: challengedFetch(), connect: (async () => { connected = true; }) as typeof import("../../src/client/connect").connectClient, scheduleRestart: () => { throw new Error("restart unavailable"); }, }, input)) as typeof import("../../src/client/link-join").joinHome, diff --git a/tests/server/port-reclaim.test.ts b/tests/server/port-reclaim.test.ts index b11f45e5af4..be63b0c449a 100644 --- a/tests/server/port-reclaim.test.ts +++ b/tests/server/port-reclaim.test.ts @@ -1,5 +1,13 @@ import { describe, expect, spyOn, test } from "bun:test"; +import { createServer } from "node:net"; +import * as childProcess from "node:child_process"; import { + listenAddressServes, + normalizeListenAddress, + parseListenEntriesFromLsof, + parseListenEntriesFromNetstat, + parseListenEntriesFromSs, + scanListenPidsForAddress, ownsIpv4LoopbackListener, parseIpv4LoopbackListenPidsFromNetstat, parseProcLoopbackListenInodes, @@ -126,6 +134,167 @@ describe("exact IPv4 loopback listener ownership", () => { }); }); +describe("listen-entry parsers keep the bound address", () => { + test("netstat entries report each listener's local address", () => { + const output = [ + "tcp 0 0 127.0.0.1:10100 0.0.0.0:* LISTEN 4242/bun", + "tcp 0 0 127.0.0.2:10100 0.0.0.0:* LISTEN 7777/foreign", + "tcp 0 0 127.0.0.1:22 0.0.0.0:* LISTEN 1/sshd", + ].join("\n"); + expect(parseListenEntriesFromNetstat(output, 10100)).toEqual([ + { pid: 4242, address: "127.0.0.1" }, + { pid: 7777, address: "127.0.0.2" }, + ]); + }); + + test("ss -Hltnp rows report address and pid; unattributed rows are dropped", () => { + const output = [ + "LISTEN 0 128 127.0.0.1:10100 0.0.0.0:* users:((\"bun\",pid=4242,fd=20))", + "LISTEN 0 128 127.0.0.2:10100 0.0.0.0:* users:((\"foreign\",pid=7777,fd=6))", + "LISTEN 0 128 127.0.0.1:10100 0.0.0.0:*", + "LISTEN 0 511 *:22 *:* users:((\"sshd\",pid=1,fd=3))", + ].join("\n"); + expect(parseListenEntriesFromSs(output, 10100)).toEqual([ + { pid: 4242, address: "127.0.0.1" }, + { pid: 7777, address: "127.0.0.2" }, + ]); + }); + + test("lsof NAME column supplies the bound address", () => { + const output = [ + "COMMAND PID USER FD TYPE DEVICE SIZE/OFF NODE NAME", + "bun 4242 devin 20u IPv4 0xdeadbeef 0t0 TCP 127.0.0.1:10100 (LISTEN)", + "other 7777 devin 21u IPv4 0xdeadbeef 0t0 TCP 127.0.0.2:10100 (LISTEN)", + ].join("\n"); + expect(parseListenEntriesFromLsof(output, 10100)).toEqual([ + { pid: 4242, address: "127.0.0.1" }, + { pid: 7777, address: "127.0.0.2" }, + ]); + }); + + test("address matching treats wildcards as serving any bound address", () => { + expect(listenAddressServes("127.0.0.1", "127.0.0.1")).toBe(true); + expect(listenAddressServes("127.0.0.2", "127.0.0.1")).toBe(false); + expect(listenAddressServes("0.0.0.0", "127.0.0.1")).toBe(true); + expect(listenAddressServes("*", "127.0.0.1")).toBe(true); + expect(listenAddressServes("::", "127.0.0.1")).toBe(true); + expect(listenAddressServes("[::1]:443", "::1")).toBe(true); + expect(normalizeListenAddress("::ffff:127.0.0.1")).toBe("127.0.0.1"); + expect(listenAddressServes("::ffff:127.0.0.1", "127.0.0.1")).toBe(true); + }); + + const multiAddressCases = [ + { + name: "Windows netstat", parse: parseListenEntriesFromNetstat, + rows: [ + "TCP 127.0.0.1:10100 0.0.0.0:0 LISTENING 4242", + "TCP 127.0.0.2:10100 0.0.0.0:0 LISTENING 4242", + "TCP [::ffff:127.0.0.1]:10100 [::]:0 LISTENING 4242", + ], + }, + { + name: "POSIX netstat", parse: parseListenEntriesFromNetstat, + rows: [ + "tcp 0 0 127.0.0.1:10100 0.0.0.0:* LISTEN 4242/ssh", + "tcp 0 0 127.0.0.2:10100 0.0.0.0:* LISTEN 4242/ssh", + "tcp 0 0 127.0.0.1:10100 0.0.0.0:* LISTEN 4242/ssh", + ], + }, + { + name: "ss", parse: parseListenEntriesFromSs, + rows: [ + 'LISTEN 0 128 127.0.0.1:10100 0.0.0.0:* users:(("ssh",pid=4242,fd=3))', + 'LISTEN 0 128 127.0.0.2:10100 0.0.0.0:* users:(("ssh",pid=4242,fd=4))', + 'LISTEN 0 128 [::ffff:127.0.0.1]:10100 [::]:* users:(("ssh",pid=4242,fd=5))', + ], + }, + { + name: "lsof", parse: parseListenEntriesFromLsof, + rows: [ + "ssh 4242 user 3u IPv4 0x1 0t0 TCP 127.0.0.1:10100 (LISTEN)", + "ssh 4242 user 4u IPv4 0x2 0t0 TCP 127.0.0.2:10100 (LISTEN)", + "ssh 4242 user 5u IPv6 0x3 0t0 TCP [::ffff:127.0.0.1]:10100 (LISTEN)", + ], + }, + ]; + for (const { name, parse, rows } of multiAddressCases) { + for (const reverse of [false, true]) { + test(`${name} retains all same-PID addresses with reverse=${reverse}`, () => { + const ordered = reverse ? [...rows].reverse() : rows; + const entries = parse([...ordered, ordered[0]].join("\n"), 10100); + expect([...entries].sort((a, b) => a.address.localeCompare(b.address))).toEqual([ + { pid: 4242, address: "127.0.0.1" }, + { pid: 4242, address: "127.0.0.2" }, + ]); + for (const address of ["127.0.0.1", "127.0.0.2"]) { + expect(entries.filter(entry => listenAddressServes(entry.address, address)).map(entry => entry.pid)).toEqual([4242]); + } + }); + } + } + + test("the PID-only netstat API still deduplicates multiple addresses", () => { + expect(parseListenPidsFromNetstat(multiAddressCases[0]!.rows.join("\n"), 10100)).toEqual([4242]); + }); + + test("the address-scoped scanner filters before deduplicating same-PID listeners", () => { + const fixture = process.platform === "win32" ? multiAddressCases[0]! : multiAddressCases[3]!; + for (const reverse of [false, true]) { + const rows = reverse ? [...fixture.rows].reverse() : fixture.rows; + const scan = spyOn(childProcess, "execFileSync").mockImplementation(() => rows.join("\n")); + try { + expect(scanListenPidsForAddress(10100)).toEqual({ ok: true, pids: [4242] }); + expect(scanListenPidsForAddress(10100, "127.0.0.1")).toEqual({ ok: true, pids: [4242] }); + expect(scanListenPidsForAddress(10100, "127.0.0.2")).toEqual({ ok: true, pids: [4242] }); + expect(scanListenPidsForAddress(10100, "127.0.0.3")).toEqual({ ok: true, pids: [] }); + expect(scanListenPidsForAddress(10100, "0.0.0.0")).toEqual({ ok: true, pids: [4242] }); + } finally { + scan.mockRestore(); + } + } + }); +}); + +describe("scanListenPidsForAddress (real scanner)", () => { + test("finds this process on its own bound port and filters other addresses", async () => { + const server = createServer(); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => resolve()); + }); + try { + const address = server.address(); + if (typeof address === "object" && address) { + const scan = scanListenPidsForAddress(address.port, "127.0.0.1"); + // Missing platform tools must report a failed scan rather than an empty result. + if (scan.ok) expect(scan.pids).toContain(process.pid); + } + } finally { + server.close(); + } + }); + + test("a listener on another loopback address does not serve 127.0.0.1", async () => { + const server = createServer(); + const bound = await new Promise(resolve => { + server.once("error", () => resolve(false)); + server.listen(0, "127.0.0.2", () => resolve(true)); + }); + if (!bound) return; + try { + const address = server.address(); + if (typeof address === "object" && address) { + const scan = scanListenPidsForAddress(address.port, "127.0.0.1"); + if (scan.ok) expect(scan.pids).not.toContain(process.pid); + const wide = scanListenPidsForAddress(address.port, "0.0.0.0"); + if (wide.ok) expect(wide.pids).toContain(process.pid); + } + } finally { + server.close(); + } + }); +}); + describe("parseTcpQuadsForLocalPort / IPv6", () => { test("collects every TCP row on the local port including non-LISTEN states", () => { const output = [