diff --git a/.github/workflows/macos-e2e.yaml b/.github/workflows/macos-e2e.yaml index e78a0e25271..56f7f1e8931 100644 --- a/.github/workflows/macos-e2e.yaml +++ b/.github/workflows/macos-e2e.yaml @@ -70,6 +70,12 @@ jobs: npm ci --ignore-scripts npm run build + - name: Run gateway lifecycle regressions + run: >- + npx vitest run --project integration + test/tunnel-gateway-port-release-runtime.test.ts + test/onboard-gateway-prelaunch-cutover.test.ts + - name: Detect Docker availability id: docker run: | diff --git a/ci/platform-matrix.json b/ci/platform-matrix.json index 67f47445ee6..c133f0b098e 100644 --- a/ci/platform-matrix.json +++ b/ci/platform-matrix.json @@ -93,7 +93,7 @@ "name": "Other OpenAI-compatible endpoint", "status": "caveated", "endpoint_type": "Custom OpenAI-compatible", - "notes": "Adapter path validated against OpenRouter as the `compatible-endpoint` provider with `openrouter/auto` (see `src/lib/inference/config.test.ts:119`); the onboarding prompt that surfaces OpenRouter as the worked example is at `src/lib/onboard.ts:3585`. Behavior on other OpenAI-compatible proxies, gateways, and self-hosted implementations may vary; this row claims the adapter, not the universe of compatible endpoints." + "notes": "Adapter path validated against OpenRouter as the `compatible-endpoint` provider with `openrouter/auto` (see `src/lib/inference/config.test.ts:119`); the onboarding prompt that surfaces OpenRouter as the worked example is in `handleRemoteProviderSelection` in `src/lib/onboard.ts`. Behavior on other OpenAI-compatible proxies, gateways, and self-hosted implementations may vary; this row claims the adapter, not the universe of compatible endpoints." }, { "name": "Anthropic", @@ -218,12 +218,12 @@ { "name": "Podman / other container runtimes", "status": "unsupported", - "notes": "Onboard surfaces an explicit unsupported-runtime error for Podman (`src/lib/onboard/fatal-runtime-preflight.ts:50` prints the rejection; `src/lib/onboard/preflight.ts:676` flags the unsupported runtime upstream). Only Docker Engine, Docker Desktop, and Colima are supported. See issue #420 (closed)." + "notes": "Onboard surfaces an explicit unsupported-runtime error for Podman (`src/lib/onboard/fatal-runtime-preflight.ts:50` prints the rejection; `src/lib/onboard/preflight.ts:677` flags the unsupported runtime upstream). Only Docker Engine, Docker Desktop, and Colima are supported. See issue #420 (closed)." }, { "name": "Intel Mac (macOS x86_64)", "status": "unsupported", - "notes": "OpenShell does not publish macOS x86_64 standalone gateway assets. Install hard-fails on x86_64 macOS (`scripts/install-openshell.sh:654`). See issue #954 (closed)." + "notes": "OpenShell does not publish macOS x86_64 standalone gateway assets. Install hard-fails on x86_64 macOS (`scripts/install-openshell.sh:663`). See issue #954 (closed)." }, { "name": "Non-Ubuntu/Debian Linux distros", diff --git a/docs/inference/inference-options.mdx b/docs/inference/inference-options.mdx index 6a66873931b..421c1556c4d 100644 --- a/docs/inference/inference-options.mdx +++ b/docs/inference/inference-options.mdx @@ -43,7 +43,7 @@ NemoClaw uses provider-specific local tokens for those routes, and rebuilds of l |----------|--------|---------------|-------| | NVIDIA Endpoints | Tested | OpenAI-compatible | Hosted models on integrate.api.nvidia.com | | OpenAI | Tested | Native OpenAI-compatible | Uses OpenAI model IDs | -| Other OpenAI-compatible endpoint | Tested with limitations | Custom OpenAI-compatible | Adapter path validated against OpenRouter as the `compatible-endpoint` provider with `openrouter/auto` (see `src/lib/inference/config.test.ts:119`); the onboarding prompt that surfaces OpenRouter as the worked example is at `src/lib/onboard.ts:3585`. Behavior on other OpenAI-compatible proxies, gateways, and self-hosted implementations may vary; this row claims the adapter, not the universe of compatible endpoints. | +| Other OpenAI-compatible endpoint | Tested with limitations | Custom OpenAI-compatible | Adapter path validated against OpenRouter as the `compatible-endpoint` provider with `openrouter/auto` (see `src/lib/inference/config.test.ts:119`); the onboarding prompt that surfaces OpenRouter as the worked example is in `handleRemoteProviderSelection` in `src/lib/onboard.ts`. Behavior on other OpenAI-compatible proxies, gateways, and self-hosted implementations may vary; this row claims the adapter, not the universe of compatible endpoints. | | Anthropic | Tested | Native Anthropic | Uses anthropic-messages | | Other Anthropic-compatible endpoint | Tested with limitations | Custom Anthropic-compatible | Adapter path validated with AWS Bedrock (`src/lib/onboard/bedrock-runtime.ts`). Behavior on other Anthropic-compatible proxies and gateways may vary; this row claims the adapter, not the universe of compatible endpoints. | | Google Gemini | Tested | OpenAI-compatible | Uses Google's OpenAI-compatible endpoint | diff --git a/docs/reference/commands-nemohermes.mdx b/docs/reference/commands-nemohermes.mdx index c7c9d81e7bc..4c211a310dc 100644 --- a/docs/reference/commands-nemohermes.mdx +++ b/docs/reference/commands-nemohermes.mdx @@ -1710,7 +1710,9 @@ Use `nemohermes channels stop ` when you only want to pause one nemohermes tunnel stop ``` -`nemohermes stop` remains as a deprecated alias that prints a warning and delegates to `tunnel stop`. +`nemohermes stop` remains as a deprecated legacy full stop. +It stops the tunnel services and also releases the managed host gateway port. +Use `nemohermes tunnel stop` when the shared gateway should remain available. ### `nemohermes tunnel status` @@ -1733,10 +1735,12 @@ This command remains as a compatibility alias to `nemohermes tunnel start`. ### `nemohermes stop` -Deprecated. Use `nemohermes tunnel stop` instead. +Deprecated legacy full stop. +Use `nemohermes tunnel stop` when the shared gateway should remain available. -This command remains as a compatibility alias to `nemohermes tunnel stop`. +This command stops tunnel services and also releases the managed host gateway port. +It is retained for compatibility with full-stop automation; unlike `nemohermes tunnel stop`, it intentionally tears down that host gateway. ### `nemohermes status` diff --git a/docs/reference/commands.mdx b/docs/reference/commands.mdx index 4432234ac24..1705e9b7f8e 100644 --- a/docs/reference/commands.mdx +++ b/docs/reference/commands.mdx @@ -2121,7 +2121,9 @@ Use `$$nemoclaw channels stop ` when you only want to pause one $$nemoclaw tunnel stop ``` -`$$nemoclaw stop` remains as a deprecated alias that prints a warning and delegates to `tunnel stop`. +`$$nemoclaw stop` remains as a deprecated legacy full stop. +It stops the tunnel services and also releases the managed host gateway port. +Use `$$nemoclaw tunnel stop` when the shared gateway should remain available. ### `$$nemoclaw tunnel status` @@ -2144,10 +2146,12 @@ This command remains as a compatibility alias to `$$nemoclaw tunnel start`. ### `$$nemoclaw stop` -Deprecated. Use `$$nemoclaw tunnel stop` instead. +Deprecated legacy full stop. +Use `$$nemoclaw tunnel stop` when the shared gateway should remain available. -This command remains as a compatibility alias to `$$nemoclaw tunnel stop`. +This command stops tunnel services and also releases the managed host gateway port. +It is retained for compatibility with full-stop automation; unlike `$$nemoclaw tunnel stop`, it intentionally tears down that host gateway. ### `$$nemoclaw status` diff --git a/docs/reference/platform-support.mdx b/docs/reference/platform-support.mdx index f4bc00f77fa..d9053575804 100644 --- a/docs/reference/platform-support.mdx +++ b/docs/reference/platform-support.mdx @@ -95,7 +95,7 @@ NemoClaw routes inference through the OpenShell gateway. Each row below is a pro |----------|--------|---------------|-------| | NVIDIA Endpoints | Tested | OpenAI-compatible | Hosted models on integrate.api.nvidia.com | | OpenAI | Tested | Native OpenAI-compatible | Uses OpenAI model IDs | -| Other OpenAI-compatible endpoint | Tested with limitations | Custom OpenAI-compatible | Adapter path validated against OpenRouter as the `compatible-endpoint` provider with `openrouter/auto` (see `src/lib/inference/config.test.ts:119`); the onboarding prompt that surfaces OpenRouter as the worked example is at `src/lib/onboard.ts:3585`. Behavior on other OpenAI-compatible proxies, gateways, and self-hosted implementations may vary; this row claims the adapter, not the universe of compatible endpoints. | +| Other OpenAI-compatible endpoint | Tested with limitations | Custom OpenAI-compatible | Adapter path validated against OpenRouter as the `compatible-endpoint` provider with `openrouter/auto` (see `src/lib/inference/config.test.ts:119`); the onboarding prompt that surfaces OpenRouter as the worked example is in `handleRemoteProviderSelection` in `src/lib/onboard.ts`. Behavior on other OpenAI-compatible proxies, gateways, and self-hosted implementations may vary; this row claims the adapter, not the universe of compatible endpoints. | | Anthropic | Tested | Native Anthropic | Uses anthropic-messages | | Other Anthropic-compatible endpoint | Tested with limitations | Custom Anthropic-compatible | Adapter path validated with AWS Bedrock (`src/lib/onboard/bedrock-runtime.ts`). Behavior on other Anthropic-compatible proxies and gateways may vary; this row claims the adapter, not the universe of compatible endpoints. | | Google Gemini | Tested | OpenAI-compatible | Uses Google's OpenAI-compatible endpoint | @@ -160,8 +160,8 @@ They are listed here so launch material, sales conversations, and support triage {/* out-of-scope:begin */} | Item | Status | Why | |------|--------|-----| -| Podman / other container runtimes | Unsupported | Onboard surfaces an explicit unsupported-runtime error for Podman (`src/lib/onboard/fatal-runtime-preflight.ts:50` prints the rejection; `src/lib/onboard/preflight.ts:676` flags the unsupported runtime upstream). Only Docker Engine, Docker Desktop, and Colima are supported. See issue #420 (closed). | -| Intel Mac (macOS x86_64) | Unsupported | OpenShell does not publish macOS x86_64 standalone gateway assets. Install hard-fails on x86_64 macOS (`scripts/install-openshell.sh:654`). See issue #954 (closed). | +| Podman / other container runtimes | Unsupported | Onboard surfaces an explicit unsupported-runtime error for Podman (`src/lib/onboard/fatal-runtime-preflight.ts:50` prints the rejection; `src/lib/onboard/preflight.ts:677` flags the unsupported runtime upstream). Only Docker Engine, Docker Desktop, and Colima are supported. See issue #420 (closed). | +| Intel Mac (macOS x86_64) | Unsupported | OpenShell does not publish macOS x86_64 standalone gateway assets. Install hard-fails on x86_64 macOS (`scripts/install-openshell.sh:663`). See issue #954 (closed). | | Non-Ubuntu/Debian Linux distros | Unsupported | Installer assumes `apt-get`. Fedora/Rocky/Alma/Arch/NixOS are not validated and the installer's package-manager probes do not cover them. See open issue #899 (Fedora hang). | | Native Kubernetes or OpenShift deployments | Unsupported | NemoClaw runs the sandbox as a Docker container, not a Kubernetes pod. The default Docker-driver topology does not embed k3s. Operator-managed K8s/OpenShift deployments are out of scope; see issue #407 (community OpenShift through agent-sandbox CRD). | | Air-gapped / offline installs | Unsupported | Onboard assumes network reachability for package fetches, container pulls, and provider validation. See open issues #4872 and #2218 (production-deployment epic covering air-gapped support, China network guidance, multi-host topology). | diff --git a/src/commands/simple-global-oclif-adapters.test.ts b/src/commands/simple-global-oclif-adapters.test.ts index f3ba7e704e9..7d79837b566 100644 --- a/src/commands/simple-global-oclif-adapters.test.ts +++ b/src/commands/simple-global-oclif-adapters.test.ts @@ -283,6 +283,12 @@ describe("simple global oclif adapters", () => { expect(mocks.runStopCommand).toHaveBeenCalledWith( expect.objectContaining({ listSandboxes: expect.any(Function), stopAll: mocks.stopAll }), ); + expect(mocks.runStopCommand.mock.calls).toEqual( + expect.arrayContaining([ + [expect.not.objectContaining({ releaseGatewayPort: true })], + [expect.objectContaining({ releaseGatewayPort: true })], + ]), + ); }); it("passes uninstall runtime dependencies to the uninstall action", async () => { diff --git a/src/commands/stop.ts b/src/commands/stop.ts index 8380d211544..93466f2cbf2 100644 --- a/src/commands/stop.ts +++ b/src/commands/stop.ts @@ -2,26 +2,27 @@ // SPDX-License-Identifier: Apache-2.0 import { NemoClawCommand } from "../lib/cli/nemoclaw-oclif-command"; - -import { stopAll } from "../lib/tunnel/services"; -import { runStopCommand } from "../lib/tunnel/service-command"; import { serviceDeps } from "../lib/tunnel/command-support"; +import { runStopCommand } from "../lib/tunnel/service-command"; +import { stopAll } from "../lib/tunnel/services"; export default class DeprecatedStopCommand extends NemoClawCommand { static id = "stop"; static strict = true; - static summary = "Deprecated alias for 'tunnel stop'"; - static description = "Deprecated alias for tunnel stop."; + static summary = "Deprecated full stop (also releases the managed gateway port)"; + static description = + "Stop tunnel services and release the managed host gateway port. Use 'tunnel stop' to preserve the shared gateway."; static usage = ["stop"]; static examples = ["<%= config.bin %> stop"]; static state = "deprecated" as const; static deprecationOptions = { - message: "Deprecated: 'nemoclaw stop' is now 'nemoclaw tunnel stop'. See 'nemoclaw help'.", + message: + "Deprecated: use 'nemoclaw tunnel stop' for tunnel-only shutdown. This legacy command also releases the managed host gateway port.", }; static flags = {}; public async run(): Promise { await this.parse(DeprecatedStopCommand); - runStopCommand({ ...serviceDeps(), stopAll }); + runStopCommand({ ...serviceDeps(), stopAll, releaseGatewayPort: true }); } } diff --git a/src/lib/onboard.ts b/src/lib/onboard.ts index 81bf251eca1..afb39d2be8a 100644 --- a/src/lib/onboard.ts +++ b/src/lib/onboard.ts @@ -67,6 +67,9 @@ const dockerGpuLocalInference: typeof import("./onboard/docker-gpu-local-inferen const dockerGpuSandboxCreate: typeof import("./onboard/docker-gpu-sandbox-create") = require("./onboard/docker-gpu-sandbox-create"); const dockerDriverGatewayLaunch: typeof import("./onboard/docker-driver-gateway-launch") = require("./onboard/docker-driver-gateway-launch"); const dockerDriverGatewayRuntime: typeof import("./onboard/docker-driver-gateway-runtime") = require("./onboard/docker-driver-gateway-runtime"); +const dockerDriverGatewayCutover: typeof import("./onboard/docker-driver-gateway-cutover") = require("./onboard/docker-driver-gateway-cutover"); +const { reapHostGatewayBeforeLaunchOrFail, reapDuplicateHostGatewaysExceptOrFail } = + require("./onboard/docker-driver-gateway-prelaunch") as typeof import("./onboard/docker-driver-gateway-prelaunch"); const { findReadableNvidiaCdiSpecFiles, parseDockerCdiSpecDirs, @@ -162,7 +165,7 @@ const pRetry = require("p-retry"); * Covers CSI (color, erase, cursor), OSC, and C1 two-byte escapes per ECMA-48. */ const ANSI_RE = /\x1B(?:\[[0-?]*[ -/]*[@-~]|\][^\x07]*(?:\x07|\x1B\\)|[@-_])/g; const runner: typeof import("./runner") = require("./runner"); -const { ROOT, SCRIPTS, redact, run, runCapture, runFile, validateName } = runner; +const { ROOT, SCRIPTS, redact, run, runCapture, runCaptureEx, runFile, validateName } = runner; const braveProviderProfile: typeof import("./onboard/brave-provider-profile") = require("./onboard/brave-provider-profile"); const { runSandboxProviderPreDeleteCleanup } = require("./onboard/sandbox-provider-cleanup") as typeof import("./onboard/sandbox-provider-cleanup"); @@ -624,6 +627,7 @@ const { clearDockerDriverGatewayRuntimeFiles, getDockerDriverGatewayEnv, getDockerDriverGatewayPid, + getDockerDriverGatewayPortListenerScan, getDockerDriverGatewayPortListenerPid, getDockerDriverGatewayRuntimeDrift, getDockerDriverGatewayRuntimeDriftFromSnapshot, @@ -643,12 +647,14 @@ const { getInstalledOpenshellVersion, isOpenshellDevVersion, runCapture, + runCaptureEx, shouldUseOpenshellDevChannel, supportedOpenshellFallbackVersion: SUPPORTED_OPENSHELL_FALLBACK_VERSION, }); import type { JsonObject as LooseObject } from "./core/json-types"; import type { PreparedSandboxBuildContext } from "./onboard/build-context-stage"; + // Non-interactive mode: set by --non-interactive flag or env var. // When active, all prompts use env var overrides or sensible defaults. let NON_INTERACTIVE = false; @@ -1255,10 +1261,8 @@ function retireLegacyGatewayForDockerDriverUpgrade(): void { } } -function restartDockerDriverGatewayProcessForDrift(pid: number, reason: string): void { +function logDockerDriverGatewayRestart(reason: string): void { console.log(` Existing OpenShell Docker-driver gateway is stale (${reason}); restarting...`); - terminateDockerDriverGatewayProcess(pid); - clearDockerDriverGatewayRuntimeFiles(); } async function refreshDockerDriverGatewayReuseState( @@ -2059,6 +2063,15 @@ async function startGatewayWithOptions( process.env.OPENSHELL_GATEWAY = GATEWAY_NAME; } +/** + * Reconcile or create the host Docker-driver gateway. The public onboard() + * entrypoint holds acquireOnboardLock()'s atomic cross-process filesystem lock + * (created with openSync("wx")) across this whole call, so separate concurrent + * `nemoclaw onboard` CLI processes cannot race creation. + * The strict post-reap bind check below remains a second boundary against + * recovery commands or external processes that do not participate in that + * lock; the OS then permits only one child to bind the port. + */ async function startDockerDriverGateway({ exitOnFailure = true, skipSandboxBridgeReachability = false, @@ -2114,93 +2127,77 @@ async function startDockerDriverGateway({ ignoreError: true, }); const activeGatewayInfo = runCaptureOpenshell(["gateway", "info"], { ignoreError: true }); - const pidFileGatewayPid = getDockerDriverGatewayPid(); - if ( - pidFileGatewayPid !== null && - isDockerDriverGatewayProcessAlive() && - isGatewayHealthy(gatewayStatus, gwInfo, activeGatewayInfo) - ) { - const drift = getDockerDriverGatewayRuntimeDrift( - pidFileGatewayPid, - driftGatewayEnv, + // Port availability and listener enumeration are not atomic. The cutover + // rechecks health before adoption, reaps every observed duplicate, and + // requires a fresh strict bind proof after reaping before launch. + const portListenerScan = getDockerDriverGatewayPortListenerScan( + await checkGatewayPortAvailable(), + { gatewayBin: identityGatewayBin }, + ); + const cutover = await dockerDriverGatewayCutover.runDockerDriverGatewayCutover( + { + gatewayBin, + identityGatewayBin, driftGatewayBin, - ); - if (drift) { - restartDockerDriverGatewayProcessForDrift(pidFileGatewayPid, drift.reason); - } else if (registerDockerDriverGatewayEndpoint() && (await isDockerDriverGatewayHttpReady())) { - await verifySandboxBridgeGatewayReachableOrExit(exitOnFailure, { - skip: skipSandboxBridgeReachability, - port: GATEWAY_PORT, - }); - console.log(" ✓ Reusing existing Docker-driver gateway"); - return; - } else { - console.log( - ` Docker-driver gateway metadata reports healthy but http://127.0.0.1:${GATEWAY_PORT}/ is not responding. Starting a fresh gateway...`, - ); - } - } - - const portCheck = await checkGatewayPortAvailable(); - const portListenerPid = getDockerDriverGatewayPortListenerPid(portCheck, { - gatewayBin: identityGatewayBin, - }); - if (portListenerPid !== null) { - const drift = getDockerDriverGatewayRuntimeDrift( - portListenerPid, driftGatewayEnv, - driftGatewayBin, - ); - if (drift) { - rememberDockerDriverGatewayPid(portListenerPid); - restartDockerDriverGatewayProcessForDrift(portListenerPid, drift.reason); - } else { - rememberDockerDriverGatewayPid(portListenerPid); - } - if (!drift && registerDockerDriverGatewayEndpoint()) { - const adoptedStatus = runCaptureOpenshell(["status"], { ignoreError: true }); - const adoptedGwInfo = runCaptureOpenshell(["gateway", "info", "-g", GATEWAY_NAME], { - ignoreError: true, - }); - const adoptedActiveGatewayInfo = runCaptureOpenshell(["gateway", "info"], { - ignoreError: true, - }); - if ( - isGatewayHealthy(adoptedStatus, adoptedGwInfo, adoptedActiveGatewayInfo) && - (await isDockerDriverGatewayHttpReady()) - ) { - await verifySandboxBridgeGatewayReachableOrExit(exitOnFailure, { - skip: skipSandboxBridgeReachability, + exitOnFailure, + skipSandboxBridgeReachability, + stateDir, + portListenerScan, + pidFileGatewayPid: getDockerDriverGatewayPid(), + initialHealth: { + status: gatewayStatus, + namedInfo: gwInfo, + activeInfo: activeGatewayInfo, + }, + }, + { + isDockerDriverGatewayProcessAlive, + isGatewayHealthy, + getDockerDriverGatewayRuntimeDrift, + logDockerDriverGatewayRestart, + registerDockerDriverGatewayEndpoint, + isDockerDriverGatewayHttpReady, + verifySandboxBridgeGatewayReachableOrExit: (fail, options) => + verifySandboxBridgeGatewayReachableOrExit(fail, { + ...options, port: GATEWAY_PORT, - }); - console.log(` ✓ Reusing existing Docker-driver gateway process (PID ${portListenerPid})`); - return; - } - } - } - if (!gatewayBin) { - console.error(" OpenShell Docker-driver gateway binary not found."); - console.error( - ` Install OpenShell v${SUPPORTED_OPENSHELL_FALLBACK_VERSION}, or set NEMOCLAW_OPENSHELL_GATEWAY_BIN.`, - ); - if (exitOnFailure) process.exit(1); - throw new Error("OpenShell gateway binary not found"); - } - - const existingPid = getDockerDriverGatewayPid() ?? portListenerPid; - if (existingPid !== null && isPidAlive(existingPid)) { - if (!isDockerDriverGatewayProcess(existingPid, identityGatewayBin)) { - clearDockerDriverGatewayRuntimeFiles(); - } else { - console.log(` Restarting unhealthy Docker-driver gateway process (PID ${existingPid})...`); - try { - process.kill(existingPid, "SIGTERM"); - sleepSeconds(1); - } catch { - /* best effort; the new process will surface any remaining port conflict */ - } - } - } + }), + readGatewayHealth: () => ({ + status: runCaptureOpenshell(["status"], { ignoreError: true }), + namedInfo: runCaptureOpenshell(["gateway", "info", "-g", GATEWAY_NAME], { + ignoreError: true, + }), + activeInfo: runCaptureOpenshell(["gateway", "info"], { ignoreError: true }), + }), + rememberDockerDriverGatewayPid, + reapDuplicateHostGatewaysExceptOrFail, + reapHostGatewayBeforeLaunchOrFail, + isGatewayPortAvailable: async () => { + const probe = await checkGatewayPortAvailable(); + return probe.ok && !probe.warning; + }, + reportUntrustedGatewayPort: (message) => { + const detail = + `Refusing to start a second OpenShell gateway: ${message}. ` + + `Inspect port ${GATEWAY_PORT} and stop only its owning process before retrying.`; + console.error(` ${detail}`); + if (exitOnFailure) process.exit(1); + throw new Error(detail); + }, + reportMissingGatewayBinary: () => { + console.error(" OpenShell Docker-driver gateway binary not found."); + console.error( + ` Install OpenShell v${SUPPORTED_OPENSHELL_FALLBACK_VERSION}, or set NEMOCLAW_OPENSHELL_GATEWAY_BIN.`, + ); + if (exitOnFailure) process.exit(1); + throw new Error("OpenShell gateway binary not found"); + }, + log: (message) => console.log(message), + }, + ); + if (cutover === "reused") return; + if (!gatewayBin) throw new Error("OpenShell gateway binary missing after cutover"); fs.mkdirSync(stateDir, { recursive: true, mode: 0o700 }); const logPath = path.join(stateDir, "openshell-gateway.log"); diff --git a/src/lib/onboard/docker-driver-gateway-cutover.ts b/src/lib/onboard/docker-driver-gateway-cutover.ts new file mode 100644 index 00000000000..69df96045b0 --- /dev/null +++ b/src/lib/onboard/docker-driver-gateway-cutover.ts @@ -0,0 +1,165 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import type { DockerDriverGatewayPortListenerScan } from "./docker-driver-gateway-port-listener"; + +interface GatewayHealthSnapshot { + status: string; + namedInfo: string; + activeInfo: string; +} + +export interface DockerDriverGatewayCutoverInput { + gatewayBin: string | null; + identityGatewayBin: string | null; + driftGatewayBin: string | null; + driftGatewayEnv: Record; + exitOnFailure: boolean; + skipSandboxBridgeReachability: boolean; + stateDir: string; + portListenerScan: DockerDriverGatewayPortListenerScan; + pidFileGatewayPid: number | null; + initialHealth: GatewayHealthSnapshot; +} + +export interface DockerDriverGatewayCutoverDeps { + isDockerDriverGatewayProcessAlive(): boolean; + isGatewayHealthy(status: string, namedInfo: string, activeInfo: string): boolean; + getDockerDriverGatewayRuntimeDrift( + pid: number, + desiredEnv: Record, + gatewayBin: string | null, + ): { reason: string } | null; + logDockerDriverGatewayRestart(reason: string): void; + registerDockerDriverGatewayEndpoint(): boolean; + isDockerDriverGatewayHttpReady(): Promise; + verifySandboxBridgeGatewayReachableOrExit( + exitOnFailure: boolean, + options: { skip: boolean }, + ): Promise; + readGatewayHealth(): GatewayHealthSnapshot; + rememberDockerDriverGatewayPid(pid: number): void; + reapDuplicateHostGatewaysExceptOrFail( + keepPid: number, + gatewayBin: string | null, + candidatePids: number[], + exitOnFailure: boolean, + ): unknown; + reapHostGatewayBeforeLaunchOrFail(options: { + stateDir: string; + gatewayBin: string | null; + extraPids: number[]; + exitOnFailure: boolean; + }): unknown; + isGatewayPortAvailable(): Promise; + reportUntrustedGatewayPort(message: string): never; + reportMissingGatewayBinary(): never; + log(message: string): void; +} + +/** + * Resolve reuse, adoption, or replacement for the host Docker-driver gateway. + * Every reuse path requires a complete listener scan; replacement reaps only + * port-observed PIDs before the fresh-launch callback is allowed to run. + */ +export async function runDockerDriverGatewayCutover( + input: DockerDriverGatewayCutoverInput, + deps: DockerDriverGatewayCutoverDeps, +): Promise<"reused" | "launch"> { + const portListenerPids = input.portListenerScan.pids; + const portListenerPid = input.portListenerScan.complete ? (portListenerPids[0] ?? null) : null; + + const pidFileGatewayAlive = + input.pidFileGatewayPid !== null && deps.isDockerDriverGatewayProcessAlive(); + const pidFileGatewayDrift = pidFileGatewayAlive + ? deps.getDockerDriverGatewayRuntimeDrift( + input.pidFileGatewayPid as number, + input.driftGatewayEnv, + input.driftGatewayBin, + ) + : null; + // PID-file state alone is never a cleanup candidate: on macOS a stale marker + // cannot distinguish PID reuse after reboot. Same-port duplicates are safe + // only when the complete listener scan observed them explicitly. + const cleanupPids = portListenerPids; + + if ( + input.portListenerScan.complete && + portListenerPids.length === 1 && + input.pidFileGatewayPid !== null && + portListenerPids[0] === input.pidFileGatewayPid && + pidFileGatewayAlive && + deps.isGatewayHealthy( + input.initialHealth.status, + input.initialHealth.namedInfo, + input.initialHealth.activeInfo, + ) + ) { + const drift = pidFileGatewayDrift; + if (drift) { + deps.logDockerDriverGatewayRestart(drift.reason); + } else if ( + deps.registerDockerDriverGatewayEndpoint() && + (await deps.isDockerDriverGatewayHttpReady()) + ) { + await deps.verifySandboxBridgeGatewayReachableOrExit(input.exitOnFailure, { + skip: input.skipSandboxBridgeReachability, + }); + deps.log(" ✓ Reusing existing Docker-driver gateway"); + return "reused"; + } else { + deps.log( + " Docker-driver gateway metadata reports healthy but its HTTP endpoint is not responding. Starting a fresh gateway...", + ); + } + } + + if (portListenerPid !== null) { + const drift = + pidFileGatewayAlive && portListenerPid === input.pidFileGatewayPid + ? pidFileGatewayDrift + : deps.getDockerDriverGatewayRuntimeDrift( + portListenerPid, + input.driftGatewayEnv, + input.driftGatewayBin, + ); + if (drift) deps.logDockerDriverGatewayRestart(drift.reason); + else deps.rememberDockerDriverGatewayPid(portListenerPid); + + if (!drift && deps.registerDockerDriverGatewayEndpoint()) { + const health = deps.readGatewayHealth(); + if ( + deps.isGatewayHealthy(health.status, health.namedInfo, health.activeInfo) && + (await deps.isDockerDriverGatewayHttpReady()) + ) { + deps.reapDuplicateHostGatewaysExceptOrFail( + portListenerPid, + input.identityGatewayBin, + cleanupPids, + input.exitOnFailure, + ); + await deps.verifySandboxBridgeGatewayReachableOrExit(input.exitOnFailure, { + skip: input.skipSandboxBridgeReachability, + }); + deps.log(` ✓ Reusing existing Docker-driver gateway process (PID ${portListenerPid})`); + return "reused"; + } + } + } + + if (!input.gatewayBin) deps.reportMissingGatewayBinary(); + deps.reapHostGatewayBeforeLaunchOrFail({ + stateDir: input.stateDir, + gatewayBin: input.identityGatewayBin, + extraPids: cleanupPids, + exitOnFailure: input.exitOnFailure, + }); + if (!(await deps.isGatewayPortAvailable())) { + deps.reportUntrustedGatewayPort( + input.portListenerScan.complete + ? "the gateway port remains occupied after scoped cleanup" + : "listener enumeration was incomplete and the gateway port remains occupied after scoped cleanup", + ); + } + return "launch"; +} diff --git a/src/lib/onboard/docker-driver-gateway-port-listener.test.ts b/src/lib/onboard/docker-driver-gateway-port-listener.test.ts new file mode 100644 index 00000000000..1444bb04fe1 --- /dev/null +++ b/src/lib/onboard/docker-driver-gateway-port-listener.test.ts @@ -0,0 +1,129 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { describe, expect, it, vi } from "vitest"; + +import { + createDockerDriverGatewayPortListenerHelpers, + type DockerDriverGatewayPortListenerDeps, +} from "./docker-driver-gateway-port-listener"; + +function makeHelpers(overrides: Partial = {}) { + const runCaptureEx = vi.fn(() => ({ stdout: "", exitCode: 1, timedOut: false })); + const deps: DockerDriverGatewayPortListenerDeps = { + gatewayPort: 18080, + runCaptureEx, + isPidAlive: () => true, + isDockerDriverGatewayProcess: () => true, + ...overrides, + }; + return { + helpers: createDockerDriverGatewayPortListenerHelpers(deps), + runCaptureEx: deps.runCaptureEx, + }; +} + +describe("Docker-driver gateway port listener discovery", () => { + it("rejects a primary listener when the injected gateway identity check fails", () => { + const { helpers } = makeHelpers(); + const isDockerDriverGatewayProcessFn = vi.fn(() => false); + + expect( + helpers.getDockerDriverGatewayPortListenerPid( + { ok: false, process: "openshell-gateway", pid: 1234 }, + { + platform: "linux", + gatewayBin: "/opt/openshell/openshell-gateway", + isPidAliveFn: () => true, + isDockerDriverGatewayProcessFn, + }, + ), + ).toBeNull(); + expect(isDockerDriverGatewayProcessFn).toHaveBeenCalledWith( + 1234, + "/opt/openshell/openshell-gateway", + ); + }); + + it("collects every verified gateway listener on the configured port", () => { + const gatewayBin = "/opt/openshell/openshell-gateway"; + const runCaptureEx = vi.fn(() => ({ + stdout: "1234\n2345\n9999\n", + exitCode: 0, + timedOut: false, + })); + const { helpers } = makeHelpers({ runCaptureEx }); + const isDockerDriverGatewayProcessFn = vi.fn( + (pid: number, candidateBin?: string | null) => + (pid === 1234 || pid === 2345) && candidateBin === gatewayBin, + ); + + expect( + helpers.getDockerDriverGatewayPortListenerScan( + { ok: false, process: "openshell-gateway", pid: 1234 }, + { + platform: "linux", + gatewayBin, + isPidAliveFn: () => true, + isDockerDriverGatewayProcessFn, + }, + ), + ).toEqual({ complete: true, pids: [1234, 2345] }); + expect(runCaptureEx).toHaveBeenCalledWith(["lsof", "-ti", ":18080", "-sTCP:LISTEN"]); + }); + + it("retains a verified primary PID when complete enumeration fails", () => { + const { helpers } = makeHelpers({ + runCaptureEx: vi.fn(() => ({ stdout: "", exitCode: 127, timedOut: false })), + }); + + expect( + helpers.getDockerDriverGatewayPortListenerScan( + { ok: false, process: "openshell-gateway", pid: 1234 }, + { + platform: "linux", + isPidAliveFn: () => true, + isDockerDriverGatewayProcessFn: () => true, + }, + ), + ).toEqual({ complete: false, pids: [1234] }); + }); + + it("treats empty lsof output as incomplete while the independent port probe is busy", () => { + const { helpers } = makeHelpers(); + + expect( + helpers.getDockerDriverGatewayPortListenerScan({ + ok: false, + pid: null, + reason: "bind probe reported EADDRINUSE", + }), + ).toEqual({ complete: false, pids: [] }); + }); + + it("marks listener enumeration incomplete when the structured runner throws", () => { + const { helpers } = makeHelpers({ + runCaptureEx: vi.fn(() => { + throw new Error("lsof unavailable"); + }), + }); + + expect(helpers.getDockerDriverGatewayPortListenerScan({ ok: true })).toEqual({ + complete: false, + pids: [], + }); + }); + + it("resolves a dynamic gateway port for every listener scan", () => { + let gatewayPort = 18080; + const runCaptureEx = vi.fn(() => ({ stdout: "", exitCode: 1, timedOut: false })); + const { helpers } = makeHelpers({ gatewayPort: () => gatewayPort, runCaptureEx }); + + helpers.getDockerDriverGatewayPortListenerScan({ ok: true }); + gatewayPort = 18081; + helpers.getDockerDriverGatewayPortListenerScan({ ok: true }); + + expect(runCaptureEx).toHaveBeenNthCalledWith(1, ["lsof", "-ti", ":18080", "-sTCP:LISTEN"]); + expect(runCaptureEx).toHaveBeenNthCalledWith(2, ["lsof", "-ti", ":18081", "-sTCP:LISTEN"]); + }); +}); diff --git a/src/lib/onboard/docker-driver-gateway-port-listener.ts b/src/lib/onboard/docker-driver-gateway-port-listener.ts new file mode 100644 index 00000000000..c53076261c5 --- /dev/null +++ b/src/lib/onboard/docker-driver-gateway-port-listener.ts @@ -0,0 +1,129 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { isLinuxDockerDriverGatewayEnabled } from "./docker-driver-platform"; +import type { PortProbeResult } from "./preflight"; + +export interface DockerDriverGatewayPortListenerOptions { + platform?: NodeJS.Platform; + arch?: NodeJS.Architecture; + gatewayBin?: string | null; + isPidAliveFn?: (pid: number) => boolean; + isDockerDriverGatewayProcessFn?: (pid: number, gatewayBin?: string | null) => boolean; +} + +export interface DockerDriverGatewayPortListenerScan { + /** Every cmdline-verified listener observed by the primary and complete scans. */ + pids: number[]; + /** False when lsof could not authoritatively enumerate the whole listener set. */ + complete: boolean; +} + +interface ListenerCaptureResult { + stdout: string; + exitCode: number | null; + timedOut: boolean; +} + +export interface DockerDriverGatewayPortListenerDeps { + gatewayPort: number | (() => number); + runCaptureEx(args: readonly string[]): ListenerCaptureResult; + isPidAlive(pid: number): boolean; + isDockerDriverGatewayProcess( + pid: number, + gatewayBin: string | null | undefined, + platform: NodeJS.Platform, + ): boolean; +} + +function parseListenerPids(output: string): number[] { + return output + .split(/\r?\n/) + .map((line) => Number.parseInt(line.trim(), 10)) + .filter((pid) => Number.isInteger(pid) && pid > 0); +} + +export function createDockerDriverGatewayPortListenerHelpers( + deps: DockerDriverGatewayPortListenerDeps, +): { + getDockerDriverGatewayPortListenerPid( + portCheck: PortProbeResult, + opts?: DockerDriverGatewayPortListenerOptions, + ): number | null; + getDockerDriverGatewayPortListenerScan( + portCheck: PortProbeResult, + opts?: DockerDriverGatewayPortListenerOptions, + ): DockerDriverGatewayPortListenerScan; + isDockerDriverGatewayPortListener( + portCheck: PortProbeResult, + opts?: DockerDriverGatewayPortListenerOptions, + ): boolean; +} { + const currentGatewayPort = () => + typeof deps.gatewayPort === "function" ? deps.gatewayPort() : deps.gatewayPort; + + function getDockerDriverGatewayPortListenerPid( + portCheck: PortProbeResult, + opts: DockerDriverGatewayPortListenerOptions = {}, + ): number | null { + if (portCheck.ok) return null; + const platform = opts.platform ?? process.platform; + if (!isLinuxDockerDriverGatewayEnabled(platform, opts.arch ?? process.arch)) return null; + const pid = Number(portCheck.pid); + if (!Number.isInteger(pid) || pid <= 0) return null; + if ( + !String(portCheck.process || "") + .toLowerCase() + .startsWith("openshell") + ) + return null; + const alive = opts.isPidAliveFn ?? deps.isPidAlive; + if (!alive(pid)) return null; + const isGateway = + opts.isDockerDriverGatewayProcessFn ?? + ((candidatePid: number, gatewayBin?: string | null) => + deps.isDockerDriverGatewayProcess(candidatePid, gatewayBin, platform)); + return isGateway(pid, opts.gatewayBin) ? pid : null; + } + + function getDockerDriverGatewayPortListenerScan( + portCheck: PortProbeResult, + opts: DockerDriverGatewayPortListenerOptions = {}, + ): DockerDriverGatewayPortListenerScan { + const candidates = new Set(); + const primaryPid = getDockerDriverGatewayPortListenerPid(portCheck, opts); + if (primaryPid !== null) candidates.add(primaryPid); + + let result: ListenerCaptureResult; + try { + result = deps.runCaptureEx(["lsof", "-ti", `:${currentGatewayPort()}`, "-sTCP:LISTEN"]); + } catch { + result = { stdout: "", exitCode: null, timedOut: false }; + } + // Status 1 means "no listeners" only when the independent port probe also + // saw a free port. EADDRINUSE plus empty lsof output is a visibility + // contradiction (commonly a root-owned listener), not a complete scan. + const complete = result.exitCode === 0 || (result.exitCode === 1 && portCheck.ok); + if (result.exitCode === 0) { + for (const pid of parseListenerPids(result.stdout)) candidates.add(pid); + } + + const platform = opts.platform ?? process.platform; + const alive = opts.isPidAliveFn ?? deps.isPidAlive; + const isGateway = + opts.isDockerDriverGatewayProcessFn ?? + ((pid: number, gatewayBin?: string | null) => + deps.isDockerDriverGatewayProcess(pid, gatewayBin, platform)); + return { + pids: Array.from(candidates).filter((pid) => alive(pid) && isGateway(pid, opts.gatewayBin)), + complete, + }; + } + + return { + getDockerDriverGatewayPortListenerPid, + getDockerDriverGatewayPortListenerScan, + isDockerDriverGatewayPortListener: (portCheck, opts) => + getDockerDriverGatewayPortListenerPid(portCheck, opts) !== null, + }; +} diff --git a/src/lib/onboard/docker-driver-gateway-prelaunch.test.ts b/src/lib/onboard/docker-driver-gateway-prelaunch.test.ts new file mode 100644 index 00000000000..0f075d273d2 --- /dev/null +++ b/src/lib/onboard/docker-driver-gateway-prelaunch.test.ts @@ -0,0 +1,214 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { describe, expect, it, vi } from "vitest"; +import { + prelaunchReapFailureMessage, + reapDuplicateHostGatewaysExcept, + reapDuplicateHostGatewaysExceptOrFail, + reapHostGatewayBeforeLaunch, + reapHostGatewayBeforeLaunchOrFail, +} from "./docker-driver-gateway-prelaunch"; +import type { StopHostGatewayOptions, StopHostGatewayResult } from "./host-gateway-process"; + +function emptyResult(overrides: Partial = {}): StopHostGatewayResult { + return { + failed: [], + skippedDeadPids: [], + skippedNonMatchingPids: [], + stopped: [], + sudoRemediationPids: [], + ...overrides, + }; +} + +// Capture the options the reaper hands to stopHostGatewayProcesses. +function stopSpy(result: StopHostGatewayResult): { + fn: typeof import("./host-gateway-process").stopHostGatewayProcesses; + lastOptions: () => StopHostGatewayOptions | undefined; + callCount: () => number; +} { + let captured: StopHostGatewayOptions | undefined; + let calls = 0; + const fn = vi.fn((_deps?: unknown, options?: StopHostGatewayOptions) => { + calls += 1; + captured = options; + return result; + }) as unknown as typeof import("./host-gateway-process").stopHostGatewayProcesses; + return { fn, lastOptions: () => captured, callCount: () => calls }; +} + +describe("reapHostGatewayBeforeLaunch (#5968)", () => { + it("reaps the recorded pid and the port listener, scoped to this port with no host-wide sweep", () => { + const stop = stopSpy(emptyResult({ stopped: [4242] })); + + const result = reapHostGatewayBeforeLaunch( + { + pidFile: "/state/openshell-docker-gateway-8090/openshell-gateway.pid", + stateDir: "/state/openshell-docker-gateway-8090", + gatewayBin: "/usr/local/bin/openshell-gateway", + extraPids: [4242], + }, + {}, + stop.fn, + ); + + expect(result.stopped).toEqual([4242]); + const options = stop.lastOptions(); + expect(options?.pids).toEqual([4242]); + expect(options?.usePidFile).toBe(false); + expect(options?.usePgrepFallback).toBe(false); + expect(options?.pidFile).toBe("/state/openshell-docker-gateway-8090/openshell-gateway.pid"); + expect(options?.stateDir).toBe("/state/openshell-docker-gateway-8090"); + expect(options?.gatewayBin).toBe("/usr/local/bin/openshell-gateway"); + }); + + it("drops null/invalid/duplicate candidate pids so a missing pid-file/listener is a quiet no-op", () => { + const stop = stopSpy(emptyResult()); + + reapHostGatewayBeforeLaunch( + { + pidFile: "/state/openshell-docker-gateway/openshell-gateway.pid", + stateDir: "/state/openshell-docker-gateway", + gatewayBin: null, + extraPids: [null, undefined, 0, -1, 7777, 7777], + }, + {}, + stop.fn, + ); + + expect(stop.lastOptions()?.pids).toEqual([7777]); + }); +}); + +describe("prelaunchReapFailureMessage (#5968)", () => { + it("returns null when no matched gateway resisted the reap", () => { + expect(prelaunchReapFailureMessage(emptyResult({ stopped: [10] }))).toBeNull(); + }); + + it("describes the unreaped gateway pids and a remediation scoped to those pids", () => { + const message = prelaunchReapFailureMessage(emptyResult({ failed: [321, 654] })); + expect(message).toContain("321, 654"); + // Scoped to the matched pids, never a host-wide `pkill -f openshell-gateway`. + expect(message).toContain("sudo kill -9 321 654"); + expect(message).not.toContain("pkill"); + }); +}); + +describe("reapHostGatewayBeforeLaunchOrFail (#5968)", () => { + const options = { + pidFile: "/state/openshell-docker-gateway-8090/openshell-gateway.pid", + stateDir: "/state/openshell-docker-gateway-8090", + gatewayBin: "/usr/local/bin/openshell-gateway", + extraPids: [4242], + }; + + it("returns the cleared result and does not exit when the port is clear", () => { + const stop = stopSpy(emptyResult({ stopped: [4242] })); + const exit = vi.fn(() => undefined as never); + + const result = reapHostGatewayBeforeLaunchOrFail(options, {}, stop.fn, exit); + + expect(result.stopped).toEqual([4242]); + expect(exit).not.toHaveBeenCalled(); + }); + + it("throws and never spawns when a matched gateway could not be stopped (exitOnFailure off)", () => { + const stop = stopSpy(emptyResult({ failed: [4242] })); + const exit = vi.fn(() => undefined as never); + + expect(() => + reapHostGatewayBeforeLaunchOrFail({ ...options, exitOnFailure: false }, {}, stop.fn, exit), + ).toThrow(/could not be stopped/); + expect(exit).not.toHaveBeenCalled(); + }); + + it("exits with code 1 when a matched gateway could not be stopped and exitOnFailure is set", () => { + const stop = stopSpy(emptyResult({ failed: [4242] })); + const exit = vi.fn((_code: number) => { + throw new Error("exit-called"); + }) as unknown as (code: number) => never; + + expect(() => + reapHostGatewayBeforeLaunchOrFail({ ...options, exitOnFailure: true }, {}, stop.fn, exit), + ).toThrow(/exit-called/); + expect(exit).toHaveBeenCalledWith(1); + }); +}); + +describe("reapDuplicateHostGatewaysExcept (#5968)", () => { + it("reaps a known stale duplicate pid while excluding the gateway being reused", () => { + const stop = stopSpy(emptyResult({ stopped: [111] })); + + const result = reapDuplicateHostGatewaysExcept( + 999, + "/usr/local/bin/openshell-gateway", + [111, 999, null, 999], + {}, + stop.fn, + ); + + expect(result.stopped).toEqual([111]); + const captured = stop.lastOptions(); + expect(captured?.pids).toEqual([111]); + expect(captured?.usePgrepFallback).toBe(false); + expect(captured?.gatewayBin).toBe("/usr/local/bin/openshell-gateway"); + // The duplicate reap must not read or clear the adopted gateway's live + // pid-file/runtime marker. + expect(captured?.usePidFile).toBe(false); + expect(captured?.clearRuntimeFiles).toBe(false); + }); + + it("never calls the stopper when the only known candidate is the reused gateway", () => { + const stop = stopSpy(emptyResult()); + + const result = reapDuplicateHostGatewaysExcept(999, null, [999, null, 0, -3], {}, stop.fn); + + expect(result).toEqual(emptyResult()); + expect(stop.callCount()).toBe(0); + }); +}); + +describe("reapDuplicateHostGatewaysExceptOrFail (#5968)", () => { + const gatewayBin = "/usr/local/bin/openshell-gateway"; + + it("returns the result and does not exit when the stale duplicate was reaped", () => { + const stop = stopSpy(emptyResult({ stopped: [111] })); + const exit = vi.fn(() => undefined as never); + + const result = reapDuplicateHostGatewaysExceptOrFail( + 999, + gatewayBin, + [111], + false, + {}, + stop.fn, + exit, + ); + + expect(result.stopped).toEqual([111]); + expect(exit).not.toHaveBeenCalled(); + }); + + it("throws and never reports reuse when a matched duplicate could not be stopped (exitOnFailure off)", () => { + const stop = stopSpy(emptyResult({ failed: [111] })); + const exit = vi.fn(() => undefined as never); + + expect(() => + reapDuplicateHostGatewaysExceptOrFail(999, gatewayBin, [111], false, {}, stop.fn, exit), + ).toThrow(/could not be stopped/); + expect(exit).not.toHaveBeenCalled(); + }); + + it("exits with code 1 when a matched duplicate could not be stopped and exitOnFailure is set", () => { + const stop = stopSpy(emptyResult({ failed: [111] })); + const exit = vi.fn((_code: number) => { + throw new Error("exit-called"); + }) as unknown as (code: number) => never; + + expect(() => + reapDuplicateHostGatewaysExceptOrFail(999, gatewayBin, [111], true, {}, stop.fn, exit), + ).toThrow(/exit-called/); + expect(exit).toHaveBeenCalledWith(1); + }); +}); diff --git a/src/lib/onboard/docker-driver-gateway-prelaunch.ts b/src/lib/onboard/docker-driver-gateway-prelaunch.ts new file mode 100644 index 00000000000..1effffe9e3f --- /dev/null +++ b/src/lib/onboard/docker-driver-gateway-prelaunch.ts @@ -0,0 +1,196 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +/** + * Pre-launch reaping for the host OpenShell Docker-driver gateway. + * + * When onboard cannot reuse an already-running gateway (its metadata reports + * unhealthy, the HTTP endpoint is unresponsive, or runtime drift forces a + * restart) it replaces that gateway with a fresh process. Historically that + * replacement only sent a single `SIGTERM` and slept one second before + * spawning — with no `SIGKILL` escalation, no wait for the old process to + * actually exit, and no sweep of a duplicate listener — so a slow-to-die + * gateway could still be alive when the new one spawned, leaving two + * host-process gateways bound to the same port (#5968: "gateway must be shared + * (exactly one instance …); got container=0 host-process=2"). + * + * This reuses the shared `stopHostGatewayProcesses` reaper (TERM→KILL with + * bounded waits, wait-for-exit, and cmdline gating on the `openshell-gateway` + * identity) so the existing gateway is *confirmed gone* before the caller + * spawns its replacement. It is scoped to the resolved per-port candidates with + * `usePgrepFallback: false` — never a host-wide sweep — so a different + * worktree's gateway on another port is never torn down. + * + * Two follow-on guards keep the linked singleton invariant on the start path: + * - `reapHostGatewayBeforeLaunchOrFail` fails closed when a matched gateway + * resists the reap (`failed` non-empty), so a replacement is never spawned + * over a still-alive gateway. + * - `reapDuplicateHostGatewaysExcept` lets a reuse path clean up a *known* + * stale duplicate (e.g. a previously recorded pid that differs from the + * adopted port listener) without tearing down the gateway being reused. + */ + +import path from "node:path"; + +import { + type HostGatewayProcessDeps, + type StopHostGatewayResult, + stopHostGatewayProcesses, +} from "./host-gateway-process"; + +export interface ReapHostGatewayBeforeLaunchOptions { + /** Per-port gateway state dir (holds the pid file and runtime marker). */ + stateDir: string; + /** Recorded gateway pid file; defaults to `/openshell-gateway.pid`. */ + pidFile?: string; + /** Canonical gateway binary; cmdline-gates which PIDs may be signalled. */ + gatewayBin: string | null; + /** Extra candidate PIDs to reap (e.g. the current port listener). */ + extraPids?: Array; +} + +// A `stopHostGatewayProcesses` result with nothing stopped — returned when there +// is no live candidate to reap so callers always get a well-formed result. +function emptyStopResult(): StopHostGatewayResult { + return { + failed: [], + skippedDeadPids: [], + skippedNonMatchingPids: [], + stopped: [], + sudoRemediationPids: [], + }; +} + +function validPids(pids: Array, exclude?: number): number[] { + return Array.from( + new Set( + pids.filter( + (pid): pid is number => + typeof pid === "number" && Number.isInteger(pid) && pid > 0 && pid !== exclude, + ), + ), + ); +} + +/** + * Reap any host `openshell-gateway` already bound to this gateway port so the + * caller can spawn exactly one replacement. Best-effort and idempotent: a quiet + * no-op when nothing matching is alive. Returns the stopper result so callers + * (and tests) can observe what was stopped. + */ +export function reapHostGatewayBeforeLaunch( + options: ReapHostGatewayBeforeLaunchOptions, + deps: Partial = {}, + stop: typeof stopHostGatewayProcesses = stopHostGatewayProcesses, +): StopHostGatewayResult { + return stop( + { env: process.env, ...deps }, + { + pids: validPids(options.extraPids ?? []), + pidFile: options.pidFile ?? path.join(options.stateDir, "openshell-gateway.pid"), + stateDir: options.stateDir, + gatewayBin: options.gatewayBin, + // PID-file state is bookkeeping, not proof that the process owns this + // port. Signal only the port-observed candidates supplied by the caller; + // a stale/recycled PID must never reap another worktree's gateway. + usePidFile: false, + usePgrepFallback: false, + }, + ); +} + +/** + * Message describing host gateways the prelaunch reap could not stop, or `null` + * when the port is clear. A non-empty `failed` means a matched gateway resisted + * TERM→KILL (e.g. a privileged process); spawning a replacement over it would + * leave two host gateways (#5968 host-process=2), so callers must fail closed. + */ +export function prelaunchReapFailureMessage(result: StopHostGatewayResult): string | null { + if (result.failed.length === 0) return null; + // Recommend killing exactly the PIDs we matched, not a host-wide + // `pkill -f openshell-gateway`: this path is deliberately scoped to this port + // (usePgrepFallback:false), so a host-wide kill could take down another + // worktree's gateway. + return ( + "Refusing to start a second OpenShell gateway: existing host gateway process " + + `${result.failed.join(", ")} could not be stopped. Run: sudo kill -9 ${result.failed.join(" ")}` + ); +} + +/** + * Reap the existing host gateway for this port, then fail closed when a matched + * gateway resisted stopping so onboard never spawns a replacement over a + * still-alive gateway. Honours `exitOnFailure` like the rest of onboard: + * `process.exit(1)` when set, otherwise throw. Returns the (cleared) stop result. + */ +export function reapHostGatewayBeforeLaunchOrFail( + options: ReapHostGatewayBeforeLaunchOptions & { exitOnFailure?: boolean }, + deps: Partial = {}, + stop: typeof stopHostGatewayProcesses = stopHostGatewayProcesses, + exit: (code: number) => never = (code) => process.exit(code) as never, +): StopHostGatewayResult { + const result = reapHostGatewayBeforeLaunch(options, deps, stop); + const failure = prelaunchReapFailureMessage(result); + if (failure) { + console.error(` ${failure}`); + if (options.exitOnFailure) exit(1); + throw new Error(failure); + } + return result; +} + +/** + * Reap KNOWN host gateways (cmdline-gated, no host-wide pgrep sweep) other than + * the gateway being reused, so a reuse path can enforce a single matching host + * gateway without tearing down the adopted one. Used when a previously recorded + * gateway pid differs from the port listener now being adopted — that stale pid + * is a duplicate orphan and is reaped here. A quiet no-op when the only known + * candidate is `keepPid`. Pid-file discovery and runtime-file cleanup are + * disabled so the adopted gateway's live state is never read as a candidate or + * cleared. + */ +export function reapDuplicateHostGatewaysExcept( + keepPid: number, + gatewayBin: string | null, + candidatePids: Array, + deps: Partial = {}, + stop: typeof stopHostGatewayProcesses = stopHostGatewayProcesses, +): StopHostGatewayResult { + const pids = validPids(candidatePids, keepPid); + if (pids.length === 0) return emptyStopResult(); + return stop( + { env: process.env, ...deps }, + { + clearRuntimeFiles: false, + pids, + gatewayBin, + usePidFile: false, + usePgrepFallback: false, + }, + ); +} + +/** + * Like `reapDuplicateHostGatewaysExcept`, but fail closed when a matched + * duplicate resisted stopping (`failed` non-empty): a reuse path must not report + * success while a second matching host gateway is still alive (#5968). Honours + * `exitOnFailure` (`process.exit(1)` when set, otherwise throw). + */ +export function reapDuplicateHostGatewaysExceptOrFail( + keepPid: number, + gatewayBin: string | null, + candidatePids: Array, + exitOnFailure?: boolean, + deps: Partial = {}, + stop: typeof stopHostGatewayProcesses = stopHostGatewayProcesses, + exit: (code: number) => never = (code) => process.exit(code) as never, +): StopHostGatewayResult { + const result = reapDuplicateHostGatewaysExcept(keepPid, gatewayBin, candidatePids, deps, stop); + const failure = prelaunchReapFailureMessage(result); + if (failure) { + console.error(` ${failure}`); + if (exitOnFailure) exit(1); + throw new Error(failure); + } + return result; +} diff --git a/src/lib/onboard/docker-driver-gateway-runtime.test.ts b/src/lib/onboard/docker-driver-gateway-runtime.test.ts index 4388bf3630d..8d656385481 100644 --- a/src/lib/onboard/docker-driver-gateway-runtime.test.ts +++ b/src/lib/onboard/docker-driver-gateway-runtime.test.ts @@ -234,28 +234,6 @@ describe("docker-driver gateway runtime helpers", () => { } }); - it("rejects an openshell port listener when the injected gateway identity check fails", () => { - const { helpers } = makeHelpers(); - const isDockerDriverGatewayProcessFn = vi.fn(() => false); - - expect( - helpers.getDockerDriverGatewayPortListenerPid( - { ok: false, process: "openshell-gateway", pid: 1234 }, - { - platform: "linux", - gatewayBin: "/opt/openshell/openshell-gateway", - isPidAliveFn: () => true, - isDockerDriverGatewayProcessFn, - }, - ), - ).toBeNull(); - - expect(isDockerDriverGatewayProcessFn).toHaveBeenCalledWith( - 1234, - "/opt/openshell/openshell-gateway", - ); - }); - it("does not match process args that only contain openshell-gateway as a suffix", () => { const pid = 12_345; const { helpers, runCapture } = makeHelpers({ diff --git a/src/lib/onboard/docker-driver-gateway-runtime.ts b/src/lib/onboard/docker-driver-gateway-runtime.ts index 0036f6a1b27..8daddc9eb6f 100644 --- a/src/lib/onboard/docker-driver-gateway-runtime.ts +++ b/src/lib/onboard/docker-driver-gateway-runtime.ts @@ -7,14 +7,23 @@ import path from "node:path"; import { resolveOpenshell } from "../adapters/openshell/resolve"; import { isErrnoException } from "../core/errno"; +import { + createDockerDriverGatewayPortListenerHelpers, + type DockerDriverGatewayPortListenerOptions, + type DockerDriverGatewayPortListenerScan, +} from "./docker-driver-gateway-port-listener"; import * as dockerDriverGatewayRuntimeMarker from "./docker-driver-gateway-runtime-marker"; -import { isLinuxDockerDriverGatewayEnabled } from "./docker-driver-platform"; import * as gatewayBinding from "./gateway-binding"; import { gatewayProcessCmdlineMatches, OPENSHELL_GATEWAY_PROCESS_NAMES, } from "./gateway-process-identity"; import type { PortProbeResult } from "./preflight"; + +// Keep the listener option type on the established runtime facade while the +// implementation remains isolated in docker-driver-gateway-port-listener.ts. +export type { DockerDriverGatewayPortListenerOptions } from "./docker-driver-gateway-port-listener"; + import * as vmDriverProcess from "./vm-driver-process"; const OPENSHELL_SUPERVISOR_MANIFEST_DIGESTS: Readonly> = { @@ -24,6 +33,11 @@ const OPENSHELL_SUPERVISOR_MANIFEST_DIGESTS: Readonly> = export type DockerDriverGatewayRuntimeDrift = { reason: string }; type RunCapture = (args: string[], opts?: { ignoreError?: boolean }) => string; +type RunCaptureEx = (args: readonly string[]) => { + stdout: string; + exitCode: number | null; + timedOut: boolean; +}; type DockerDriverGatewayEnvModule = typeof import("./docker-driver-gateway-env"); // Source boundary: OpenShell does not currently expose an authoritative local @@ -42,6 +56,7 @@ export interface DockerDriverGatewayRuntimeDeps { isOpenshellDevVersion(versionOutput: string | null | undefined): boolean; loadDockerDriverGatewayEnv?(): DockerDriverGatewayEnvModule; runCapture: RunCapture; + runCaptureEx?: RunCaptureEx; shouldUseOpenshellDevChannel(): boolean; supportedOpenshellFallbackVersion: string; } @@ -54,15 +69,18 @@ export function createDockerDriverGatewayRuntimeHelpers(deps: DockerDriverGatewa ): Record; getDockerDriverGatewayPid(): number | null; getDockerDriverGatewayPidFile(): string; + getDockerDriverGatewayPortListenerScan( + portCheck: PortProbeResult, + opts?: DockerDriverGatewayPortListenerOptions, + ): DockerDriverGatewayPortListenerScan; + /** Compatibility view for callers that only need the verified PID list. */ + getDockerDriverGatewayPortListenerPids( + portCheck: PortProbeResult, + opts?: DockerDriverGatewayPortListenerOptions, + ): number[]; getDockerDriverGatewayPortListenerPid( portCheck: PortProbeResult, - opts?: { - platform?: NodeJS.Platform; - arch?: NodeJS.Architecture; - gatewayBin?: string | null; - isPidAliveFn?: (pid: number) => boolean; - isDockerDriverGatewayProcessFn?: (pid: number, gatewayBin?: string | null) => boolean; - }, + opts?: DockerDriverGatewayPortListenerOptions, ): number | null; getDockerDriverGatewayRuntimeDrift( pid: number, @@ -411,52 +429,42 @@ export function createDockerDriverGatewayRuntimeHelpers(deps: DockerDriverGatewa ); } - function getDockerDriverGatewayPortListenerPid( - portCheck: PortProbeResult, - opts: { - platform?: NodeJS.Platform; - arch?: NodeJS.Architecture; - gatewayBin?: string | null; - isPidAliveFn?: (pid: number) => boolean; - isDockerDriverGatewayProcessFn?: (pid: number, gatewayBin?: string | null) => boolean; - } = {}, - ): number | null { - if (portCheck.ok) return null; - if ( - !isLinuxDockerDriverGatewayEnabled( - opts.platform ?? process.platform, - opts.arch ?? process.arch, - ) - ) - return null; - const pid = Number(portCheck.pid); - if (!Number.isInteger(pid) || pid <= 0) return null; - const proc = String(portCheck.process || "").toLowerCase(); - if (!proc.startsWith("openshell")) return null; - const alive = opts.isPidAliveFn ?? isPidAlive; - if (!alive(pid)) return null; - const isGateway = - opts.isDockerDriverGatewayProcessFn ?? - ((candidatePid: number, gatewayBin?: string | null) => - isDockerDriverGatewayProcess(candidatePid, gatewayBin, { - requireDockerDriverEnv: shouldRequireDockerDriverEnv(opts.platform ?? process.platform), - })); - if (!isGateway(pid, opts.gatewayBin)) return null; - return pid; - } - - function isDockerDriverGatewayPortListener( + // Bind listener discovery to this factory's liveness and process-identity + // dependencies. Returning the configured methods keeps onboard on one + // authoritative runtime instance rather than constructing a second factory. + const { + getDockerDriverGatewayPortListenerPid, + getDockerDriverGatewayPortListenerScan, + isDockerDriverGatewayPortListener, + } = createDockerDriverGatewayPortListenerHelpers({ + gatewayPort: currentGatewayPort, + runCaptureEx: + deps.runCaptureEx ?? + ((args) => { + try { + return { stdout: deps.runCapture([...args]), exitCode: 0, timedOut: false }; + } catch { + return { stdout: "", exitCode: null, timedOut: false }; + } + }), + isPidAlive, + isDockerDriverGatewayProcess: (pid, gatewayBin, platform) => + isDockerDriverGatewayProcess(pid, gatewayBin, { + requireDockerDriverEnv: shouldRequireDockerDriverEnv(platform), + }), + }); + const getDockerDriverGatewayPortListenerPids = ( portCheck: PortProbeResult, - opts: Parameters[1] = {}, - ): boolean { - return getDockerDriverGatewayPortListenerPid(portCheck, opts) !== null; - } + opts: DockerDriverGatewayPortListenerOptions = {}, + ): number[] => getDockerDriverGatewayPortListenerScan(portCheck, opts).pids; return { clearDockerDriverGatewayRuntimeFiles, getDockerDriverGatewayEnv, getDockerDriverGatewayPid, getDockerDriverGatewayPidFile, + getDockerDriverGatewayPortListenerScan, + getDockerDriverGatewayPortListenerPids, getDockerDriverGatewayPortListenerPid, getDockerDriverGatewayRuntimeDrift, getDockerDriverGatewayRuntimeDriftFromSnapshot, diff --git a/src/lib/onboard/host-gateway-process.test.ts b/src/lib/onboard/host-gateway-process.test.ts index b77892c4f30..1fb74cce07b 100644 --- a/src/lib/onboard/host-gateway-process.test.ts +++ b/src/lib/onboard/host-gateway-process.test.ts @@ -8,9 +8,9 @@ import path from "node:path"; import { describe, expect, it, vi } from "vitest"; import { - stopHostGatewayProcesses, type HostGatewayProcessDeps, type RunResult, + stopHostGatewayProcesses, } from "./host-gateway-process"; interface RunArgs { @@ -196,7 +196,7 @@ describe("stopHostGatewayProcesses", () => { expect(result.stopped).toEqual([]); expect(warn).toHaveBeenCalledWith( "pgrep not found; could not scan for orphan host openshell-gateway processes. " + - "If port 8080 is still bound, run: sudo pkill -f openshell-gateway", + "Inspect any remaining listener and stop only the matching gateway process.", ); expect(log).not.toHaveBeenCalledWith("No host openshell-gateway processes found"); }); @@ -253,7 +253,7 @@ describe("stopHostGatewayProcesses", () => { expect(result.failed).toEqual([9999042]); expect(result.sudoRemediationPids).toEqual([9999042]); expect(warn).toHaveBeenCalledWith( - "Cannot stop root-owned host openshell-gateway process 9999042. Run: sudo pkill -f openshell-gateway", + "Cannot stop root-owned host openshell-gateway process 9999042. Run: sudo kill -9 9999042", ); }); diff --git a/src/lib/onboard/host-gateway-process.ts b/src/lib/onboard/host-gateway-process.ts index 2fcd6589697..d0bb226d6be 100644 --- a/src/lib/onboard/host-gateway-process.ts +++ b/src/lib/onboard/host-gateway-process.ts @@ -1,7 +1,7 @@ // SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. // SPDX-License-Identifier: Apache-2.0 -import { spawnSync, type SpawnSyncOptions } from "node:child_process"; +import { type SpawnSyncOptions, spawnSync } from "node:child_process"; import fs from "node:fs"; import os from "node:os"; import path from "node:path"; @@ -26,6 +26,8 @@ export interface HostGatewayProcessDeps { } export interface StopHostGatewayOptions { + /** Whether successful stops may clear the pid file/runtime marker. */ + clearRuntimeFiles?: boolean; gatewayBin?: string | null; killWaitMs?: number; logNoProcesses?: boolean; @@ -34,6 +36,8 @@ export interface StopHostGatewayOptions { pollIntervalMs?: number; stateDir?: string; termWaitMs?: number; + /** Whether to read and act on the resolved pid file. */ + usePidFile?: boolean; usePgrepFallback?: boolean; } @@ -78,6 +82,9 @@ function defaultKill(pid: number, signal?: NodeJS.Signals | number): boolean { } function defaultCommandExists(command: string, env: NodeJS.ProcessEnv): boolean { + // `command` is always an internal, trusted literal ("pgrep"); it is never + // user-supplied. It is also JSON.stringify-quoted, so the `sh -c` here carries + // no shell-injection surface. return ( defaultRun("sh", ["-c", `command -v ${JSON.stringify(command)} >/dev/null 2>&1`], { env, @@ -207,7 +214,7 @@ function warnSudoRemediation(pid: number, deps: HostGatewayProcessDeps): void { const ownerLabel = owner ? `${owner}-owned` : "privileged"; warn( `Cannot stop ${ownerLabel} host openshell-gateway process ${pid}. ` + - "Run: sudo pkill -f openshell-gateway", + `Run: sudo kill -9 ${pid}`, ); } @@ -238,6 +245,7 @@ export function stopHostGatewayProcesses( const deps = defaultDeps(depsOverrides); const stateDir = options.stateDir ?? resolveDockerDriverGatewayStateDir(deps.env); const pidFile = options.pidFile ?? path.join(stateDir, "openshell-gateway.pid"); + const clearRuntimeState = options.clearRuntimeFiles ?? true; const candidates = new Map>(); const result: StopHostGatewayResult = { failed: [], @@ -247,11 +255,13 @@ export function stopHostGatewayProcesses( sudoRemediationPids: [], }; - const pidFromFile = readPidFile(pidFile); - if (pidFromFile !== null) { - addPid(candidates, pidFromFile, "pid-file"); - } else if (fs.existsSync(pidFile)) { - clearRuntimeFiles(pidFile, stateDir); + if (options.usePidFile ?? true) { + const pidFromFile = readPidFile(pidFile); + if (pidFromFile !== null) { + addPid(candidates, pidFromFile, "pid-file"); + } else if (clearRuntimeState && fs.existsSync(pidFile)) { + clearRuntimeFiles(pidFile, stateDir); + } } const explicitPids = Array.from(options.pids ?? []).filter( @@ -281,7 +291,7 @@ export function stopHostGatewayProcesses( for (const [pid, sources] of candidates) { if (!pidExists(pid, deps)) { result.skippedDeadPids.push(pid); - if (sources.has("pid-file") && !clearedRuntimeFiles) { + if (clearRuntimeState && sources.has("pid-file") && !clearedRuntimeFiles) { clearRuntimeFiles(pidFile, stateDir); clearedRuntimeFiles = true; } @@ -289,7 +299,7 @@ export function stopHostGatewayProcesses( } if (!hostGatewayCmdlineMatches(processArgs(pid, deps), options.gatewayBin)) { result.skippedNonMatchingPids.push(pid); - if (sources.has("pid-file") && !clearedRuntimeFiles) { + if (clearRuntimeState && sources.has("pid-file") && !clearedRuntimeFiles) { clearRuntimeFiles(pidFile, stateDir); clearedRuntimeFiles = true; } @@ -298,7 +308,7 @@ export function stopHostGatewayProcesses( if (tryStopPid(pid, deps, waitOptions) === "stopped") { result.stopped.push(pid); - if (!clearedRuntimeFiles) { + if (clearRuntimeState && !clearedRuntimeFiles) { clearRuntimeFiles(pidFile, stateDir); clearedRuntimeFiles = true; } @@ -317,7 +327,7 @@ export function stopHostGatewayProcesses( const warn = deps.warn ?? ((message: string) => console.warn(message)); warn( "pgrep not found; could not scan for orphan host openshell-gateway processes. " + - "If port 8080 is still bound, run: sudo pkill -f openshell-gateway", + "Inspect any remaining listener and stop only the matching gateway process.", ); } else { const log = deps.log ?? ((message: string) => console.log(message)); diff --git a/src/lib/state/onboard-session-cross-process-lock.test.ts b/src/lib/state/onboard-session-cross-process-lock.test.ts new file mode 100644 index 00000000000..eb97bab7d0d --- /dev/null +++ b/src/lib/state/onboard-session-cross-process-lock.test.ts @@ -0,0 +1,73 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { spawn } from "node:child_process"; +import { once } from "node:events"; +import fs from "node:fs"; +import { createRequire } from "node:module"; +import os from "node:os"; +import path from "node:path"; + +import { afterEach, beforeEach, describe, expect, it } from "vitest"; + +const require = createRequire(import.meta.url); +const sessionPath = require.resolve("./onboard-session"); +const originalHome = process.env.HOME; +type OnboardSessionModule = typeof import("./onboard-session"); +let session: OnboardSessionModule; +let tempHome: string; + +function restoreHome(): boolean { + return originalHome === undefined + ? Reflect.deleteProperty(process.env, "HOME") + : Reflect.set(process.env, "HOME", originalHome); +} + +beforeEach(() => { + tempHome = fs.mkdtempSync(path.join(os.tmpdir(), "nemoclaw-onboard-lock-process-")); + process.env.HOME = tempHome; + delete require.cache[sessionPath]; + session = require("./onboard-session"); + session.releaseOnboardLock(); +}); + +afterEach(() => { + session.releaseOnboardLock(); + delete require.cache[sessionPath]; + fs.rmSync(tempHome, { recursive: true, force: true }); + restoreHome(); +}); + +describe("cross-process onboard lock", () => { + it("rejects a concurrent CLI process before gateway creation", async () => { + const childScript = ` + const fs = require("node:fs"); + const path = require("node:path"); + const lockFile = process.argv[1]; + fs.mkdirSync(path.dirname(lockFile), { recursive: true }); + const fd = fs.openSync(lockFile, "wx", 0o600); + fs.writeSync(fd, JSON.stringify({ + pid: process.pid, + startedAt: new Date().toISOString(), + command: "separate nemoclaw onboard process", + })); + process.stdout.write("locked\\n"); + setInterval(() => {}, 1000); + `; + const child = spawn(process.execPath, ["-e", childScript, session.LOCK_FILE], { + stdio: ["ignore", "pipe", "inherit"], + }); + await once(child.stdout, "data"); + + try { + const acquired = session.acquireOnboardLock("competing nemoclaw onboard"); + expect(acquired.acquired).toBe(false); + expect(acquired.holderPid).toBe(child.pid); + expect(acquired.holderCommand).toBe("separate nemoclaw onboard process"); + } finally { + const exited = once(child, "exit"); + child.kill(); + await exited; + } + }); +}); diff --git a/src/lib/tunnel/gateway-port-confirmation.test.ts b/src/lib/tunnel/gateway-port-confirmation.test.ts new file mode 100644 index 00000000000..d0c635731be --- /dev/null +++ b/src/lib/tunnel/gateway-port-confirmation.test.ts @@ -0,0 +1,48 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { describe, expect, it, vi } from "vitest"; + +import { confirmGatewayPortReleased } from "./gateway-port-confirmation"; + +describe("confirmGatewayPortReleased", () => { + it("caps failed listener inspections at twenty without spawning a bind probe", () => { + let clock = 0; + const listeningPids = vi.fn(() => null); + const probePortFree = vi.fn(() => true); + + const result = confirmGatewayPortReleased({ + port: 8080, + timeoutMs: 100_000, + pollIntervalMs: 1, + now: () => clock++, + sleep: () => {}, + probePortFree, + listeningPids, + }); + + expect(result.released).toBe(false); + expect(listeningPids).toHaveBeenCalledTimes(20); + expect(probePortFree).not.toHaveBeenCalled(); + }); + + it("runs the independent bind probe once after listeners clear", () => { + let clock = 0; + const listeningPids = vi.fn().mockReturnValueOnce([4242]).mockReturnValue([]); + const probePortFree = vi.fn(() => true); + + const result = confirmGatewayPortReleased({ + port: 8080, + timeoutMs: 100_000, + pollIntervalMs: 1, + now: () => clock++, + sleep: () => {}, + probePortFree, + listeningPids, + }); + + expect(result).toEqual({ released: true, remaining: [] }); + expect(listeningPids).toHaveBeenCalledTimes(2); + expect(probePortFree).toHaveBeenCalledTimes(1); + }); +}); diff --git a/src/lib/tunnel/gateway-port-confirmation.ts b/src/lib/tunnel/gateway-port-confirmation.ts new file mode 100644 index 00000000000..5328e6c7fd5 --- /dev/null +++ b/src/lib/tunnel/gateway-port-confirmation.ts @@ -0,0 +1,93 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { spawnSync } from "node:child_process"; + +import { waitUntil } from "../core/wait"; + +const PORT_FREE_PROBE_SCRIPT = ` +const net = require("node:net"); +const port = Number(process.argv[1]); +const server = net.createServer(); +let done = false; +const finish = (code) => { + if (done) return; + done = true; + server.close(() => process.exit(code)); +}; +// Bind errors are asynchronous in Node, so exit nonzero from the error event. +server.once("error", () => process.exit(1)); +// The listening callback is the proof that this child acquired the port. +server.listen(port, "127.0.0.1", () => finish(0)); +`; + +export interface ConfirmGatewayPortOptions { + port: number; + timeoutMs: number; + pollIntervalMs: number; + now: () => number; + sleep?: (ms: number) => void; + probePortFree: (port: number) => boolean; + /** Optional authoritative listener scan; null means the scan itself failed. */ + listeningPids?: () => number[] | null; +} + +export interface ConfirmGatewayPortResult { + released: boolean; + remaining: number[]; +} + +/** + * Bind loopback in a child so this synchronous stop path can prove the port is + * free. Node's in-process net.Server reports bind success/failure + * asynchronously; using it here would require making the full stop API async. + * The child performs one bind and confirmGatewayPortReleased invokes it only + * once, after any authoritative listener scan has cleared. + */ +export function defaultProbePortFree(port: number): boolean { + try { + return ( + spawnSync(process.execPath, ["-e", PORT_FREE_PROBE_SCRIPT, String(port)], { + stdio: "ignore", + timeout: 2000, + }).status === 0 + ); + } catch { + return false; + } +} + +/** + * Confirm both observation layers agree: lsof sees no listener (when + * available) and an independent bind succeeds. The bind is required even + * after an empty lsof result because unprivileged lsof can hide root-owned + * listeners. Listener polling is bounded by both deadline and attempt count; + * the independent bind subprocess runs exactly once. + */ +export function confirmGatewayPortReleased( + options: ConfirmGatewayPortOptions, +): ConfirmGatewayPortResult { + let remaining: number[] = []; + const listeningPids = options.listeningPids; + const listenersReleased = listeningPids + ? waitUntil( + () => { + const pids = listeningPids(); + if (pids === null) return false; + remaining = pids; + return pids.length === 0; + }, + { + deadlineMs: options.now() + options.timeoutMs, + maxAttempts: 20, + initialIntervalMs: options.pollIntervalMs, + maxIntervalMs: options.pollIntervalMs, + backoffFactor: 1, + now: options.now, + ...(options.sleep ? { sleep: options.sleep } : {}), + }, + ) + : true; + if (!listenersReleased) return { released: false, remaining }; + return { released: options.probePortFree(options.port), remaining }; +} diff --git a/src/lib/tunnel/gateway-port-listeners.test.ts b/src/lib/tunnel/gateway-port-listeners.test.ts new file mode 100644 index 00000000000..29deb995e56 --- /dev/null +++ b/src/lib/tunnel/gateway-port-listeners.test.ts @@ -0,0 +1,33 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +import { afterEach, describe, expect, it } from "vitest"; + +import { defaultGatewayReleaseCommandExists } from "./gateway-port-listeners"; + +const tempDirs: string[] = []; + +afterEach(() => { + for (const directory of tempDirs.splice(0)) { + fs.rmSync(directory, { force: true, recursive: true }); + } +}); + +describe("defaultGatewayReleaseCommandExists", () => { + it("finds an executable directly on the configured PATH", () => { + const directory = fs.mkdtempSync(path.join(os.tmpdir(), "nemoclaw-command-path-")); + tempDirs.push(directory); + const executable = path.join(directory, "lsof"); + fs.writeFileSync(executable, "#!/bin/sh\nexit 0\n", { mode: 0o700 }); + + expect(defaultGatewayReleaseCommandExists("lsof", { PATH: directory })).toBe(true); + }); + + it("does not invoke a shell or trust an empty PATH entry", () => { + expect(defaultGatewayReleaseCommandExists("lsof; exit 0", { PATH: "" })).toBe(false); + }); +}); diff --git a/src/lib/tunnel/gateway-port-listeners.ts b/src/lib/tunnel/gateway-port-listeners.ts new file mode 100644 index 00000000000..5322c0792d1 --- /dev/null +++ b/src/lib/tunnel/gateway-port-listeners.ts @@ -0,0 +1,58 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { type SpawnSyncOptions, spawnSync } from "node:child_process"; +import fs from "node:fs"; +import path from "node:path"; + +import type { HostGatewayProcessDeps, RunResult } from "../onboard/host-gateway-process"; + +export function defaultGatewayReleaseRun( + command: string, + args: string[], + options: SpawnSyncOptions = {}, +): RunResult { + const result = spawnSync(command, args, { encoding: "utf-8", ...options }); + return { + status: result.status, + stdout: typeof result.stdout === "string" ? result.stdout : String(result.stdout ?? ""), + stderr: typeof result.stderr === "string" ? result.stderr : String(result.stderr ?? ""), + }; +} + +export function defaultGatewayReleaseCommandExists( + command: string, + env: NodeJS.ProcessEnv, +): boolean { + // Resolve the internal literal (currently only "lsof") directly from PATH. + // Ignore empty entries instead of treating the working directory as trusted. + return (env.PATH ?? "") + .split(path.delimiter) + .filter(Boolean) + .some((directory) => { + try { + fs.accessSync(path.join(directory, command), fs.constants.X_OK); + return true; + } catch { + return false; + } + }); +} + +export function listeningGatewayPids( + port: number, + run: NonNullable, + env: NodeJS.ProcessEnv, + warn: (message: string) => void, +): number[] | null { + const result = run("lsof", ["-ti", `:${port}`, "-sTCP:LISTEN"], { env }); + if (result.status !== 0 && result.status !== 1) { + const detail = result.stderr.trim() || `status ${String(result.status)}`; + warn(`lsof failed while scanning gateway port ${port}: ${detail}`); + return null; + } + return result.stdout + .split(/\r?\n/) + .map((line) => Number.parseInt(line.trim(), 10)) + .filter((pid) => Number.isInteger(pid) && pid > 0); +} diff --git a/src/lib/tunnel/gateway-port-release-fail-closed.test.ts b/src/lib/tunnel/gateway-port-release-fail-closed.test.ts new file mode 100644 index 00000000000..b0fa727eca5 --- /dev/null +++ b/src/lib/tunnel/gateway-port-release-fail-closed.test.ts @@ -0,0 +1,282 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { describe, expect, it, vi } from "vitest"; + +import { DEFAULT_GATEWAY_PORT } from "../core/ports"; +import type { HostGatewayProcessDeps } from "../onboard/host-gateway-process"; +import { releaseManagedGatewayPort } from "./gateway-port-release"; +import { + baseDeps, + emptyStopResult, + lsofResponder, + ok, + stopSpy, +} from "./gateway-port-release-test-helpers"; + +describe("releaseManagedGatewayPort fail-closed behavior (#5968)", () => { + it("does not fall back to the default port when the persisted gateway binding is invalid", () => { + // Source-of-truth guard: a corrupt registry entry must NOT cause + // default-port cleanup or any stopHostGatewayProcesses invocation. + const lsof = lsofResponder(ok("999\n")); + const stop = stopSpy(emptyStopResult()); + const warn = vi.fn(); + + const result = releaseManagedGatewayPort( + { sandboxName: "nemoclaw-5968" }, + { + ...baseDeps(), + warn, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => ({ gatewayPort: 0 }), + }, + ); + + expect(result.skipped).toBe(true); + expect(result.released).toBe(false); + expect(result.port).toBe(null); + expect(stop.lastOptions()).toBeUndefined(); + expect(lsof.calls).toBe(0); + expect(warn.mock.calls.map((c) => c[0]).join("\n")).toContain( + "no valid gateway binding is registered", + ); + }); + + it("skips default-port cleanup for a named sandbox whose registry entry is absent", () => { + // A named stop with no registry entry must not scan or signal the + // process-wide default gateway, which could belong to another worktree. + const lsof = lsofResponder(ok("777\n")); + const stop = stopSpy(emptyStopResult()); + const warn = vi.fn(); + + const result = releaseManagedGatewayPort( + { sandboxName: "no-such-sandbox" }, + { + ...baseDeps(), + warn, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.skipped).toBe(true); + expect(result.released).toBe(false); + expect(result.port).toBe(null); + expect(stop.lastOptions()).toBeUndefined(); + expect(lsof.calls).toBe(0); + expect(warn.mock.calls.map((c) => c[0]).join("\n")).toContain( + "no valid gateway binding is registered", + ); + }); + + it("emits a NODE_DEBUG=nemoclaw:gateway diagnostic when the fail-closed path is taken", () => { + // The default warning stays concise; NODE_DEBUG adds the underlying cause. + const errorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); + const stop = stopSpy(emptyStopResult()); + + releaseManagedGatewayPort( + { sandboxName: "alpha" }, + { + ...baseDeps(), + env: { HOME: "/home/tester", NODE_DEBUG: "nemoclaw:gateway" } as NodeJS.ProcessEnv, + run: lsofResponder(ok("999\n")).run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => { + throw new Error("corrupt registry"); + }, + }, + ); + + expect(errorSpy.mock.calls.map((c) => String(c[0])).join("\n")).toContain( + "[nemoclaw:gateway] registry lookup for sandbox", + ); + errorSpy.mockRestore(); + }); + + it("skips the destructive path when the registry lookup throws", () => { + const lsof = lsofResponder(ok("888\n")); + const stop = stopSpy(emptyStopResult()); + const warn = vi.fn(); + + const result = releaseManagedGatewayPort( + { sandboxName: "nemoclaw-5968" }, + { + ...baseDeps(), + warn, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => { + throw new Error("corrupt registry"); + }, + }, + ); + + expect(result.skipped).toBe(true); + expect(result.released).toBe(false); + expect(stop.lastOptions()).toBeUndefined(); + expect(lsof.calls).toBe(0); + expect(warn.mock.calls.map((call) => String(call[0])).join("\n")).toContain( + "Registry lookup failed for sandbox", + ); + }); + + it("warns and refuses unsafe pid-file cleanup when lsof exits with a real failure", () => { + // lsof status > 1 is a genuine error (not "no listeners"); surface it and + // do not treat unverified PID-file contents as signal-safe candidates. + const stop = stopSpy(emptyStopResult()); + const warn = vi.fn(); + const run: NonNullable = (command) => + command === "lsof" ? { status: 2, stdout: "", stderr: "lsof: boom" } : ok(); + + const result = releaseManagedGatewayPort( + {}, + { + ...baseDeps(), + warn, + run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.scanned).toBe(false); + // A genuine lsof error (vs lsof simply being absent) means we confirmed + // nothing, so the port must not be reported as released (#5968): this is what + // lets stopAll surface its unconfirmed-release warning. + expect(result.released).toBe(false); + expect(stop.lastOptions()?.pids).toEqual([]); + expect(stop.lastOptions()?.usePidFile).toBe(false); + expect(warn.mock.calls.map((c) => c[0]).join("\n")).toContain("lsof failed while scanning"); + }); + + it("does not report released when the confirmation probe itself fails", () => { + // Port is bound on the initial scan (so the stop path runs), but lsof + // errors on every confirmation probe. A failed probe is not proof the port + // is free, so released must stay false rather than coercing null -> []. + const lsof = lsofResponder(ok("555\n"), { status: 2, stdout: "", stderr: "boom" }); + const stop = stopSpy(emptyStopResult({ stopped: [555] })); + + const result = releaseManagedGatewayPort( + { confirmTimeoutMs: 10 }, + { + ...baseDeps(), + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.released).toBe(false); + }); + + it("skips unsafe pid-file cleanup and relies on the bind proof when lsof is absent", () => { + const stop = stopSpy(emptyStopResult()); + const run = vi.fn(() => ok()); + const probePortFree = vi.fn(() => true); + + const result = releaseManagedGatewayPort( + { sandboxName: "nemoclaw-5968" }, + { + ...baseDeps(), + commandExists: () => false, + probePortFree, + run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => ({ gatewayPort: DEFAULT_GATEWAY_PORT }), + }, + ); + + expect(result.scanned).toBe(false); + expect(result.released).toBe(true); + expect(stop.lastOptions()?.pids).toEqual([]); + expect(stop.lastOptions()?.usePidFile).toBe(false); + expect(run).not.toHaveBeenCalled(); + expect(probePortFree).toHaveBeenCalledWith(DEFAULT_GATEWAY_PORT); + }); + + it("does not report release without lsof when an unrecorded listener still owns the port", () => { + const stop = stopSpy(emptyStopResult({ stopped: [111] })); + const log = vi.fn(); + + const result = releaseManagedGatewayPort( + { confirmTimeoutMs: 10 }, + { + ...baseDeps(), + commandExists: () => false, + probePortFree: () => false, + log, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.scanned).toBe(false); + expect(result.released).toBe(false); + expect(result.stopped).toEqual([111]); + expect(log).not.toHaveBeenCalledWith( + expect.stringContaining(`Released NemoClaw gateway port ${DEFAULT_GATEWAY_PORT}`), + ); + }); + + it("runs one bind proof for one managed gateway release", () => { + let clock = 0; + const probePortFree = vi.fn(() => false); + + const result = releaseManagedGatewayPort( + { confirmTimeoutMs: 100_000, confirmPollIntervalMs: 1 }, + { + ...baseDeps(), + commandExists: () => false, + now: () => clock++, + probePortFree, + stopHostGatewayProcesses: stopSpy(emptyStopResult()).fn, + }, + ); + + expect(result.released).toBe(false); + expect(probePortFree).toHaveBeenCalledTimes(1); + }); + + it("does not trust empty lsof output when a hidden listener prevents rebinding", () => { + const stop = stopSpy(emptyStopResult()); + const probePortFree = vi.fn(() => false); + const lsof = lsofResponder({ status: 1, stdout: "", stderr: "" }); + + const result = releaseManagedGatewayPort( + { confirmTimeoutMs: 10 }, + { + ...baseDeps(), + probePortFree, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.scanned).toBe(true); + expect(result.released).toBe(false); + expect(probePortFree).toHaveBeenCalledWith(DEFAULT_GATEWAY_PORT); + }); + + it("never reports release when a matched gateway could not be stopped", () => { + const stop = stopSpy(emptyStopResult({ failed: [777] })); + const probePortFree = vi.fn(() => true); + + const result = releaseManagedGatewayPort( + {}, + { + ...baseDeps(), + probePortFree, + run: lsofResponder({ status: 1, stdout: "", stderr: "" }).run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.released).toBe(false); + expect(result.remaining).toEqual([777]); + expect(probePortFree).not.toHaveBeenCalled(); + }); +}); diff --git a/src/lib/tunnel/gateway-port-release-lifecycle.test.ts b/src/lib/tunnel/gateway-port-release-lifecycle.test.ts new file mode 100644 index 00000000000..e2068aabbd7 --- /dev/null +++ b/src/lib/tunnel/gateway-port-release-lifecycle.test.ts @@ -0,0 +1,186 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import path from "node:path"; +import { describe, expect, it, vi } from "vitest"; + +import { DEFAULT_GATEWAY_PORT } from "../core/ports"; +import type { HostGatewayProcessDeps } from "../onboard/host-gateway-process"; +import { releaseManagedGatewayPort } from "./gateway-port-release"; +import { + baseDeps, + emptyStopResult, + lsofResponder, + ok, + stopSpy, +} from "./gateway-port-release-test-helpers"; + +describe("releaseManagedGatewayPort lifecycle (#5968)", () => { + it("stops lsof-discovered gateways, then reports the port released", () => { + const lsof = lsofResponder(ok("111\n222\n"), ok("")); + const stop = stopSpy(emptyStopResult({ stopped: [111, 222] })); + + const log = vi.fn(); + const result = releaseManagedGatewayPort( + { sandboxName: "nemoclaw-5968", confirmTimeoutMs: 1000 }, + { + ...baseDeps(), + log, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => ({ gatewayPort: DEFAULT_GATEWAY_PORT }), + }, + ); + + expect(result.released).toBe(true); + expect(result.port).toBe(DEFAULT_GATEWAY_PORT); + expect(result.stopped).toEqual([111, 222]); + + expect(stop.fn).toHaveBeenCalledTimes(1); + const stopOptions = stop.lastOptions(); + expect(stopOptions?.pids).toEqual([111, 222]); + expect(stopOptions?.usePgrepFallback).toBe(false); + expect(stopOptions?.usePidFile).toBe(false); + expect(stopOptions?.stateDir).toBe( + path.join("/home/tester", ".local", "state", "nemoclaw", "openshell-docker-gateway"), + ); + expect(log.mock.calls.map((c) => c[0]).join("\n")).toContain( + `Released NemoClaw gateway port ${DEFAULT_GATEWAY_PORT}`, + ); + }); + + it("scopes the sweep to the sandbox's own gateway port so another worktree's gateway is untouched", () => { + // Cross-worktree isolation: a stop for sandbox A (port 8090) must only ever + // probe :8090 and target the 8090 state dir, and must never run a host-wide + // pgrep sweep — so sandbox B's gateway on a different port is never reaped. + const calls: string[][] = []; + const run: NonNullable = (command, args) => { + calls.push([command, ...args]); + return ok("8190\n"); + }; + const stop = stopSpy(emptyStopResult({ stopped: [8190] })); + + releaseManagedGatewayPort( + { sandboxName: "alpha", confirmTimeoutMs: 5 }, + { + ...baseDeps(), + run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => ({ gatewayPort: 8090 }), + }, + ); + + const lsofCalls = calls.filter((c) => c[0] === "lsof"); + expect(lsofCalls.length).toBeGreaterThan(0); + expect(lsofCalls.every((c) => c.includes(":8090"))).toBe(true); + expect(lsofCalls.some((c) => c.includes(":8091"))).toBe(false); + expect(stop.lastOptions()?.usePgrepFallback).toBe(false); + expect(stop.lastOptions()?.stateDir).toContain("openshell-docker-gateway-8090"); + }); + + it("targets the per-port state dir for a non-default gateway port", () => { + const lsof = lsofResponder(ok("")); + const stop = stopSpy(emptyStopResult()); + + releaseManagedGatewayPort( + { sandboxName: "nemoclaw-5968" }, + { + ...baseDeps(), + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => ({ gatewayPort: 8090 }), + }, + ); + + expect(stop.lastOptions()?.stateDir).toBe( + path.join("/home/tester", ".local", "state", "nemoclaw", "openshell-docker-gateway-8090"), + ); + }); + + it("is a quiet no-op when nothing is bound to the gateway port", () => { + const lsof = lsofResponder(ok("")); + const stop = stopSpy(emptyStopResult()); + const log = vi.fn(); + const warn = vi.fn(); + + const result = releaseManagedGatewayPort( + {}, + { + ...baseDeps(), + log, + warn, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.released).toBe(true); + expect(log).not.toHaveBeenCalled(); + expect(warn).not.toHaveBeenCalled(); + }); + + it("does not trust a per-port pid file as proof that its process owns the port", () => { + const lsof = lsofResponder(ok("222\n"), ok("")); + const stop = stopSpy(emptyStopResult({ stopped: [222] })); + + releaseManagedGatewayPort( + { sandboxName: "alpha" }, + { + ...baseDeps(), + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => ({ gatewayPort: 8090 }), + }, + ); + + expect(stop.lastOptions()?.pids).toEqual([222]); + expect(stop.lastOptions()?.usePidFile).toBe(false); + }); + + it("warns with sudo remediation when the port stays bound after stop", () => { + // lsof keeps reporting a listener even after the stop attempt — the orphan + // could not be reaped (e.g. a privileged process). + const lsof = lsofResponder(ok("333\n")); + const stop = stopSpy(emptyStopResult({ failed: [333] })); + const warn = vi.fn(); + + const result = releaseManagedGatewayPort( + { confirmTimeoutMs: 10 }, + { + ...baseDeps(), + warn, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.released).toBe(false); + expect(result.remaining).toEqual([333]); + expect(warn.mock.calls.map((c) => c[0]).join("\n")).toContain("sudo kill -9 333"); + }); + + it("leaves a non-matching listener alone without sudo pkill remediation", () => { + // lsof reports a PID the stopper classifies as non-matching (e.g. a + // Docker-published port held by docker-proxy). No matched gateway failed, + // so no scary remediation hint. + const lsof = lsofResponder(ok("444\n"), ok("444\n")); + const stop = stopSpy(emptyStopResult({ skippedNonMatchingPids: [444] })); + const warn = vi.fn(); + + const result = releaseManagedGatewayPort( + { confirmTimeoutMs: 10 }, + { + ...baseDeps(), + warn, + run: lsof.run, + stopHostGatewayProcesses: stop.fn, + getSandbox: () => null, + }, + ); + + expect(result.released).toBe(false); + expect(warn).not.toHaveBeenCalled(); + }); +}); diff --git a/src/lib/tunnel/gateway-port-release-test-helpers.ts b/src/lib/tunnel/gateway-port-release-test-helpers.ts new file mode 100644 index 00000000000..035d49b8ee4 --- /dev/null +++ b/src/lib/tunnel/gateway-port-release-test-helpers.ts @@ -0,0 +1,100 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { vi } from "vitest"; + +import type { + HostGatewayProcessDeps, + RunResult, + StopHostGatewayOptions, + StopHostGatewayResult, +} from "../onboard/host-gateway-process"; +import type { ReleaseGatewayPortDeps } from "./gateway-port-release"; + +export function emptyStopResult( + overrides: Partial = {}, +): StopHostGatewayResult { + return { + failed: [], + skippedDeadPids: [], + skippedNonMatchingPids: [], + stopped: [], + sudoRemediationPids: [], + ...overrides, + }; +} + +export function ok(stdout = ""): RunResult { + return { status: 0, stdout, stderr: "" }; +} + +type StopFn = ( + depsOverrides?: Partial, + options?: StopHostGatewayOptions, +) => StopHostGatewayResult; + +// Build a host-gateway stopper mock that records the options it was called +// with. The explicit StopFn type keeps it assignable to the real +// (optional-param) signature, and capturing in a closure avoids fragile tuple +// indexing. +export function stopSpy(result: StopHostGatewayResult): { + fn: StopFn; + lastOptions: () => StopHostGatewayOptions | undefined; +} { + let captured: StopHostGatewayOptions | undefined; + const fn: StopFn = vi.fn( + (_deps?: Partial, options?: StopHostGatewayOptions) => { + captured = options; + return result; + }, + ); + return { fn, lastOptions: () => captured }; +} + +// A queued `lsof` responder so a test can model the port being held on the +// first probe and free on the confirmation probe. +export function lsofResponder(...responses: RunResult[]): { + run: NonNullable; + calls: number; +} { + const state = { calls: 0 }; + const run: NonNullable = (command) => { + const isLsof = command === "lsof"; + const idx = Math.min(state.calls, responses.length - 1); + const response = isLsof ? (responses[idx] ?? ok()) : ok(); + state.calls += isLsof ? 1 : 0; + return response; + }; + return { + run, + get calls() { + return state.calls; + }, + }; +} + +// Advancing fake clock so the confirmation poll's deadline is always reached — +// a constant clock would make `waitUntil` spin forever when the port never +// frees. +function clock(step = 1): () => number { + let t = 0; + return () => { + const v = t; + t += step; + return v; + }; +} + +export function baseDeps(): ReleaseGatewayPortDeps { + return { + env: { HOME: "/home/tester" } as NodeJS.ProcessEnv, + homeDir: "/home/tester", + commandExists: () => true, + kill: () => true, + now: clock(), + sleep: () => {}, + probePortFree: () => true, + log: () => {}, + warn: () => {}, + }; +} diff --git a/src/lib/tunnel/gateway-port-release.test.ts b/src/lib/tunnel/gateway-port-release.test.ts new file mode 100644 index 00000000000..7a512c6dd1c --- /dev/null +++ b/src/lib/tunnel/gateway-port-release.test.ts @@ -0,0 +1,60 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { describe, expect, it } from "vitest"; + +import { DEFAULT_GATEWAY_PORT } from "../core/ports"; +import { resolveStopGatewayPort } from "./gateway-port-release"; + +describe("resolveStopGatewayPort (#5968)", () => { + it("prefers an explicit port override", () => { + expect(resolveStopGatewayPort({ port: 9090 }, () => null)).toBe(9090); + }); + + it("fails closed (null) for an explicit but invalid port override", () => { + // An out-of-range override must not silently fall through to the sandbox + // binding or the default port — it is a caller error, so skip. + expect(resolveStopGatewayPort({ port: 70000 }, () => ({ gatewayPort: 8090 }))).toBe(null); + expect(resolveStopGatewayPort({ port: 0, sandboxName: "alpha" }, () => null)).toBe(null); + }); + + it("derives the port from the sandbox's persisted gateway binding", () => { + const port = resolveStopGatewayPort({ sandboxName: "alpha" }, () => ({ gatewayPort: 8090 })); + expect(port).toBe(8090); + }); + + it("fails closed (null) when a named sandbox has no registry entry", () => { + // A named stop whose registry entry is absent must not fall back to + // default-port cleanup: an unknown name could otherwise tear down a + // different sandbox's / worktree's default gateway. + expect(resolveStopGatewayPort({ sandboxName: "alpha" }, () => null)).toBe(null); + }); + + it("falls back to the default gateway port for a call with no sandbox name", () => { + // A direct "release the default gateway" request (no sandbox identity). + expect(resolveStopGatewayPort({}, () => null)).toBe(DEFAULT_GATEWAY_PORT); + }); + + it("falls back to the default gateway port for a legacy entry with no gateway fields", () => { + // A real legacy entry (e.g. `{}`) maps to the base `nemoclaw` name and + // resolves to the default port, keeping single-sandbox deployments working. + expect(resolveStopGatewayPort({ sandboxName: "alpha" }, () => ({}))).toBe(DEFAULT_GATEWAY_PORT); + }); + + it("fails closed (null) when the persisted gateway binding is invalid", () => { + // An out-of-range gatewayPort is a corrupt/tampered binding; + // resolveSandboxGatewayName throws and we must not coerce to the default. + expect(resolveStopGatewayPort({ sandboxName: "alpha" }, () => ({ gatewayPort: 70000 }))).toBe( + null, + ); + }); + + it("fails closed (null) when the registry lookup itself throws", () => { + // A corrupt registry that throws on read must not be treated as a clean + // "no entry" and fall back to the default port. + const port = resolveStopGatewayPort({ sandboxName: "alpha" }, () => { + throw new Error("corrupt registry"); + }); + expect(port).toBe(null); + }); +}); diff --git a/src/lib/tunnel/gateway-port-release.ts b/src/lib/tunnel/gateway-port-release.ts new file mode 100644 index 00000000000..d51fe772a81 --- /dev/null +++ b/src/lib/tunnel/gateway-port-release.ts @@ -0,0 +1,166 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +/** Port-scoped host gateway release for `nemoclaw stop` (#5968). */ + +import os from "node:os"; + +import type { SandboxGatewayBinding } from "../onboard/gateway-binding"; +import { + type HostGatewayProcessDeps, + type StopHostGatewayResult, + stopHostGatewayProcesses, +} from "../onboard/host-gateway-process"; +import { getSandbox as getRegisteredSandbox } from "../state/registry"; +import { confirmGatewayPortReleased, defaultProbePortFree } from "./gateway-port-confirmation"; +import { + defaultGatewayReleaseCommandExists, + defaultGatewayReleaseRun, + listeningGatewayPids, +} from "./gateway-port-listeners"; +import { + makeGatewayDebug, + resolveGatewayReleaseStateDir, + resolveStopGatewayPort, +} from "./gateway-port-resolution"; + +export { resolveStopGatewayPort }; + +export interface ReleaseGatewayPortDeps extends Partial { + homeDir?: string; + now?: () => number; + sleep?: (ms: number) => void; + stopHostGatewayProcesses?: typeof stopHostGatewayProcesses; + getSandbox?: (name: string) => SandboxGatewayBinding | null; + probePortFree?: (port: number) => boolean; +} + +export interface ReleaseGatewayPortOptions { + sandboxName?: string; + port?: number; + confirmTimeoutMs?: number; + confirmPollIntervalMs?: number; +} + +export interface ReleaseGatewayPortResult { + port: number | null; + released: boolean; + stopped: number[]; + remaining: number[]; + scanned: boolean; + skipped: boolean; +} + +/** + * Stop cmdline-verified gateway listeners on the selected port and prove the + * port can be rebound. PID-file contents are never signal candidates on their + * own: only lsof-observed PIDs are passed to the stopper, preventing a stale or + * recycled PID from killing another worktree's same-named gateway. + */ +export function releaseManagedGatewayPort( + options: ReleaseGatewayPortOptions = {}, + depsOverrides: ReleaseGatewayPortDeps = {}, +): ReleaseGatewayPortResult { + const env = depsOverrides.env ?? process.env; + const homeDir = depsOverrides.homeDir ?? env.HOME ?? os.homedir(); + const run = depsOverrides.run ?? defaultGatewayReleaseRun; + const log = depsOverrides.log ?? ((message: string) => console.log(message)); + const warn = depsOverrides.warn ?? ((message: string) => console.warn(message)); + const commandExists = + depsOverrides.commandExists ?? + ((command: string) => defaultGatewayReleaseCommandExists(command, env)); + const stop = depsOverrides.stopHostGatewayProcesses ?? stopHostGatewayProcesses; + const getSandbox = depsOverrides.getSandbox ?? getRegisteredSandbox; + const probePortFree = depsOverrides.probePortFree ?? defaultProbePortFree; + + const port = resolveStopGatewayPort(options, getSandbox, makeGatewayDebug(env), warn); + if (port === null) { + warn( + `Skipping gateway port release for sandbox ${JSON.stringify(options.sandboxName)}: ` + + "no valid gateway binding is registered for it (the entry is missing, " + + "invalid, or unreadable). Resolve the registry entry, then re-run stop.", + ); + return { + port: null, + released: false, + stopped: [], + remaining: [], + scanned: false, + skipped: true, + }; + } + + const stateDir = resolveGatewayReleaseStateDir(port, env, homeDir); + let lsofPids: number[] = []; + let scanned = false; + // The two lsof failure stages fail closed differently. An initial failure + // leaves the destructive candidate scan incomplete, so confirmation is + // skipped entirely: a later successful bind cannot make that scan complete. + // When this initial scan succeeds, a later confirmation failure is retried + // by confirmGatewayPortReleased and never treated as an empty listener set. + let scanFailed = false; + if (commandExists("lsof")) { + const result = listeningGatewayPids(port, run, env, warn); + if (result === null) scanFailed = true; + else { + lsofPids = result; + scanned = true; + } + } + + const hostDeps: Partial = { env }; + if (depsOverrides.run) hostDeps.run = depsOverrides.run; + if (depsOverrides.kill) hostDeps.kill = depsOverrides.kill; + if (depsOverrides.commandExists) hostDeps.commandExists = depsOverrides.commandExists; + if (depsOverrides.log) hostDeps.log = depsOverrides.log; + if (depsOverrides.warn) hostDeps.warn = depsOverrides.warn; + + const stopResult: StopHostGatewayResult = stop(hostDeps, { + stateDir, + pids: lsofPids, + // A per-port PID file is bookkeeping, not proof that its PID owns this + // port. Only the lsof-observed, cmdline-gated candidates are signal-safe. + usePidFile: false, + usePgrepFallback: false, + }); + + // Stage 1: scanFailed=true selects the fallback result below, so an initial + // lsof error never falls through to bind-only confirmation. Stage 2: after a + // successful initial scan, listeningGatewayPids() returning null makes each + // confirmation attempt false; exhaustion returns released=false. Neither + // failure is ever coerced to an empty listener set. + const confirmation = + !scanFailed && stopResult.failed.length === 0 + ? confirmGatewayPortReleased({ + port, + timeoutMs: options.confirmTimeoutMs ?? 2000, + pollIntervalMs: options.confirmPollIntervalMs ?? 100, + now: depsOverrides.now ?? Date.now, + ...(depsOverrides.sleep ? { sleep: depsOverrides.sleep } : {}), + probePortFree, + ...(scanned ? { listeningPids: () => listeningGatewayPids(port, run, env, warn) } : {}), + }) + : { released: false, remaining: stopResult.failed }; + + if (confirmation.released && stopResult.stopped.length > 0) { + log( + `Released NemoClaw gateway port ${port} (stopped host process ${stopResult.stopped.join(", ")}).`, + ); + } + if (stopResult.failed.length > 0) { + warn( + `NemoClaw gateway port ${port} is still in use after stop ` + + `(host process ${stopResult.failed.join(", ")} could not be stopped). ` + + `Run: sudo kill -9 ${stopResult.failed.join(" ")}`, + ); + } + + return { + port, + released: confirmation.released, + stopped: stopResult.stopped, + remaining: confirmation.remaining, + scanned, + skipped: false, + }; +} diff --git a/src/lib/tunnel/gateway-port-resolution.ts b/src/lib/tunnel/gateway-port-resolution.ts new file mode 100644 index 00000000000..e99fb2788d8 --- /dev/null +++ b/src/lib/tunnel/gateway-port-resolution.ts @@ -0,0 +1,78 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import path from "node:path"; + +import { GATEWAY_PORT } from "../core/ports"; +import { + resolveGatewayPortFromName, + resolveGatewayStateDirName, + resolveSandboxGatewayName, + type SandboxGatewayBinding, +} from "../onboard/gateway-binding"; +import type { ReleaseGatewayPortOptions } from "./gateway-port-release"; + +function isValidPort(value: number | undefined): value is number { + return typeof value === "number" && Number.isInteger(value) && value >= 1 && value <= 65535; +} + +export function makeGatewayDebug(env: NodeJS.ProcessEnv): (message: string) => void { + const enabled = (env.NODE_DEBUG ?? "").includes("nemoclaw:gateway"); + return enabled ? (message: string) => console.error(`[nemoclaw:gateway] ${message}`) : () => {}; +} + +/** + * Resolve the selected sandbox's persisted gateway port. Because the caller + * will signal processes, every missing, unreadable, or invalid named-sandbox + * binding fails closed instead of falling back to another sandbox's default + * port. A no-name call and a valid legacy row may still use GATEWAY_PORT. + */ +export function resolveStopGatewayPort( + options: ReleaseGatewayPortOptions, + getSandbox: (name: string) => SandboxGatewayBinding | null, + debug: (message: string) => void = () => {}, + warn: (message: string) => void = () => {}, +): number | null { + if (options.port !== undefined) return isValidPort(options.port) ? options.port : null; + if (!options.sandboxName) return GATEWAY_PORT; + + let entry: SandboxGatewayBinding | null; + try { + entry = getSandbox(options.sandboxName); + } catch (error) { + // Source boundary: the registry write path should guarantee readable data. + // Keep this guard until that path also validates/heals pre-existing rows. + warn( + `Registry lookup failed for sandbox ${JSON.stringify(options.sandboxName)}; ` + + "skipping gateway release. Run with NODE_DEBUG=nemoclaw:gateway for details.", + ); + debug( + `registry lookup for sandbox ${JSON.stringify(options.sandboxName)} threw; ` + + `skipping gateway release: ${(error as Error).message ?? String(error)}`, + ); + return null; + } + if (!entry) return null; + + try { + return resolveGatewayPortFromName(resolveSandboxGatewayName(entry)); + } catch (error) { + // Source boundary: onboard/registry writes validate new bindings, but old + // or tampered rows can still exist. Never coerce one to the default port. + debug( + `persisted gateway binding for sandbox ${JSON.stringify(options.sandboxName)} is invalid; ` + + `skipping gateway release: ${(error as Error).message ?? String(error)}`, + ); + return null; + } +} + +export function resolveGatewayReleaseStateDir( + port: number, + env: NodeJS.ProcessEnv, + homeDir: string, +): string { + const configured = env.NEMOCLAW_OPENSHELL_GATEWAY_STATE_DIR; + if (configured && configured.trim()) return path.resolve(configured.trim()); + return path.join(homeDir, ".local", "state", "nemoclaw", resolveGatewayStateDirName(port)); +} diff --git a/src/lib/tunnel/gateway-stop.ts b/src/lib/tunnel/gateway-stop.ts new file mode 100644 index 00000000000..c3df9d5db9e --- /dev/null +++ b/src/lib/tunnel/gateway-stop.ts @@ -0,0 +1,118 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { resolveGatewayPortFromName, resolveSandboxGatewayName } from "../onboard/gateway-binding"; +import * as registry from "../state/registry"; +import * as gatewayPortRelease from "./gateway-port-release"; + +type Log = (message: string) => void; + +export interface GatewayStopDeps { + env?: NodeJS.ProcessEnv; + info?: Log; + warn?: Log; + listSandboxes?: typeof registry.listSandboxes; + releaseManagedGatewayPort?: typeof gatewayPortRelease.releaseManagedGatewayPort; +} + +type SharedGatewayOwner = { + name: string; + port: number; +}; + +/** + * Find another registered sandbox that may still own the selected sandbox's + * host gateway. The registry is intentionally a conservative ownership + * signal, matching destroy's last-sandbox gate: a stale registration can keep + * a gateway alive, but tearing it down while a registered peer is live would + * break that peer. + * + * A missing selected entry is left to releaseManagedGatewayPort(), whose + * sandbox-specific resolver already fails closed. Invalid or unreadable + * registry state throws into the best-effort catch below so teardown is + * skipped rather than guessed. + */ +function findSharedGatewayOwner( + sandboxName: string, + listSandboxes: typeof registry.listSandboxes, +): SharedGatewayOwner | null { + const sandboxes = listSandboxes().sandboxes; + const selected = sandboxes.find((sandbox) => sandbox.name === sandboxName); + if (!selected) return null; + + const gatewayName = resolveSandboxGatewayName(selected); + const port = resolveGatewayPortFromName(gatewayName); + if (port === null) { + throw new Error(`Could not resolve gateway port for registered sandbox ${sandboxName}`); + } + + for (const sandbox of sandboxes) { + if (sandbox.name === sandboxName) continue; + try { + if (resolveSandboxGatewayName(sandbox) === gatewayName) { + return { name: sandbox.name, port }; + } + } catch (error) { + throw new Error( + `Invalid persisted sandbox gateway for peer '${sandbox.name}': ` + + `${(error as Error).message ?? String(error)}`, + ); + } + } + return null; +} + +/** + * Release the selected sandbox's host gateway only when no registered peer + * shares it. A missing sandbox name is a deliberate no-op: falling back to the + * process-wide default port could tear down another worktree's gateway. + */ +export function releaseGatewayPortForStop( + sandboxName: string | undefined, + deps: GatewayStopDeps = {}, +): void { + if (!sandboxName) return; + + const env = deps.env ?? process.env; + const info = deps.info ?? console.log; + const warn = deps.warn ?? console.warn; + const listSandboxes = deps.listSandboxes ?? registry.listSandboxes; + const releaseManagedGatewayPort = + deps.releaseManagedGatewayPort ?? gatewayPortRelease.releaseManagedGatewayPort; + + try { + const sharedOwner = findSharedGatewayOwner(sandboxName, listSandboxes); + if (sharedOwner) { + info( + `Keeping shared NemoClaw gateway port ${sharedOwner.port} running for ` + + `registered sandbox '${sharedOwner.name}'.`, + ); + return; + } + + const release = releaseManagedGatewayPort({ sandboxName }); + // The release helper reports invalid bindings itself. For an attempted but + // unconfirmed release, do not recommend killing raw lsof PIDs: an unrelated + // listener may be one the scoped stopper deliberately left alone. + if (!release.released && !release.skipped) { + warn( + `NemoClaw gateway port ${release.port ?? "?"} was not confirmed released. ` + + "Inspect the remaining listener and stop it only if it is the matching gateway process.", + ); + } + } catch (error) { + // A corrupt peer registry entry makes gateway ownership ambiguous. Do not + // block the selected sandbox's non-gateway stop work, but skip destructive + // release so a potentially shared gateway is never torn down by guessing. + warn( + `Could not release the NemoClaw gateway port: ${(error as Error).message ?? String(error)}. ` + + "Gateway ownership is ambiguous; repair the sandbox registry and retry. " + + "Run with NODE_DEBUG=nemoclaw:gateway for details.", + ); + // Best-effort by design: keep normal output concise, with the full stack + // available only to an operator explicitly debugging gateway teardown. + if ((env.NODE_DEBUG ?? "").includes("nemoclaw:gateway")) { + console.error((error as Error).stack ?? String(error)); + } + } +} diff --git a/src/lib/tunnel/service-command.test.ts b/src/lib/tunnel/service-command.test.ts index 365b9397e31..93eb8d12ef3 100644 --- a/src/lib/tunnel/service-command.test.ts +++ b/src/lib/tunnel/service-command.test.ts @@ -78,4 +78,14 @@ describe("services command", () => { }); expect(stopAll).toHaveBeenCalledWith({ sandboxName: undefined }); }); + + it("opts the legacy full-stop command into managed gateway release", () => { + const stopAll = vi.fn(); + runStopCommand({ + listSandboxes: () => ({ defaultSandbox: "alpha" }), + stopAll, + releaseGatewayPort: true, + }); + expect(stopAll).toHaveBeenCalledWith({ sandboxName: "alpha", releaseGatewayPort: true }); + }); }); diff --git a/src/lib/tunnel/service-command.ts b/src/lib/tunnel/service-command.ts index 1b5f250ee50..2053aec8c97 100644 --- a/src/lib/tunnel/service-command.ts +++ b/src/lib/tunnel/service-command.ts @@ -12,7 +12,9 @@ export interface StartCommandDeps { export interface StopCommandDeps { listSandboxes: () => SandboxSummary; - stopAll: (options: { sandboxName?: string }) => void; + stopAll: (options: { sandboxName?: string; releaseGatewayPort?: boolean }) => void; + /** Legacy `nemoclaw stop` tears down the managed host gateway too. */ + releaseGatewayPort?: boolean; } const SAFE_SANDBOX_RE = /^[a-zA-Z0-9][a-zA-Z0-9._-]*$/; @@ -33,5 +35,9 @@ export async function runStartCommand(deps: StartCommandDeps): Promise { } export function runStopCommand(deps: StopCommandDeps): void { - deps.stopAll({ sandboxName: resolveDefaultSandboxName(deps.listSandboxes) }); + const options: { sandboxName?: string; releaseGatewayPort?: boolean } = { + sandboxName: resolveDefaultSandboxName(deps.listSandboxes), + }; + if (deps.releaseGatewayPort) options.releaseGatewayPort = true; + deps.stopAll(options); } diff --git a/src/lib/tunnel/services-gateway-ownership.test.ts b/src/lib/tunnel/services-gateway-ownership.test.ts new file mode 100644 index 00000000000..99a37393b11 --- /dev/null +++ b/src/lib/tunnel/services-gateway-ownership.test.ts @@ -0,0 +1,225 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, describe, expect, it, vi } from "vitest"; + +import type { SandboxEntry } from "../state/registry"; +import type { ReleaseGatewayPortResult } from "./gateway-port-release"; +import type { GatewayStopDeps } from "./gateway-stop"; +import * as gatewayStop from "./gateway-stop"; +import { stopAll } from "./services"; + +vi.mock("../adapters/docker", () => ({ + dockerSpawnSync: vi.fn(() => ({ status: 1, stdout: "", stderr: "" })), +})); + +vi.mock("../adapters/openshell/resolve", () => ({ + resolveOpenshell: vi.fn(() => null), +})); + +function sandboxList(sandboxes: SandboxEntry[]): NonNullable { + return vi.fn(() => ({ sandboxes, defaultSandbox: sandboxes[0]?.name ?? null })); +} + +function releaseResult( + overrides: Partial = {}, +): ReleaseGatewayPortResult { + return { + port: 8080, + released: true, + stopped: [], + remaining: [], + scanned: true, + skipped: false, + ...overrides, + }; +} + +function gatewayRelease( + result: ReleaseGatewayPortResult = releaseResult(), +): NonNullable { + return vi.fn(() => result); +} + +describe("releaseGatewayPortForStop", () => { + it("keeps the host gateway when another registered sandbox shares its port", () => { + const release = gatewayRelease(); + const info = vi.fn<(message: string) => void>(); + + gatewayStop.releaseGatewayPortForStop("alpha", { + listSandboxes: sandboxList([ + { name: "alpha", gatewayName: "nemoclaw", gatewayPort: 8080 }, + { name: "beta", gatewayName: "nemoclaw", gatewayPort: 8080 }, + ]), + releaseManagedGatewayPort: release, + info, + }); + + expect(release).not.toHaveBeenCalled(); + expect(info).toHaveBeenCalledWith( + "Keeping shared NemoClaw gateway port 8080 running for registered sandbox 'beta'.", + ); + }); + + it("releases the host gateway for the only registered sandbox", () => { + const release = gatewayRelease(); + + gatewayStop.releaseGatewayPortForStop("alpha", { + listSandboxes: sandboxList([{ name: "alpha", gatewayName: "nemoclaw", gatewayPort: 8080 }]), + releaseManagedGatewayPort: release, + }); + + expect(release).toHaveBeenCalledTimes(1); + expect(release).toHaveBeenCalledWith({ sandboxName: "alpha" }); + }); + + it("releases only the selected port when another sandbox uses a different gateway", () => { + const release = gatewayRelease(); + + gatewayStop.releaseGatewayPortForStop("alpha", { + listSandboxes: sandboxList([ + { name: "alpha", gatewayName: "nemoclaw", gatewayPort: 8080 }, + { name: "beta", gatewayName: "nemoclaw-18080", gatewayPort: 18080 }, + ]), + releaseManagedGatewayPort: release, + }); + + expect(release).toHaveBeenCalledTimes(1); + expect(release).toHaveBeenCalledWith({ sandboxName: "alpha" }); + }); + + it("does not resolve or release a process-wide default without a sandbox name", () => { + const listSandboxes = sandboxList([]); + const release = gatewayRelease(); + + gatewayStop.releaseGatewayPortForStop(undefined, { + listSandboxes, + releaseManagedGatewayPort: release, + }); + + expect(listSandboxes).not.toHaveBeenCalled(); + expect(release).not.toHaveBeenCalled(); + }); + + it("warns without failing stop when gateway release throws", () => { + const release = vi.fn(() => { + throw new Error("registry boom"); + }); + const warn = vi.fn<(message: string) => void>(); + + expect(() => + gatewayStop.releaseGatewayPortForStop("alpha", { + listSandboxes: sandboxList([{ name: "alpha", gatewayPort: 8080 }]), + releaseManagedGatewayPort: release, + warn, + }), + ).not.toThrow(); + + const output = warn.mock.calls.map((call) => call[0]).join("\n"); + expect(output).toContain("Could not release the NemoClaw gateway port: registry boom"); + expect(output).toContain("repair the sandbox registry and retry"); + expect(output).toContain("NODE_DEBUG=nemoclaw:gateway"); + }); + + it("uses inspect-only guidance when release cannot confirm the port is free", () => { + const warn = vi.fn<(message: string) => void>(); + + gatewayStop.releaseGatewayPortForStop("alpha", { + listSandboxes: sandboxList([{ name: "alpha", gatewayPort: 8080 }]), + releaseManagedGatewayPort: gatewayRelease( + releaseResult({ released: false, remaining: [4242] }), + ), + warn, + }); + + const output = warn.mock.calls.map((call) => call[0]).join("\n"); + expect(output).toContain("gateway port 8080 was not confirmed released"); + expect(output).not.toContain("4242"); + expect(output).not.toContain("pkill"); + expect(output).toContain("only if it is the matching gateway process"); + }); + + it("does not duplicate the release helper warning for an invalid binding", () => { + const warn = vi.fn<(message: string) => void>(); + + gatewayStop.releaseGatewayPortForStop("alpha", { + listSandboxes: sandboxList([{ name: "alpha", gatewayPort: 8080 }]), + releaseManagedGatewayPort: gatewayRelease( + releaseResult({ port: null, released: false, scanned: false, skipped: true }), + ), + warn, + }); + + expect(warn).not.toHaveBeenCalled(); + }); + + it("fails closed when a peer has an invalid gateway binding", () => { + const release = gatewayRelease(); + const warn = vi.fn<(message: string) => void>(); + + gatewayStop.releaseGatewayPortForStop("alpha", { + listSandboxes: sandboxList([ + { name: "alpha", gatewayPort: 8080 }, + { name: "beta", gatewayPort: 0 }, + ]), + releaseManagedGatewayPort: release, + warn, + }); + + expect(release).not.toHaveBeenCalled(); + const output = warn.mock.calls.map((call) => call[0]).join("\n"); + expect(output).toContain("Invalid persisted sandbox gateway for peer 'beta'"); + expect(output).toContain("repair the sandbox registry and retry"); + expect(output).toContain("NODE_DEBUG=nemoclaw:gateway"); + }); +}); + +describe("stopAll gateway-stop wiring", () => { + afterEach(() => { + vi.unstubAllEnvs(); + vi.restoreAllMocks(); + }); + + it("passes the resolved sandbox and service reporters to the focused stop module", () => { + const pidDir = mkdtempSync(join(tmpdir(), "nemoclaw-gateway-stop-wiring-")); + vi.stubEnv("PATH", ""); + const releaseForStop = vi + .spyOn(gatewayStop, "releaseGatewayPortForStop") + .mockImplementation(() => {}); + const logSpy = vi.spyOn(console, "log").mockImplementation(() => {}); + + try { + stopAll({ pidDir, sandboxName: "alpha", releaseGatewayPort: true }); + } finally { + rmSync(pidDir, { recursive: true, force: true }); + } + + expect(releaseForStop).toHaveBeenCalledTimes(1); + expect(releaseForStop).toHaveBeenCalledWith("alpha", { + info: expect.any(Function), + warn: expect.any(Function), + }); + expect(logSpy.mock.calls.map((call) => String(call[0] ?? "")).join("\n")).toContain( + "All services stopped", + ); + }); + + it("preserves the shared gateway for canonical tunnel-only stop", () => { + const pidDir = mkdtempSync(join(tmpdir(), "nemoclaw-tunnel-stop-wiring-")); + vi.stubEnv("PATH", ""); + const releaseForStop = vi + .spyOn(gatewayStop, "releaseGatewayPortForStop") + .mockImplementation(() => {}); + + try { + stopAll({ pidDir, sandboxName: "alpha" }); + } finally { + rmSync(pidDir, { recursive: true, force: true }); + } + + expect(releaseForStop).not.toHaveBeenCalled(); + }); +}); diff --git a/src/lib/tunnel/services.ts b/src/lib/tunnel/services.ts index cd458a51a25..3e4dcd20835 100644 --- a/src/lib/tunnel/services.ts +++ b/src/lib/tunnel/services.ts @@ -23,6 +23,7 @@ import { isRecord } from "../core/json-types"; import { DASHBOARD_PORT } from "../core/ports"; import { buildSubprocessEnv } from "../subprocess-env"; import { registerTunnelOrigin } from "./allowed-origins"; +import * as gatewayStop from "./gateway-stop"; // --------------------------------------------------------------------------- // Types @@ -39,6 +40,8 @@ export interface ServiceOptions { pidDir?: string; /** Cloudflare named tunnel token. Falls back to CLOUDFLARE_TUNNEL_TOKEN. */ cloudflareTunnelToken?: string; + /** Also release the managed host gateway port (legacy full-stop only). */ + releaseGatewayPort?: boolean; } export interface ServiceStatus { @@ -620,6 +623,11 @@ export function stopAll(opts: ServiceOptions = {}): void { // Stop host-side services. stopService(pidDir, "cloudflared"); + + if (opts.releaseGatewayPort) { + gatewayStop.releaseGatewayPortForStop(sandboxName, { info, warn }); + } + info("All services stopped."); } diff --git a/test/cli/tunnel-command.test.ts b/test/cli/tunnel-command.test.ts index 17768e9a65f..ed742cda875 100644 --- a/test/cli/tunnel-command.test.ts +++ b/test/cli/tunnel-command.test.ts @@ -66,10 +66,11 @@ describe("tunnel CLI dispatch", () => { expect(r.out).toContain("tunnel status"); }); - it("deprecated stop --help exits 0 and shows alias usage", () => { + it("deprecated stop --help exits 0 and explains legacy full-stop behavior", () => { const r = run("stop --help"); expect(r.code).toBe(0); expect(r.out).toContain("stop"); - expect(r.out).toContain("Deprecated alias"); + expect(r.out).toContain("Deprecated full stop"); + expect(r.out).toContain("releases the managed host gateway port"); }); }); diff --git a/test/onboard-gateway-prelaunch-cutover.test.ts b/test/onboard-gateway-prelaunch-cutover.test.ts new file mode 100644 index 00000000000..06383a7c37d --- /dev/null +++ b/test/onboard-gateway-prelaunch-cutover.test.ts @@ -0,0 +1,201 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +import { describe, expect, it } from "vitest"; + +import { + type DockerDriverGatewayCutoverDeps, + type DockerDriverGatewayCutoverInput, + runDockerDriverGatewayCutover, +} from "../src/lib/onboard/docker-driver-gateway-cutover"; + +type Event = { + type: string; + extraPids?: number[]; + keepPid?: number; + pid?: number; + message?: string; +}; + +interface HarnessOptions { + listenerPids: number[]; + scanComplete?: boolean; + postReapPortAvailable?: boolean; + pidFileGatewayPid?: number | null; + driftPids?: number[]; + prelaunchError?: string; + duplicateError?: string; +} + +function throwHarnessError(message: string): never { + throw new Error(message); +} + +function makeHarness(options: HarnessOptions) { + const events: Event[] = []; + const input: DockerDriverGatewayCutoverInput = { + gatewayBin: "/test/bin/openshell-gateway", + identityGatewayBin: "/test/bin/openshell-gateway", + driftGatewayBin: "/test/bin/openshell-gateway", + driftGatewayEnv: { OPENSHELL_DRIVERS: "docker" }, + exitOnFailure: false, + skipSandboxBridgeReachability: false, + stateDir: "/test/state", + portListenerScan: { + complete: options.scanComplete ?? true, + pids: options.listenerPids, + }, + pidFileGatewayPid: options.pidFileGatewayPid === undefined ? 4242 : options.pidFileGatewayPid, + initialHealth: { + status: "Gateway: nemoclaw\nConnected", + namedInfo: "Gateway: nemoclaw", + activeInfo: "Gateway: nemoclaw", + }, + }; + const driftPids = new Set(options.driftPids ?? []); + const deps: DockerDriverGatewayCutoverDeps = { + isDockerDriverGatewayProcessAlive: () => true, + isGatewayHealthy: () => true, + getDockerDriverGatewayRuntimeDrift: (pid) => + driftPids.has(pid) ? { reason: "test runtime drift" } : null, + logDockerDriverGatewayRestart: (message) => events.push({ type: "restart", message }), + registerDockerDriverGatewayEndpoint: () => true, + isDockerDriverGatewayHttpReady: async () => { + events.push({ type: "http-ready" }); + return true; + }, + verifySandboxBridgeGatewayReachableOrExit: async () => { + events.push({ type: "verify-sandbox-bridge" }); + }, + readGatewayHealth: () => ({ + status: "Gateway: nemoclaw\nConnected", + namedInfo: "Gateway: nemoclaw", + activeInfo: "Gateway: nemoclaw", + }), + rememberDockerDriverGatewayPid: (pid) => events.push({ type: "remember-pid", pid }), + reapDuplicateHostGatewaysExceptOrFail: (keepPid, _gatewayBin, extraPids) => { + events.push({ type: "duplicate-reap", keepPid, extraPids }); + options.duplicateError && throwHarnessError(options.duplicateError); + }, + reapHostGatewayBeforeLaunchOrFail: ({ extraPids }) => { + events.push({ type: "prelaunch-reap", extraPids }); + options.prelaunchError && throwHarnessError(options.prelaunchError); + }, + isGatewayPortAvailable: async () => options.postReapPortAvailable ?? true, + reportUntrustedGatewayPort: (message) => { + throw new Error(message); + }, + reportMissingGatewayBinary: () => { + throw new Error("missing gateway binary"); + }, + log: (message) => events.push({ type: "log", message }), + }; + + return { + events, + async run(): Promise<"reused" | "launch"> { + const action = await runDockerDriverGatewayCutover(input, deps); + action === "launch" && events.push({ type: "spawn-fresh" }); + return action; + }, + }; +} + +describe("Docker-driver gateway prelaunch cutover (#5968)", () => { + it("reaps stale port listeners before allowing a fresh launch", async () => { + const harness = makeHarness({ + listenerPids: [4242, 4343], + driftPids: [4242], + }); + + await expect(harness.run()).resolves.toBe("launch"); + const reapIndex = harness.events.findIndex((event) => event.type === "prelaunch-reap"); + const launchIndex = harness.events.findIndex((event) => event.type === "spawn-fresh"); + expect(harness.events[reapIndex]?.extraPids).toEqual([4242, 4343]); + expect(reapIndex).toBeGreaterThanOrEqual(0); + expect(launchIndex).toBeGreaterThan(reapIndex); + }); + + it("bypasses sole-binder reuse and reaps the duplicate when an extra listener exists", async () => { + const harness = makeHarness({ listenerPids: [4242, 4343] }); + + await expect(harness.run()).resolves.toBe("reused"); + expect(harness.events).toContainEqual({ + type: "duplicate-reap", + keepPid: 4242, + extraPids: [4242, 4343], + }); + expect(harness.events.some((event) => event.type === "spawn-fresh")).toBe(false); + }); + + it("does not reuse a healthy pid-file gateway when listener enumeration is incomplete", async () => { + const harness = makeHarness({ listenerPids: [4242], scanComplete: false }); + + await expect(harness.run()).resolves.toBe("launch"); + expect(harness.events).toContainEqual({ type: "prelaunch-reap", extraPids: [4242] }); + expect(harness.events.some((event) => event.type === "http-ready")).toBe(false); + }); + + it("fails closed when no listener is attributable and the port remains occupied", async () => { + const harness = makeHarness({ + listenerPids: [], + scanComplete: true, + pidFileGatewayPid: null, + postReapPortAvailable: false, + }); + + await expect(harness.run()).rejects.toThrow("gateway port remains occupied"); + expect(harness.events).toContainEqual({ type: "prelaunch-reap", extraPids: [] }); + expect(harness.events.some((event) => event.type === "http-ready")).toBe(false); + expect(harness.events.some((event) => event.type === "spawn-fresh")).toBe(false); + }); + + it("never includes an unobserved pid-file process in port-scoped cleanup", async () => { + const harness = makeHarness({ listenerPids: [4343], pidFileGatewayPid: 4242 }); + + await expect(harness.run()).resolves.toBe("reused"); + expect(harness.events).toContainEqual({ + type: "duplicate-reap", + keepPid: 4343, + extraPids: [4343], + }); + }); + + it("also excludes a drifted pid-file process from port-scoped cleanup", async () => { + const harness = makeHarness({ + listenerPids: [4343], + pidFileGatewayPid: 4242, + driftPids: [4242], + }); + + await expect(harness.run()).resolves.toBe("reused"); + expect(harness.events).toContainEqual({ + type: "duplicate-reap", + keepPid: 4343, + extraPids: [4343], + }); + }); + + it("does not launch when the scoped prelaunch reaper fails", async () => { + const harness = makeHarness({ + listenerPids: [4242], + driftPids: [4242], + prelaunchError: "__prelaunch_reap_failed__", + }); + + await expect(harness.run()).rejects.toThrow("__prelaunch_reap_failed__"); + expect(harness.events.some((event) => event.type === "spawn-fresh")).toBe(false); + }); + + it("does not report adopted reuse when duplicate cleanup fails", async () => { + const harness = makeHarness({ + listenerPids: [4343, 4242], + pidFileGatewayPid: null, + duplicateError: "__duplicate_reap_failed__", + }); + + await expect(harness.run()).rejects.toThrow("__duplicate_reap_failed__"); + expect(harness.events.some((event) => event.type === "verify-sandbox-bridge")).toBe(false); + expect(harness.events.some((event) => event.type === "spawn-fresh")).toBe(false); + }); +}); diff --git a/test/tunnel-gateway-port-release-runtime.test.ts b/test/tunnel-gateway-port-release-runtime.test.ts new file mode 100644 index 00000000000..2ded9088b6b --- /dev/null +++ b/test/tunnel-gateway-port-release-runtime.test.ts @@ -0,0 +1,160 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Runtime validation for the #5968 gateway port release. The unit suite in +// src/lib/tunnel/gateway-port-release.test.ts mocks lsof/the stopper to cover +// branch decisions; this test exercises the REAL release path end-to-end: +// it starts an actual process whose argv0 basename is `openshell-gateway` +// (the identity the host-gateway stopper cmdline-gates on), bound to an +// isolated non-default port with an isolated HOME/state dir, then runs the +// real releaseManagedGatewayPort and proves a fresh process can immediately +// rebind the freed port. Nothing here touches a real user gateway. +// +// The fake gateway is launched through a short-lived launcher that exits +// immediately, so the gateway is orphaned to init rather than parented by the +// (synchronous, event-loop-blocked) test process — otherwise a killed child +// would linger as an unreaped zombie that `ps` still reports as alive. + +import { spawn, spawnSync } from "node:child_process"; +import fs from "node:fs"; +import net from "node:net"; +import os from "node:os"; +import path from "node:path"; + +import { afterEach, describe, expect, it } from "vitest"; + +import { waitUntil } from "../src/lib/core/wait"; +import { resolveGatewayStateDirName } from "../src/lib/onboard/gateway-binding"; +import { releaseManagedGatewayPort } from "../src/lib/tunnel/gateway-port-release"; + +// POSIX-only: the release path relies on lsof/ps/POSIX signals and the +// cmdline gate reads /proc or `ps -o args=`. Windows has no equivalent and is +// not a NemoClaw host target for the gateway. +const posix = process.platform !== "win32"; +const hasLsof = posix && !spawnSync("lsof", ["-v"], { stdio: "ignore" }).error; + +let gatewayPid = 0; +let tmpHome: string | null = null; + +function killQuietly(pid: number): void { + try { + pid > 0 && process.kill(pid, "SIGKILL"); + } catch { + /* already gone */ + } +} + +afterEach(() => { + killQuietly(gatewayPid); + gatewayPid = 0; + tmpHome && fs.rmSync(tmpHome, { recursive: true, force: true }); + tmpHome = null; +}); + +// Reserve a free localhost TCP port by binding :0, then releasing it. +function reserveFreePort(): Promise { + return new Promise((resolve, reject) => { + const probe = net.createServer(); + probe.once("error", reject); + probe.listen(0, "127.0.0.1", () => { + const address = probe.address(); + const port = typeof address === "object" && address ? address.port : 0; + probe.close(() => resolve(port)); + }); + }); +} + +// Resolve true when a fresh server can bind the port, false otherwise. +function canBind(port: number): Promise { + return new Promise((resolve) => { + const server = net.createServer(); + server.once("error", () => resolve(false)); + server.listen(port, "127.0.0.1", () => { + server.close(() => resolve(true)); + }); + }); +} + +function readPidQuietly(pidFile: string): number { + try { + return Number.parseInt(fs.readFileSync(pidFile, "utf-8").trim() || "0", 10) || 0; + } catch { + return 0; + } +} + +describe("releaseManagedGatewayPort runtime validation (#5968)", () => { + it.skipIf(!posix || !hasLsof)( + "stops a real openshell-gateway process and frees the port for immediate rebind", + async () => { + const port = await reserveFreePort(); + tmpHome = fs.mkdtempSync(path.join(os.tmpdir(), "nemoclaw-gw-rt-")); + const argv0Path = path.join(tmpHome, "openshell-gateway"); + + // Persist realistic per-port bookkeeping, then rely on real lsof to prove + // this PID owns the selected port before the stopper may signal it. + const stateDir = path.join( + tmpHome, + ".local", + "state", + "nemoclaw", + resolveGatewayStateDirName(port), + ); + fs.mkdirSync(stateDir, { recursive: true }); + const pidFile = path.join(stateDir, "openshell-gateway.pid"); + + // The gateway binds the port and records its own pid; the launcher spawns + // it detached (argv0 basename `openshell-gateway`) and exits, orphaning it. + const gatewayFile = path.join(tmpHome, "gateway.cjs"); + fs.writeFileSync( + gatewayFile, + `const net=require("node:net");const fs=require("node:fs");` + + `const server=net.createServer();` + + `server.listen(${String(port)},"127.0.0.1",()=>fs.writeFileSync(${JSON.stringify(pidFile)},String(process.pid)));` + + `process.on("SIGTERM",()=>process.exit(0));`, + ); + const launcherScript = + `const {spawn}=require("node:child_process");` + + `spawn(process.argv[1],[process.argv[2]],{argv0:process.argv[3],detached:true,stdio:"ignore"}).unref();`; + spawn(process.execPath, ["-e", launcherScript, process.execPath, gatewayFile, argv0Path], { + stdio: "ignore", + }); + + // Wait until the orphaned gateway has recorded its pid and bound the port. + const pidRecorded = waitUntil( + () => { + gatewayPid = readPidQuietly(pidFile); + return gatewayPid > 0; + }, + { + deadlineMs: Date.now() + 10_000, + initialIntervalMs: 25, + maxIntervalMs: 25, + backoffFactor: 1, + }, + ); + expect(pidRecorded).toBe(true); + expect(gatewayPid).toBeGreaterThan(0); + await expect(canBind(port)).resolves.toBe(false); + + // Run the REAL release path (real spawnSync/ps/kill/stopper); only the + // registry lookup and HOME are isolated so no real gateway is touched. + const result = releaseManagedGatewayPort( + { sandboxName: "nemoclaw-5968-runtime", confirmTimeoutMs: 8000 }, + { + homeDir: tmpHome, + env: { ...process.env, HOME: tmpHome }, + getSandbox: () => ({ gatewayPort: port }), + }, + ); + + expect(result.port).toBe(port); + expect(result.stopped).toContain(gatewayPid); + expect(result.released).toBe(true); + + // Ground truth: a fresh process can rebind the freed port immediately. + await expect(canBind(port)).resolves.toBe(true); + }, + 30000, + ); +});