From 09cee3720b30ecba72ff8a54236ee1ec97e804f5 Mon Sep 17 00:00:00 2001 From: Epinephrine Date: Sat, 26 Sep 2026 09:34:18 +0000 Subject: [PATCH 1/6] fix(link): gate join credential on tunnel survival --- src/client/link-join.ts | 18 ++++++++++++++---- structure/runtime.md | 2 +- tests/server/link-join-route.test.ts | 23 ++++++++++++++++++++--- 3 files changed, 35 insertions(+), 8 deletions(-) diff --git a/src/client/link-join.ts b/src/client/link-join.ts index b782ef42259..9951a323380 100644 --- a/src/client/link-join.ts +++ b/src/client/link-join.ts @@ -22,6 +22,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 LINK_ID = /^lnk_[0-9a-f]{16}$/; @@ -203,6 +204,7 @@ async function compensateStaleSidecar(deps: ClientLinkJoinDeps): Promise { async function waitForReady( deps: ClientLinkJoinDeps, + tunnel: ClientLinkTunnelHandle, port: number, key: string, ): Promise { @@ -210,11 +212,19 @@ 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)), + ]); for (;;) { try { - const response = await fetchImpl(`http://127.0.0.1:${port}/readyz`, { - headers: { "x-opencodex-api-key": key }, - }); + const response = await Promise.race([ + tunnelExited, + 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"); } catch (error) { @@ -287,7 +297,7 @@ 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); diff --git a/structure/runtime.md b/structure/runtime.md index fd99beb93a0..da09d6a66f4 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -77,7 +77,7 @@ verified matching processes regardless of the advisory freshness result. ## Hub management dashboard address -When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. +When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. Client-initiated Remote Link enrollment in `src/client/link-join.ts` observes the SSH tunnel through a spawn grace and every credential-bearing readiness request. A tunnel that exits cannot deliver its issued data key to an unrelated loopback listener or commit the connection. ## Codex desktop process membership diff --git a/tests/server/link-join-route.test.ts b/tests/server/link-join-route.test.ts index ca57e13df0c..03d3dda894e 100644 --- a/tests/server/link-join-route.test.ts +++ b/tests/server/link-join-route.test.ts @@ -32,7 +32,7 @@ function runnerFor(calls: string[][], issueResult = true): SshRunner { function tunnelFor(order: string[]) { return { pid: 123, - exited: Promise.resolve(0), + exited: new Promise(() => {}), stop: async () => { order.push("stop-tunnel"); }, }; } @@ -201,7 +201,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({ @@ -213,6 +213,23 @@ describe("client initiated link join", () => { } }); + 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(calls.filter(argv => argv.some(value => value.includes("revoke")))).toHaveLength(1); + }); + test("rolls back on connect failure and never exposes the issued key", async () => { const calls: string[][] = []; const logs = spyOn(console, "log").mockImplementation(() => {}); @@ -263,7 +280,7 @@ 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 () => {} }), + spawnTunnel: () => ({ pid: 1, exited: new Promise(() => {}), stop: async () => {} }), fetchImpl: async () => new Response(null, { status: 200 }), connect: (async () => { throw new Error("connect failed"); }) as typeof import("../../src/client/connect").connectClient, }); From 9e2538c6034264c8bf9346cfd45ae52d3172dbed Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sat, 26 Sep 2026 12:40:17 +0000 Subject: [PATCH 2/6] fix(link): challenge-gate the join key and watch the tunnel through connect Co-Authored-By: Epinephrine --- src/client/link-join.ts | 55 +++++++++++++++--------- structure/runtime.md | 2 +- tests/server/link-join-route.test.ts | 63 ++++++++++++++++++++++++---- 3 files changed, 91 insertions(+), 29 deletions(-) diff --git a/src/client/link-join.ts b/src/client/link-join.ts index 9951a323380..eeb6c6835dc 100644 --- a/src/client/link-join.ts +++ b/src/client/link-join.ts @@ -219,14 +219,23 @@ async function waitForReady( ]); for (;;) { try { - const response = await Promise.race([ + // The tunnel can be alive without owning the port yet: an unrelated + // loopback listener must never receive the issued key, so the keyed + // request only goes to a listener answering the link-auth challenge. + const probe = await Promise.race([ tunnelExited, - fetchImpl(`http://127.0.0.1:${port}/readyz`, { - headers: { "x-opencodex-api-key": key }, - }), + fetchImpl(`http://127.0.0.1:${port}/readyz`), ]); - if (response.status === 200) return; - if (response.status === 401) throw new ClientLinkJoinError("admission_failed"); + if (probe.status === 401) { + const response = await Promise.race([ + tunnelExited, + 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"); + } } catch (error) { if (error instanceof ClientLinkJoinError) throw error; } @@ -305,22 +314,28 @@ export async function joinHome(deps: ClientLinkJoinDeps, input: { alias: string } 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, - }); - } catch { + // 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. + await Promise.race([ + tunnel.exited.then(() => { throw new ClientLinkJoinError("join_tunnel_failed"); }), + 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, + }), + ]); + } catch (error) { await rollback(deps, issued.linkId, tunnel); - throw new ClientLinkJoinError("join_connect_failed"); + throw new ClientLinkJoinError(error instanceof ClientLinkJoinError ? error.code : "join_connect_failed"); } await stopTunnel(tunnel); diff --git a/structure/runtime.md b/structure/runtime.md index da09d6a66f4..b40c30f7741 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -77,7 +77,7 @@ verified matching processes regardless of the advisory freshness result. ## Hub management dashboard address -When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. Client-initiated Remote Link enrollment in `src/client/link-join.ts` observes the SSH tunnel through a spawn grace and every credential-bearing readiness request. A tunnel that exits cannot deliver its issued data key to an unrelated loopback listener or commit the connection. +When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. Client-initiated Remote Link enrollment in `src/client/link-join.ts` watches the SSH tunnel from spawn grace through the connection commit and precedes every credential-bearing readiness request with an unauthenticated probe that must answer the link-auth challenge. An exited tunnel can neither deliver the issued data key to an unrelated loopback listener nor commit the connection. ## Codex desktop process membership diff --git a/tests/server/link-join-route.test.ts b/tests/server/link-join-route.test.ts index 03d3dda894e..3ff92f6c0ed 100644 --- a/tests/server/link-join-route.test.ts +++ b/tests/server/link-join-route.test.ts @@ -37,6 +37,16 @@ function tunnelFor(order: string[]) { }; } +// The link listener answers an unauthenticated /readyz with the 401 challenge; the keyed +// request earns 200. A foreign listener (any other status on the probe) must never see the 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 }); + }; +} + function joinDeps(overrides: Partial = {}): ClientLinkJoinDeps { const calls = overrides.runner ? [] : []; return { @@ -158,10 +168,7 @@ 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 }); - }, + fetchImpl: challengedFetch(order), connect: (async () => { order.push("connect"); }) as typeof import("../../src/client/connect").connectClient, scheduleRestart: () => { order.push("restart"); }, }, { alias: "home" })), @@ -171,7 +178,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(calls[0]?.some(value => value.includes("issue"))).toBe(true); expect(calls[0]?.some(value => value.includes("--json"))).toBe(true); }); @@ -213,6 +220,46 @@ 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(calls.filter(argv => argv.some(value => value.includes("revoke")))).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; + await expect(joinHome(joinDeps({ + runner: runnerFor(calls), + writeState: () => {}, + clearState: () => {}, + spawnTunnel: () => ({ pid: 1, exited, stop: async () => { stopped += 1; } }), + fetchImpl: challengedFetch(), + connect: (async () => { releaseExit(255); await new Promise(() => {}); connectCommitted = true; }) as typeof import("../../src/client/connect").connectClient, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(connectCommitted).toBe(false); + expect(stopped).toBe(1); + expect(calls.filter(argv => argv.some(value => value.includes("revoke")))).toHaveLength(1); + }); + test("does not disclose the issued key when the tunnel exits during its spawn grace", async () => { const calls: string[][] = []; let fetches = 0; @@ -239,7 +286,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 { @@ -281,7 +328,7 @@ describe("client initiated link join", () => { writeState: value => { sidecarPresent = true; Object.assign(sidecar, value); }, clearState: () => { sidecarPresent = false; }, spawnTunnel: () => ({ pid: 1, exited: new Promise(() => {}), stop: async () => {} }), - fetchImpl: async () => new Response(null, { status: 200 }), + 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 }); @@ -314,7 +361,7 @@ describe("client initiated link join", () => { writeState: state => { sidecar = { ...state }; }, clearState: () => { cleared = true; }, spawnTunnel: () => tunnelFor([]), - fetchImpl: async () => new Response(null, { status: 200 }), + 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, From 25bedeeee5ef8ef425f0daa39d616294150bb14e Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sat, 26 Sep 2026 13:53:00 +0000 Subject: [PATCH 3/6] fix(link-join): require tunnel-owned port and no redirects for keyed readyz A port squatter could answer the unkeyed /readyz probe with the 401 challenge and receive the following keyed request; readiness now only runs while the LISTEN owner of the tunnel port is the spawned ssh process (unverifiable scans stay not-ready), and both probes use redirect: manual so a redirecting occupant cannot reroute the challenge or the credential-bearing request. Co-Authored-By: Epinephrine --- src/client/link-join.ts | 39 ++++++++++++------- structure/runtime.md | 2 +- tests/server/link-join-route.test.ts | 56 ++++++++++++++++++++++++++++ 3 files changed, 82 insertions(+), 15 deletions(-) diff --git a/src/client/link-join.ts b/src/client/link-join.ts index eeb6c6835dc..b872194f196 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 { findAvailablePort } from "../server/ports"; +import { scanListenPids, type ListenPidScan } from "../server/port-reclaim"; import { isLinkPort } from "../link/ports"; import { buildExecArgv } from "../link/ssh-argv"; import type { SshRunner } from "../link/ssh-runner"; @@ -72,6 +73,8 @@ export interface ClientLinkJoinDeps { hostname?: () => string; randomBytes?: (size: number) => Uint8Array; fetchImpl?: typeof fetch; + /** LISTEN-owner probe for the tunnel port; defaults to the netstat/lsof scan. */ + scanListenPids?: (port: number) => ListenPidScan; spawnTunnel?: (spec: { linkId: string; alias: string; @@ -217,24 +220,32 @@ async function waitForReady( tunnelExited, new Promise(resolve => setTimeout(resolve, JOIN_TUNNEL_SPAWN_GRACE_MS)), ]); + const listenPids = deps.scanListenPids ?? scanListenPids; for (;;) { try { - // The tunnel can be alive without owning the port yet: an unrelated - // loopback listener must never receive the issued key, so the keyed - // request only goes to a listener answering the link-auth challenge. - const probe = await Promise.race([ - tunnelExited, - fetchImpl(`http://127.0.0.1:${port}/readyz`), - ]); - if (probe.status === 401) { - const response = await Promise.race([ + // 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); + 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`, { - headers: { "x-opencodex-api-key": key }, - }), + fetchImpl(`http://127.0.0.1:${port}/readyz`, { redirect: "manual" }), ]); - if (response.status === 200) return; - if (response.status === 401) throw new ClientLinkJoinError("admission_failed"); + if (probe.status === 401) { + 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; diff --git a/structure/runtime.md b/structure/runtime.md index b40c30f7741..1d3ee041555 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -77,7 +77,7 @@ verified matching processes regardless of the advisory freshness result. ## Hub management dashboard address -When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. Client-initiated Remote Link enrollment in `src/client/link-join.ts` watches the SSH tunnel from spawn grace through the connection commit and precedes every credential-bearing readiness request with an unauthenticated probe that must answer the link-auth challenge. An exited tunnel can neither deliver the issued data key to an unrelated loopback listener nor commit the connection. +When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. Client-initiated Remote Link enrollment in `src/client/link-join.ts` watches the SSH tunnel from spawn grace through the connection commit. Readiness is accepted only while the LISTEN owner of the tunnel port is the spawned ssh process — a live tunnel does not prove it owns the socket, and a foreign listener answering the link-auth challenge would otherwise collect the issued key — and both the probe and the keyed request run with `redirect: "manual"` so a redirecting occupant cannot reroute the challenge. An exited tunnel can neither deliver the issued data key to an unrelated loopback listener nor commit the connection. ## Codex desktop process membership diff --git a/tests/server/link-join-route.test.ts b/tests/server/link-join-route.test.ts index 3ff92f6c0ed..f9f09aa4346 100644 --- a/tests/server/link-join-route.test.ts +++ b/tests/server/link-join-route.test.ts @@ -1,5 +1,6 @@ import { describe, expect, test, spyOn } from "bun:test"; import { ClientLinkJoinError, joinHome, type ClientLinkJoinDeps } from "../../src/client/link-join"; +import { spawnClientLinkTunnel } from "../../src/client/link-tunnel"; import { handleLinkRoutes, type LinkRouteState } from "../../src/server/management/link-routes"; import type { ManagementContext } from "../../src/server/management/context"; import type { SshRunner } from "../../src/link/ssh-runner"; @@ -49,6 +50,10 @@ function challengedFetch(order?: string[]) { function joinDeps(overrides: Partial = {}): ClientLinkJoinDeps { const calls = overrides.runner ? [] : []; + // The readiness gate only trusts the port when the LISTEN pid is the spawned tunnel's; + // wrap whichever spawnTunnel is under test so the default scan reports that pid. + let tunnelPid = 0; + const spawn = overrides.spawnTunnel ?? spawnClientLinkTunnel; return { runner: overrides.runner ?? runnerFor(calls), knownHostsFile: "/tmp/ocx-known-hosts", @@ -60,6 +65,12 @@ function joinDeps(overrides: Partial = {}): ClientLinkJoinDe readSidecar: () => null, readConnectionState: () => ({ kind: "disconnected" }), ...overrides, + spawnTunnel: (spec, spawnDeps) => { + const handle = spawn(spec, spawnDeps); + tunnelPid = handle.pid; + return handle; + }, + scanListenPids: overrides.scanListenPids ?? (() => ({ ok: true, pids: [tunnelPid] })), }; } @@ -168,6 +179,7 @@ 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); }, + scanListenPids: () => ({ ok: true, pids: [123] }), fetchImpl: challengedFetch(order), connect: (async () => { order.push("connect"); }) as typeof import("../../src/client/connect").connectClient, scheduleRestart: () => { order.push("restart"); }, @@ -241,6 +253,49 @@ describe("client initiated link join", () => { expect(calls.filter(argv => argv.some(value => value.includes("revoke")))).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 foreign process holds the port: a live ssh does not prove it owns the socket. + 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(calls.filter(argv => argv.some(value => value.includes("revoke")))).toHaveLength(1); + }); + + test("a redirect on the readiness probe is never followed with the issued key", async () => { + const calls: string[][] = []; + 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) => { + if (new Headers(init?.headers).has("x-opencodex-api-key")) keyedFetches += 1; + // A port occupant redirecting the probe used to let the challenge pass at a foreign URL. + return new Response(null, { status: 302, headers: { location: "http://169.254.1.1/fake-readyz" } }); + }, + }), { alias: "home" })).rejects.toMatchObject({ code: "join_tunnel_failed" }); + expect(keyedFetches).toBe(0); + expect(calls.filter(argv => argv.some(value => value.includes("revoke")))).toHaveLength(1); + }); + test("a tunnel that exits during connect cannot commit the connection", async () => { const calls: string[][] = []; let releaseExit!: (code: number) => void; @@ -361,6 +416,7 @@ describe("client initiated link join", () => { writeState: state => { sidecar = { ...state }; }, clearState: () => { cleared = true; }, spawnTunnel: () => tunnelFor([]), + scanListenPids: () => ({ ok: true, pids: [123] }), fetchImpl: challengedFetch(), connect: (async () => { connected = true; }) as typeof import("../../src/client/connect").connectClient, scheduleRestart: () => { throw new Error("restart unavailable"); }, From c5862e89206883a1744ac555ec43021c607036cc Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sat, 26 Sep 2026 15:08:54 +0000 Subject: [PATCH 4/6] fix(client): scope tunnel-port ownership to the bound address and add ss fallback MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three join-readiness hardening fixes: - scanListenEntries now keeps each listener's bound address and the readiness check only counts sockets that serve the tunnel's 127.0.0.1 bind — a listener on 127.0.0.2 or another interface no longer stalls enrollment until the issued link is revoked. - The POSIX scanner chain gains ss -Hltnp between lsof and netstat, so minimal Linux installs with only iproute2 can still verify ownership instead of failing every probe as unavailable. - Ownership is re-verified in the same iteration immediately before the keyed request, narrowing the scan-to-request takeover window that could have delivered the issued key to a port flipper. Co-Authored-By: Epinephrine --- src/client/link-join.ts | 24 +++- src/server/port-reclaim.ts | 202 ++++++++++++++++++++++----- structure/runtime.md | 2 +- tests/server/link-join-route.test.ts | 44 ++++++ tests/server/port-reclaim.test.ts | 105 +++++++++++++- 5 files changed, 338 insertions(+), 39 deletions(-) diff --git a/src/client/link-join.ts b/src/client/link-join.ts index b872194f196..04b15f86fa2 100644 --- a/src/client/link-join.ts +++ b/src/client/link-join.ts @@ -1,7 +1,7 @@ import { randomBytes } from "node:crypto"; import { hostname } from "node:os"; import { findAvailablePort } from "../server/ports"; -import { scanListenPids, type ListenPidScan } from "../server/port-reclaim"; +import { scanListenPidsForAddress, type ListenPidScan } from "../server/port-reclaim"; import { isLinkPort } from "../link/ports"; import { buildExecArgv } from "../link/ssh-argv"; import type { SshRunner } from "../link/ssh-runner"; @@ -73,8 +73,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 scan. */ - scanListenPids?: (port: number) => ListenPidScan; + /** + * 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; @@ -220,14 +224,17 @@ async function waitForReady( tunnelExited, new Promise(resolve => setTimeout(resolve, JOIN_TUNNEL_SPAWN_GRACE_MS)), ]); - const listenPids = deps.scanListenPids ?? scanListenPids; + 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 { // 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); + 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. @@ -236,6 +243,13 @@ async function waitForReady( 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 the keyed request if the port is no longer solely the tunnel's. + const recheck = listenPids(port, tunnelAddress); + if (!recheck.ok || recheck.pids.length !== 1 || recheck.pids[0] !== tunnel.pid) { + continue; + } const response = await Promise.race([ tunnelExited, fetchImpl(`http://127.0.0.1:${port}/readyz`, { diff --git a/src/server/port-reclaim.ts b/src/server/port-reclaim.ts index 0bd0660ff7b..b10300486d4 100644 --- a/src/server/port-reclaim.ts +++ b/src/server/port-reclaim.ts @@ -16,6 +16,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 @@ -58,12 +68,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 listener's bound local address. 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(); @@ -86,9 +133,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) { + entries.set(pid, { pid, address: normalizeListenAddress(parts[localIdx]) }); + } + } + return [...entries.values()]; +} + +/** + * Parse `netstat -ano` (Windows) / `netstat -anlp` listen lines for a port. + * Exported for unit tests. + */ +export function parseListenPidsFromNetstat(output: string, port: number): number[] { + return parseListenEntriesFromNetstat(output, port).map(entry => entry.pid); +} + +/** + * Parse `ss -Hltnp` rows for a port, keeping the bound local address. 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) { + entries.set(pid, { pid, address: normalizeListenAddress(parts[localIdx]) }); + } } - return [...pids]; + return [...entries.values()]; +} + +/** + * Parse `lsof -nP -iTCP: -sTCP:LISTEN` output (without -t). The NAME column is + * the last address token, optionally followed by the `(LISTEN)` state; skip the + * header and any line whose pid is not numeric. Exported for unit 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; + entries.set(pid, { pid, address: normalizeListenAddress(address) }); + } + return [...entries.values()]; } function normalizeListenPidScan(result: ListenPidScan | number[]): ListenPidScan { @@ -119,50 +223,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: string): 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/runtime.md b/structure/runtime.md index 1d3ee041555..f8d4ea8358d 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -77,7 +77,7 @@ verified matching processes regardless of the advisory freshness result. ## Hub management dashboard address -When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. Client-initiated Remote Link enrollment in `src/client/link-join.ts` watches the SSH tunnel from spawn grace through the connection commit. Readiness is accepted only while the LISTEN owner of the tunnel port is the spawned ssh process — a live tunnel does not prove it owns the socket, and a foreign listener answering the link-auth challenge would otherwise collect the issued key — and both the probe and the keyed request run with `redirect: "manual"` so a redirecting occupant cannot reroute the challenge. An exited tunnel can neither deliver the issued data key to an unrelated loopback listener nor commit the connection. +When hub management ingress is enabled, `src/cli/dispatch.ts` opens the dashboard on the literal IPv4 loopback address and configured ingress port, matching the listener in `src/server/index.ts`. Other dashboard address selection is unchanged. Client-initiated Remote Link enrollment in `src/client/link-join.ts` watches the SSH tunnel from spawn grace through the connection commit. Readiness is accepted only while the LISTEN owner of the tunnel port is the spawned ssh process — a live tunnel does not prove it owns the socket, and a foreign listener answering the link-auth challenge would otherwise collect the issued key — and ownership is judged only among listeners that serve the tunnel's 127.0.0.1 bind, so an occupant on another loopback or interface address can neither satisfy the check nor block it. The socket scan (`src/server/port-reclaim.ts`) runs `lsof`, then `ss`, then `netstat`, so minimal Linux installs with only iproute2 still enumerate listeners. The ownership check repeats immediately before the keyed request, narrowing the takeover window between the 401 probe and the request that carries the issued key, and both requests run with `redirect: "manual"` so a redirecting occupant cannot reroute the challenge. An exited tunnel can neither deliver the issued data key to an unrelated loopback listener nor commit the connection. ## Codex desktop process membership diff --git a/tests/server/link-join-route.test.ts b/tests/server/link-join-route.test.ts index f9f09aa4346..f6c9859d667 100644 --- a/tests/server/link-join-route.test.ts +++ b/tests/server/link-join-route.test.ts @@ -276,6 +276,50 @@ describe("client initiated link join", () => { expect(calls.filter(argv => argv.some(value => value.includes("revoke")))).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), + // Only sockets serving 127.0.0.1 count: the real scanner reports listeners on + // other loopback/interface addresses too, and the dep must scope them out. + 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 () => {} }), + // First scan names the tunnel; by the time the 401 arrives a squatter holds the port. + 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(calls.filter(argv => argv.some(value => value.includes("revoke")))).toHaveLength(1); + }); + test("a redirect on the readiness probe is never followed with the issued key", async () => { const calls: string[][] = []; let keyedFetches = 0; diff --git a/tests/server/port-reclaim.test.ts b/tests/server/port-reclaim.test.ts index 930a7866c9c..5c6dfbd041b 100644 --- a/tests/server/port-reclaim.test.ts +++ b/tests/server/port-reclaim.test.ts @@ -1,5 +1,15 @@ import { describe, expect, spyOn, test } from "bun:test"; -import { reclaimListenPort, type ReclaimListenPortOptions } from "../../src/server/port-reclaim"; +import { createServer } from "node:net"; +import { + listenAddressServes, + normalizeListenAddress, + parseListenEntriesFromLsof, + parseListenEntriesFromNetstat, + parseListenEntriesFromSs, + reclaimListenPort, + scanListenPidsForAddress, + type ReclaimListenPortOptions, +} from "../../src/server/port-reclaim"; import { isBareIpv6Address, parseTcpQuadsForLocalPort, @@ -41,6 +51,99 @@ describe("parseListenPidsFromNetstat", () => { }); }); +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); + }); +}); + +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"); + // The default scanner is whatever the platform ships (netstat/lsof/ss); when none + // is installed the scan reports a probe failure rather than an empty list. + 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; // platform does not allow the second loopback address + 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 = [ From b38ea5a99e0989f2cc092dbccdc2ae9be2731818 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Sun, 27 Sep 2026 06:28:20 +0000 Subject: [PATCH 5/6] fix(link): cancel and drain enrollment before tunnel rollback --- src/client/connect.ts | 34 ++++++++++++++++------- src/client/link-join.ts | 26 +++++++++++++---- structure/remote-link.md | 4 +++ tests/clients/client-link-connect.test.ts | 28 +++++++++++++++++++ tests/server/link-join-route.test.ts | 11 +++++++- 5 files changed, 87 insertions(+), 16 deletions(-) diff --git a/src/client/connect.ts b/src/client/connect.ts index 67ea8d0968b..c50e47ebb1e 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 ? async (input, init = {}) => { + 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" }); + } : 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 ad2fe285de9..0527fc21114 100644 --- a/src/client/link-join.ts +++ b/src/client/link-join.ts @@ -360,14 +360,22 @@ export async function joinHome(deps: ClientLinkJoinDeps, input: { alias: string 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; // 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. - await Promise.race([ - tunnel.exited.then(() => { throw new ClientLinkJoinError("join_tunnel_failed"); }), - connect({ + const tunnelFailure = tunnel.exited.then(() => { + const error = new ClientLinkJoinError("join_tunnel_failed"); + if (!enrollmentFinished) enrollmentAbort.abort(error); + throw error; + }); + 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 }, @@ -378,11 +386,19 @@ export async function joinHome(deps: ClientLinkJoinDeps, input: { alias: string }, { 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(error instanceof ClientLinkJoinError ? error.code : "join_connect_failed"); + } finally { + enrollmentFinished = true; } await stopTunnel(tunnel); diff --git a/structure/remote-link.md b/structure/remote-link.md index fc9d1c9897e..61b8ec4b693 100644 --- a/structure/remote-link.md +++ b/structure/remote-link.md @@ -61,3 +61,7 @@ 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/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 9ea2b8c26c1..8c722c3226b 100644 --- a/tests/server/link-join-route.test.ts +++ b/tests/server/link-join-route.test.ts @@ -477,14 +477,23 @@ describe("client initiated link join", () => { 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 () => { releaseExit(255); await new Promise(() => {}); connectCommitted = true; }) as typeof import("../../src/client/connect").connectClient, + 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); From 5672d3bc4891f418d747854faa5d761c505afd0d Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Sun, 27 Sep 2026 06:33:49 +0000 Subject: [PATCH 6/6] fix(link): preserve Bun fetch interface on cancellation wrapper --- src/client/connect.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/client/connect.ts b/src/client/connect.ts index c50e47ebb1e..7aa3e65e005 100644 --- a/src/client/connect.ts +++ b/src/client/connect.ts @@ -543,12 +543,12 @@ export async function connectClient( ): Promise { deps.signal?.throwIfAborted(); const rawFetch = deps.fetchImpl ?? fetch; - const fetchImpl: typeof fetch = deps.signal ? async (input, init = {}) => { + 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" }); - } : rawFetch; + }, { preconnect: rawFetch.preconnect }) : rawFetch; const assertActiveConnectingState = (fingerprint?: string) => { deps.signal?.throwIfAborted(); assertConnectingState(fingerprint);