From f32d9faabc190e3e906245040be582725eb1d7b4 Mon Sep 17 00:00:00 2001 From: JUN Date: Sun, 27 Sep 2026 00:31:31 +0900 Subject: [PATCH] fix(link): the Child tunnel reconnects on its own after sleep, outages and crashes Root cause: the Child's own ssh -L tunnel treated 'failed' as terminal. Five minutes down (sleep, a Home outage) or one auth/forward stderr left it failed forever. A respawn from reconnecting was never promoted (only 'connecting' had the 5 s alive grace), so at since+5 min the tick killed the healthy ssh and cut in-flight SSE. On macOS every leftover pidfile was trusted: 'unresolved' and 'owned' set connected with no child, and no ssh was ever spawned again. One unreadable config read ended the link and recycled the Child to standalone. Requests during a reconnect or right after the restart failed at once, and a join took an OS ephemeral tunnel port that another socket can hold after a reboot. Fix: - src/link/tunnel-state.ts: an opt-in retry policy. With CLIENT_TUNNEL_RETRY_POLICY failed carries retryAt: timeout and forward retry after 60 s, auth after 5 min (so at most 12 auth attempts an hour), a changed host key never. The retry runs while the state reads failed (inFlight), and an exit keeps the one-a-minute cadence. Without a policy (the Home's -R supervisor) nothing changes; a guard test pins it. - src/client/link-tunnel.ts: promotion only on a keyed GET /readyz through the tunnel that proves the link (connecting, reconnecting or a retry from failed): a 200, or a 503 whose body (read up to 4 KiB) carries service: "opencodex". The Home's link listener answers 401 before /readyz, so that 503 only means the Home's own start-up readiness is pending or failed, which relayed requests do not depend on; it shows as probe 'home_not_ready' instead of holding every request and killing a working ssh at the 5-minute mark. A display-only 30 s probe while connected (401/403 -> 'unauthorized', readiness 503 -> 'home_not_ready', else 'home_unreachable', shown as the child reason). One probe runs at a time, detached from the 1 s check; stop() and the end of the probed tunnel abort it, so a hanging probe never delays a disconnect, a sidecar removal or shutdown. While a request is held the probe runs on every check instead of backing off. macOS orphans are verified with ps: a dead or reused pid is a stale pidfile, an exact-argv child of launchd is reaped like Linux, anything else is adopted, probed, never signalled, and replaced once it dies. A leftover pidfile is settled before the first spawn; an unusable start-up read defers that to the first matching check, so nothing spawns over an unreaped orphan. An unreadable sidecar or connection state acts only after 3 checks in a row; an explicit disconnect or another link id still ends the link on the next check, once. - src/client/link-relay.ts: a bounded hold. Only while the tunnel is connecting or reconnecting a request waits for connected, at most 15 s from its first wait and 64 at once, then is forwarded once. A refused connection (nothing sent, body untouched) may be resent inside the same window; a reset or any failure after the body started is never replayed. - src/client/runtime.ts and machine-listener.ts: the supervisor is the relay's tunnel gate, and one cached key source serves the relay and the probe. - src/client/link-join.ts: a join picks its tunnel port at random from 20000-29999, below the OS ephemeral ranges. Persisted ports never change. Performance: standalone, hub and Home paths are untouched. A connected Child adds one pending() call per relayed request and one keyed probe every 30 s; holds and their timers exist only while the tunnel is actually down. The supervisor's 1 s check is now an unref'd interval that stats the sidecar and config.json and parses them only after a change (it parsed config.json twice a second before), and it never awaits the network. Probes back off 1-5 s while connecting and up to 30 s once the link reads failed; only while a request is held do they run once a second. Only a 503 body is read, capped at 4 KiB. The key is read once; no request or probe touches the disk. Security: no new admission surface. Probes send the link key as a header to the same 127.0.0.1 tunnel the data requests use, with no-store, and never log it. A readiness 503 is accepted only with the Home's service marker, which proves no more than the 200 already did. Host-key trust is unchanged (StrictHostKeyChecking=yes; a retry only re-checks and never trusts a new key); host-key failures are not retried. Auth retries are capped by their 5-minute spacing. On macOS a process is signalled only on an exact argv match with launchd as parent; nothing unverified is ever signalled. A 401 never disconnects the link. Decision K19 in the remote link devlog records the change to D13 for Child-owned tunnels. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../003_decisions.md | 1 + .../src/content/docs/fr/guides/remote-link.md | 4 +- .../src/content/docs/guides/remote-link.md | 4 +- .../src/content/docs/ja/guides/remote-link.md | 4 +- .../src/content/docs/ko/guides/remote-link.md | 4 +- .../src/content/docs/ru/guides/remote-link.md | 4 +- .../src/content/docs/tr/guides/remote-link.md | 4 +- .../content/docs/zh-cn/guides/remote-link.md | 4 +- .../content/docs/zh-tw/guides/remote-link.md | 4 +- src/client/link-join.ts | 25 +- src/client/link-relay.ts | 91 ++- src/client/link-status.ts | 5 +- src/client/link-tunnel.ts | 513 +++++++++++--- src/client/machine-listener.ts | 9 +- src/client/runtime.ts | 12 +- src/link/ports.ts | 9 + src/link/tunnel-state.ts | 67 +- structure/remote-link.md | 10 +- tests/clients/client-link-relay.test.ts | 200 ++++++ tests/clients/client-link-status.test.ts | 12 + tests/clients/client-link-tunnel.test.ts | 623 +++++++++++++++++- tests/clients/link-supervisor.test.ts | 33 + tests/clients/link-tunnel-state.test.ts | 68 ++ tests/server/link-join-route.test.ts | 32 +- 24 files changed, 1604 insertions(+), 138 deletions(-) diff --git a/devlog/_plan/260925_remote_home_child_link/003_decisions.md b/devlog/_plan/260925_remote_home_child_link/003_decisions.md index 80a96b212d1..8507a8c0972 100644 --- a/devlog/_plan/260925_remote_home_child_link/003_decisions.md +++ b/devlog/_plan/260925_remote_home_child_link/003_decisions.md @@ -22,3 +22,4 @@ decade 문서(020-060) 작성 중 올라온 계약 편차와 열린 질문에 | K16 | status/candidates 응답 모양 (050) | `GET /api/link/status` → `{role: "standalone"|"home"|"child", listener: {state: "off"|"listening"|"failed", port: number|null}, links: [{id, alias, direction, state: "connecting"|"connected"|"reconnecting"|"failed"|"idle", since: string, reason: string|null, tunnelPort: number}], child: null | {alias, state, since, reason}}`. `GET /api/link/candidates` → `{candidates: [{alias, source: "ssh_config"|"tailscale"}]}`. `POST /api/link/probe` → `{alias, fingerprint, keyType}`. `POST /api/link/confirm-host` → `{alias, fingerprint, ocxVersion}`. `POST /api/link/apply` → 202 `{linkId}` 후 status 폴링. 오류는 `{error: {code, message}}` | wp4, wp5, wp6 | | K17 | GUI locale 수 (050) | gui/src/i18n의 실제 locale 파일 전부(050 확인: 10개). docs-site는 영어 원문 + 기존 번역 locale에 같은 페이지를 추가하되 번역이 없으면 영어 원문 링크를 두지 않고 해당 locale 생략 | wp5 | | K18 | link 데이터 키 메모리 제로화 (wp3 감사 Leibniz) | 요구하지 않는다. JS 문자열은 지울 수 없고 키는 0600 토큰 파일에 저장되는 장기 비밀이며 기존 hub connect와 같은 처리다. 대신 stdin 4 KiB 상한, 원본 입력 버퍼 0 채움, 키가 로그·오류·config·journal·status에 나타나지 않음을 테스트로 고정 | wp3, wp6 | +| K19 | D13 "failed는 조치 필요"를 Child가 소유한 터널에도 적용할지 (PR-D, 자동 복구 요구) | Child가 소유한 `-L` 터널에만 재시도 정책을 둔다. timeout·forward는 약 60초마다, auth는 5분마다(시간당 최대 12회) 다시 시도하고, 바뀐 호스트 키는 계속 failed로 남는다. 재시도 중에도 상태는 failed로 보인다. 재연결은 링크 키를 담은 `/readyz`가 링크를 증명할 때만 connected로 올린다. 200, 또는 `service: "opencodex"` 본문을 가진 503이다(Home의 link listener는 `/readyz` 앞에서 401을 주므로 이 503은 Home 자체 startup readiness가 pending/failed라는 뜻일 뿐이고 표시용 `home_not_ready`로 남긴다). Home의 `-R` supervisor는 바꾸지 않는다(PR-F 몫) | wp6 | diff --git a/docs-site/src/content/docs/fr/guides/remote-link.md b/docs-site/src/content/docs/fr/guides/remote-link.md index c760fe5105b..6a619f8d5dd 100644 --- a/docs-site/src/content/docs/fr/guides/remote-link.md +++ b/docs-site/src/content/docs/fr/guides/remote-link.md @@ -41,8 +41,8 @@ Le rôle **Child** n’est disponible que lorsque OpenCodex tourne sur son port ## État de la liaison - **Connected** signifie que le tunnel SSH est prêt et que Child peut utiliser la liaison Home. -- **Reconnecting** signifie que le tunnel est réessayé. Les requêtes peuvent temporairement renvoyer `503` avec `Retry-After`. -- **Failed** signifie que la liaison nécessite une intervention. Vérifiez l’authentification SSH, la clé d’hôte confirmée, la redirection ou le délai indiqué. +- **Reconnecting** signifie que le tunnel est réessayé. Les requêtes peuvent temporairement renvoyer `503` avec `Retry-After`. Sur un Child connecté depuis son propre tableau de bord, une requête attend d’abord jusqu’à 15 secondes le retour du tunnel. +- **Failed** signifie que la liaison nécessite une intervention. Vérifiez l’authentification SSH, la clé d’hôte confirmée, la redirection ou le délai indiqué. Un Child connecté depuis son propre tableau de bord réessaie de lui-même après une mise en veille, une panne ou un redémarrage : environ une fois par minute après un délai dépassé ou une erreur de redirection, et toutes les cinq minutes après une erreur d’authentification. Une clé d’hôte modifiée n’est jamais réessayée. Une liaison en échec ne bascule pas silencieusement vers un fournisseur local. diff --git a/docs-site/src/content/docs/guides/remote-link.md b/docs-site/src/content/docs/guides/remote-link.md index 7031673bf32..2b8b911b78d 100644 --- a/docs-site/src/content/docs/guides/remote-link.md +++ b/docs-site/src/content/docs/guides/remote-link.md @@ -41,8 +41,8 @@ The **Child** role is available only while OpenCodex runs on its configured port ## Link status - **Connected** means the SSH tunnel is ready and the Child can use the Home link. -- **Reconnecting** means the tunnel is being retried. Requests can temporarily return `503` with `Retry-After` while the retry is in progress. -- **Failed** means the link needs attention. Check SSH authentication, the confirmed host key, forwarding, or the timeout reason shown in the dashboard. +- **Reconnecting** means the tunnel is being retried. Requests can temporarily return `503` with `Retry-After` while the retry is in progress. On a Child that connected from its own dashboard, a request first waits up to 15 seconds for the tunnel to come back. +- **Failed** means the link needs attention. Check SSH authentication, the confirmed host key, forwarding, or the timeout reason shown in the dashboard. A Child that connected from its own dashboard keeps retrying by itself, after sleep, an outage or a restart: about once a minute after a timeout or forwarding error, and every five minutes after an authentication error. A changed host key is never retried. A failed link does not silently switch to a local provider. diff --git a/docs-site/src/content/docs/ja/guides/remote-link.md b/docs-site/src/content/docs/ja/guides/remote-link.md index 84f36392202..d8878f2fbd8 100644 --- a/docs-site/src/content/docs/ja/guides/remote-link.md +++ b/docs-site/src/content/docs/ja/guides/remote-link.md @@ -41,8 +41,8 @@ Home のプロバイダーを使うコンピューターで次の操作を行い ## リンクの状態 - **Connected** は SSH トンネルが準備でき、Child が Home のリンクを使える状態です。 -- **Reconnecting** はトンネルを再試行している状態です。再試行中はリクエストが `Retry-After` 付きの `503` を一時的に返すことがあります。 -- **Failed** は対応が必要な状態です。SSH 認証、確認済みのホストキー、転送、タイムアウトの理由を確認してください。 +- **Reconnecting** はトンネルを再試行している状態です。再試行中はリクエストが `Retry-After` 付きの `503` を一時的に返すことがあります。自分のダッシュボードから接続した Child では、リクエストはまずトンネルの復帰を最大 15 秒待ちます。 +- **Failed** は対応が必要な状態です。SSH 認証、確認済みのホストキー、転送、タイムアウトの理由を確認してください。自分のダッシュボードから接続した Child は、スリープ、障害、再起動の後も自動で再試行します。タイムアウトや転送エラーの後は約 1 分ごと、認証エラーの後は 5 分ごとです。ホストキーが変わった場合は再試行しません。 リンクが失敗しても、ローカルプロバイダーへ自動的に切り替わることはありません。 diff --git a/docs-site/src/content/docs/ko/guides/remote-link.md b/docs-site/src/content/docs/ko/guides/remote-link.md index 9cad2362ed5..6ac312feeb9 100644 --- a/docs-site/src/content/docs/ko/guides/remote-link.md +++ b/docs-site/src/content/docs/ko/guides/remote-link.md @@ -41,8 +41,8 @@ Home의 프로바이더를 사용할 컴퓨터에서 다음을 진행합니다. ## 링크 상태 - **Connected**는 SSH 터널이 준비되어 Child가 Home 링크를 사용할 수 있다는 뜻입니다. -- **Reconnecting**은 터널을 다시 연결하는 중이라는 뜻입니다. 재시도 중에는 요청이 일시적으로 `Retry-After`와 함께 `503`을 반환할 수 있습니다. -- **Failed**는 조치가 필요하다는 뜻입니다. SSH 인증, 확인한 호스트 키, 포워딩 또는 타임아웃 사유를 확인하세요. +- **Reconnecting**은 터널을 다시 연결하는 중이라는 뜻입니다. 재시도 중에는 요청이 일시적으로 `Retry-After`와 함께 `503`을 반환할 수 있습니다. 자기 대시보드에서 연결한 Child에서는 요청이 먼저 터널이 돌아오기를 최대 15초 기다립니다. +- **Failed**는 조치가 필요하다는 뜻입니다. SSH 인증, 확인한 호스트 키, 포워딩 또는 타임아웃 사유를 확인하세요. 자기 대시보드에서 연결한 Child는 절전, 장애, 재시작 뒤에도 스스로 다시 시도합니다. 타임아웃이나 포워딩 오류 뒤에는 약 1분마다, 인증 오류 뒤에는 5분마다 시도합니다. 호스트 키가 바뀐 경우에는 다시 시도하지 않습니다. 링크가 실패해도 로컬 프로바이더로 조용히 전환하지 않습니다. diff --git a/docs-site/src/content/docs/ru/guides/remote-link.md b/docs-site/src/content/docs/ru/guides/remote-link.md index fb00add6694..c53f8b63a5f 100644 --- a/docs-site/src/content/docs/ru/guides/remote-link.md +++ b/docs-site/src/content/docs/ru/guides/remote-link.md @@ -41,8 +41,8 @@ SSH с паролем и Windows сейчас не поддерживаются. ## Состояние связи - **Connected** означает, что SSH-туннель готов и Child может использовать связь Home. -- **Reconnecting** означает, что туннель переподключается. Во время повторных попыток запросы могут временно получать `503` с `Retry-After`. -- **Failed** означает, что связь требует действий. Проверьте SSH-аутентификацию, подтверждённый ключ хоста, перенаправление или причину тайм-аута. +- **Reconnecting** означает, что туннель переподключается. Во время повторных попыток запросы могут временно получать `503` с `Retry-After`. На Child, подключённом из собственной панели, запрос сначала до 15 секунд ждёт восстановления туннеля. +- **Failed** означает, что связь требует действий. Проверьте SSH-аутентификацию, подтверждённый ключ хоста, перенаправление или причину тайм-аута. Child, подключённый из собственной панели, сам повторяет попытки после сна, сбоя или перезапуска: примерно раз в минуту после тайм-аута или ошибки перенаправления и каждые пять минут после ошибки аутентификации. Изменённый ключ хоста никогда не повторяется. При сбое связи система молча не переключается на локального провайдера. diff --git a/docs-site/src/content/docs/tr/guides/remote-link.md b/docs-site/src/content/docs/tr/guides/remote-link.md index df5c33a0201..3fe13ce1e89 100644 --- a/docs-site/src/content/docs/tr/guides/remote-link.md +++ b/docs-site/src/content/docs/tr/guides/remote-link.md @@ -41,8 +41,8 @@ Bağlanmak bu bilgisayardaki OpenCodex'i yeniden başlatır. Zaten çalışan Co ## Bağlantı durumu - **Connected**, SSH tünelinin hazır ve Child'ın Home bağlantısını kullanabilir olduğu anlamına gelir. -- **Reconnecting**, tünelin yeniden denendiği anlamına gelir. Yeniden deneme sırasında istekler geçici olarak `Retry-After` ile birlikte `503` döndürebilir. -- **Failed**, bağlantının ilgilenilmesi gerektiği anlamına gelir. SSH kimlik doğrulamasını, onaylanan ana bilgisayar anahtarını, yönlendirmeyi veya zaman aşımı nedenini kontrol edin. +- **Reconnecting**, tünelin yeniden denendiği anlamına gelir. Yeniden deneme sırasında istekler geçici olarak `Retry-After` ile birlikte `503` döndürebilir. Kendi panosundan bağlanan bir Child üzerinde istek önce tünelin geri gelmesi için en fazla 15 saniye bekler. +- **Failed**, bağlantının ilgilenilmesi gerektiği anlamına gelir. SSH kimlik doğrulamasını, onaylanan ana bilgisayar anahtarını, yönlendirmeyi veya zaman aşımı nedenini kontrol edin. Kendi panosundan bağlanan bir Child; uyku, kesinti veya yeniden başlatmadan sonra kendiliğinden yeniden dener: zaman aşımı veya yönlendirme hatasından sonra yaklaşık dakikada bir, kimlik doğrulama hatasından sonra beş dakikada bir. Değişmiş bir ana bilgisayar anahtarı asla yeniden denenmez. Bağlantı başarısız olduğunda sistem sessizce yerel bir sağlayıcıya geçmez. diff --git a/docs-site/src/content/docs/zh-cn/guides/remote-link.md b/docs-site/src/content/docs/zh-cn/guides/remote-link.md index 1440a27c3fd..6b092f1b9ab 100644 --- a/docs-site/src/content/docs/zh-cn/guides/remote-link.md +++ b/docs-site/src/content/docs/zh-cn/guides/remote-link.md @@ -41,8 +41,8 @@ description: 通过 SSH 将 OpenCodex 主机与子机连接起来。 ## 链接状态 - **Connected** 表示 SSH 隧道已就绪,子机可以使用主机链接。 -- **Reconnecting** 表示正在重试隧道。重试期间请求可能暂时返回带有 `Retry-After` 的 `503`。 -- **Failed** 表示链接需要处理。请检查 SSH 身份验证、已确认的主机密钥、转发或超时原因。 +- **Reconnecting** 表示正在重试隧道。重试期间请求可能暂时返回带有 `Retry-After` 的 `503`。在从自己的仪表板连接的 Child 上,请求会先最多等待 15 秒让隧道恢复。 +- **Failed** 表示链接需要处理。请检查 SSH 身份验证、已确认的主机密钥、转发或超时原因。从自己的仪表板连接的 Child 会在睡眠、故障或重启后自动重试:超时或转发错误后大约每分钟一次,身份验证错误后每五分钟一次。主机密钥变更时不会重试。 链接失败时不会静默切换到本地提供商。 diff --git a/docs-site/src/content/docs/zh-tw/guides/remote-link.md b/docs-site/src/content/docs/zh-tw/guides/remote-link.md index 02258a4dc16..f30a8987c43 100644 --- a/docs-site/src/content/docs/zh-tw/guides/remote-link.md +++ b/docs-site/src/content/docs/zh-tw/guides/remote-link.md @@ -41,8 +41,8 @@ description: 透過 SSH 連接 OpenCodex Home 電腦與 Child 電腦。 ## 連結狀態 - **Connected** 表示 SSH 通道已準備好,Child 可以使用 Home 連結。 -- **Reconnecting** 表示正在重試通道。重試期間請求可能暫時回傳帶有 `Retry-After` 的 `503`。 -- **Failed** 表示連結需要處理。請檢查 SSH 驗證、已確認的主機金鑰、轉送或逾時原因。 +- **Reconnecting** 表示正在重試通道。重試期間請求可能暫時回傳帶有 `Retry-After` 的 `503`。在從自己的儀表板連線的 Child 上,請求會先最多等待 15 秒讓通道恢復。 +- **Failed** 表示連結需要處理。請檢查 SSH 驗證、已確認的主機金鑰、轉送或逾時原因。從自己的儀表板連線的 Child 會在睡眠、故障或重新啟動後自動重試:逾時或轉送錯誤後大約每分鐘一次,驗證錯誤後每五分鐘一次。主機金鑰變更時不會重試。 連結失敗時不會靜默切換到本機供應商。 diff --git a/src/client/link-join.ts b/src/client/link-join.ts index 8d6f1c8c979..5059a4a5515 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 { isLinkPort } from "../link/ports"; +import { isPortAvailable } from "../server/ports"; +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"; import { connectClient, type ClientConnectDeps } from "./connect"; @@ -24,6 +24,7 @@ const JOIN_TUNNEL_READY_TIMEOUT_MS = 15_000; const JOIN_TUNNEL_POLL_MS = 100; const JOIN_REVOKE_TIMEOUT_MS = 30_000; const JOIN_CONFIRM_TTL_MS = 5 * 60_000; +const JOIN_PORT_ATTEMPTS = 32; const LINK_ID = /^lnk_[0-9a-f]{16}$/; const API_KEY_ID = /^[A-Za-z0-9][A-Za-z0-9_.:-]{0,255}$/; const DATA_KEY = /^ocx_data_[0-9a-f]{40}$/; @@ -120,6 +121,24 @@ function parseIssuedLink(stdout: string): IssuedLink | null { }; } +/** + * A free loopback port in the join range (`JOIN_TUNNEL_PORT_MIN`-`JOIN_TUNNEL_PORT_MAX`), tried at + * random so a fixed-port service on this computer is not hit every time. + */ +export async function chooseJoinTunnelPort(deps: { + isAvailable?: (port: number) => Promise; + random?: () => number; +} = {}): Promise { + const isAvailable = deps.isAvailable ?? (port => isPortAvailable(port, "127.0.0.1")); + const random = deps.random ?? Math.random; + const span = JOIN_TUNNEL_PORT_MAX - JOIN_TUNNEL_PORT_MIN + 1; + for (let attempt = 0; attempt < JOIN_PORT_ATTEMPTS; attempt += 1) { + const port = JOIN_TUNNEL_PORT_MIN + Math.min(span - 1, Math.floor(random() * span)); + if (await isAvailable(port)) return port; + } + throw new Error("no free port in the join tunnel range"); +} + function localAlias(deps: ClientLinkJoinDeps): string { const raw = (deps.hostname ?? hostname)().trim(); const normalized = raw.replace(/[^A-Za-z0-9_\.\-]/g, "-").replace(/^-+/, "").slice(0, 253); @@ -242,7 +261,7 @@ export async function joinHome(deps: ClientLinkJoinDeps, input: { alias: string await compensateStaleSidecar(deps); let tunnelPort: number; try { - tunnelPort = await (deps.choosePort ?? (() => findAvailablePort(0, "127.0.0.1")))(); + tunnelPort = await (deps.choosePort ?? (() => chooseJoinTunnelPort()))(); if (!isLinkPort(tunnelPort)) throw new Error("invalid link port"); } catch (error) { void error; diff --git a/src/client/link-relay.ts b/src/client/link-relay.ts index 683e9106b02..911bc9a08b3 100644 --- a/src/client/link-relay.ts +++ b/src/client/link-relay.ts @@ -22,6 +22,21 @@ export interface LinkRelayClock { clearTimeout: typeof clearTimeout; } +/** + * The Child's tunnel as the relay sees it (the client link supervisor implements it). A request + * waits on it only while the tunnel is being (re)established; a connected tunnel costs one + * `pending()` call per request and nothing else. + */ +export interface LinkTunnelGate { + /** True while the tunnel is connecting or reconnecting. */ + pending(): boolean; + /** + * Resolves true once the tunnel is connected, and false after `timeoutMs`, when `signal` aborts, + * when the tunnel stops being re-established (failed or stopped), or when too many requests wait. + */ + waitForConnected(timeoutMs: number, signal?: AbortSignal): Promise; +} + export interface LinkRelayDeps { fetchImpl?: typeof fetch; clock?: LinkRelayClock; @@ -30,6 +45,11 @@ export interface LinkRelayDeps { sseIdleTimeoutMs?: number; /** Byte cap for the streamed request body and for a non-SSE response body. */ bodyLimitBytes?: number; + /** The Child's tunnel; without one a request is forwarded at once, as before. */ + tunnel?: LinkTunnelGate; + /** The longest a request waits for a reconnecting tunnel, from its first wait. */ + holdMs?: number; + now?: () => number; } export const LINK_RELAY_RETRY_AFTER_SECONDS = 1; @@ -41,6 +61,12 @@ export const LINK_RELAY_SSE_IDLE_TIMEOUT_MS = 300_000; export const LINK_RELAY_HEADER_TIMEOUT_MS = 300_000; /** The data-plane default: the same inbound limit a standalone listener admits. */ export const LINK_RELAY_BODY_MAX_BYTES = resolveInboundBodyLimitBytes(undefined); +/** + * How long a request waits for a tunnel that is connecting or reconnecting (right after the + * restart into a Child, after sleep, after an ssh exit) before it is answered 503. Codex's own + * retries cover only a few seconds. + */ +export const LINK_RELAY_HOLD_MS = 15_000; const defaultClock: LinkRelayClock = { setTimeout: globalThis.setTimeout, @@ -126,6 +152,7 @@ function byteCappedRequestBody( body: ReadableStream, limit: number, onOverflow: () => void, + onCancel: () => void, ): ReadableStream { const reader = body.getReader(); let bytes = 0; @@ -151,11 +178,21 @@ function byteCappedRequestBody( } }, async cancel(reason) { + onCancel(); try { await reader.cancel(reason); } catch { /* best effort */ } }, }); } +/** + * The connection itself was refused: nothing reached the Home, so the request may be sent again. + * Any other failure (a reset, a timeout, an error after the body started) may have reached it. + */ +function connectionRefused(error: unknown): boolean { + const code = (error as { code?: unknown } | null)?.code; + return code === "ConnectionRefused" || code === "ECONNREFUSED"; +} + function idleBoundedStream( body: ReadableStream, signal: AbortSignal, @@ -283,28 +320,52 @@ export async function relayLinkDataRequest( if (req.signal.aborted) onClientAbort(); let bodyOverflow = false; + let bodyCancelled = false; const body = req.method === "GET" || req.method === "HEAD" || !req.body ? null - : byteCappedRequestBody(req.body, bodyLimit, () => { bodyOverflow = true; }); + : byteCappedRequestBody(req.body, bodyLimit, () => { bodyOverflow = true; }, () => { bodyCancelled = true; }); // A streamed body keeps the caller's Content-Length, so the Home sees the same framing a // buffered body produced instead of a chunked upload. if (body && declaredLength !== null) headers.set("content-length", declaredLength); - let upstream: Response; - try { - const init: RequestInit & { duplex?: "half" } = { - method: req.method, - headers, - redirect: "manual", - signal: relayAbort.signal, - ...(body ? { body, duplex: "half" } : {}), - }; - upstream = await (deps.fetchImpl ?? fetch)(destination, init); - } catch { + // The hold: only while the tunnel is connecting or reconnecting, and at most `holdMs` from the + // first wait. A connected tunnel skips it after one `pending()` call. + const tunnel = deps.tunnel; + let holdUntil: number | undefined; + const holdForTunnel = async (gate: LinkTunnelGate): Promise => { + const clockNow = deps.now ?? Date.now; + holdUntil ??= clockNow() + (positive(deps.holdMs) ?? LINK_RELAY_HOLD_MS); + const remaining = holdUntil - clockNow(); + return remaining > 0 && await gate.waitForConnected(remaining, relayAbort.signal); + }; + const tunnelUnavailable = async (): Promise => { cleanup(); - return bodyOverflow - ? jsonError(413, "link relay request body too large") - : jsonError(503, "link tunnel unavailable", true); + try { await body?.cancel(); } catch { /* best effort */ } + return jsonError(503, "link tunnel unavailable", true); + }; + if (tunnel?.pending() && !await holdForTunnel(tunnel)) return await tunnelUnavailable(); + + let upstream: Response | undefined; + while (!upstream) { + try { + const init: RequestInit & { duplex?: "half" } = { + method: req.method, + headers, + redirect: "manual", + signal: relayAbort.signal, + ...(body ? { body, duplex: "half" } : {}), + }; + upstream = await (deps.fetchImpl ?? fetch)(destination, init); + } catch (error) { + if (bodyOverflow) { + cleanup(); + return jsonError(413, "link relay request body too large"); + } + // A refused connection sent nothing, and an untouched body can go again once the tunnel + // is back inside the same hold window. Anything else is never replayed. + const resendable = connectionRefused(error) && (body === null || (!body.locked && !bodyCancelled)); + if (!tunnel || !resendable || !tunnel.pending() || !await holdForTunnel(tunnel)) return await tunnelUnavailable(); + } } if (relayAbort.signal.aborted) { cleanup(); diff --git a/src/client/link-status.ts b/src/client/link-status.ts index a3245dee9cb..41c78fcef4f 100644 --- a/src/client/link-status.ts +++ b/src/client/link-status.ts @@ -31,6 +31,9 @@ function childState(sidecar: ClientLinkSidecarRead, supervisor: ClientLinkSuperv alias, state: tunnel.kind, since: new Date(tunnel.kind === "idle" ? now : tunnel.since).toISOString(), - reason: tunnel.kind === "failed" ? tunnel.reason : null, + // The keyed probe's finding (the Home refused the key, did not answer, or reports its own + // readiness as not ready) names the cause better than the tunnel state, so it wins; otherwise + // a failed tunnel shows its own reason. + reason: supervisor.probe ?? (tunnel.kind === "failed" ? tunnel.reason : null), }; } diff --git a/src/client/link-tunnel.ts b/src/client/link-tunnel.ts index d41d898a1eb..60bc2af21fe 100644 --- a/src/client/link-tunnel.ts +++ b/src/client/link-tunnel.ts @@ -1,19 +1,25 @@ -import { chmodSync, mkdirSync, readFileSync, unlinkSync } from "node:fs"; +import { chmodSync, mkdirSync, readFileSync, statSync, unlinkSync } from "node:fs"; import { dirname, join } from "node:path"; import { atomicWriteFile, isMissingPathError } from "../config/atomic-write"; +import { getConfigPath } from "../config/paths"; +import { readBoundedResponseBytes } from "../lib/bounded-body"; import { linkDir, linkKnownHostsPath } from "../link/paths"; import { buildTunnelArgv } from "../link/ssh-argv"; import { createSshRunner, type SshChild, type SshRunner } from "../link/ssh-runner"; import { + CLIENT_TUNNEL_RETRY_POLICY, classifySshStderr, dueForSpawn, + failedTunnel, IDLE, reduceTunnel, + type StderrClass, type TunnelState, } from "../link/tunnel-state"; import { isLinkPort } from "../link/ports"; import { isLinkConnection, readClientConnectionState } from "./state"; import { clientLinkStatePath, readClientLinkState, type ClientLinkState } from "./link-state"; +import type { LinkTunnelGate } from "./link-relay"; /** * The client-owned `ssh -N -L 127.0.0.1::127.0.0.1: ` @@ -45,21 +51,36 @@ export interface ClientLinkTunnelDeps { export type OrphanTunnelResult = | { tunnel: "reaped" } | { tunnel: "absent" } - | { tunnel: "owned" } + | { tunnel: "owned"; pid?: number } | { tunnel: "unresolved"; pid: number }; +/** A macOS process as `ps -o ppid= -o args=` shows it: the parent pid and the space-joined argv. */ +export interface DarwinProcessInfo { + ppid: number; + args: string; +} + export interface OrphanReapDeps { configDir?: string; platform?: NodeJS.Platform; readProcessArgv?: (pid: number) => readonly string[] | null; + /** macOS: the parent pid and argv of `pid`, or null when they cannot be read. */ + readProcessInfo?: (pid: number) => DarwinProcessInfo | null; isAlive?: (pid: number) => boolean; signal?: (pid: number, signal: NodeJS.Signals) => void; sleep?: (ms: number) => Promise; } +/** + * What the periodic keyed probe last saw, for display only: it never changes a connected tunnel. + * `home_not_ready` is a Home that admitted the key but reports its own startup readiness as pending + * or failed; its link still works. + */ +export type ClientTunnelProbeReason = "unauthorized" | "home_unreachable" | "home_not_ready"; + export type ClientLinkSupervisorStatus = | { kind: "stopped" } - | { kind: "tunnel"; linkId: string; state: TunnelState; pid: number | null } + | { kind: "tunnel"; linkId: string; state: TunnelState; pid: number | null; probe?: ClientTunnelProbeReason } | { kind: "failed"; reason: "sidecar_invalid" }; export interface ClientLinkTunnelStatusProjection { @@ -69,7 +90,8 @@ export interface ClientLinkTunnelStatusProjection { reason: "sidecar_invalid"; } -export interface ClientLinkSupervisor { +/** The supervisor is also the relay's tunnel gate: requests wait on it while it reconnects. */ +export interface ClientLinkSupervisor extends LinkTunnelGate { start(): void; /** Stops the tunnel (TERM, up to 5 s, KILL). The runtime calls this before stopping its listener. */ stop(): Promise; @@ -89,12 +111,21 @@ export function clientLinkTunnelStatus( } } +/** `connectedLinkId` when the connection state could not be read (a write in flight, a bad file). */ +export const CONNECTION_UNREADABLE = "unreadable"; + export interface ClientLinkSupervisorDeps extends ClientLinkTunnelDeps, OrphanReapDeps { readSidecar?: () => ClientLinkState | null; - /** Current link id of a connected link-transport client, or null when that no longer holds. */ - connectedLinkId?: () => string | null; + /** + * Current link id of a connected link-transport client, null when that no longer holds, or + * `CONNECTION_UNREADABLE` when the connection state could not be read. + */ + connectedLinkId?: () => string | null | typeof CONNECTION_UNREADABLE; /** Called once after the tunnel stopped because the link ended (the runtime recycles here). */ onLinkEnded?: () => void; + /** The link key for the keyed readiness probe. The runtime passes its cached key source. */ + linkKey?: () => string | null; + fetchImpl?: typeof fetch; now?: () => number; random?: () => number; warn?: (message: string) => void; @@ -118,8 +149,18 @@ export interface ClientTunnelPidfile { } const STOP_TIMEOUT_MS = 5_000; +const REAP_POLL_MS = 100; const TIMER_MS = 1_000; -const SPAWN_GRACE_MS = 5_000; +/** Consecutive unreadable reads before the supervisor acts on them. */ +export const INVALID_READ_TICKS = 3; +/** How often a connected tunnel is probed. The result is display-only. */ +export const CONNECTED_PROBE_MS = 30_000; +const PROBE_TIMEOUT_MS = 5_000; +const PROBE_BACKOFF_MAX_MS = 5_000; +/** A Home `/readyz` body is a few hundred bytes; a larger one is not the Home's. */ +const PROBE_BODY_MAX_BYTES = 4_096; +/** Requests that may wait on a reconnecting tunnel at once; more are answered 503 at once. */ +export const CLIENT_LINK_MAX_HOLDS = 64; function sameArgv(left: readonly string[], right: readonly string[]): boolean { return left.length === right.length && left.every((value, index) => value === right[index]); @@ -196,6 +237,22 @@ function linuxProcessArgv(pid: number): readonly string[] | null { } } +/** `ps -ww -o ppid= -o args= -p `: one line, the right-aligned parent pid, one space, the argv. */ +function darwinProcessInfo(pid: number): DarwinProcessInfo | null { + try { + const result = Bun.spawnSync(["/bin/ps", "-ww", "-o", "ppid=", "-o", "args=", "-p", String(pid)], { + stdin: "ignore", + stdout: "pipe", + stderr: "ignore", + }); + if (result.exitCode !== 0) return null; + const match = /^\s*(\d+) (.+)$/.exec(result.stdout.toString().replace(/\r?\n$/, "")); + return match ? { ppid: Number(match[1]), args: match[2]! } : null; + } catch { + return null; + } +} + function timerDeps(deps: ClientLinkTunnelDeps): Required> { return { setTimer: deps.setTimer ?? ((callback, ms) => setTimeout(callback, ms)), @@ -265,44 +322,158 @@ export function spawnClientLinkTunnel(spec: ClientLinkTunnelSpec, deps: ClientLi return handle; } +type TunnelIdentity = + | { kind: "gone" } + | { kind: "other" } + | { kind: "ours"; orphaned: boolean } + | { kind: "unknown" }; + +/** + * Whether the pidfile's process is still our tunnel. Linux compares `/proc//cmdline` with + * the recorded argv. macOS compares `ps` args with the argv joined by spaces and reports the + * process orphaned only when launchd (pid 1) is its parent. Elsewhere a live process is unknown. + */ +function tunnelIdentity( + pidfile: ClientTunnelPidfile, + platform: NodeJS.Platform, + deps: OrphanReapDeps, + isAlive: (pid: number) => boolean, +): TunnelIdentity { + if (platform === "linux") { + const actualArgv = (deps.readProcessArgv ?? linuxProcessArgv)(pidfile.pid); + if (!actualArgv) return { kind: "gone" }; + return sameArgv(actualArgv, pidfile.argv) ? { kind: "ours", orphaned: true } : { kind: "other" }; + } + if (!isAlive(pidfile.pid)) return { kind: "gone" }; + if (platform !== "darwin") return { kind: "unknown" }; + const info = (deps.readProcessInfo ?? darwinProcessInfo)(pidfile.pid); + if (!info) return { kind: "unknown" }; + if (info.args !== pidfile.argv.join(" ")) return { kind: "other" }; + return { kind: "ours", orphaned: info.ppid === 1 }; +} + +/** + * Settle a leftover tunnel pidfile before a new tunnel starts. + * + * - A pidfile whose process is gone, or is provably another program (a reused pid, as after a + * reboot), is stale: it is removed and the result is `absent`, so a new tunnel starts. + * - While the owner lives, a tunnel that is still ours (or cannot be told apart) is `owned`. + * - After the owner exits, a proven orphan is reaped: Linux on an exact `/proc` argv match, macOS + * on an exact `ps` argv match with launchd as the parent. TERM, up to five seconds, KILL. + * - Anything else is `unresolved`: never signalled; the caller watches it. + */ export async function reapOrphanTunnel(deps: OrphanReapDeps = {}): Promise { const path = clientTunnelPidfilePath(deps.configDir); const pidfile = readPidfile(path); if (!pidfile) return { tunnel: "absent" }; const isAlive = deps.isAlive ?? defaultIsAlive; - if (isAlive(pidfile.ownerPid)) return { tunnel: "owned" }; const platform = deps.platform ?? process.platform; - if (platform !== "linux") return { tunnel: "unresolved", pid: pidfile.pid }; - const readProcessArgv = deps.readProcessArgv ?? linuxProcessArgv; - const actualArgv = readProcessArgv(pidfile.pid); - if (!actualArgv || !sameArgv(actualArgv, pidfile.argv)) { - try { unlinkSync(path); } catch (error) { if (!isMissingPathError(error)) throw error; } + const identity = tunnelIdentity(pidfile, platform, deps, isAlive); + if (identity.kind === "gone" || identity.kind === "other") { + removePidfileIfPid(path, pidfile.pid); return { tunnel: "absent" }; } + if (isAlive(pidfile.ownerPid)) return { tunnel: "owned", pid: pidfile.pid }; + if (identity.kind !== "ours" || !identity.orphaned) return { tunnel: "unresolved", pid: pidfile.pid }; const signal = deps.signal ?? defaultSignal; const sleep = deps.sleep ?? ((ms: number) => new Promise(resolve => setTimeout(resolve, ms))); try { signal(pidfile.pid, "SIGTERM"); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ESRCH") throw error; } - await sleep(STOP_TIMEOUT_MS); + for (let waited = 0; waited < STOP_TIMEOUT_MS && isAlive(pidfile.pid); waited += REAP_POLL_MS) await sleep(REAP_POLL_MS); if (isAlive(pidfile.pid)) { try { signal(pidfile.pid, "SIGKILL"); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ESRCH") throw error; } } - try { unlinkSync(path); } catch (error) { if (!isMissingPathError(error)) throw error; } + removePidfileIfPid(path, pidfile.pid); return { tunnel: "reaped" }; } -function defaultConnectedLinkId(): string | null { +function defaultConnectedLinkId(): string | null | typeof CONNECTION_UNREADABLE { const state = readClientConnectionState(); + if (state.kind === "invalid" || state.kind === "mismatched") return CONNECTION_UNREADABLE; if (state.kind !== "connected" || !isLinkConnection(state.value)) return null; return state.value.link?.linkId ?? null; } +function fileSignature(path: string): string | null { + try { + const stat = statSync(path); + return `${stat.dev}:${stat.ino}:${stat.size}:${stat.mtimeMs}:${stat.ctimeMs}`; + } catch (error) { + return isMissingPathError(error) ? "missing" : null; + } +} + +/** + * `read` again only when the file at `path` changed, so a steady link costs one stat per check + * instead of a parse. A read that throws, a value `keep` refuses, or a failed stat is never cached. + */ +export function readWhenFileChanges(path: () => string, read: () => T, keep: (value: T) => boolean = () => true): () => T { + let cached: { signature: string; value: T } | null = null; + return () => { + const signature = fileSignature(path()); + if (signature !== null && cached?.signature === signature) return cached.value; + cached = null; + const value = read(); + if (signature !== null && keep(value)) cached = { signature, value }; + return value; + }; +} + +/** `ready` and `home_not_ready` both prove the link works; the rest do not. */ +type ProbeResult = "ready" | ClientTunnelProbeReason; + +/** Whether a 503 `/readyz` body is the Home's own readiness answer (`service: "opencodex"`). */ +async function isOpencodexReadiness(response: Response, signal: AbortSignal): Promise { + try { + const { bytes, oversized } = await readBoundedResponseBytes(response, { maxBytes: PROBE_BODY_MAX_BYTES, signal }); + if (oversized) return false; + const body = JSON.parse(new TextDecoder().decode(bytes)) as unknown; + return typeof body === "object" && body !== null && (body as { service?: unknown }).service === "opencodex"; + } catch { + return false; + } +} + +/** + * `GET /readyz` through the tunnel with the link key; the key is sent as a header only. The Home's + * link listener answers 401 before it reaches `/readyz`, so both a 200 and a 503 carrying the + * Home's readiness body prove that the forward reaches that listener and that the key is admitted. + * The 503 only means the Home's own startup readiness is pending or failed, which does not stop + * relayed requests; it is reported as `home_not_ready`, for display. + */ +async function probeTunnel(fetchImpl: typeof fetch, tunnelPort: number, key: string, stop: AbortSignal): Promise { + const signal = AbortSignal.any([stop, AbortSignal.timeout(PROBE_TIMEOUT_MS)]); + try { + const response = await fetchImpl(`http://127.0.0.1:${tunnelPort}/readyz`, { + headers: { "x-opencodex-api-key": key }, + cache: "no-store", + redirect: "manual", + signal, + }); + if (response.status === 503) return await isOpencodexReadiness(response, signal) ? "home_not_ready" : "home_unreachable"; + try { await response.body?.cancel(); } catch { /* the body is not needed */ } + if (response.status === 200) return "ready"; + return response.status === 401 || response.status === 403 ? "unauthorized" : "home_unreachable"; + } catch { + return "home_unreachable"; + } +} + export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): ClientLinkSupervisor { - const readSidecar = deps.readSidecar ?? (() => readClientLinkState(clientLinkStatePath(deps.configDir))); - const connectedLinkId = deps.connectedLinkId ?? defaultConnectedLinkId; + const sidecarPath = (): string => clientLinkStatePath(deps.configDir); + const readSidecar = deps.readSidecar ?? readWhenFileChanges(sidecarPath, () => readClientLinkState(sidecarPath())); + const connectedLinkId = deps.connectedLinkId + ?? readWhenFileChanges(getConfigPath, defaultConnectedLinkId, value => value !== CONNECTION_UNREADABLE); const now = deps.now ?? (() => Date.now()); const random = deps.random ?? Math.random; - const setSupervisorTimer = deps.setTimer - ?? ((callback: () => void, ms: number) => setInterval(callback, ms) as unknown as ReturnType); + const policy = CLIENT_TUNNEL_RETRY_POLICY; + const isAlive = deps.isAlive ?? defaultIsAlive; + const linkKey = deps.linkKey ?? (() => null); + const fetchImpl = deps.fetchImpl ?? fetch; + const setSupervisorTimer = deps.setTimer ?? ((callback: () => void, ms: number) => { + const interval = setInterval(callback, ms); + (interval as { unref?: () => void }).unref?.(); + return interval as unknown as ReturnType; + }); const clearSupervisorTimer = deps.clearTimer ?? ((timer: ReturnType) => clearInterval(timer as unknown as ReturnType)); let timer: ReturnType | undefined; @@ -312,27 +483,80 @@ export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): let initializing = false; let onLinkEndedCalled = false; let child: ClientLinkTunnelHandle | undefined; + /** A leftover tunnel this supervisor may not signal: watched, and replaced once it dies. */ + let adopted: number | null = null; let state: TunnelState = IDLE; let linkId: string | null = null; let failure: ClientLinkSupervisorStatus | undefined; let tickFlight: Promise | undefined; + /** The one keyed probe in flight, which tunnel it probes, and the abort that stop() fires. */ + let probeFlight: { child: ClientLinkTunnelHandle | undefined; adopted: number | null; abort: AbortController } | undefined; + let probe: ClientTunnelProbeReason | null = null; + let nextProbeAt = 0; + let probeFailures = 0; + let invalidReads = 0; + /** A leftover tunnel pidfile has not been settled yet; nothing may spawn until it is. */ + let reapPending = true; + const waiters = new Set<(connected: boolean) => void>(); + + /** The tunnel is being (re)established, so a request may wait for it instead of failing. */ + const pending = (): boolean => !stopping && failure === undefined + && ((started && !initialized) || state.kind === "connecting" || state.kind === "reconnecting"); + + const settleWaiters = (): void => { + if (waiters.size === 0) return; + const connected = !stopping && failure === undefined && state.kind === "connected"; + if (!connected && pending()) return; + for (const settle of [...waiters]) settle(connected); + }; + + const setState = (next: TunnelState): void => { + state = next; + settleWaiters(); + }; + + const awaitingReady = (): boolean => state.kind === "connecting" || state.kind === "reconnecting" + || (state.kind === "failed" && state.inFlight === true); const readCurrent = (): { sidecar: ClientLinkState | null; invalid: boolean } => { try { return { sidecar: readSidecar(), invalid: false }; - } catch (error) { + } catch { deps.warn?.("client link sidecar could not be read"); return { sidecar: null, invalid: true }; } }; - const stopTunnel = async (): Promise => { + const readConnectedLinkId = (): string | null | typeof CONNECTION_UNREADABLE => { + try { + return connectedLinkId(); + } catch { + return CONNECTION_UNREADABLE; + } + }; + + /** Drops the probe in flight; its answer, if it still comes, is ignored. */ + const abortProbe = (): void => { + const flight = probeFlight; + probeFlight = undefined; + flight?.abort.abort(); + }; + + /** Stops our own child only; the tunnel state is left for the caller to decide. */ + const killChild = async (): Promise => { const current = child; child = undefined; - state = reduceTunnel(state, { type: "stop" }); + probe = null; + abortProbe(); if (current) await current.stop(); }; + const stopTunnel = async (): Promise => { + adopted = null; + setState(reduceTunnel(state, { type: "stop" })); + await killChild(); + }; + const endLink = async (): Promise => { await stopTunnel(); if (onLinkEndedCalled || stopping) return; @@ -341,8 +565,14 @@ export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): deps.onLinkEnded?.(); }; + const onExit = (stderrClass: StderrClass): void => { + setState(reduceTunnel(state, { type: "exit", now: now(), stderrClass }, random, policy)); + }; + const spawn = (sidecar: ClientLinkState): void => { - if (stopping || child || state.kind === "failed") return; + if (stopping || child || adopted !== null || reapPending) return; + const timestamp = now(); + if (state.kind === "failed" && !dueForSpawn(state, timestamp)) return; try { child = spawnClientLinkTunnel({ linkId: sidecar.linkId, @@ -350,53 +580,158 @@ export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): tunnelPort: sidecar.tunnelPort, peerListenerPort: sidecar.peerListenerPort, }, deps); - linkId = sidecar.linkId; - state = reduceTunnel(state, { type: "spawn", now: now() }, random); - const current = child; - void current.exited.then(async () => { - if (child !== current) return; - child = undefined; - const stderr = (current as ClientLinkTunnelHandle & { stderr?: Promise }).stderr - ? await (current as ClientLinkTunnelHandle & { stderr?: Promise }).stderr!.catch(() => "") - : ""; - const next = reduceTunnel(state, { type: "exit", now: now(), stderrClass: classifySshStderr(stderr) }, random); - state = next; - }).catch(() => { - if (child !== current) return; - child = undefined; - state = reduceTunnel(state, { type: "exit", now: now(), stderrClass: "network" }, random); - }); - } catch (error) { - state = { kind: "failed", since: now(), reason: "forward" }; + } catch { + setState(failedTunnel("forward", timestamp, policy, state.kind === "failed" ? state.since : timestamp)); deps.warn?.("client link tunnel could not be started"); + return; + } + linkId = sidecar.linkId; + probe = null; + probeFailures = 0; + nextProbeAt = timestamp; + setState(reduceTunnel(state, { type: "spawn", now: timestamp }, random, policy)); + const current = child; + void current.exited.then(async () => { + if (child !== current) return; + child = undefined; + probe = null; + const stderr = (current as ClientLinkTunnelHandle & { stderr?: Promise }).stderr + ? await (current as ClientLinkTunnelHandle & { stderr?: Promise }).stderr!.catch(() => "") + : ""; + onExit(classifySshStderr(stderr)); + }).catch(() => { + if (child !== current) return; + child = undefined; + probe = null; + onExit("network"); + }); + }; + + /** + * Applies one keyed probe's answer. While the tunnel is being established an answer that proves + * the link (`ready`, or `home_not_ready`) promotes it to connected; otherwise the next probe backs + * off from one to five seconds, or up to 30 seconds once the link reads failed (a revoked key or a + * stopped Home costs one probe per 30 s). While connected the answer is display-only, and the next + * probe is 30 seconds out. + */ + const applyProbe = (result: ProbeResult, wasConnected: boolean): void => { + const timestamp = now(); + const reason = result === "ready" ? null : result; + if (wasConnected) { + if (state.kind === "connected") probe = reason; + nextProbeAt = timestamp + CONNECTED_PROBE_MS; + return; + } + if ((result === "ready" || result === "home_not_ready") && awaitingReady()) { + probe = reason; + probeFailures = 0; + nextProbeAt = timestamp + CONNECTED_PROBE_MS; + setState(reduceTunnel(state, { type: "ready", now: timestamp }, random, policy)); + return; + } + probe = reason; + probeFailures += 1; + const cap = state.kind === "failed" ? CONNECTED_PROBE_MS : PROBE_BACKOFF_MAX_MS; + nextProbeAt = timestamp + Math.min(cap, 1_000 * 2 ** Math.min(probeFailures - 1, 16)); + }; + + /** + * Starts one keyed probe without waiting for it, so a slow Home never delays the check that + * notices a disconnect, and stop() aborts it instead of waiting up to the probe timeout. + */ + const startProbe = (tunnelPort: number, key: string): void => { + const flight = { child, adopted, abort: new AbortController() }; + probeFlight = flight; + const wasConnected = state.kind === "connected"; + void probeTunnel(fetchImpl, tunnelPort, key, flight.abort.signal).then(result => { + // An aborted probe, or one whose tunnel exited meanwhile, says nothing about the next tunnel. + if (probeFlight !== flight) return; + probeFlight = undefined; + if (stopping || child !== flight.child || adopted !== flight.adopted) return; + applyProbe(result, wasConnected); + }, () => { + if (probeFlight === flight) probeFlight = undefined; + }); + }; + + /** + * Settles a leftover tunnel pidfile once, before this supervisor's first spawn: a proven orphan + * is reaped and a live tunnel that may not be ours is adopted. It needs a valid, matching read + * after the reap; without one it stays pending and the next check that has one runs it again. + */ + const settleLeftover = async (): Promise => { + const orphan = await reapOrphanTunnel(deps); + if (stopping) return; + const afterReap = readCurrent(); + if (afterReap.invalid || !afterReap.sidecar || readConnectedLinkId() !== afterReap.sidecar.linkId) return; + reapPending = false; + if ((orphan.tunnel === "owned" || orphan.tunnel === "unresolved") && orphan.pid !== undefined && isAlive(orphan.pid)) { + // A leftover tunnel that may not be ours to stop: watch it and probe through it, and + // start our own once it dies. It is never signalled. + adopted = orphan.pid; + linkId = afterReap.sidecar.linkId; + nextProbeAt = now(); + setState({ kind: "connecting", since: now() }); + return; } + spawn(afterReap.sidecar); }; const tick = async (): Promise => { if (stopping || !initialized) return; const current = readCurrent(); - if (current.invalid) { - failure = { kind: "failed", reason: "sidecar_invalid" }; - await stopTunnel(); + const connected = current.invalid ? null : readConnectedLinkId(); + if (current.invalid || connected === CONNECTION_UNREADABLE) { + // A single unreadable read (a write in progress, a transient I/O error) must not end a + // healthy link; only INVALID_READ_TICKS in a row do. + invalidReads += 1; + if (invalidReads < INVALID_READ_TICKS) return; + if (current.invalid) { + failure = { kind: "failed", reason: "sidecar_invalid" }; + await stopTunnel(); + } else if (child || linkId || adopted !== null) { + await endLink(); + } + settleWaiters(); return; } - const connected = connectedLinkId(); - if (!current.sidecar) { - if (child || linkId) await endLink(); + invalidReads = 0; + if (!current.sidecar || connected !== current.sidecar.linkId) { + if (child || linkId || adopted !== null) await endLink(); return; } - if (connected !== current.sidecar.linkId) { - if (child || linkId) await endLink(); + failure = undefined; + const sidecar = current.sidecar; + if (reapPending) { + // The start-up read was unreadable or did not match, so a leftover tunnel was never + // settled. Settle it now, before anything spawns, so an orphan cannot keep the port. + await settleLeftover(); + settleWaiters(); return; } - failure = undefined; - const timestamp = now(); - state = reduceTunnel(state, { type: "tick", now: timestamp }, random); - if (child && state.kind === "connecting" && timestamp - state.since >= SPAWN_GRACE_MS) { - state = reduceTunnel(state, { type: "ready", now: timestamp }, random); + if (adopted !== null && !isAlive(adopted)) { + // The leftover tunnel is gone: start our own on this tick. + adopted = null; + probe = null; + setState(IDLE); } - if (state.kind === "failed" && child) await stopTunnel(); - else if (!child && (state.kind === "idle" || dueForSpawn(state, timestamp))) spawn(current.sidecar); + setState(reduceTunnel(state, { type: "tick", now: now() }, random, policy)); + // An adopted tunnel is the attempt in flight: it is never killed, so a timeout leaves it + // probed, and it can still be promoted, until it dies. + if (state.kind === "failed" && !state.inFlight && adopted !== null) setState({ ...state, inFlight: true }); + if (state.kind === "failed" && !state.inFlight && child) await killChild(); + // A probe still running for a tunnel that has since exited or died answers nothing useful. + if (probeFlight && (probeFlight.child !== child || probeFlight.adopted !== adopted)) abortProbe(); + // While requests are held the backoff does not apply: one probe per check (one a second), so a + // forward that just came up releases them at once. Held requests exist only while it is down. + if ((child || adopted !== null) && (state.kind === "connected" || awaitingReady()) && !probeFlight + && (now() >= nextProbeAt || waiters.size > 0)) { + // Without the committed key nothing can be proven, and the relay refuses requests anyway. + const key = linkKey(); + if (key) startProbe(sidecar.tunnelPort, key); + } + if (!child && adopted === null && (state.kind === "idle" || dueForSpawn(state, now()))) spawn(sidecar); + settleWaiters(); }; const runTick = (): void => { @@ -407,27 +742,21 @@ export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): const initialize = async (): Promise => { if (initializing || initialized || stopping) return; initializing = true; - const current = readCurrent(); - if (current.invalid) { - failure = { kind: "failed", reason: "sidecar_invalid" }; - initialized = true; - initializing = false; - return; - } - if (current.sidecar && connectedLinkId() === current.sidecar.linkId) { - const orphan = await reapOrphanTunnel(deps); - if (orphan.tunnel === "owned" || orphan.tunnel === "unresolved") { - linkId = current.sidecar.linkId; - state = { kind: "connected", since: now() }; - initialized = true; - initializing = false; + try { + const current = readCurrent(); + if (current.invalid) { + failure = { kind: "failed", reason: "sidecar_invalid" }; return; } - const afterReap = readCurrent(); - if (!afterReap.invalid && afterReap.sidecar && connectedLinkId() === afterReap.sidecar.linkId) spawn(afterReap.sidecar); + // An unreadable or mismatched read leaves the leftover tunnel for the first check that + // reads a matching link (reapPending), so no tunnel is ever spawned over an unreaped one. + if (!current.sidecar || readConnectedLinkId() !== current.sidecar.linkId) return; + await settleLeftover(); + } finally { + initialized = true; + initializing = false; + settleWaiters(); } - initialized = true; - initializing = false; }; return { @@ -440,6 +769,7 @@ export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): failure = { kind: "failed", reason: "sidecar_invalid" }; initialized = true; initializing = false; + settleWaiters(); }); }, async stop(): Promise { @@ -448,6 +778,8 @@ export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): return; } stopping = true; + abortProbe(); + settleWaiters(); if (timer !== undefined) { clearSupervisorTimer(timer); timer = undefined; @@ -458,8 +790,35 @@ export function createClientLinkSupervisor(deps: ClientLinkSupervisorDeps = {}): }, status(): ClientLinkSupervisorStatus { if (failure) return failure; - if (!child && !linkId) return { kind: "stopped" }; - return { kind: "tunnel", linkId: linkId ?? "", state, pid: child?.pid ?? null }; + if (!child && !linkId && adopted === null) return { kind: "stopped" }; + return { + kind: "tunnel", + linkId: linkId ?? "", + state, + pid: child?.pid ?? adopted, + ...(probe ? { probe } : {}), + }; + }, + pending, + waitForConnected(timeoutMs: number, signal?: AbortSignal): Promise { + if (!stopping && failure === undefined && state.kind === "connected") return Promise.resolve(true); + if (!pending() || waiters.size >= CLIENT_LINK_MAX_HOLDS || signal?.aborted || !(timeoutMs > 0)) { + return Promise.resolve(false); + } + return new Promise(resolve => { + let holdTimer: ReturnType | undefined; + const onAbort = (): void => settle(false); + const settle = (connected: boolean): void => { + if (!waiters.delete(settle)) return; + if (holdTimer !== undefined) clearTimeout(holdTimer); + signal?.removeEventListener("abort", onAbort); + resolve(connected); + }; + waiters.add(settle); + holdTimer = setTimeout(() => settle(false), timeoutMs); + (holdTimer as { unref?: () => void }).unref?.(); + signal?.addEventListener("abort", onAbort, { once: true }); + }); }, }; } diff --git a/src/client/machine-listener.ts b/src/client/machine-listener.ts index 093fa6d7a71..3edf50973ac 100644 --- a/src/client/machine-listener.ts +++ b/src/client/machine-listener.ts @@ -17,6 +17,7 @@ import { handleMachineApi, type HubReachability, type MachineApiDeps } from "./m import { MACHINE_GUI_ORIGIN_HEADER, requireMachineAuth } from "./machine-auth"; import { HUB_RELAY_REQUEST_BODY_MAX_BYTES, relayHubManagementRequest } from "./hub-relay"; import { createLinkKeySource, handleLinkIngress, type LinkIngress, type LinkKeySourceDeps } from "./link-ingress"; +import type { LinkTunnelGate } from "./link-relay"; import { readClientLinkState, type ClientLinkState } from "./link-state"; import { projectClientLinkChild, type ClientLinkSidecarRead } from "./link-status"; import type { ClientLinkSupervisorStatus } from "./link-tunnel"; @@ -42,6 +43,10 @@ export interface MachineListenerDeps { readSidecar?: () => ClientLinkState | null; /** Link mode: how the link key is read; the listener reads it once and caches it. */ linkKey?: LinkKeySourceDeps; + /** Link mode: a key source the runtime already holds (shared with the tunnel supervisor). */ + linkKeySource?: () => string | null; + /** Link mode: the tunnel supervisor; relayed requests wait on it while it reconnects. */ + linkTunnel?: LinkTunnelGate; /** Link mode: relay seams (deadline clock and byte cap). */ linkRelay?: LinkIngress["relay"]; serve?: (options: Parameters[0]) => Server; @@ -112,8 +117,8 @@ export function startMachineListener( ? { tunnelPort: connection.link!.tunnelPort, policy: requestPolicyView(config, "127.0.0.1"), - linkKey: createLinkKeySource(connection.tokenFingerprint, deps.linkKey), - relay: { fetchImpl: deps.fetchImpl, bodyLimitBytes: inboundBodyLimit, ...deps.linkRelay }, + linkKey: deps.linkKeySource ?? createLinkKeySource(connection.tokenFingerprint, deps.linkKey), + relay: { fetchImpl: deps.fetchImpl, bodyLimitBytes: inboundBodyLimit, tunnel: deps.linkTunnel, ...deps.linkRelay }, } : null; const readSidecar = deps.readSidecar ?? (() => readClientLinkState()); diff --git a/src/client/runtime.ts b/src/client/runtime.ts index ed7079d7917..c5b10c7f4dc 100644 --- a/src/client/runtime.ts +++ b/src/client/runtime.ts @@ -7,6 +7,7 @@ import { installCrashGuards } from "../lib/crash-guard"; import { selfLaunchArgv } from "../lib/self-launch-argv"; import { loadServiceTokenFromFile, serviceApiTokenFingerprint } from "../lib/service-secrets"; import { findAvailablePort, PortUnavailableError } from "../server/ports"; +import { createLinkKeySource } from "./link-ingress"; import { createClientLinkSupervisor, type ClientLinkSupervisor } from "./link-tunnel"; import { clientLinkStatePath } from "./link-state"; import { startMachineListener } from "./machine-listener"; @@ -121,16 +122,21 @@ export async function startClientRuntime( } throw error; } - // Created before the listener so its status route can read the tunnel state; it starts no - // process and no timer until start(). + // One cached key source serves the relay and the supervisor's keyed probe: the token file is + // read once here, not per request or per probe. + const linkKey = linkMode ? createLinkKeySource(state.value.tokenFingerprint) : undefined; + // Created before the listener so its status route can read the tunnel state and relayed + // requests can wait on it; it starts no process and no timer until start(). const supervisor = linkMode && existsSync(clientLinkStatePath()) ? createClientLinkSupervisor({ onLinkEnded: () => scheduleStandaloneRecycle(state.value.tokenFingerprint), + linkKey, }) : null; const server = startMachineListener(port, { state: state.value, - ...(linkMode ? { linkStatus: () => supervisor?.status() ?? { kind: "stopped" as const } } : {}), + ...(linkMode ? { linkStatus: () => supervisor?.status() ?? { kind: "stopped" as const }, linkKeySource: linkKey } : {}), + ...(supervisor ? { linkTunnel: supervisor } : {}), }); const boundPort = server.port ?? port; activeServer = server; diff --git a/src/link/ports.ts b/src/link/ports.ts index 5dbd56ab366..8995e23c130 100644 --- a/src/link/ports.ts +++ b/src/link/ports.ts @@ -7,6 +7,15 @@ export const MIN_LINK_PORT = 1024; export const MAX_LINK_PORT = 65535; +/** + * A Child joining from its dashboard picks its tunnel port from this range. It sits below the + * common OS ephemeral ranges (macOS and Windows 49152-65535, Linux 32768-60999), so an outgoing + * connection that happens to hold the port rarely blocks the tunnel when it comes back after a + * reboot. The persisted port of an existing link is never rewritten. + */ +export const JOIN_TUNNEL_PORT_MIN = 20_000; +export const JOIN_TUNNEL_PORT_MAX = 29_999; + export function isLinkPort(value: unknown): value is number { return typeof value === "number" && Number.isInteger(value) && value >= MIN_LINK_PORT && value <= MAX_LINK_PORT; } diff --git a/src/link/tunnel-state.ts b/src/link/tunnel-state.ts index f8416366d4d..1df2c703009 100644 --- a/src/link/tunnel-state.ts +++ b/src/link/tunnel-state.ts @@ -5,8 +5,11 @@ * - connected: the forward is up. * - reconnecting: a transient failure; requests through the link fail with 503 meanwhile, and a * new attempt is due at `retryAt`. - * - failed: needs the user. Auth, host key and forward failures are not retried, and neither is a - * link that stayed down for FAILED_AFTER_MS. + * - failed: auth, host key and forward failures, and a link that stayed down for FAILED_AFTER_MS. + * Without a retry policy (the Home's `-R` supervisor) it needs the user. With one (the Child's + * own `-L` tunnel) a reason that has a delay is tried again at `retryAt`, and the retry attempt + * runs with `inFlight` while the state still reads failed; a reason without a delay stays + * terminal. */ export type TunnelFailure = "auth" | "hostkey" | "forward" | "timeout"; @@ -17,7 +20,7 @@ export type TunnelState = | { kind: "connecting"; since: number } | { kind: "connected"; since: number } | { kind: "reconnecting"; since: number; attempt: number; retryAt: number; inFlight: boolean } - | { kind: "failed"; since: number; reason: TunnelFailure }; + | { kind: "failed"; since: number; reason: TunnelFailure; retryAt?: number; inFlight?: boolean }; export type TunnelEvent = | { type: "spawn"; now: number } @@ -26,10 +29,25 @@ export type TunnelEvent = | { type: "tick"; now: number } | { type: "stop" }; +/** Opt-in retry of a failed tunnel: the delay before each failure reason is tried again. */ +export interface TunnelRetryPolicy { + retryFailedAfterMs: Partial>; +} + export const FAILED_AFTER_MS = 5 * 60_000; export const BASE_DELAY_MS = 1_000; export const MAX_DELAY_MS = 30_000; +/** + * The client-owned tunnel's policy. Timeout and forward failures retry about once a minute. Auth + * retries every five minutes, and because each auth failure schedules the next attempt five + * minutes out, no more than 12 attempts reach the Home's sshd in an hour. A changed host key is a + * security signal and is never retried. + */ +export const CLIENT_TUNNEL_RETRY_POLICY: TunnelRetryPolicy = { + retryFailedAfterMs: { timeout: 60_000, forward: 60_000, auth: 5 * 60_000 }, +}; + export const IDLE: TunnelState = { kind: "idle" }; /** Capped exponential backoff with ±20% jitter. `random` is injectable for tests. */ @@ -40,22 +58,52 @@ export function nextDelayMs(attempt: number, random: () => number = Math.random) return Math.round(Math.min(MAX_DELAY_MS, base * jitter)); } -export function reduceTunnel(state: TunnelState, event: TunnelEvent, random?: () => number): TunnelState { +/** A failed state; under a policy that retries `reason` it carries the next attempt's time. */ +export function failedTunnel( + reason: TunnelFailure, + now: number, + policy?: TunnelRetryPolicy, + since: number = now, +): TunnelState { + const delay = policy?.retryFailedAfterMs[reason]; + return delay === undefined + ? { kind: "failed", since, reason } + : { kind: "failed", since, reason, retryAt: now + delay, inFlight: false }; +} + +function failureOf(cls: StderrClass): TunnelFailure | null { + return cls === "auth" || cls === "hostkey" || cls === "forward" ? cls : null; +} + +export function reduceTunnel( + state: TunnelState, + event: TunnelEvent, + random?: () => number, + policy?: TunnelRetryPolicy, +): TunnelState { if (event.type === "stop") return IDLE; switch (event.type) { case "spawn": + if (state.kind === "failed" && state.retryAt !== undefined) return { ...state, inFlight: true }; if (state.kind === "idle" || state.kind === "failed") return { kind: "connecting", since: event.now }; if (state.kind === "reconnecting") return { ...state, inFlight: true }; return state; case "ready": if (state.kind === "connecting" || state.kind === "reconnecting") return { kind: "connected", since: event.now }; + if (state.kind === "failed" && state.inFlight) return { kind: "connected", since: event.now }; return state; case "exit": { - if (state.kind === "idle" || state.kind === "failed") return state; - const cls = event.stderrClass; - if (cls === "auth" || cls === "hostkey" || cls === "forward") return { kind: "failed", since: event.now, reason: cls }; + if (state.kind === "idle") return state; + if (state.kind === "failed") { + if (!state.inFlight) return state; + // A retry of a failed link failed again. It stays failed from the original moment; a + // transient exit keeps the slow cadence (timeout) instead of restarting fast backoff. + return failedTunnel(failureOf(event.stderrClass) ?? "timeout", event.now, policy, state.since); + } + const failure = failureOf(event.stderrClass); + if (failure) return failedTunnel(failure, event.now, policy); const since = state.kind === "reconnecting" || state.kind === "connecting" ? state.since : event.now; - if (event.now - since >= FAILED_AFTER_MS) return { kind: "failed", since: event.now, reason: "timeout" }; + if (event.now - since >= FAILED_AFTER_MS) return failedTunnel("timeout", event.now, policy); const attempt = state.kind === "reconnecting" ? state.attempt + 1 : 1; return { kind: "reconnecting", since, attempt, retryAt: event.now + nextDelayMs(attempt, random), inFlight: false }; } @@ -63,7 +111,7 @@ export function reduceTunnel(state: TunnelState, event: TunnelEvent, random?: () // An attempt in flight does not pause the clock: a first attempt or a retry that hangs past // the limit still fails the link, and the supervisor kills the child on seeing `failed`. if ((state.kind === "reconnecting" || state.kind === "connecting") && event.now - state.since >= FAILED_AFTER_MS) { - return { kind: "failed", since: event.now, reason: "timeout" }; + return failedTunnel("timeout", event.now, policy); } return state; } @@ -71,6 +119,7 @@ export function reduceTunnel(state: TunnelState, event: TunnelEvent, random?: () /** Whether the supervisor should start a new ssh attempt now. */ export function dueForSpawn(state: TunnelState, now: number): boolean { + if (state.kind === "failed") return state.retryAt !== undefined && !state.inFlight && now >= state.retryAt; return state.kind === "reconnecting" && !state.inFlight && now >= state.retryAt; } diff --git a/structure/remote-link.md b/structure/remote-link.md index 71151482b90..96f1a58f395 100644 --- a/structure/remote-link.md +++ b/structure/remote-link.md @@ -8,7 +8,7 @@ Every remote `ocx` call goes through `remoteOcxArgv`, which runs `sh -c` with a `src/link/ssh-config.ts` lists host candidates from `~/.ssh/config`. Arguments are split with the rules of OpenSSH's `argv_split`. Pattern hosts, `Match` blocks and aliases that fail the alias check produce no candidates, and only top-level `Include` directives are followed, because an include inside a `Host` or `Match` block is conditional. A candidate is an offer, not trust. -`src/link/tunnel-state.ts` is the tunnel lifecycle reducer: connecting, connected, reconnecting with capped jittered backoff, and failed for auth, host-key and forward errors or after five minutes without a connection, whether or not an attempt is in flight. `src/link/supervisor.ts` keeps a spawned tunnel in connecting until its link key authenticates a catalog request or the child remains alive for five seconds. +`src/link/tunnel-state.ts` is the tunnel lifecycle reducer: connecting, connected, reconnecting with capped jittered backoff, and failed for auth, host-key and forward errors or after five minutes without a connection, whether or not an attempt is in flight. Without a retry policy, failed is terminal; that is how `src/link/supervisor.ts` runs the Home's `-R` tunnels, and it keeps a spawned tunnel in connecting until its link key authenticates a catalog request or the child remains alive for five seconds. `CLIENT_TUNNEL_RETRY_POLICY`, which only the Child's own tunnel uses, gives failed a `retryAt`: timeout and forward failures are retried after 60 seconds and auth failures after five minutes, so no more than 12 auth attempts reach the Home's sshd in an hour; a host-key failure is never retried. A retry runs while the state still reads failed (`inFlight`), a ready event promotes it to connected, and an exit returns it to failed with the original `since` and the next `retryAt`, so a long outage keeps the one-a-minute cadence instead of restarting fast backoff. `src/link/routes.ts` holds the one route and method table for linked traffic; the hub-link listener and the client relay both decide admission from it. @@ -18,11 +18,11 @@ 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 }`. -`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. +`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. Its periodic state check 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. Normal shutdown, including recycle, stops the client supervisor before the client listener; it sends TERM, waits at most five seconds, then sends KILL. +`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. -The client tunnel pidfile is `/link/client-tunnel.pid` with `{ version: 1, linkId, pid, argv, ownerPid }` and private permissions. `reapOrphanTunnel` leaves a tunnel alone while `ownerPid` is alive and reports `tunnel: "owned"`. After the owner exits, Linux reaps only a process whose `/proc//cmdline` argv exactly matches the pidfile: TERM is followed by at most five seconds and then KILL, the pidfile is removed, and the result is `reaped`. A missing process or argv mismatch removes only the pidfile and reports `absent`; macOS and Windows leave the process and pidfile untouched and report `unresolved`. +The client tunnel pidfile is `/link/client-tunnel.pid` with `{ version: 1, linkId, pid, argv, ownerPid }` and private permissions. `reapOrphanTunnel` first checks the recorded process: on Linux its `/proc//cmdline` argv, on macOS `ps -ww -o ppid= -o args=`, compared with the recorded argv joined by spaces. A process that is gone, or provably another program (a pid reused after a reboot), leaves a stale pidfile, which is removed without signalling anything, and the result is `absent`. While `ownerPid` is alive a remaining tunnel is `owned`. After the owner exits, Linux reaps an exact argv match and macOS an exact argv match whose parent is launchd (pid 1): TERM, at most five seconds, then KILL, the pidfile is removed, and the result is `reaped`. Anything else (another parent, an unreadable `ps`, Windows) is `unresolved` and never signalled. The client supervisor adopts an `owned` or `unresolved` tunnel that is still alive: it probes through it, never signals it, and starts its own tunnel on the check after it dies. It settles the pidfile before its first spawn; when the start-up read is unreadable or does not match, the first check that reads a matching link settles it instead, and nothing spawns before that. `ocx disconnect` on the client tears down a matching sidecar link by attempting one SSH `ocx link revoke --link-id ` on Home, then disconnecting the client state and deleting the sidecar after rechecking ownership. A revoke failure still completes local disconnect and prints `Home revoke failed; run ocx link revoke --link-id on the home.`; orphan cleanup runs through the same `reapOrphanTunnel` rule before sidecar parsing, including when the sidecar is corrupt or mismatched. After `connectClient` commits during a join, a restart scheduling failure leaves the connection and sidecar intact and returns `join_restart_failed`; the operator restarts OpenCodex to finish connecting as a Child. @@ -50,6 +50,6 @@ Codex keeps the standalone loopback routing: `routingTarget` in `src/client/conn `src/client/link-ingress.ts` is the link-mode data plane of the machine listener. A `/v1/responses` WebSocket upgrade answers `426 upgrade_required`, which codex-rs maps to its HTTP fallback, and no upgrade is ever relayed. Every relayed route first passes the standalone loopback Host and Origin gate (`isAllowedRequestOrigin` in `src/server/auth-cors.ts`), so a rebinding or cross-site page gets `403 origin_rejected` and nothing is fetched upstream. `/readyz` is answered locally. The link key is read from the service token file once, when the listener starts, and held in memory; the file must hold the key whose fingerprint the connection committed. While it does not, relayed routes answer `503 link_credential_unavailable` without an upstream fetch, and the file is read again at most once a second; a valid key is never re-read. The key is never logged or returned. -`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. 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 binds the configured port or fails to start, turns the management relay off, and refuses key rotation and revocation, which belong to the hub. +`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 binds the configured port or fails to start, 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`. diff --git a/tests/clients/client-link-relay.test.ts b/tests/clients/client-link-relay.test.ts index f92b6435011..5384c8d981f 100644 --- a/tests/clients/client-link-relay.test.ts +++ b/tests/clients/client-link-relay.test.ts @@ -7,10 +7,12 @@ import { forwardLinkRequestHeaders, LINK_RELAY_BODY_MAX_BYTES, LINK_RELAY_HEADER_TIMEOUT_MS, + LINK_RELAY_HOLD_MS, LINK_RELAY_SSE_IDLE_TIMEOUT_MS, relayLinkDataRequest, sanitizeLinkResponseHeaders, type LinkRelayClock, + type LinkTunnelGate, } from "../../src/client/link-relay"; import { HUB_RELAY_REQUEST_BODY_MAX_BYTES } from "../../src/client/hub-relay"; import { startMachineListener } from "../../src/client/machine-listener"; @@ -420,3 +422,201 @@ describe("client link HTTP relay", () => { expect(received).toMatchObject({ authorization: `Bearer ${LINK_KEY}`, contentLength: null, transferEncoding: "chunked", body: '{"input":"chunked"}' }); }); }); + +/** A tunnel gate the test opens by hand. */ +function manualGate(initiallyPending = true) { + const waits: number[] = []; + const state = { pending: initiallyPending }; + let open: (connected: boolean) => void = () => {}; + const tunnel: LinkTunnelGate = { + pending: () => state.pending, + waitForConnected: timeoutMs => { + waits.push(timeoutMs); + return new Promise(resolve => { open = resolve; }); + }, + }; + return { + tunnel, + waits, + state, + release(connected: boolean) { + state.pending = false; + open(connected); + }, + }; +} + +describe("client link relay while the tunnel reconnects", () => { + test("holds a request while the tunnel reconnects and forwards it once when it connects", async () => { + const gate = manualGate(); + const sent: string[] = []; + const pending = relayLinkDataRequest(relayRequest({ + method: "POST", headers: { "Content-Type": "application/json" }, body: '{"input":"held"}', + }), target, { + tunnel: gate.tunnel, + fetchImpl: (async (_input, init) => { + sent.push(await new Response(init!.body).text()); + return Response.json({ ok: true }); + }) as typeof fetch, + }); + await Bun.sleep(5); + expect(sent).toEqual([]); + expect(gate.waits).toEqual([LINK_RELAY_HOLD_MS]); + gate.release(true); + const response = await pending; + expect(response.status).toBe(200); + expect(sent).toEqual(['{"input":"held"}']); + }); + + test("answers 503 with Retry-After when the hold ends without a tunnel", async () => { + const gate = manualGate(); + let calls = 0; + const pending = relayLinkDataRequest(relayRequest({ method: "POST", body: "{}" }), target, { + tunnel: gate.tunnel, + fetchImpl: (async () => { calls += 1; return new Response(); }) as typeof fetch, + }); + await Bun.sleep(1); + gate.release(false); + const response = await pending; + expect(response.status).toBe(503); + expect(response.headers.get("retry-after")).toBe("1"); + expect(await response.json()).toEqual({ error: "link tunnel unavailable" }); + expect(calls).toBe(0); + }); + + test("a connected tunnel adds no wait", async () => { + let waited = 0; + const tunnel: LinkTunnelGate = { pending: () => false, waitForConnected: async () => { waited += 1; return true; } }; + const response = await relayLinkDataRequest(relayRequest({ method: "POST", body: "{}" }), target, { + tunnel, + fetchImpl: (async () => Response.json({ ok: true })) as typeof fetch, + }); + expect(response.status).toBe(200); + expect(waited).toBe(0); + }); + + test("never replays a request after a reset, even while the tunnel reconnects", async () => { + let pending = false; + let calls = 0; + let waited = 0; + const tunnel: LinkTunnelGate = { pending: () => pending, waitForConnected: async () => { waited += 1; return true; } }; + const response = await relayLinkDataRequest(relayRequest({ method: "POST", body: '{"input":"once"}' }), target, { + tunnel, + fetchImpl: (async (_input, init) => { + calls += 1; + await new Response(init!.body).text(); + pending = true; + throw Object.assign(new Error("socket reset"), { code: "ECONNRESET" }); + }) as typeof fetch, + }); + expect(response.status).toBe(503); + expect(response.headers.get("retry-after")).toBe("1"); + expect(calls).toBe(1); + expect(waited).toBe(0); + }); + + test("resends a refused request once the tunnel is back, inside the same hold window", async () => { + let pending = false; + let calls = 0; + let clock = 1_000; + const waits: number[] = []; + const tunnel: LinkTunnelGate = { + pending: () => pending, + waitForConnected: async timeoutMs => { + waits.push(timeoutMs); + clock += 14_000; + pending = false; + return true; + }, + }; + const bodies: string[] = []; + const response = await relayLinkDataRequest(relayRequest({ + method: "POST", headers: { "Content-Type": "application/json" }, body: '{"input":"again"}', + }), target, { + tunnel, + now: () => clock, + fetchImpl: (async (_input, init) => { + calls += 1; + // Refused twice: the second refusal finds the 15 s window spent and gets 503. + if (calls <= 2) { + pending = true; + throw Object.assign(new Error("Unable to connect"), { code: "ConnectionRefused" }); + } + bodies.push(await new Response(init!.body).text()); + return Response.json({ ok: true }); + }) as typeof fetch, + }); + expect(waits).toEqual([LINK_RELAY_HOLD_MS, LINK_RELAY_HOLD_MS - 14_000]); + expect(calls).toBe(3); + expect(response.status).toBe(200); + expect(bodies).toEqual(['{"input":"again"}']); + + clock = 0; + calls = 0; + const spent = await relayLinkDataRequest(relayRequest({ method: "POST", body: "{}" }), target, { + tunnel: { pending: () => true, waitForConnected: async () => { clock += LINK_RELAY_HOLD_MS; return true; } }, + now: () => clock, + fetchImpl: (async () => { + calls += 1; + throw Object.assign(new Error("Unable to connect"), { code: "ConnectionRefused" }); + }) as typeof fetch, + }); + expect(spent.status).toBe(503); + expect(calls).toBe(1); + }); + + test("a refused connection to a real closed port leaves the streamed body intact for the resend", async () => { + const probe = Bun.serve({ hostname: "127.0.0.1", port: 0, fetch: () => new Response() }); + const port = probe.port!; + probe.stop(true); + let pending = false; + let home: Server | undefined; + const received: string[] = []; + try { + const response = await relayLinkDataRequest(relayRequest({ + method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ input: "x".repeat(100_000) }), + }), { tunnelPort: port, admissionKey: LINK_KEY }, { + tunnel: { + pending: () => pending, + waitForConnected: async () => { + home = Bun.serve({ hostname: "127.0.0.1", port, async fetch(req) { received.push(await req.text()); return Response.json({ ok: true }); } }); + pending = false; + return true; + }, + }, + fetchImpl: (async (input, init) => { + try { + return await fetch(input, init); + } catch (error) { + pending = true; + throw error; + } + }) as typeof fetch, + }); + expect(response.status).toBe(200); + expect(received).toEqual([JSON.stringify({ input: "x".repeat(100_000) })]); + } finally { + home?.stop(true); + } + }); + + test("the Child's listener holds a relayed request on its tunnel gate", async () => { + writeFileSync(join(root, "service-api-token"), `${LINK_KEY}\n`, { mode: 0o600 }); + const hub = Bun.serve({ hostname: "127.0.0.1", port: 0, fetch: () => Response.json({ relayed: true }) }); + servers.push(hub); + const gate = manualGate(); + const machine = startMachineListener(0, { state: linkConnection(hub.port!), linkTunnel: gate.tunnel }); + servers.push(machine); + let settled = false; + const pending = fetch(new URL("/v1/responses", machine.url), { + method: "POST", headers: { "Content-Type": "application/json" }, body: "{}", + }).then(response => { settled = true; return response; }); + await Bun.sleep(50); + expect(settled).toBe(false); + expect(gate.waits).toEqual([LINK_RELAY_HOLD_MS]); + gate.release(true); + const response = await pending; + expect(response.status).toBe(200); + expect(await response.json()).toEqual({ relayed: true }); + }); +}); diff --git a/tests/clients/client-link-status.test.ts b/tests/clients/client-link-status.test.ts index 4e497192911..741138a76cf 100644 --- a/tests/clients/client-link-status.test.ts +++ b/tests/clients/client-link-status.test.ts @@ -52,6 +52,18 @@ describe("a Child's own link status", () => { .toEqual({ ...invalid, alias: "home-mac" } as never); }); + test("reports what the keyed probe found as the reason, whatever the tunnel state", () => { + const connected: ClientLinkSupervisorStatus = { ...tunnel({ kind: "connected", since: SINCE }), probe: "unauthorized" }; + expect(projectClientLinkChild(sidecar, connected, NOW).child).toMatchObject({ state: "connected", reason: "unauthorized" }); + const retrying: ClientLinkSupervisorStatus = { + ...tunnel({ kind: "failed", since: SINCE, reason: "timeout", retryAt: NOW + 60_000, inFlight: true }), + probe: "home_unreachable", + }; + expect(projectClientLinkChild(sidecar, retrying, NOW).child).toEqual({ + alias: "home-mac", state: "failed", since: new Date(SINCE).toISOString(), reason: "home_unreachable", + }); + }); + test("a Home-initiated Child without a sidecar reports no child row", () => { expect(projectClientLinkChild(null, { kind: "stopped" }, NOW).child).toBeNull(); }); diff --git a/tests/clients/client-link-tunnel.test.ts b/tests/clients/client-link-tunnel.test.ts index 2b41415c323..9a177500ca4 100644 --- a/tests/clients/client-link-tunnel.test.ts +++ b/tests/clients/client-link-tunnel.test.ts @@ -1,13 +1,17 @@ import { expect, test } from "bun:test"; -import { existsSync, mkdirSync, readFileSync, statSync, writeFileSync, mkdtempSync, rmSync } from "node:fs"; +import { existsSync, mkdirSync, readFileSync, renameSync, statSync, writeFileSync, mkdtempSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { + CLIENT_LINK_MAX_HOLDS, clientLinkTunnelStatus, clientTunnelPidfilePath, + CONNECTION_UNREADABLE, createClientLinkSupervisor, + readWhenFileChanges, reapOrphanTunnel, spawnClientLinkTunnel, + type ClientLinkSupervisorDeps, } from "../../src/client/link-tunnel"; import type { ClientLinkState } from "../../src/client/link-state"; import type { SshChild, SshRunner } from "../../src/link/ssh-runner"; @@ -114,7 +118,7 @@ test("sends TERM and then KILL after the bounded stop wait", async () => { } }); -test("reaps only a dead owner's exact Linux tunnel and preserves non-Linux ambiguity", async () => { +test("reaps only a dead owner's exact Linux tunnel and drops a stale pidfile", async () => { const configDir = tempConfigDir(); const path = clientTunnelPidfilePath(configDir); const argv = ["ssh", "-N", "-T"]; @@ -124,9 +128,17 @@ test("reaps only a dead owner's exact Linux tunnel and preserves non-Linux ambig }; try { writePidfile(41, 42); - expect(await reapOrphanTunnel({ configDir, isAlive: pid => pid === 41, platform: "linux" })).toEqual({ tunnel: "owned" }); + expect(await reapOrphanTunnel({ + configDir, isAlive: pid => pid === 41 || pid === 42, platform: "linux", readProcessArgv: () => argv, + })).toEqual({ tunnel: "owned", pid: 42 }); expect(existsSync(path)).toBe(true); + // A live owner whose tunnel is gone (pids reused after a reboot) leaves only a stale pidfile. + expect(await reapOrphanTunnel({ + configDir, isAlive: pid => pid === 41, platform: "linux", readProcessArgv: () => null, + })).toEqual({ tunnel: "absent" }); + expect(existsSync(path)).toBe(false); + const signals: NodeJS.Signals[] = []; writePidfile(41, 42); expect(await reapOrphanTunnel({ @@ -148,15 +160,89 @@ test("reaps only a dead owner's exact Linux tunnel and preserves non-Linux ambig readProcessArgv: () => argv, })).toEqual({ tunnel: "absent" }); expect(existsSync(path)).toBe(false); + } finally { + rmSync(configDir, { recursive: true, force: true }); + } +}); - writePidfile(41, 42); - expect(await reapOrphanTunnel({ configDir, isAlive: () => false, platform: "darwin" })).toEqual({ tunnel: "unresolved", pid: 42 }); - expect(existsSync(path)).toBe(true); +test("macOS reaps an exact-argv orphan of launchd and never signals a process it cannot prove is its own", async () => { + const configDir = tempConfigDir(); + const path = clientTunnelPidfilePath(configDir); + const argv = ["ssh", "-N", "-T", "-L", "127.0.0.1:19002:127.0.0.1:19001", "--", "home.example.test"]; + const writePidfile = () => { + mkdirSync(join(configDir, "link"), { recursive: true }); + writeFileSync(path, JSON.stringify({ version: 1, linkId: sidecar().linkId, pid: 42, argv, ownerPid: 41 })); + }; + const signals: Array<[number, NodeJS.Signals]> = []; + const record = (pid: number, signal: NodeJS.Signals) => { signals.push([pid, signal]); }; + try { + // Owner and tunnel both gone: the pidfile is stale. + writePidfile(); + expect(await reapOrphanTunnel({ configDir, platform: "darwin", isAlive: () => false, signal: record })).toEqual({ tunnel: "absent" }); + expect(existsSync(path)).toBe(false); + + // The exact argv under launchd after the owner died: TERM, and it is gone. + writePidfile(); + let alive = true; + expect(await reapOrphanTunnel({ + configDir, + platform: "darwin", + isAlive: pid => pid === 42 && alive, + readProcessInfo: () => ({ ppid: 1, args: argv.join(" ") }), + signal: (pid, signal) => { record(pid, signal); alive = false; }, + sleep: async () => {}, + })).toEqual({ tunnel: "reaped" }); + expect(signals).toEqual([[42, "SIGTERM"]]); + expect(existsSync(path)).toBe(false); + + // Our argv under another parent, or a process ps cannot read: watched, never signalled. + for (const info of [{ ppid: 500, args: argv.join(" ") }, null]) { + signals.length = 0; + writePidfile(); + expect(await reapOrphanTunnel({ + configDir, platform: "darwin", isAlive: pid => pid === 42, readProcessInfo: () => info, signal: record, sleep: async () => {}, + })).toEqual({ tunnel: "unresolved", pid: 42 }); + expect(signals).toEqual([]); + expect(existsSync(path)).toBe(true); + } + + // The pid now runs another program (reused after a reboot): stale, and never signalled. + expect(await reapOrphanTunnel({ + configDir, platform: "darwin", isAlive: pid => pid === 42, + readProcessInfo: () => ({ ppid: 1, args: "/usr/libexec/some-daemon --agent" }), signal: record, + })).toEqual({ tunnel: "absent" }); + expect(signals).toEqual([]); + expect(existsSync(path)).toBe(false); } finally { rmSync(configDir, { recursive: true, force: true }); } }); +test.skipIf(process.platform !== "darwin")("the macOS process reader identifies a live child by its exact argv and parent", async () => { + const configDir = tempConfigDir(); + const path = clientTunnelPidfilePath(configDir); + const sleeper = Bun.spawn(["sleep", "30"], { stdout: "ignore", stderr: "ignore" }); + const signals: NodeJS.Signals[] = []; + const writePidfile = (argv: string[]) => { + mkdirSync(join(configDir, "link"), { recursive: true }); + // Owner 999999 does not exist, so only the identity check stands between the child and a signal. + writeFileSync(path, JSON.stringify({ version: 1, linkId: sidecar().linkId, pid: sleeper.pid, argv, ownerPid: 999_999 })); + }; + try { + // Real `ps`: the argv matches, but the parent is this test, not launchd. + writePidfile(["sleep", "30"]); + expect(await reapOrphanTunnel({ configDir, platform: "darwin", signal: (_pid, signal) => { signals.push(signal); } })) + .toEqual({ tunnel: "unresolved", pid: sleeper.pid }); + writePidfile(["sleep", "31"]); + expect(await reapOrphanTunnel({ configDir, platform: "darwin", signal: (_pid, signal) => { signals.push(signal); } })) + .toEqual({ tunnel: "absent" }); + expect(signals).toEqual([]); + } finally { + sleeper.kill(); + rmSync(configDir, { recursive: true, force: true }); + } +}); + test("supervisor starts once while connected and ends once after a missing or mismatched sidecar", async () => { const configDir = tempConfigDir(); const fake = fakeRunner(); @@ -249,3 +335,528 @@ test("projects an invalid sidecar as a failed child status", () => { rmSync(configDir, { recursive: true, force: true }); } }); + +const LINK_KEY = `ocx_data_${"f".repeat(40)}`; +const HOST_KEY_CHANGED = "Host key verification failed."; + +/** ssh children that stay up until the test ends them with an exit code and stderr. */ +function scriptedRunner() { + const children: Array<{ child: SshChild; alive: boolean; signals: NodeJS.Signals[]; exit(code: number, stderr?: string): void }> = []; + const runner: SshRunner = { + async run() { return { code: 0, stdout: "", stderr: "" }; }, + spawnTunnel(argv) { + const exited = deferred(); + let stderrText = ""; + const item = { + child: undefined as unknown as SshChild, + alive: true, + signals: [] as NodeJS.Signals[], + exit(code: number, stderr = "") { + if (!item.alive) return; + item.alive = false; + stderrText = stderr; + exited.resolve(code); + }, + }; + item.child = { + pid: 31_000 + children.length, + argv: [...argv], + exited: exited.promise, + stderr: exited.promise.then(() => stderrText), + kill(signal = "SIGTERM") { item.signals.push(signal); item.exit(143); }, + }; + children.push(item); + return item.child; + }, + }; + return { runner, children, live: () => children.filter(item => item.alive) }; +} + +/** + * A client supervisor on a scripted ssh, an injected clock and a scripted Home `/readyz`. + * `step` moves the clock and runs one supervisor tick to completion. + */ +function supervisorHarness(overrides: Partial = {}) { + const configDir = tempConfigDir(); + const ssh = scriptedRunner(); + const intervals: Array<() => void> = []; + const probes: Array<{ url: string; keyed: boolean; cache: RequestCache | undefined }> = []; + const view = { + readyz: 200 as number | "refused", + readyzBody: null as string | null, + sidecar: sidecar() as ClientLinkState | null | "unreadable", + connected: sidecar().linkId as string | null, + clock: 0, + ended: 0, + }; + const supervisor = createClientLinkSupervisor({ + configDir, + runner: ssh.runner, + readSidecar: () => { + if (view.sidecar === "unreadable") throw new Error("client-link.json is not valid JSON"); + return view.sidecar; + }, + connectedLinkId: () => view.connected, + setTimer: (callback, ms) => { + if (ms === 1_000) intervals.push(callback); + return intervals.length as unknown as ReturnType; + }, + clearTimer: () => {}, + now: () => view.clock, + random: () => 0.5, + linkKey: () => LINK_KEY, + fetchImpl: (async (input, init) => { + probes.push({ url: String(input), keyed: new Headers(init?.headers).get("x-opencodex-api-key") === LINK_KEY, cache: init?.cache }); + if (view.readyz === "refused") throw Object.assign(new Error("Unable to connect"), { code: "ConnectionRefused" }); + return new Response(view.readyzBody, { status: view.readyz }); + }) as typeof fetch, + onLinkEnded: () => { view.ended += 1; }, + ...overrides, + }); + const settle = async () => { + for (let turn = 0; turn < 3; turn += 1) await new Promise(resolve => setTimeout(resolve, 0)); + }; + return { + configDir, + ssh, + supervisor, + probes, + view, + settle, + state: () => { + const status = supervisor.status(); + return status.kind === "tunnel" ? status.state.kind : status.kind; + }, + async start() { + supervisor.start(); + await settle(); + }, + async step(ms = 1_000) { + view.clock += ms; + intervals[0]!(); + await settle(); + }, + async close() { + await supervisor.stop(); + rmSync(configDir, { recursive: true, force: true }); + }, + }; +} + +test("after a six-minute outage the tunnel is back within a minute of the Home returning", async () => { + const h = supervisorHarness(); + try { + await h.start(); + await h.step(); + expect(h.state()).toBe("connected"); + + h.view.readyz = "refused"; + h.ssh.children[0]!.exit(255, "ssh: connect to host home port 22: Connection refused"); + let spawnsInLastMinute = 0; + for (let elapsed = 1_000; elapsed <= 6 * 60_000; elapsed += 1_000) { + for (const item of h.ssh.live()) item.exit(255, "ssh: connect to host home port 22: Operation timed out"); + const before = h.ssh.children.length; + await h.step(); + if (elapsed > 5 * 60_000) spawnsInLastMinute += h.ssh.children.length - before; + } + expect(h.state()).toBe("failed"); + // Past the five-minute window the retries slow to about one a minute. + expect(spawnsInLastMinute).toBeLessThanOrEqual(1); + + h.view.readyz = 200; + let reconnectedAfter = -1; + for (let elapsed = 1_000; elapsed <= 90_000; elapsed += 1_000) { + await h.step(); + if (h.state() === "connected") { + reconnectedAfter = elapsed; + break; + } + } + expect(reconnectedAfter).toBeGreaterThan(0); + expect(reconnectedAfter).toBeLessThanOrEqual(62_000); + } finally { + await h.close(); + } +}, 30_000); + +test("a reconnect is promoted by a keyed readyz and the healthy ssh is never cut afterwards", async () => { + const h = supervisorHarness(); + try { + await h.start(); + await h.step(); + expect(h.state()).toBe("connected"); + h.ssh.children[0]!.exit(255, "client_loop: send disconnect: Broken pipe"); + await h.settle(); + expect(h.state()).toBe("reconnecting"); + await h.step(); + expect(h.ssh.children).toHaveLength(2); + await h.step(); + expect(h.state()).toBe("connected"); + const probesWhenConnected = h.probes.length; + + for (let minute = 1; minute <= 30; minute += 1) { + for (let second = 0; second < 60; second += 5) await h.step(5_000); + if (minute === 5 || minute === 30) { + expect(h.state()).toBe("connected"); + expect(h.ssh.children[1]!.signals).toEqual([]); + expect(h.ssh.children[1]!.alive).toBe(true); + } + } + expect(h.ssh.children).toHaveLength(2); + // While connected the only background cost is one keyed probe every 30 seconds. + expect(h.probes.length - probesWhenConnected).toBeLessThanOrEqual(60); + expect(h.probes.every(probe => probe.keyed && probe.cache === "no-store" && probe.url === "http://127.0.0.1:19002/readyz")).toBe(true); + } finally { + await h.close(); + } +}, 30_000); + +test("the keyed probe reports an unauthorized key or an unreachable Home for display only", async () => { + const h = supervisorHarness(); + try { + await h.start(); + await h.step(); + expect(h.supervisor.status()).toMatchObject({ kind: "tunnel", state: { kind: "connected" } }); + h.view.readyz = 401; + await h.step(30_000); + expect(h.supervisor.status()).toMatchObject({ state: { kind: "connected" }, probe: "unauthorized" }); + h.view.readyz = 503; + await h.step(30_000); + expect(h.supervisor.status()).toMatchObject({ state: { kind: "connected" }, probe: "home_unreachable" }); + h.view.readyz = 200; + await h.step(30_000); + const status = h.supervisor.status(); + expect(status).toMatchObject({ state: { kind: "connected" } }); + expect("probe" in status).toBe(false); + expect(h.ssh.children[0]!.signals).toEqual([]); + expect(h.view.ended).toBe(0); + } finally { + await h.close(); + } +}); + +test("a Home that admits the key but reports its own readiness failed still carries the link, and its ssh is never cut", async () => { + const h = supervisorHarness(); + try { + // The Home's link listener answers 401 before /readyz, so its 503 readiness body proves the + // forward and the key; only the Home's own start-up sync failed, which relayed requests ignore. + h.view.readyz = 503; + h.view.readyzBody = JSON.stringify({ service: "opencodex", status: "failed" }); + await h.start(); + const held = h.supervisor.waitForConnected(15_000); + await h.step(); + expect(h.state()).toBe("connected"); + expect(await held).toBe(true); + expect(h.supervisor.status()).toMatchObject({ state: { kind: "connected" }, probe: "home_not_ready" }); + for (let elapsed = 0; elapsed < 12 * 60_000; elapsed += 5_000) await h.step(5_000); + expect(h.state()).toBe("connected"); + expect(h.ssh.children).toHaveLength(1); + expect(h.ssh.children[0]!.signals).toEqual([]); + + // A 503 that is not the Home's readiness answer proves nothing, so a reconnect waits for one. + h.view.readyzBody = "Service Unavailable"; + h.ssh.children[0]!.exit(255, "Connection reset by peer"); + await h.settle(); + for (let tick = 0; tick < 10; tick += 1) await h.step(); + expect(h.ssh.children).toHaveLength(2); + expect(h.supervisor.status()).toMatchObject({ state: { kind: "reconnecting" }, probe: "home_unreachable" }); + h.view.readyzBody = JSON.stringify({ service: "opencodex", status: "pending" }); + for (let tick = 0; tick < 6 && h.state() !== "connected"; tick += 1) await h.step(); + expect(h.state()).toBe("connected"); + expect(h.ssh.children[1]!.signals).toEqual([]); + } finally { + await h.close(); + } +}, 30_000); + +test("a probe that hangs never delays noticing a disconnect, and stop() aborts it instead of waiting", async () => { + let hang = false; + const probeSignals: AbortSignal[] = []; + const fetchImpl = (async (_input: string | URL | Request, init?: RequestInit) => { + if (!hang) return new Response(null, { status: 200 }); + const signal = init!.signal!; + probeSignals.push(signal); + return await new Promise((_resolve, reject) => { + signal.addEventListener("abort", () => reject(signal.reason), { once: true }); + }); + }) as typeof fetch; + const race = (work: Promise) => + Promise.race([work.then(() => "done"), new Promise(resolve => setTimeout(() => resolve("waiting"), 500))]); + + const h = supervisorHarness({ fetchImpl }); + try { + await h.start(); + await h.step(); + expect(h.state()).toBe("connected"); + hang = true; + await h.step(30_000); + expect(probeSignals).toHaveLength(1); + h.view.connected = null; + await h.step(); + expect(h.view.ended).toBe(1); + expect(h.ssh.children[0]!.signals).toContain("SIGTERM"); + expect(probeSignals[0]!.aborted).toBe(true); + } finally { + expect(await race(h.close())).toBe("done"); + } + + const stopping = supervisorHarness({ fetchImpl }); + try { + await stopping.start(); + await stopping.step(); + expect(probeSignals).toHaveLength(2); + expect(await race(stopping.supervisor.stop())).toBe("done"); + expect(probeSignals[1]!.aborted).toBe(true); + expect(stopping.ssh.children[0]!.signals).toContain("SIGTERM"); + } finally { + await stopping.close(); + } +}); + +test("while a request is held the tunnel is probed on every check, and otherwise it backs off", async () => { + const h = supervisorHarness(); + try { + h.view.readyz = "refused"; + await h.start(); + for (let tick = 0; tick < 5; tick += 1) await h.step(); + // Nothing held: probes at one, two and four seconds, then the next is due at eight. + expect(h.probes).toHaveLength(3); + const held = h.supervisor.waitForConnected(15_000); + await h.step(); + expect(h.probes).toHaveLength(4); + h.view.readyz = 200; + await h.step(); + expect(await held).toBe(true); + expect(h.state()).toBe("connected"); + const probesWhenConnected = h.probes.length; + for (let tick = 0; tick < 29; tick += 1) await h.step(); + expect(h.probes.length).toBe(probesWhenConnected); + } finally { + await h.close(); + } +}); + +test("a start-up read that cannot be used still reaps a leftover tunnel before the first spawn", async () => { + const signals: Array<[number, NodeJS.Signals]> = []; + let orphanAlive = true; + const argv = buildTunnelArgv({ alias: "home.example.test", direction: "L", bindPort: 19002, targetPort: 19001, knownHostsFile: "/tmp/ocx-known-hosts" }); + const h = supervisorHarness({ + platform: "linux", + readProcessArgv: pid => (pid === 42 && orphanAlive ? argv : null), + isAlive: pid => pid === 42 && orphanAlive, + signal: (pid, signal) => { signals.push([pid, signal]); orphanAlive = false; }, + sleep: async () => {}, + }); + mkdirSync(join(h.configDir, "link"), { recursive: true }); + writeFileSync(clientTunnelPidfilePath(h.configDir), JSON.stringify({ version: 1, linkId: sidecar().linkId, pid: 42, argv, ownerPid: 41 })); + try { + h.view.connected = CONNECTION_UNREADABLE as never; + await h.start(); + expect(h.ssh.children).toHaveLength(0); + h.view.connected = sidecar().linkId; + await h.step(); + expect(signals).toEqual([[42, "SIGTERM"]]); + expect(h.ssh.children).toHaveLength(1); + await h.step(); + expect(h.state()).toBe("connected"); + } finally { + await h.close(); + } +}); + +test("a macOS pidfile whose owner and ssh are both dead gives way to a new tunnel at once", async () => { + const h = supervisorHarness({ platform: "darwin", isAlive: () => false }); + mkdirSync(join(h.configDir, "link"), { recursive: true }); + writeFileSync(clientTunnelPidfilePath(h.configDir), JSON.stringify({ version: 1, linkId: sidecar().linkId, pid: 42, argv: ["ssh", "-N"], ownerPid: 41 })); + try { + await h.start(); + expect(h.ssh.children).toHaveLength(1); + expect(h.state()).toBe("connecting"); + await h.step(); + expect(h.state()).toBe("connected"); + } finally { + await h.close(); + } +}); + +test("a macOS orphan with the exact argv under launchd is reaped and replaced", async () => { + const signals: Array<[number, NodeJS.Signals]> = []; + let orphanAlive = true; + const argv = buildTunnelArgv({ alias: "home.example.test", direction: "L", bindPort: 19002, targetPort: 19001, knownHostsFile: "/tmp/ocx-known-hosts" }); + const h = supervisorHarness({ + platform: "darwin", + isAlive: pid => pid === 42 && orphanAlive, + readProcessInfo: pid => (pid === 42 ? { ppid: 1, args: argv.join(" ") } : null), + signal: (pid, signal) => { signals.push([pid, signal]); orphanAlive = false; }, + sleep: async () => {}, + }); + mkdirSync(join(h.configDir, "link"), { recursive: true }); + writeFileSync(clientTunnelPidfilePath(h.configDir), JSON.stringify({ version: 1, linkId: sidecar().linkId, pid: 42, argv, ownerPid: 41 })); + try { + await h.start(); + expect(signals).toEqual([[42, "SIGTERM"]]); + expect(h.ssh.children).toHaveLength(1); + } finally { + await h.close(); + } +}); + +test("a macOS tunnel that cannot be verified is watched, never signalled, and replaced once it dies", async () => { + const signals: Array<[number, NodeJS.Signals]> = []; + let orphanAlive = true; + const h = supervisorHarness({ + platform: "darwin", + isAlive: pid => pid === 42 && orphanAlive, + readProcessInfo: () => null, + signal: (pid, signal) => { signals.push([pid, signal]); }, + }); + mkdirSync(join(h.configDir, "link"), { recursive: true }); + writeFileSync(clientTunnelPidfilePath(h.configDir), JSON.stringify({ version: 1, linkId: sidecar().linkId, pid: 42, argv: ["ssh", "-N"], ownerPid: 41 })); + try { + h.view.readyz = "refused"; + await h.start(); + expect(h.ssh.children).toHaveLength(0); + expect(h.supervisor.status()).toMatchObject({ kind: "tunnel", state: { kind: "connecting" }, pid: 42 }); + // Past the five-minute window the link reads failed, but the adopted tunnel is still probed. + for (let elapsed = 0; elapsed < 6 * 60_000; elapsed += 5_000) await h.step(5_000); + expect(h.state()).toBe("failed"); + expect(h.ssh.children).toHaveLength(0); + h.view.readyz = 200; + for (let elapsed = 0; elapsed < 31_000 && h.state() !== "connected"; elapsed += 1_000) await h.step(); + expect(h.state()).toBe("connected"); + expect(h.ssh.children).toHaveLength(0); + orphanAlive = false; + await h.step(); + expect(h.ssh.children).toHaveLength(1); + await h.step(); + expect(h.state()).toBe("connected"); + expect(signals).toEqual([]); + } finally { + await h.close(); + } +}); + +test("an unreadable read does not end the link; three in a row or an explicit disconnect do, once", async () => { + const h = supervisorHarness(); + try { + await h.start(); + await h.step(); + h.view.connected = CONNECTION_UNREADABLE as never; + await h.step(); + await h.step(); + expect(h.view.ended).toBe(0); + expect(h.ssh.children[0]!.signals).toEqual([]); + h.view.connected = sidecar().linkId; + await h.step(); + h.view.sidecar = "unreadable"; + await h.step(); + await h.step(); + expect(h.ssh.children[0]!.signals).toEqual([]); + expect(h.supervisor.status()).toMatchObject({ kind: "tunnel", state: { kind: "connected" } }); + h.view.sidecar = sidecar(); + await h.step(); + h.view.connected = CONNECTION_UNREADABLE as never; + for (let tick = 0; tick < 3; tick += 1) await h.step(); + expect(h.view.ended).toBe(1); + expect(h.ssh.children[0]!.signals).toContain("SIGTERM"); + } finally { + await h.close(); + } + + const disconnected = supervisorHarness(); + try { + await disconnected.start(); + await disconnected.step(); + disconnected.view.connected = null; + await disconnected.step(); + expect(disconnected.view.ended).toBe(1); + await disconnected.step(); + await disconnected.step(); + expect(disconnected.view.ended).toBe(1); + } finally { + await disconnected.close(); + } +}); + +test("requests wait on a reconnect only while it lasts and are released when it connects or fails", async () => { + const h = supervisorHarness(); + try { + expect(h.supervisor.pending()).toBe(false); + await h.start(); + expect(h.supervisor.pending()).toBe(true); + await h.step(); + expect(h.supervisor.pending()).toBe(false); + expect(await h.supervisor.waitForConnected(15_000)).toBe(true); + + h.view.readyz = "refused"; + h.ssh.children[0]!.exit(255, "Connection reset by peer"); + await h.settle(); + expect(h.supervisor.pending()).toBe(true); + let released: boolean | undefined; + const waiting = h.supervisor.waitForConnected(15_000).then(value => { released = value; return value; }); + await h.step(); + expect(h.ssh.children).toHaveLength(2); + expect(released).toBeUndefined(); + h.view.readyz = 200; + await h.step(); + expect(await waiting).toBe(true); + + h.ssh.children[1]!.exit(255, "Connection reset by peer"); + await h.settle(); + const failing = h.supervisor.waitForConnected(15_000); + await h.step(); + h.ssh.children[2]!.exit(255, HOST_KEY_CHANGED); + await h.settle(); + expect(await failing).toBe(false); + expect(h.supervisor.status()).toMatchObject({ state: { kind: "failed", reason: "hostkey" } }); + expect(h.supervisor.pending()).toBe(false); + expect(await h.supervisor.waitForConnected(15_000)).toBe(false); + } finally { + await h.close(); + } +}); + +test("the number of requests waiting on a reconnect is capped", async () => { + const h = supervisorHarness(); + try { + h.view.readyz = "refused"; + await h.start(); + const held = Array.from({ length: CLIENT_LINK_MAX_HOLDS }, () => h.supervisor.waitForConnected(15_000)); + expect(await h.supervisor.waitForConnected(15_000)).toBe(false); + h.view.readyz = 200; + await h.step(); + expect(await Promise.all(held)).toEqual(Array.from({ length: CLIENT_LINK_MAX_HOLDS }, () => true)); + } finally { + await h.close(); + } +}); + +test("the supervisor's own reads parse a file again only after it changed", () => { + const dir = tempConfigDir(); + const path = join(dir, "config.json"); + try { + writeFileSync(path, "one"); + let reads = 0; + const read = readWhenFileChanges(() => path, () => { reads += 1; return readFileSync(path, "utf8"); }); + expect(read()).toBe("one"); + expect(read()).toBe("one"); + expect(reads).toBe(1); + writeFileSync(`${path}.next`, "two"); + renameSync(`${path}.next`, path); + expect(read()).toBe("two"); + expect(reads).toBe(2); + + let refused = 0; + const unreadable = readWhenFileChanges(() => path, () => { refused += 1; return CONNECTION_UNREADABLE; }, value => value !== CONNECTION_UNREADABLE); + unreadable(); + unreadable(); + expect(refused).toBe(2); + let thrown = 0; + const throwing = readWhenFileChanges(() => path, () => { thrown += 1; throw new Error("bad"); }); + expect(() => throwing()).toThrow("bad"); + expect(() => throwing()).toThrow("bad"); + expect(thrown).toBe(2); + } finally { + rmSync(dir, { recursive: true, force: true }); + } +}); diff --git a/tests/clients/link-supervisor.test.ts b/tests/clients/link-supervisor.test.ts index 19b3d80317b..3a2c7eb41af 100644 --- a/tests/clients/link-supervisor.test.ts +++ b/tests/clients/link-supervisor.test.ts @@ -141,6 +141,39 @@ test("auth and host ownership failures do not retry, and client links stay clien await supervisor.stop(); }); +test("the Home's -R supervisor keeps the old policy: a failed tunnel is not retried and a live one is promoted after five seconds", async () => { + const fake = fakeRunner(); + const timers: Array<() => void> = []; + let current = 0; + const store = baseStore(record("hub-initiated", "lnk_0123456789abcdef")); + const supervisor = createLinkSupervisor({ + readStore: () => store, + runner: fake.runner, + pidfileDir: "/tmp/opencodex-link-supervisor-test", + now: () => current, + setTimer: callback => { timers.push(callback); return 1 as unknown as ReturnType; }, + clearTimer: () => {}, + random: () => 0.5, + }); + supervisor.start(); + current = 5_000; + timers[0]!(); + expect(supervisor.status()[0]!.state).toEqual({ kind: "connected", since: 5_000 }); + fake.children[0]!.stderr = "Error: remote port forwarding failed for listen port 19002"; + fake.children[0]!.resolve(255); + await Promise.resolve(); + await Promise.resolve(); + const failed = supervisor.status()[0]!.state; + expect(failed).toEqual({ kind: "failed", since: 5_000, reason: "forward" }); + for (let second = 0; second < 2 * 60 * 60; second += 10) { + current += 10_000; + timers[0]!(); + } + expect(fake.children).toHaveLength(1); + expect(supervisor.status()[0]!.state).toEqual(failed); + await supervisor.stop(); +}); + test("reaps a pidfile only after an exact Linux argv match", () => { const fake = fakeRunner(); const store = baseStore(record("hub-initiated", "lnk_0123456789abcdef")); diff --git a/tests/clients/link-tunnel-state.test.ts b/tests/clients/link-tunnel-state.test.ts index dd2844e7ebc..dd2a3d3db0e 100644 --- a/tests/clients/link-tunnel-state.test.ts +++ b/tests/clients/link-tunnel-state.test.ts @@ -1,6 +1,7 @@ import { expect, test } from "bun:test"; import { classifySshStderr, + CLIENT_TUNNEL_RETRY_POLICY, dueForSpawn, FAILED_AFTER_MS, IDLE, @@ -9,6 +10,9 @@ import { type TunnelState, } from "../../src/link/tunnel-state"; +const POLICY = CLIENT_TUNNEL_RETRY_POLICY; +const HOUR_MS = 60 * 60_000; + test("tunnel lifecycle reaches connected from idle", () => { let state = IDLE; state = reduceTunnel(state, { type: "spawn", now: 100 }); @@ -88,3 +92,67 @@ test("ssh stderr is classified by its retry policy", () => { ] as const; for (const [stderr, expected] of cases) expect(classifySshStderr(stderr)).toBe(expected); }); + +test("the client policy retries timeout and forward after a minute and auth after five minutes, never a host key", () => { + const forward = reduceTunnel({ kind: "connected", since: 0 }, { type: "exit", now: 1_000, stderrClass: "forward" }, () => 0.5, POLICY); + expect(forward).toEqual({ kind: "failed", since: 1_000, reason: "forward", retryAt: 61_000, inFlight: false }); + expect(dueForSpawn(forward, 60_999)).toBe(false); + expect(dueForSpawn(forward, 61_000)).toBe(true); + + const reconnecting = reduceTunnel({ kind: "connected", since: 0 }, { type: "exit", now: 0, stderrClass: "network" }, () => 0.5, POLICY); + const timeout = reduceTunnel(reconnecting, { type: "tick", now: FAILED_AFTER_MS }, () => 0.5, POLICY); + expect(timeout).toEqual({ kind: "failed", since: FAILED_AFTER_MS, reason: "timeout", retryAt: FAILED_AFTER_MS + 60_000, inFlight: false }); + expect(dueForSpawn(timeout, FAILED_AFTER_MS + 60_000)).toBe(true); + + const auth = reduceTunnel({ kind: "connecting", since: 0 }, { type: "exit", now: 2_000, stderrClass: "auth" }, () => 0.5, POLICY); + expect(auth).toMatchObject({ kind: "failed", reason: "auth", retryAt: 302_000 }); + expect(dueForSpawn(auth, 301_999)).toBe(false); + expect(dueForSpawn(auth, 302_000)).toBe(true); + + const hostkey = reduceTunnel({ kind: "connected", since: 0 }, { type: "exit", now: 3_000, stderrClass: "hostkey" }, () => 0.5, POLICY); + expect(hostkey).toEqual({ kind: "failed", since: 3_000, reason: "hostkey" }); + expect(dueForSpawn(hostkey, 3_000 + 24 * HOUR_MS)).toBe(false); +}); + +test("a retry runs while the link still reads failed, connects on ready and falls back to the slow cadence", () => { + const failed: TunnelState = { kind: "failed", since: 100, reason: "timeout", retryAt: 60_100, inFlight: false }; + let state = reduceTunnel(failed, { type: "spawn", now: 60_100 }, () => 0.5, POLICY); + expect(state).toEqual({ ...failed, inFlight: true }); + expect(dueForSpawn(state, 10 * HOUR_MS)).toBe(false); + // A transient exit of the retry keeps the original failure time and waits another minute. + state = reduceTunnel(state, { type: "exit", now: 61_000, stderrClass: "network" }, () => 0.5, POLICY); + expect(state).toEqual({ kind: "failed", since: 100, reason: "timeout", retryAt: 121_000, inFlight: false }); + // A hung retry is not timed out by the reducer; ssh's own keepalive bounds it. + state = reduceTunnel(state, { type: "spawn", now: 121_000 }, () => 0.5, POLICY); + expect(reduceTunnel(state, { type: "tick", now: 121_000 + 10 * FAILED_AFTER_MS }, () => 0.5, POLICY)).toBe(state); + expect(reduceTunnel(state, { type: "ready", now: 122_000 }, () => 0.5, POLICY)).toEqual({ kind: "connected", since: 122_000 }); + // A retry that meets a changed host key stops retrying. + expect(reduceTunnel(state, { type: "exit", now: 122_000, stderrClass: "hostkey" }, () => 0.5, POLICY)) + .toEqual({ kind: "failed", since: 100, reason: "hostkey" }); +}); + +test("auth retries reach the Home's sshd at most 12 times in any hour", () => { + let state: TunnelState = reduceTunnel({ kind: "connecting", since: 0 }, { type: "exit", now: 0, stderrClass: "auth" }, () => 0.5, POLICY); + const attempts: number[] = []; + for (let now = 0; now <= 3 * HOUR_MS; now += 1_000) { + if (!dueForSpawn(state, now)) continue; + state = reduceTunnel(state, { type: "spawn", now }, () => 0.5, POLICY); + attempts.push(now); + state = reduceTunnel(state, { type: "exit", now: now + 50, stderrClass: "auth" }, () => 0.5, POLICY); + } + expect(attempts.length).toBeGreaterThan(30); + for (const start of attempts) { + expect(attempts.filter(at => at >= start && at < start + HOUR_MS).length).toBeLessThanOrEqual(12); + } +}); + +test("without a retry policy (the Home's -R supervisor) a failed tunnel is never due", () => { + for (const stderrClass of ["auth", "forward"] as const) { + const state = reduceTunnel({ kind: "connected", since: 0 }, { type: "exit", now: 5, stderrClass }); + expect(state).toEqual({ kind: "failed", since: 5, reason: stderrClass }); + expect(dueForSpawn(state, 24 * HOUR_MS)).toBe(false); + } + const timeout = reduceTunnel({ kind: "connecting", since: 0 }, { type: "tick", now: FAILED_AFTER_MS }); + expect(timeout).toEqual({ kind: "failed", since: FAILED_AFTER_MS, reason: "timeout" }); + expect(dueForSpawn(timeout, 24 * HOUR_MS)).toBe(false); +}); diff --git a/tests/server/link-join-route.test.ts b/tests/server/link-join-route.test.ts index 4f432eff68a..653e749bf37 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 { chooseJoinTunnelPort, ClientLinkJoinError, joinHome, type ClientLinkJoinDeps } from "../../src/client/link-join"; +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"; import type { SshRunner } from "../../src/link/ssh-runner"; @@ -288,6 +289,35 @@ describe("client initiated link join", () => { expect(cleared).toBe(true); }); + test("picks the join tunnel port from 20000-29999, outside the OS ephemeral ranges", async () => { + const calls: string[][] = []; + let sidecarPort = 0; + await joinHome(joinDeps({ + runner: runnerFor(calls), + choosePort: undefined, + writeState: state => { sidecarPort = state.tunnelPort; }, + spawnTunnel: () => tunnelFor([]), + fetchImpl: async () => new Response(null, { status: 200 }), + connect: (async () => {}) as typeof import("../../src/client/connect").connectClient, + scheduleRestart: () => {}, + }), { alias: "home" }); + const issue = calls.find(isWrappedIssue)?.at(-1) ?? ""; + const port = Number(/'--tunnel-port' '(\d+)'/.exec(issue)?.[1]); + expect(port).toBe(sidecarPort); + expect(port).toBeGreaterThanOrEqual(JOIN_TUNNEL_PORT_MIN); + expect(port).toBeLessThanOrEqual(JOIN_TUNNEL_PORT_MAX); + expect(isLinkPort(port)).toBe(true); + + const tried: number[] = []; + const picked = await chooseJoinTunnelPort({ + isAvailable: async candidate => { tried.push(candidate); return tried.length === 3; }, + random: () => 0.999_999_9, + }); + expect(picked).toBe(JOIN_TUNNEL_PORT_MAX); + expect(tried).toHaveLength(3); + await expect(chooseJoinTunnelPort({ isAvailable: async () => false })).rejects.toThrow("join tunnel range"); + }); + test("rolls back on readiness timeout and admission rejection", async () => { for (const readiness of ["timeout", "unauthorized"] as const) { const calls: string[][] = [];