Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion ci/source-architecture-budget.json
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@
"src/lib/actions/uninstall/run-plan.ts": 26,
"src/lib/inference/onboard-probes.ts": 20,
"src/lib/inference/vllm.ts": 21,
"src/lib/onboard.ts": 212,
"src/lib/onboard.ts": 211,
"src/lib/onboard/machine/handlers/sandbox.ts": 21,
"src/lib/sandbox/config.ts": 22,
"src/lib/shields/index.ts": 23
Expand Down
80 changes: 19 additions & 61 deletions src/lib/onboard.ts
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,6 @@ const dockerGpuRoute: typeof import("./onboard/docker-gpu-route") = require("./o
const sandboxGpuCreateFlow: typeof import("./onboard/sandbox-gpu-create-flow") = require("./onboard/sandbox-gpu-create-flow");
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 gatewayService: typeof import("./onboard/docker-driver-gateway-service") = require("./onboard/docker-driver-gateway-service");
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");
Expand Down Expand Up @@ -769,6 +768,25 @@ const { getGatewayReuseSnapshot, selectNamedGatewayForReuseIfNeeded } =
cliDisplayName,
});

const { refreshDockerDriverGatewayReuseState } =
gatewayReuse.createDockerDriverGatewayReuseApplication({
gatewayName: () => GATEWAY_NAME,
getGatewayCompatContainerName: () =>
gatewayBinding.resolveGatewayCompatContainerName(GATEWAY_PORT),
isDockerDriverGatewayEnabled: isLinuxDockerDriverGatewayEnabled,
resolveOpenShellGatewayBinary,
getDockerDriverGatewayEnv,
runCaptureOpenshell,
getDockerDriverGatewayStateDir,
resolveOpenShellSandboxBinary,
getDockerDriverGatewayPid,
isDockerDriverGatewayProcessAlive,
getDockerDriverGatewayReuseDrift: getGatewayReuseDrift,
checkGatewayPortAvailable,
getDockerDriverGatewayPortListenerPid,
rememberDockerDriverGatewayPid,
});

// biome-ignore format: keep src/lib/onboard.ts net-neutral for growth guardrail.
const { getSandboxReuseState, getSandboxRecreateObservation } = sandboxReuse.createSandboxReuseHelpers({ runCaptureOpenshell, getSandboxStateFromOutputs, getGatewayName: () => GATEWAY_NAME });

Expand Down Expand Up @@ -1223,66 +1241,6 @@ function logDockerDriverGatewayRestart(reason: string): void {
console.log(` Existing OpenShell Docker-driver gateway is stale (${reason}); restarting...`);
}

async function refreshDockerDriverGatewayReuseState(
gatewayReuseState: GatewayReuseState,
): Promise<GatewayReuseState> {
if (!isLinuxDockerDriverGatewayEnabled() || gatewayReuseState !== "healthy") {
return gatewayReuseState;
}
const gatewayBin = resolveOpenShellGatewayBinary();
const baseDesiredEnv = getDockerDriverGatewayEnv(
runCaptureOpenshell(["--version"], { ignoreError: true }),
);
const runtimeIdentity = gatewayBin
? dockerDriverGatewayLaunch.buildDockerDriverGatewayRuntimeIdentity({
gatewayBin,
gatewayEnv: baseDesiredEnv,
stateDir: getDockerDriverGatewayStateDir(),
sandboxBin: resolveOpenShellSandboxBinary(),
gatewayName: GATEWAY_NAME,
compatContainerName: gatewayBinding.resolveGatewayCompatContainerName(GATEWAY_PORT),
})
: null;
const desiredEnv = runtimeIdentity?.desiredEnv ?? baseDesiredEnv;
const driftBin = dockerDriverGatewayLaunch.resolveDriftGatewayBin(runtimeIdentity, gatewayBin);
const identityBin = runtimeIdentity?.identityGatewayBin ?? gatewayBin;
const managedServicePid = gatewayService.getTrustedActiveOpenShellGatewayUserServicePid();
const pid = getDockerDriverGatewayPid();
if (pid !== null && isDockerDriverGatewayProcessAlive()) {
const drift = getGatewayReuseDrift(pid, desiredEnv, driftBin, managedServicePid);
if (drift) {
console.log(
` Existing OpenShell Docker-driver gateway is stale (${drift.reason}); it will be recreated.`,
);
return "stale";
}
return gatewayReuseState;
}

const portCheck = await checkGatewayPortAvailable();
const dockerGatewayPid = getDockerDriverGatewayPortListenerPid(portCheck, {
gatewayBin: identityBin,
});
if (dockerGatewayPid !== null) {
const drift = getGatewayReuseDrift(dockerGatewayPid, desiredEnv, driftBin, managedServicePid);
if (dockerGatewayPid !== managedServicePid) rememberDockerDriverGatewayPid(dockerGatewayPid);
if (drift) {
console.log(
` Existing OpenShell Docker-driver gateway is stale (${drift.reason}); it will be recreated.`,
);
return "stale";
}
return "healthy";
}

// `openshell status` already proved the selected gateway is reachable. If
// the port probe cannot identify the owning PID, avoid tearing down a live
// gateway solely because the pid file is stale.
if (!portCheck.ok && !portCheck.pid) return "healthy";

return "stale";
}

