import { spawn } from "node:child_process"; import { isRecord, isStringRecord } from "@openclaw/normalization-core/record-coerce"; import { withContainerEnvFile } from "../infra/container-env-file.js"; import { attachChildProcessBridge } from "../process/child-process-bridge.js"; import { runCommandWithTimeout } from "../process/exec.js"; import { buildCellCreateArgs, buildCellRunArgs, validateCellContainerProfile, validateFleetImage, type CellContainerProfile, type FleetContainerRuntimeName, } from "./cell-profile.js"; import { createRedactingStreamWriter } from "./containers.redaction.js"; type FleetContainerCommandOptions = { allowFailure?: boolean; redactValues?: readonly string[]; }; type FleetContainerCommandResult = { stdout: string; stderr: string; code: number; }; type FleetContainerCommandExecutor = ( runtime: FleetContainerRuntimeName, args: string[], options: FleetContainerCommandOptions, ) => Promise; type FleetContainerStreamResult = { code: number | null; signal: NodeJS.Signals | null; }; // Mirrors the termination set attachChildProcessBridge forwards, plus SIGPIPE for a // downstream reader (e.g. `| head`) closing the inherited stdout mid-stream. const DELIBERATE_STREAM_STOP_SIGNALS = new Set([ "SIGINT", "SIGTERM", "SIGHUP", "SIGQUIT", "SIGBREAK", "SIGPIPE", ]); type FleetContainerStreamExecutor = ( runtime: FleetContainerRuntimeName, args: string[], options: { redactValues: readonly string[] }, ) => Promise; type FleetContainerLogsOptions = { follow?: boolean; timestamps?: boolean; tail?: number; since?: string; redactValues: readonly string[]; }; export type FleetContainerInspectResult = | { kind: "ok"; containerId: string; state: string; running: boolean; labels: Record; environment: Record; imageId: string; memory: string; cpus: string; pidsLimit: number | undefined; storageOpt: Record; capDrop: string[]; // Podman-only top-level inspect field: null means every capability is // dropped, a list means caps remain, and Docker omits the field entirely. effectiveCaps: string[] | undefined; securityOpt: string[]; init: boolean | undefined; restartPolicy: string | undefined; portBindings: Array<{ containerPort: string; hostIp: string; hostPort: string }>; user?: string; usernsMode?: string; } | { kind: "missing"; state: "missing" } | { kind: "unavailable"; state: "unknown"; error: string }; export type FleetNetworkInspectResult = | { kind: "ok"; labels: Record; attachedContainers: Array<{ id: string; name?: string }>; internal: boolean; } | { kind: "missing" } | { kind: "unavailable"; error: string }; export type FleetContainerRuntime = ReturnType; const COMMAND_TIMEOUT_MS = 10 * 60_000; const COMMAND_MAX_OUTPUT_BYTES = 4 * 1024 * 1024; class InvalidInspectOutputError extends Error { constructor() { super("container inspect returned an invalid response"); } } function requireRecord(value: unknown): Record { if (!isRecord(value)) { throw new InvalidInspectOutputError(); } return value; } function requireString(value: unknown): string { if (typeof value !== "string" || value.length === 0) { throw new InvalidInspectOutputError(); } return value; } function requireBoolean(value: unknown): boolean { if (typeof value !== "boolean") { throw new InvalidInspectOutputError(); } return value; } function requireNonNegativeNumber(value: unknown): number { if (typeof value !== "number" || !Number.isFinite(value) || value < 0) { throw new InvalidInspectOutputError(); } return value; } function readOptionalInspectString(value: unknown): string | undefined { if (value === undefined || value === null || value === "") { return undefined; } if (typeof value !== "string") { throw new InvalidInspectOutputError(); } return value; } function readStringRecord(value: unknown): Record { if (value === undefined || value === null) { return {}; } if (!isStringRecord(value)) { throw new InvalidInspectOutputError(); } return value; } function readStringArray(value: unknown): string[] { if (value === undefined || value === null) { return []; } if (!Array.isArray(value) || !value.every((entry) => typeof entry === "string")) { throw new InvalidInspectOutputError(); } return value; } function readOptionalBoolean(value: unknown): boolean | undefined { if (value === undefined || value === null) { return undefined; } return requireBoolean(value); } function readRestartPolicy(value: unknown): string | undefined { if (value === undefined || value === null) { return undefined; } return readOptionalInspectString(requireRecord(value).Name); } function readPortBindings( value: unknown, ): Array<{ containerPort: string; hostIp: string; hostPort: string }> { if (value === undefined || value === null) { return []; } const bindings: Array<{ containerPort: string; hostIp: string; hostPort: string }> = []; for (const [containerPort, rawEntries] of Object.entries(requireRecord(value))) { if (!containerPort || !Array.isArray(rawEntries)) { throw new InvalidInspectOutputError(); } for (const rawEntry of rawEntries) { const entry = requireRecord(rawEntry); bindings.push({ containerPort, hostIp: requireString(entry.HostIp), hostPort: requireString(entry.HostPort), }); } } return bindings; } function readNetworkAttachments(value: unknown): Array<{ id: string; name?: string }> { if (value === undefined || value === null) { return []; } const record = requireRecord(value); return Object.entries(record) .map(([id, rawAttachment]) => { if (!id) { throw new InvalidInspectOutputError(); } const attachment = requireRecord(rawAttachment); const name = readOptionalInspectString(attachment.Name ?? attachment.name); const normalized: { id: string; name?: string } = { id }; if (name) { normalized.name = name; } return normalized; }) .toSorted((left, right) => left.id.localeCompare(right.id)); } function readEnvironment(value: unknown): Record { const environment: Record = {}; for (const assignment of readStringArray(value)) { const separator = assignment.indexOf("="); if (separator <= 0) { throw new InvalidInspectOutputError(); } environment[assignment.slice(0, separator)] = assignment.slice(separator + 1); } return environment; } function readPidsLimit(value: unknown): number | undefined { if (value === undefined || value === null || value === 0) { return undefined; } if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) { throw new InvalidInspectOutputError(); } return value; } function parseInspectRecord(stdout: string): Record { let parsed: unknown; try { parsed = JSON.parse(stdout); } catch { throw new InvalidInspectOutputError(); } if (!Array.isArray(parsed) || parsed.length !== 1) { throw new InvalidInspectOutputError(); } return requireRecord(parsed[0]); } function parseInspectOutput(stdout: string): Extract { const inspected = parseInspectRecord(stdout); const state = requireRecord(inspected.State); const config = requireRecord(inspected.Config); const hostConfig = requireRecord(inspected.HostConfig); const nanoCpus = requireNonNegativeNumber(hostConfig.NanoCpus); const user = readOptionalInspectString(config.User); const usernsMode = readOptionalInspectString(hostConfig.UsernsMode); return { kind: "ok", containerId: requireString(inspected.Id), state: requireString(state.Status), running: requireBoolean(state.Running), // Preserve ordinary-object assignment semantics for JSON "__proto__" labels. labels: Object.assign({}, readStringRecord(config.Labels)), environment: readEnvironment(config.Env), imageId: requireString(inspected.Image), memory: String(requireNonNegativeNumber(hostConfig.Memory)), cpus: String(nanoCpus / 1_000_000_000), pidsLimit: readPidsLimit(hostConfig.PidsLimit), storageOpt: readStringRecord(hostConfig.StorageOpt), capDrop: readStringArray(hostConfig.CapDrop), effectiveCaps: inspected.EffectiveCaps === undefined ? undefined : readStringArray(inspected.EffectiveCaps), securityOpt: readStringArray(hostConfig.SecurityOpt), init: readOptionalBoolean(hostConfig.Init), restartPolicy: readRestartPolicy(hostConfig.RestartPolicy), portBindings: readPortBindings(hostConfig.PortBindings), ...(user ? { user } : {}), ...(usernsMode ? { usernsMode } : {}), }; } function parseNetworkInspectOutput( stdout: string, ): Extract { const inspected = parseInspectRecord(stdout); return { kind: "ok", labels: Object.assign({}, readStringRecord(inspected.Labels ?? inspected.labels)), attachedContainers: readNetworkAttachments(inspected.Containers ?? inspected.containers), internal: readOptionalBoolean(inspected.Internal ?? inspected.internal) ?? false, }; } function parseDockerContextEndpoint(stdout: string): string { try { const context = parseInspectRecord(stdout); const endpoints = requireRecord(context.Endpoints); const docker = requireRecord(endpoints.docker); return requireString(docker.Host); } catch { throw new Error("docker context inspect returned an invalid response"); } } function isLocalDockerEndpoint(endpoint: string): boolean { const normalized = endpoint.toLowerCase(); const unixPrefix = "unix:///"; if (normalized.startsWith(unixPrefix)) { return normalized.length > unixPrefix.length; } const windowsPipePrefix = "npipe:////./pipe/"; return ( process.platform === "win32" && normalized.startsWith(windowsPipePrefix) && normalized.length > windowsPipePrefix.length ); } function parsePodmanServiceIsRemote(stdout: string): boolean { try { const info = requireRecord(JSON.parse(stdout)); const host = requireRecord(info.host); return requireBoolean(host.serviceIsRemote); } catch { throw new Error("podman info returned an invalid response"); } } function parseDockerRootlessInfo(stdout: string): boolean { let parsed: unknown; try { parsed = JSON.parse(stdout); } catch { throw new Error("docker info returned an invalid security-options response"); } if (!Array.isArray(parsed) || !parsed.every((entry) => typeof entry === "string")) { throw new Error("docker info returned an invalid security-options response"); } return parsed.some((entry) => entry.split(",").includes("name=rootless")); } function readEnvironmentValues(args: readonly string[], extraValues?: readonly string[]): string[] { const values: string[] = []; for (let index = 0; index < args.length; index += 1) { const arg = args[index] ?? ""; let assignment: string | undefined; if (arg === "--env" || arg === "-e") { assignment = args[index + 1]; index += 1; } else if (arg.startsWith("--env=")) { assignment = arg.slice("--env=".length); } else if (arg.startsWith("-e=")) { assignment = arg.slice("-e=".length); } if (!assignment) { continue; } const separator = assignment.indexOf("="); if (separator >= 0 && separator < assignment.length - 1) { values.push(assignment.slice(separator + 1)); } } for (const value of extraValues ?? []) { if (value) { values.push(value); } } return [...new Set(values)]; } function redactEnvironmentValues( text: string, args: readonly string[], extraValues?: readonly string[], ): string { let redacted = text; const values = readEnvironmentValues(args, extraValues).toSorted( (left, right) => right.length - left.length, ); for (const value of values) { redacted = redacted.replaceAll(value, ""); } return redacted; } function formatExecutorError( error: unknown, runtime: FleetContainerRuntimeName, args: readonly string[], extraValues?: readonly string[], ): Error { const detail = error instanceof Error ? redactEnvironmentValues(error.message, args, extraValues).trim() : ""; return new Error(detail || `${runtime} container command failed`); } function commandFailureError( runtime: FleetContainerRuntimeName, args: readonly string[], result: FleetContainerCommandResult, extraValues?: readonly string[], ): Error { const detail = redactEnvironmentValues(result.stderr, args, extraValues).trim(); return new Error(detail || `${runtime} container command failed with exit code ${result.code}`); } const defaultFleetContainerCommandExecutor: FleetContainerCommandExecutor = async ( runtime, args, options, ) => { const result = await runCommandWithTimeout([runtime, ...args], { timeoutMs: COMMAND_TIMEOUT_MS, maxOutputBytes: COMMAND_MAX_OUTPUT_BYTES, }); const normalized = { stdout: result.stdout, stderr: redactEnvironmentValues(result.stderr, args, options.redactValues), code: result.code ?? 1, }; if (normalized.code !== 0 && !options.allowFailure) { throw commandFailureError(runtime, args, normalized, options.redactValues); } return normalized; }; const defaultFleetContainerStreamExecutor: FleetContainerStreamExecutor = ( runtime, args, options, ) => new Promise((resolve, reject) => { // Pipe instead of inheriting stdio so the cell's Gateway token can be // redacted from live log content; argv itself carries no secrets. const child = spawn(runtime, args, { stdio: ["ignore", "pipe", "pipe"] }); // Forward every termination signal (SIGINT/SIGTERM/SIGHUP/...), not just Ctrl-C: // a supervisor signaling only the CLI PID must never orphan a follow stream. // The bridge detaches itself on child exit. attachChildProcessBridge(child); const stdout = createRedactingStreamWriter(process.stdout, options.redactValues); const stderr = createRedactingStreamWriter(process.stderr, options.redactValues); // A downstream reader closing our stdout (e.g. `| head`) surfaces as a // stream error here rather than SIGPIPE on the child; end the child so a // follow stream terminates instead of writing into a broken pipe forever. const onTargetError = (): void => { child.kill("SIGTERM"); }; process.stdout.once("error", onTargetError); process.stderr.once("error", onTargetError); // Honor target backpressure: pause the child stream while the terminal or // pipe drains so a noisy follow stream cannot buffer without bound. const pipeWithBackpressure = ( source: typeof child.stdout, targetStream: NodeJS.WriteStream, writer: ReturnType, ): void => { source?.on("data", (chunk: Buffer) => { if (!writer.write(chunk)) { source.pause(); targetStream.once("drain", () => source.resume()); } }); }; pipeWithBackpressure(child.stdout, process.stdout, stdout); pipeWithBackpressure(child.stderr, process.stderr, stderr); child.on("error", (error) => { if (child.pid === undefined) { reject(error); } }); child.once("close", (code, signal) => { stdout.flush(); stderr.flush(); process.stdout.removeListener("error", onTargetError); process.stderr.removeListener("error", onTargetError); resolve({ code, signal }); }); }); function isMissingContainerError(stderr: string): boolean { const normalized = stderr.toLowerCase(); return ( normalized.includes("no such object") || normalized.includes("no such container") || normalized.includes("no container with name or id") || /container .+ does not exist/u.test(normalized) ); } function isMissingNetworkError(stderr: string): boolean { const normalized = stderr.toLowerCase(); return ( normalized.includes("no such network") || normalized.includes("network not found") || /network .+ not found/u.test(normalized) || /network .+ does not exist/u.test(normalized) ); } function validateNetworkName(networkName: string): string { const normalized = networkName.trim(); if (!normalized || normalized.startsWith("-")) { throw new Error("Fleet network name is invalid."); } return normalized; } function validateContainerName(containerName: string): string { const normalized = containerName.trim(); if (!normalized || normalized.startsWith("-")) { throw new Error("Fleet container name is invalid."); } return normalized; } function buildLogsArgs(containerName: string, options: FleetContainerLogsOptions): string[] { const args = ["logs"]; if (options.follow) { args.push("--follow"); } if (options.timestamps) { args.push("--timestamps"); } if (options.tail !== undefined) { if (!Number.isSafeInteger(options.tail) || options.tail < 1) { throw new Error("Fleet logs --tail must be a positive integer."); } args.push("--tail", String(options.tail)); } if (options.since !== undefined) { if (!options.since || options.since.startsWith("-") || /[\s\p{Cc}]/u.test(options.since)) { throw new Error( "Fleet logs --since must be non-empty, must not start with '-', and must not contain whitespace or control characters.", ); } args.push("--since", options.since); } args.push(validateContainerName(containerName)); return args; } export function createFleetContainerRuntime( executor: FleetContainerCommandExecutor = defaultFleetContainerCommandExecutor, streamExecutor: FleetContainerStreamExecutor = defaultFleetContainerStreamExecutor, ) { const execute = async ( runtime: FleetContainerRuntimeName, args: string[], options: FleetContainerCommandOptions = {}, ): Promise => { try { const result = await executor(runtime, args, options); if (result.code !== 0 && !options.allowFailure) { throw commandFailureError(runtime, args, result, options.redactValues); } return { ...result, stderr: redactEnvironmentValues(result.stderr, args, options.redactValues), }; } catch (error) { throw formatExecutorError(error, runtime, args, options.redactValues); } }; return { async assertLocal(runtime: FleetContainerRuntimeName): Promise { if (runtime === "podman") { const result = await execute("podman", ["info", "--format", "json"]); if (parsePodmanServiceIsRemote(result.stdout)) { throw new Error("Fleet requires local Podman; remote cells are not supported."); } return; } // Docker exposes env-selected hosts as a virtual current context, including DOCKER_HOST. const result = await execute("docker", ["context", "inspect"]); const endpoint = parseDockerContextEndpoint(result.stdout); if (!isLocalDockerEndpoint(endpoint)) { throw new Error("Fleet requires a local Docker endpoint; remote cells are not supported."); } }, async inspect( runtime: FleetContainerRuntimeName, containerName: string, ): Promise { const args = ["container", "inspect", validateContainerName(containerName)]; let result: FleetContainerCommandResult; try { result = await execute(runtime, args, { allowFailure: true }); } catch (error) { return { kind: "unavailable", state: "unknown", error: formatExecutorError(error, runtime, args).message, }; } if (result.code !== 0) { if (isMissingContainerError(result.stderr)) { return { kind: "missing", state: "missing" }; } return { kind: "unavailable", state: "unknown", error: result.stderr.trim() || `${runtime} container inspect failed`, }; } try { return parseInspectOutput(result.stdout); } catch { return { kind: "unavailable", state: "unknown", error: "container inspect returned an invalid response", }; } }, async inspectNetwork( runtime: FleetContainerRuntimeName, networkName: string, ): Promise { const args = ["network", "inspect", validateNetworkName(networkName)]; let result: FleetContainerCommandResult; try { result = await execute(runtime, args, { allowFailure: true }); } catch (error) { return { kind: "unavailable", error: formatExecutorError(error, runtime, args).message, }; } if (result.code !== 0) { if (isMissingNetworkError(result.stderr)) { return { kind: "missing" }; } return { kind: "unavailable", error: result.stderr.trim() || `${runtime} network inspect failed`, }; } try { return parseNetworkInspectOutput(result.stdout); } catch { return { kind: "unavailable", error: "network inspect returned an invalid response", }; } }, async isDockerRootless(): Promise { const result = await execute("docker", ["info", "--format", "{{json .SecurityOptions}}"]); return parseDockerRootlessInfo(result.stdout); }, async run(profile: CellContainerProfile, start: boolean): Promise { validateCellContainerProfile(profile); await withContainerEnvFile(profile.environment, async (environmentFile) => { const args = start ? buildCellRunArgs(profile, { environmentFile }) : buildCellCreateArgs(profile, { environmentFile }); await execute(profile.runtime, args, { redactValues: Object.values(profile.environment), }); }); }, async pull(runtime: FleetContainerRuntimeName, image: string): Promise { await execute(runtime, ["pull", validateFleetImage(image)]); }, async createNetwork( runtime: FleetContainerRuntimeName, networkName: string, labels: Readonly>, options: { internal: boolean }, ): Promise { const labelArgs = Object.entries(labels) .toSorted(([left], [right]) => left.localeCompare(right)) .flatMap(([key, value]) => ["--label", `${key}=${value}`]); await execute(runtime, [ "network", "create", "--driver", "bridge", ...(options.internal ? ["--internal"] : []), ...labelArgs, validateNetworkName(networkName), ]); }, async removeNetwork(runtime: FleetContainerRuntimeName, networkName: string): Promise { await execute(runtime, ["network", "rm", validateNetworkName(networkName)]); }, async start(runtime: FleetContainerRuntimeName, containerName: string): Promise { await execute(runtime, ["start", validateContainerName(containerName)]); }, async stop(runtime: FleetContainerRuntimeName, containerName: string): Promise { await execute(runtime, ["stop", validateContainerName(containerName)]); }, async restart(runtime: FleetContainerRuntimeName, containerName: string): Promise { await execute(runtime, ["restart", validateContainerName(containerName)]); }, async logs( runtime: FleetContainerRuntimeName, containerName: string, options: FleetContainerLogsOptions, ): Promise { const result = await streamExecutor(runtime, buildLogsArgs(containerName, options), { redactValues: options.redactValues, }); if (result.code === 0) { return; } // A follow stream ended by a deliberate stop — a forwarded termination signal, // a closed downstream pipe, or docker/podman's own Ctrl-C translation to 130 — // is not a failure. Crash signals (SIGSEGV, SIGKILL, ...) stay errors. const deliberateStop = result.signal !== null && DELIBERATE_STREAM_STOP_SIGNALS.has(result.signal); if (options.follow && (deliberateStop || result.code === 130)) { return; } throw new Error( `${runtime} logs failed with ${ result.signal ? `signal ${result.signal}` : `exit code ${result.code ?? 1}` }.`, ); }, async remove( runtime: FleetContainerRuntimeName, containerName: string, force: boolean, ): Promise { await execute(runtime, [ "rm", ...(force ? ["--force"] : []), validateContainerName(containerName), ]); }, }; }