function destroyGateway(
clearRegistry: () => void = registry.clearAll,
isDockerDriverGatewayEnabledForDestroy: () => boolean = isLinuxDockerDriverGatewayEnabled,
Expand Down
131 changes: 130 additions & 1 deletion src/lib/onboard/gateway-reuse.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,11 @@
import { describe, expect, it, vi } from "vitest";

import { OPENSHELL_PROBE_TIMEOUT_MS } from "../adapters/openshell/timeouts";
import { createGatewayReuseHelpers } from "./gateway-reuse";
import {
createDockerDriverGatewayReuseApplication,
type DockerDriverGatewayReuseApplicationDeps,
createGatewayReuseHelpers,
} from "./gateway-reuse";

describe("gateway reuse snapshot", () => {
it("bounds OpenShell gateway inspection probes (#6752)", () => {
Expand Down Expand Up @@ -84,3 +88,128 @@ describe("gateway reuse snapshot", () => {
expect(helpers.getGatewayReuseSnapshot().gatewayReuseState).toBe("missing");
});
});

function createDockerDriverReuseApplication(
overrides: Partial<DockerDriverGatewayReuseApplicationDeps> = {},
) {
return createDockerDriverGatewayReuseApplication({
gatewayName: () => "nemoclaw",
getGatewayCompatContainerName: () => "openshell-gateway-nemoclaw",
isDockerDriverGatewayEnabled: () => true,
resolveOpenShellGatewayBinary: () => "/opt/openshell-gateway",
getDockerDriverGatewayEnv: () => ({ OPENSHELL_DRIVERS: "docker" }),
runCaptureOpenshell: vi.fn(() => "openshell 0.0.99"),
getDockerDriverGatewayStateDir: () => "/tmp/nemoclaw-gateway",
resolveOpenShellSandboxBinary: () => "/opt/openshell-sandbox",
getDockerDriverGatewayPid: () => 42,
isDockerDriverGatewayProcessAlive: () => true,
getDockerDriverGatewayReuseDrift: vi.fn(() => null),
checkGatewayPortAvailable: vi.fn(async () => ({ ok: true })),
getDockerDriverGatewayPortListenerPid: vi.fn(() => null),
rememberDockerDriverGatewayPid: vi.fn(),
buildDockerDriverGatewayRuntimeIdentity: vi.fn(() => ({
launch: null,
desiredEnv: { OPENSHELL_DRIVERS: "docker" },
driftGatewayBin: "/opt/openshell-gateway",
identityGatewayBin: "/opt/openshell-gateway",
})),
resolveDriftGatewayBin: vi.fn((runtimeIdentity, gatewayBin) =>
runtimeIdentity ? runtimeIdentity.driftGatewayBin : gatewayBin,
),
getTrustedActiveOpenShellGatewayUserServicePid: vi.fn(() => null),
log: vi.fn(),
...overrides,
});
}

describe("Docker-driver gateway reuse application", () => {
it("keeps reuse state unchanged when Docker-driver inspection does not apply (#7695)", async () => {
const isDockerDriverGatewayEnabled = vi.fn(() => false);
const checkGatewayPortAvailable = vi.fn(async () => ({ ok: true }));
const application = createDockerDriverReuseApplication({
isDockerDriverGatewayEnabled,
checkGatewayPortAvailable,
});

await expect(application.refreshDockerDriverGatewayReuseState("healthy")).resolves.toBe(
"healthy",
);
isDockerDriverGatewayEnabled.mockReturnValue(true);
await expect(application.refreshDockerDriverGatewayReuseState("stale")).resolves.toBe("stale");
expect(checkGatewayPortAvailable).not.toHaveBeenCalled();
});

it("marks a running Docker-driver gateway stale when runtime identity drifts", async () => {
const log = vi.fn();
const checkGatewayPortAvailable = vi.fn(async () => ({ ok: true }));
const application = createDockerDriverReuseApplication({
getDockerDriverGatewayReuseDrift: vi.fn(() => ({
reason: "runtime environment changed",
})),
checkGatewayPortAvailable,
log,
});

await expect(application.refreshDockerDriverGatewayReuseState("healthy")).resolves.toBe(
"stale",
);
expect(checkGatewayPortAvailable).not.toHaveBeenCalled();
expect(log).toHaveBeenCalledWith(
" Existing OpenShell Docker-driver gateway is stale (runtime environment changed); it will be recreated.",
);
});

it("adopts a matching gateway port listener when the PID file is absent", async () => {
const rememberDockerDriverGatewayPid = vi.fn();
const getDockerDriverGatewayReuseDrift = vi.fn(() => null);
const application = createDockerDriverReuseApplication({
getDockerDriverGatewayPid: () => null,
isDockerDriverGatewayProcessAlive: () => false,
checkGatewayPortAvailable: vi.fn(async () => ({ ok: false, pid: 731 })),
getDockerDriverGatewayPortListenerPid: vi.fn(() => 731),
getDockerDriverGatewayReuseDrift,
getTrustedActiveOpenShellGatewayUserServicePid: vi.fn(() => 900),
rememberDockerDriverGatewayPid,
});

await expect(application.refreshDockerDriverGatewayReuseState("healthy")).resolves.toBe(
"healthy",
);
expect(getDockerDriverGatewayReuseDrift).toHaveBeenCalledWith(
731,
{ OPENSHELL_DRIVERS: "docker" },
"/opt/openshell-gateway",
900,
);
expect(rememberDockerDriverGatewayPid).toHaveBeenCalledWith(731);
});

it("preserves a reachable selected gateway when the port owner is ambiguous", async () => {
const rememberDockerDriverGatewayPid = vi.fn();
const application = createDockerDriverReuseApplication({
getDockerDriverGatewayPid: () => null,
isDockerDriverGatewayProcessAlive: () => false,
checkGatewayPortAvailable: vi.fn(async () => ({ ok: false, pid: null })),
getDockerDriverGatewayPortListenerPid: vi.fn(() => null),
rememberDockerDriverGatewayPid,
});

await expect(application.refreshDockerDriverGatewayReuseState("healthy")).resolves.toBe(
"healthy",
);
expect(rememberDockerDriverGatewayPid).not.toHaveBeenCalled();
});

it("marks a gateway stale when no Docker-driver process owns the available port", async () => {
const application = createDockerDriverReuseApplication({
getDockerDriverGatewayPid: () => null,
isDockerDriverGatewayProcessAlive: () => false,
checkGatewayPortAvailable: vi.fn(async () => ({ ok: true })),
getDockerDriverGatewayPortListenerPid: vi.fn(() => null),
});

await expect(application.refreshDockerDriverGatewayReuseState("healthy")).resolves.toBe(
"stale",
);
});
});
128 changes: 127 additions & 1 deletion src/lib/onboard/gateway-reuse.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,14 @@
// SPDX-License-Identifier: Apache-2.0

import { OPENSHELL_PROBE_TIMEOUT_MS } from "../adapters/openshell/timeouts";
import { getGatewayReuseState, shouldSelectNamedGatewayForReuse } from "../state/gateway";
import {
getGatewayReuseState,
type GatewayReuseState,
shouldSelectNamedGatewayForReuse,
} from "../state/gateway";
import * as dockerDriverGatewayLaunch from "./docker-driver-gateway-launch";
import * as gatewayService from "./docker-driver-gateway-service";
import type { PortProbeResult } from "./preflight";

export type GatewayReuseSnapshot = {
gatewayStatus: string;
Expand All @@ -23,6 +30,125 @@ export interface GatewayReuseHelpers {
selectNamedGatewayForReuseIfNeeded(snapshot: GatewayReuseSnapshot): GatewayReuseSnapshot;
}

export interface DockerDriverGatewayReuseApplicationDeps {
gatewayName(): string;
getGatewayCompatContainerName(): string;
isDockerDriverGatewayEnabled(): boolean;
resolveOpenShellGatewayBinary(): string | null;
getDockerDriverGatewayEnv(versionOutput?: string | null): Record<string, string>;
runCaptureOpenshell(args: string[], opts?: { ignoreError?: boolean }): string;
getDockerDriverGatewayStateDir(): string;
resolveOpenShellSandboxBinary(): string | null;
getDockerDriverGatewayPid(): number | null;
isDockerDriverGatewayProcessAlive(): boolean;
getDockerDriverGatewayReuseDrift(
pid: number,
desiredEnv: Record<string, string>,
gatewayBin?: string | null,
trustedServicePid?: number | null,
): { reason: string } | null;
checkGatewayPortAvailable(): Promise<PortProbeResult>;
getDockerDriverGatewayPortListenerPid(
portCheck: PortProbeResult,
opts?: { gatewayBin?: string | null },
): number | null;
rememberDockerDriverGatewayPid(pid: number): void;
buildDockerDriverGatewayRuntimeIdentity?: typeof dockerDriverGatewayLaunch.buildDockerDriverGatewayRuntimeIdentity;
resolveDriftGatewayBin?: typeof dockerDriverGatewayLaunch.resolveDriftGatewayBin;
getTrustedActiveOpenShellGatewayUserServicePid?: typeof gatewayService.getTrustedActiveOpenShellGatewayUserServicePid;
log?(message: string): void;
}

export interface DockerDriverGatewayReuseApplication {
refreshDockerDriverGatewayReuseState(state: GatewayReuseState): Promise<GatewayReuseState>;
}

export function createDockerDriverGatewayReuseApplication(
deps: DockerDriverGatewayReuseApplicationDeps,
): DockerDriverGatewayReuseApplication {
const buildRuntimeIdentity =
deps.buildDockerDriverGatewayRuntimeIdentity ??
dockerDriverGatewayLaunch.buildDockerDriverGatewayRuntimeIdentity;
const resolveDriftGatewayBin =
deps.resolveDriftGatewayBin ?? dockerDriverGatewayLaunch.resolveDriftGatewayBin;
const getTrustedServicePid =
deps.getTrustedActiveOpenShellGatewayUserServicePid ??
gatewayService.getTrustedActiveOpenShellGatewayUserServicePid;
const log = deps.log ?? console.log;

async function refreshDockerDriverGatewayReuseState(
state: GatewayReuseState,
): Promise<GatewayReuseState> {
if (!deps.isDockerDriverGatewayEnabled() || state !== "healthy") return state;

const gatewayBin = deps.resolveOpenShellGatewayBinary();
const baseDesiredEnv = deps.getDockerDriverGatewayEnv(
deps.runCaptureOpenshell(["--version"], { ignoreError: true }),
);
const runtimeIdentity = gatewayBin
? buildRuntimeIdentity({
gatewayBin,
gatewayEnv: baseDesiredEnv,
stateDir: deps.getDockerDriverGatewayStateDir(),
sandboxBin: deps.resolveOpenShellSandboxBinary(),
gatewayName: deps.gatewayName(),
compatContainerName: deps.getGatewayCompatContainerName(),
})
: null;
const desiredEnv = runtimeIdentity?.desiredEnv ?? baseDesiredEnv;
const driftBin = resolveDriftGatewayBin(runtimeIdentity, gatewayBin);
const identityBin = runtimeIdentity?.identityGatewayBin ?? gatewayBin;
const managedServicePid = getTrustedServicePid();
const pid = deps.getDockerDriverGatewayPid();
if (pid !== null && deps.isDockerDriverGatewayProcessAlive()) {
const drift = deps.getDockerDriverGatewayReuseDrift(
pid,
desiredEnv,
driftBin,
managedServicePid,
);
if (drift) {
log(
` Existing OpenShell Docker-driver gateway is stale (${drift.reason}); it will be recreated.`,
);
return "stale";
}
return state;
}

const portCheck = await deps.checkGatewayPortAvailable();
const dockerGatewayPid = deps.getDockerDriverGatewayPortListenerPid(portCheck, {
gatewayBin: identityBin,
});
if (dockerGatewayPid !== null) {
const drift = deps.getDockerDriverGatewayReuseDrift(
dockerGatewayPid,
desiredEnv,
driftBin,
managedServicePid,
);
if (dockerGatewayPid !== managedServicePid) {
deps.rememberDockerDriverGatewayPid(dockerGatewayPid);
}
if (drift) {
log(
` Existing OpenShell Docker-driver gateway is stale (${drift.reason}); it will be recreated.`,
);
return "stale";
}
return "healthy";
}

// OpenShell status already proved the selected gateway is reachable. Preserve it when
// the port probe cannot identify an owner, instead of deleting a potentially live gateway.
if (!portCheck.ok && !portCheck.pid) return "healthy";

return "stale";
}

return { refreshDockerDriverGatewayReuseState };
}

export function createGatewayReuseHelpers(deps: GatewayReuseDeps): GatewayReuseHelpers {
const currentGatewayName = () =>
typeof deps.gatewayName === "function" ? deps.gatewayName() : deps.gatewayName;
Expand Down
Loading