diff --git a/companion/src/wire.ts b/companion/src/wire.ts index 5e2a749941..7cc5272e2e 100644 --- a/companion/src/wire.ts +++ b/companion/src/wire.ts @@ -10,14 +10,23 @@ // point this becomes a no-op rather than a lie, which is the right way for a // sidecar to depend on someone else's API: assume nothing, and be correct // either way. +// +// `sshAlias` is the same story with a different payload: the harness's config +// status echoes the self-hosted VPS alias — a label naming one of the user's +// servers — inside `vps`, on both GET /api/config and the `config` SSE frame. +// The phone only ever renders configured-or-not, so it gets exactly that: +// `{configured: true}` survives, the host label does not. + +/** Keys that are the harness's business, never a device's. */ +const WITHHELD_KEYS = new Set(["resumeCursors", "sshAlias"]); -/** Recursively drop `resumeCursors`, wherever it appears. */ +/** Recursively drop the withheld keys, wherever they appear. */ export function scrub(value: T): T { if (Array.isArray(value)) return value.map(scrub) as unknown as T; if (value && typeof value === "object") { const out: Record = {}; for (const [key, inner] of Object.entries(value as Record)) { - if (key === "resumeCursors") continue; + if (WITHHELD_KEYS.has(key)) continue; out[key] = scrub(inner); } return out as T; diff --git a/companion/test/wire.test.ts b/companion/test/wire.test.ts index e0a1c405cd..a920780b23 100644 --- a/companion/test/wire.test.ts +++ b/companion/test/wire.test.ts @@ -30,6 +30,21 @@ describe("scrub", () => { }); }); + it("withholds the VPS host label but keeps the configured signal", () => { + // GET /api/config and the `config` SSE frame echo the VPS SSH alias — a + // label naming one of the user's servers. The phone renders + // configured-or-not, so that is all it may receive. + const status = { + box: { configured: false }, + vps: { configured: true, sshAlias: "prod-vps" }, + }; + const cleaned = scrub(status); + + expect(JSON.stringify(cleaned)).not.toContain("sshAlias"); + expect(JSON.stringify(cleaned)).not.toContain("prod-vps"); + expect(cleaned).toEqual({ box: { configured: false }, vps: { configured: true } }); + }); + it("leaves values it does not own alone", () => { expect(scrub(null)).toBe(null); expect(scrub(42)).toBe(42); diff --git a/docs/byo-vps.md b/docs/byo-vps.md new file mode 100644 index 0000000000..7c588ee54e --- /dev/null +++ b/docs/byo-vps.md @@ -0,0 +1,111 @@ +# Bring Your Own VPS + +OpenMausBot can turn a Linux server you already own into a bot's computer. The agent process stays on your +machine; Docker's own SSH transport reaches the daemon on the VPS, and each bot gets one managed, hardened +Cua container there — a Linux desktop it can see and control. SSH is the only credential involved and the +only surface exposed: OpenMausBot never opens a port on the VPS, never stores a key or password, and never +runs an agent remotely. + +## What works + +- A per-bot Linux desktop in a managed container on your VPS, driven through the official Cua tools. +- Live screen preview in the Computer panel and in transcripts, same as a Box. +- Explicit **Cloud** with the **Self-hosted VPS** backend provisions or starts the container; **Auto** only + reuses one that is already running and verified. + +Deliberately not offered: an interactive desktop tunnel. There is no "Open desktop" for a VPS bot — the +container publishes no ports, so there is nothing to tunnel to, by design. + +## Prerequisites + +- **Locally:** a `docker` CLI, version 18.09 or newer (that is when `docker -H ssh://` shipped). The Docker + daemon does not need to run locally — only the CLI is used. +- **On the VPS:** a running Docker daemon (`dockerd`) on x86_64 Linux. +- **The SSH user** must be in the `docker` group on the VPS, so `docker info` works without sudo. + +Be clear-eyed about that last point: membership in the `docker` group is root-equivalent on that machine. +Anyone who can talk to the daemon can mount the host filesystem into a container. Using this feature means +trusting the VPS — and whoever else can reach its Docker daemon — completely. Give bots a dedicated server, +not one that also holds things you would not hand to the agent. + +## The required SSH config alias + +OpenMausBot connects only through a named alias in your `~/.ssh/config` — you type the alias into +App Settings → Connections, nothing else. The alias block is load-bearing, not a convenience: every bot +action becomes a `docker exec` over SSH, and without multiplexing each one pays a full SSH handshake; without +keepalives and a connect timeout, a VPS that drops off the network hangs the bot's turn instead of failing it. +Set the block up like this: + +``` +Host my-vps + HostName 203.0.113.7 + User deploy + IdentityFile ~/.ssh/id_ed25519 + ControlMaster auto + ControlPath ~/.ssh/cm-%r@%h-%p + ControlPersist 60m + ServerAliveInterval 15 + ServerAliveCountMax 3 + ConnectTimeout 10 +``` + +- `ControlMaster`/`ControlPath`/`ControlPersist` — every action is a docker-over-SSH exec; multiplexing turns + per-command connects into milliseconds over one persistent connection. +- `ServerAliveInterval`/`ServerAliveCountMax`/`ConnectTimeout` — a dropped VPS must fail fast (under a + minute, and ten seconds to connect), not hang a turn waiting on a dead TCP session. + +**Host key first, by hand.** Connect once manually before pointing OpenMausBot at the alias: + +```sh +ssh my-vps true +``` + +That puts the host key in `known_hosts` on your terms. The app never auto-accepts a host key — an alias whose +host is unknown simply fails until you have done this once. + +## Security + +- **No public ports.** The managed container is created with no published ports, and OpenMausBot refuses to + use a container that publishes any — the check runs before every attach, not just at creation. +- **Firewall the VPS to SSH only**, ideally from your IP. Nothing OpenMausBot does needs any other inbound + port open, so anything else open is pure attack surface. +- **Nothing sensitive is stored.** The only thing OpenMausBot persists is the alias name itself + (`~/.openmausbot/config.json`); keys, passphrases, and agent state stay with SSH. The alias is also kept + off paired phones — the companion reports configured-or-not, never the name. +- The container itself runs hardened: capabilities dropped, private network/IPC/cgroup namespaces, no host + mounts, and memory/CPU/pid limits. A container missing any of that — including one someone created under + the managed name — is refused, not repaired. + +## Container lifecycle + +Each bot owns one container on the VPS, named `openmausbot-vps--` — stable across restarts and +independent of the bot's display name. + +- **Provision** (choosing **Cloud** for the bot, or the panel's button): builds the pinned Cua image on the + VPS if needed, creates the container if missing, starts it if stopped, and waits until the desktop answers. +- **Start** only wakes an existing stopped container; it never creates one. +- **Sleep** stops the container. The VPS stops spending CPU on it; the filesystem stays put. +- **Remove** is yours, done by hand when a bot no longer needs the server: + `docker -H ssh://my-vps rm -f `. OpenMausBot never deletes a container on its own. + +What survives what: sleep/start preserves the container's filesystem; removal — including the recreate that +follows a Cua image upgrade, since a container pinned to an old image is refused rather than reused — wipes +it. Treat the container filesystem as **disposable**: anything a bot must keep should leave the VPS (pushed, +uploaded, or pasted back into chat) before the container is removed. + +A bot set to **Auto** never touches this lifecycle. It attaches only when the container is already running +and verified; otherwise it behaves as if no cloud computer existed. + +## Troubleshooting + +Work up the same path the app takes, cheapest signal first: + +1. **The alias works by hand:** `ssh my-vps true` returns silently. A password prompt means the key/agent is + not set up; a host-key prompt means the first manual connect has not happened yet. +2. **Docker over SSH reaches the daemon:** `docker -H ssh://my-vps info` prints server details. A permission + error means the SSH user is not in the `docker` group. +3. **Provision:** choose **Cloud** with the **Self-hosted VPS** backend in the bot's Computer panel. The + first provision pulls and builds the Cua image on the VPS, which can take minutes; later ones are fast. +4. **Read the status states.** The panel surfaces exactly what the server found, in check order: alias not + configured → daemon unreachable → image missing → container missing / stopped → container unmanaged or + unsafe (ports, mounts, hardening) → desktop not ready. Each message names the step to fix. diff --git a/docs/linux-desktop.md b/docs/linux-desktop.md index 1bd53b3fa2..be829fdc76 100644 --- a/docs/linux-desktop.md +++ b/docs/linux-desktop.md @@ -1,7 +1,8 @@ # Ubuntu Desktop OpenMausBot has an Ubuntu 24.04 LTS x86_64 desktop beta. The Electron package embeds the harness server, so -installed builds do not require Node, pnpm, Swift, or a terminal at runtime. +installed builds do not require Node, pnpm, Swift, or a terminal at runtime. For giving a bot the same kind +of Linux desktop on your own server instead of this machine, see [byo-vps.md](byo-vps.md). ## What works diff --git a/ios/App/ComputerView.swift b/ios/App/ComputerView.swift index f012c7557a..ae63af55bd 100644 --- a/ios/App/ComputerView.swift +++ b/ios/App/ComputerView.swift @@ -61,7 +61,10 @@ struct ComputerView: View { } } .safeAreaInset(edge: .bottom) { - if current.computer == "cloud" { + // A VPS-backed bot is "cloud" too, but the server refuses to mint + // an interactive desktop for it — no button beats a dead one. An + // older harness never sends cloudBackend, so nil keeps the button. + if current.computer == "cloud" && current.cloudBackend != "vps" { VStack(spacing: 8) { if let desktopError { Text(desktopError) diff --git a/ios/Sources/CompanionCore/Models.swift b/ios/Sources/CompanionCore/Models.swift index ede207ef08..dc726330e3 100644 --- a/ios/Sources/CompanionCore/Models.swift +++ b/ios/Sources/CompanionCore/Models.swift @@ -153,6 +153,10 @@ public struct Bot: Codable, Hashable, Identifiable, Sendable { public var autoApprove: Bool? public var alwaysAllow: [String]? public var computer: String? + /// Which cloud computer backs `computer == "cloud"`. Absent (older + /// harnesses included) means the hosted Box; "vps" means the user's own + /// server, which has no interactive desktop to offer a phone. + public var cloudBackend: String? public var speakReplies: Bool? public var voice: String? public var mascotExpression: String? diff --git a/ios/Tests/CompanionCoreTests/DecodingTests.swift b/ios/Tests/CompanionCoreTests/DecodingTests.swift index 8bf0c752e0..11e20cdccc 100644 --- a/ios/Tests/CompanionCoreTests/DecodingTests.swift +++ b/ios/Tests/CompanionCoreTests/DecodingTests.swift @@ -56,6 +56,34 @@ final class DecodingTests: XCTestCase { XCTAssertNil(fleet.bots.first?.hasMore) } + func testDecodesTheCloudBackendAndItsAbsence() throws { + // The cloud-desktop button hides on cloudBackend == "vps", so both + // sides of that gate must decode: a harness that sends the field, and + // an older one that has never heard of it (nil keeps the button). + let json = """ + { + "bots": [ + { + "id":"b1","threadId":"t1","name":"Scout","title":"","description":"", + "notifications":true,"color":"green","unread":false, + "modelSelection":{"instanceId":"i1","model":"m1"},"createdAt":1, + "computer":"cloud","cloudBackend":"vps" + }, + { + "id":"b2","threadId":"t2","name":"Rio","title":"","description":"", + "notifications":true,"color":"blue","unread":false, + "modelSelection":{"instanceId":"i1","model":"m1"},"createdAt":2, + "computer":"cloud" + } + ], + "groups": [] + } + """ + let fleet = try JSONDecoder().decode(Fleet.self, from: Data(json.utf8)) + XCTAssertEqual(fleet.bots.first?.cloudBackend, "vps") + XCTAssertNil(fleet.bots.last?.cloudBackend) + } + func testOneMalformedBotDoesNotHideTheRestOfTheFleet() throws { let json = """ { diff --git a/scripts/bundle-server.mjs b/scripts/bundle-server.mjs index 77474f654e..3b9d45ad84 100644 --- a/scripts/bundle-server.mjs +++ b/scripts/bundle-server.mjs @@ -29,6 +29,7 @@ const ENTRY_POINTS = [ "index.ts", "computer-proxy.ts", "container-mcp.ts", + "vps-container-mcp.ts", "permission-proxy.ts", "connector-proxy.ts", "drivers/agents-proxy.ts", diff --git a/server/branching.test.ts b/server/branching.test.ts index 1fd2c15a7e..cf33c672b8 100644 --- a/server/branching.test.ts +++ b/server/branching.test.ts @@ -187,6 +187,12 @@ posixOnly("conversation branching e2e (fake ACP fleet)", () => { expect((await api("POST", `/api/bots/${created.id}/messages`, { text: "first try" })).status).toBe(202); await waitFor(async () => (await getBot(created.id)).busy === true, "the hung turn to start"); + const backendBefore = (await getBot(created.id)).cloudBackend; + const backendChange = await api("PATCH", `/api/bots/${created.id}`, { cloudBackend: "vps" }); + expect(backendChange.status).toBe(409); + expect(backendChange.body.error).toContain("stop the active turn"); + expect((await getBot(created.id)).cloudBackend).toBe(backendBefore); + // a second send while busy queues (steer-queue) — never a parallel // turn: the words land in the transcript, the live turn keeps running const parallel = await api("POST", `/api/bots/${created.id}/messages`, { text: "sneaky second" }); diff --git a/server/cloud-backend.test.ts b/server/cloud-backend.test.ts new file mode 100644 index 0000000000..7807553b1a --- /dev/null +++ b/server/cloud-backend.test.ts @@ -0,0 +1,29 @@ +import { describe, expect, it } from "vitest"; + +import { + CLOUD_BACKEND_CHANGE_ERROR, + VPS_ALIAS_CHANGE_ERROR, + cloudBackendChangeError, + vpsAliasChangeError, +} from "./cloud-backend.ts"; + +describe("cloud backend switching", () => { + const activeTurnCases: Array<[string, boolean, boolean]> = [ + ["a busy bot", true, false], + ["an active VPS thread", false, true], + ]; + + it.each(activeTurnCases)("rejects changes during %s", (_reason, busy, activeVpsThread) => { + expect(cloudBackendChangeError(busy, activeVpsThread)).toBe(CLOUD_BACKEND_CHANGE_ERROR); + }); + + it("allows changes while idle", () => { + expect(cloudBackendChangeError(false, false)).toBeNull(); + }); + + it("keeps an active VPS turn on its original SSH host", () => { + expect(vpsAliasChangeError("old-vps", "new-vps", true)).toBe(VPS_ALIAS_CHANGE_ERROR); + expect(vpsAliasChangeError("old-vps", "old-vps", true)).toBeNull(); + expect(vpsAliasChangeError("old-vps", "new-vps", false)).toBeNull(); + }); +}); diff --git a/server/cloud-backend.ts b/server/cloud-backend.ts new file mode 100644 index 0000000000..f8be215da1 --- /dev/null +++ b/server/cloud-backend.ts @@ -0,0 +1,10 @@ +export const CLOUD_BACKEND_CHANGE_ERROR = "stop the active turn before changing the cloud backend"; +export const VPS_ALIAS_CHANGE_ERROR = "stop the active VPS turn before changing the SSH config alias"; + +export function cloudBackendChangeError(botBusy: boolean, activeVpsThread: boolean): string | null { + return botBusy || activeVpsThread ? CLOUD_BACKEND_CHANGE_ERROR : null; +} + +export function vpsAliasChangeError(currentAlias: string | null, nextAlias: string | null, activeVpsThread: boolean): string | null { + return activeVpsThread && currentAlias !== nextAlias ? VPS_ALIAS_CHANGE_ERROR : null; +} diff --git a/server/config.test.ts b/server/config.test.ts index 3da54f2ff7..9c62e49899 100644 --- a/server/config.test.ts +++ b/server/config.test.ts @@ -2,8 +2,10 @@ import { describe, expect, it } from "vitest"; import { instanceConfigs, + isValidSshAlias, parseConfigPatch, parseStoredConfig, + vpsSshAlias, withInstanceCli, type AppConfig, } from "./config.ts"; @@ -27,6 +29,17 @@ describe("configuration boundaries", () => { expect(() => parseConfigPatch({ opencodeGo: { apiKey: 42 } })).toThrow("opencodeGo.apiKey"); expect(() => parseConfigPatch({ profile: [] })).toThrow("profile"); }); + + it("accepts only a simple VPS SSH config alias and exposes no credentials", () => { + expect(isValidSshAlias("production-vps")).toBe(true); + expect(isValidSshAlias("prod; reboot")).toBe(false); + expect(() => parseConfigPatch({ vps: { sshAlias: "prod; reboot" } })).toThrow("vps.sshAlias"); + expect(parseConfigPatch({ vps: { sshAlias: "production-vps" } })).toEqual({ + vps: { sshAlias: "production-vps" }, + }); + expect(vpsSshAlias({ vps: { sshAlias: "production-vps" } })).toBe("production-vps"); + expect(vpsSshAlias({ vps: { sshAlias: "-bad" } })).toBeNull(); + }); }); describe("default fleet", () => { diff --git a/server/config.ts b/server/config.ts index c21d5ace40..10557d67d8 100644 --- a/server/config.ts +++ b/server/config.ts @@ -11,6 +11,31 @@ import type { InstanceConfigMap } from "./contracts.ts"; import { parseJson, schemaIssue, type JsonObject, type JsonValue } from "./schema.ts"; const optionalText = z.string().optional(); +const SSH_ALIAS = /^[A-Za-z0-9][A-Za-z0-9_.-]{0,127}$/; + +export function isValidSshAlias(value: unknown): value is string { + return typeof value === "string" && SSH_ALIAS.test(value); +} + +/** Keep the persisted VPS shape deliberately smaller than an SSH connection. */ +export function normalizeVpsConfig(raw: unknown): { sshAlias?: string } { + if (raw === undefined || raw === null) return {}; + if (!raw || typeof raw !== "object" || Array.isArray(raw)) { + throw new Error("vps must be an object containing an SSH config alias"); + } + const alias = (raw as Record).sshAlias; + if (alias === undefined || alias === "") return {}; + if (!isValidSshAlias(alias)) { + throw new Error("vps.sshAlias must be a simple SSH config alias (letters, numbers, dot, dash, or underscore)"); + } + return { sshAlias: alias }; +} + +const vpsConfigSchema = z.object({ + sshAlias: z.string().refine((value) => value === "" || isValidSshAlias(value), { + message: "must be a simple SSH config alias", + }).optional(), +}); const instanceConfigSchema = z.object({ driver: z.string().min(1), displayName: optionalText, @@ -26,6 +51,7 @@ const appConfigSchema = z.object({ * are non-secret local identifiers used to reuse one Composio Session. */ composio: z.object({ apiKey: optionalText, userId: optionalText, sessionId: optionalText }).optional(), box: z.object({ token: optionalText }).optional(), + vps: vpsConfigSchema.optional(), /** OpenCode Go key; persisted write-only and passed only to its child. */ opencodeGo: z.object({ apiKey: optionalText }).optional(), /** Voice credentials and the selected voice id. */ @@ -41,6 +67,8 @@ export interface AppConfig { xai?: { key?: string; url?: string }; composio?: { apiKey?: string; userId?: string; sessionId?: string }; box?: { token?: string }; + /** A named host from the user's SSH config. Authentication stays with SSH. */ + vps?: { sshAlias?: string }; opencodeGo?: { apiKey?: string }; tts?: { key?: string; voice?: string }; profile?: { name?: string; email?: string }; @@ -62,6 +90,10 @@ export function parseConfigPatch(value: JsonValue): ConfigPatch { return parsed.data; } +export function vpsSshAlias(cfg: AppConfig): string | null { + return isValidSshAlias(cfg.vps?.sshAlias) ? cfg.vps.sshAlias : null; +} + // OMB_DATA_DIR isolates test/soak rigs from the user's real fleet. export const DATA_DIR = process.env.OMB_DATA_DIR ?? join(homedir(), ".openmausbot"); const LEGACY_DATA_DIR = join(homedir(), ".opengrokbot"); @@ -117,6 +149,7 @@ export function saveConfig(patch: Partial): void { Object.assign(merged, section); disk[key] = merged; } + if (checkedPatch.vps !== undefined) disk.vps = normalizeVpsConfig(checkedPatch.vps); if (checkedPatch.instances) { const currentInstances = jsonObjectSchema.safeParse(disk.instances); const diskInstances: JsonObject = currentInstances.success ? currentInstances.data : {}; diff --git a/server/container-computer.test.ts b/server/container-computer.test.ts index c80fee7f55..077d5d291d 100644 --- a/server/container-computer.test.ts +++ b/server/container-computer.test.ts @@ -88,6 +88,8 @@ function readyInspect(overrides: Record = {}) { }, State: { Running: true }, Image: "sha256:managed-image-id", + // the full hardened HostConfig the stricter shared check now demands: + // unprivileged, private IPC/cgroup namespaces, pinned shm, no devices HostConfig: { Memory: 4 * 1024 * 1024 * 1024, MemorySwap: 4 * 1024 * 1024 * 1024, @@ -95,6 +97,11 @@ function readyInspect(overrides: Record = {}) { PidsLimit: 512, CapDrop: ["ALL"], CapAdd: ["CAP_SETUID", "CAP_SETGID"], + Privileged: false, + IpcMode: "private", + CgroupnsMode: "private", + ShmSize: 512 * 1024 * 1024, + RestartPolicy: { Name: "no", MaximumRetryCount: 0 }, PortBindings: { "6901/tcp": [{ HostIp: "127.0.0.1" }] }, }, Mounts: [ @@ -197,6 +204,32 @@ describe("containerComputerStatus", () => { expect(status.ready).toBe(false); }); + it("rejects a privileged or host-namespaced Local VM even with correct limits", async () => { + // pins the stricter shared hardening check: resource limits alone are + // not hardening — privilege and namespace escapes disqualify the VM too + const base = JSON.parse(readyInspect())[0].HostConfig; + for (const override of [ + { Privileged: true }, + { IpcMode: "host" }, + { PidMode: "host" }, + { CgroupnsMode: "host" }, + { SecurityOpt: ["seccomp=unconfined"] }, + { DeviceRequests: [{ Driver: "nvidia" }] }, + { RestartPolicy: { Name: "always", MaximumRetryCount: 0 } }, + ]) { + const fake = runner({ + "/usr/bin/which docker": "docker\n", + "/usr/bin/which podman": new Error("missing"), + "docker info --format {{.ServerVersion}}": "29\n", + [`docker image inspect ${IMAGE}`]: preparedImageInspect(), + [`docker inspect ${CONTAINER}`]: readyInspect({ HostConfig: { ...base, ...override } }), + }); + const status = await containerComputerStatus(fake.run, "linux"); + expect(status.security, JSON.stringify(override)).toBe("unsafe"); + expect(status.ready).toBe(false); + } + }); + it("rejects missing or unexpected host mounts instead of exposing them to the bot", async () => { const fake = runner({ "/usr/bin/which docker": "docker\n", @@ -479,6 +512,7 @@ describe("setupCommands", () => { const command = setupCommands("docker", "linux").run!; expect(command).toContain("--memory 4g --memory-swap 4g"); expect(command).toContain("--cpus 2 --pids-limit 512"); + expect(command).toContain("--ipc private --cgroupns private"); expect(command).toContain("--cap-drop ALL --cap-add SETUID --cap-add SETGID"); expect(command).toContain(`--label ${MANAGED_LABEL}=1`); expect(command).toContain(`--label ${DRIVER_LABEL}=${CUA_DRIVER_VERSION}`); diff --git a/server/container-computer.ts b/server/container-computer.ts index 61b50770b7..1b6f97c13c 100644 --- a/server/container-computer.ts +++ b/server/container-computer.ts @@ -59,6 +59,7 @@ const HOST_VIEWER_PORT = 6080; const MEMORY_BYTES = 4 * 1024 * 1024 * 1024; const NANO_CPUS = 2_000_000_000; const PIDS_LIMIT = 512; +const SHM_BYTES = 512 * 1024 * 1024; const LINUX_WHEELS = { x86_64: { @@ -239,7 +240,9 @@ function statusProblem(status: ContainerComputerStatus): string | null { return null; } -function imageLabelsMatch(labels: Record | undefined): boolean { +/** Shared with the BYO-VPS backend (vps-computer.ts): both containers are + * built from the same pinned derivative, so image compatibility is one rule. */ +export function imageLabelsMatch(labels: Record | undefined): boolean { return ( labels?.[MANAGED_LABEL] === "1" && labels?.[DRIVER_LABEL] === CUA_DRIVER_VERSION && @@ -289,10 +292,16 @@ function viewerUrl(password: string | null): string { return `${base}#${fragment.toString()}`; } -function cuaExecArgs(args: string[], interactive = false): string[] { +/** The one authoritative `exec … cua-driver` argv. Shared with the BYO-VPS + * backend and both MCP bridge entry points so the identity, env, and + * telemetry knobs can never drift between the Local VM and a VPS container. */ +export function cuaExecArgs( + args: string[], + options: { container?: string; interactive?: boolean } = {}, +): string[] { return [ "exec", - ...(interactive ? ["-i"] : []), + ...(options.interactive ? ["-i"] : []), "-u", "cua", "-e", @@ -303,7 +312,7 @@ function cuaExecArgs(args: string[], interactive = false): string[] { "CUA_DRIVER_INSTALL_CHANNEL=python_package", "-e", "CUA_DRIVER_RS_TELEMETRY_ENABLED=0", - CONTAINER, + options.container ?? CONTAINER, CUA_EXECUTABLE, ...args, ]; @@ -390,14 +399,8 @@ export async function containerComputerStatus( } else { const inspected = JSON.parse(stdout) as Array<{ Config?: { Image?: string; Labels?: Record; Env?: string[] }; - HostConfig?: { + HostConfig?: DockerHardeningConfig & { PortBindings?: Record | null>; - Memory?: number; - MemorySwap?: number; - NanoCpus?: number; - PidsLimit?: number | null; - CapDrop?: string[] | null; - CapAdd?: string[] | null; }; Mounts?: Array<{ Type?: string; @@ -556,31 +559,71 @@ function appleWorkspaceMountIsSafe( ); } -function dockerSecurityIsHardened( - config: - | { - Memory?: number; - MemorySwap?: number; - NanoCpus?: number; - PidsLimit?: number | null; - CapDrop?: string[] | null; - CapAdd?: string[] | null; - } - | undefined, +/** The Docker/Podman HostConfig surface the hardening check reads. */ +export interface DockerHardeningConfig { + Memory?: number; + MemorySwap?: number; + NanoCpus?: number; + PidsLimit?: number | null; + CapDrop?: string[] | null; + CapAdd?: string[] | null; + Privileged?: boolean; + PidMode?: string; + IpcMode?: string; + UTSMode?: string; + ShmSize?: number; + Devices?: unknown[] | null; + DeviceRequests?: unknown[] | null; + SecurityOpt?: string[] | null; + UsernsMode?: string; + CgroupnsMode?: string; + OomKillDisable?: boolean | null; + AutoRemove?: boolean; + RestartPolicy?: { Name?: string; MaximumRetryCount?: number }; +} + +/** One hardening contract for both managed containers (Local VM here, the + * BYO-VPS backend in vps-computer.ts): exact resource limits, no privilege, + * no host namespaces or devices, no disabled security profiles. The only + * knob the callers legitimately disagree on is the restart policy — the VPS + * container must survive a reboot nobody is watching ("unless-stopped"), + * while the Local VM must NOT auto-resume: its desktop leaves a stale X lock + * on stop, so a restarted container is a broken one. */ +export function dockerSecurityIsHardened( + config: DockerHardeningConfig | undefined, + options: { restartPolicy?: "no" | "unless-stopped" } = {}, ): boolean { if (!config) return false; const capDrop = (config.CapDrop ?? []).map((cap) => cap.toLowerCase()); const capAdd = (config.CapAdd ?? []) .map((cap) => cap.toLowerCase().replace(/^cap_/, "")) .sort(); + const unsafeSecurityOption = (config.SecurityOpt ?? []).some((option) => /(?:^|=)(?:unconfined|disable)$/i.test(option)); + const restartPolicy = config.RestartPolicy?.Name; + const restartPolicyOk = + options.restartPolicy === "unless-stopped" + ? restartPolicy === "unless-stopped" + : restartPolicy === undefined || restartPolicy === "" || restartPolicy === "no"; return ( - (config.Memory ?? 0) >= MEMORY_BYTES && + config.Memory === MEMORY_BYTES && (config.MemorySwap ?? 0) === MEMORY_BYTES && (config.NanoCpus ?? 0) === NANO_CPUS && - (config.PidsLimit ?? 0) > 0 && - (config.PidsLimit ?? Infinity) <= PIDS_LIMIT && + config.PidsLimit === PIDS_LIMIT && capDrop.includes("all") && - capAdd.join(",") === "setgid,setuid" + capAdd.join(",") === "setgid,setuid" && + config.Privileged === false && + !config.PidMode && + config.IpcMode === "private" && + !config.UTSMode && + config.ShmSize === SHM_BYTES && + (!config.Devices || config.Devices.length === 0) && + (!config.DeviceRequests || config.DeviceRequests.length === 0) && + !unsafeSecurityOption && + !config.UsernsMode && + config.CgroupnsMode === "private" && + config.OomKillDisable !== true && + config.AutoRemove !== true && + restartPolicyOk ); } @@ -626,6 +669,14 @@ export function containerRunArgs(runtime: Runtime, password = "CHANGE_ME"): stri "2", "--pids-limit", String(PIDS_LIMIT), + // Pinned explicitly rather than trusting daemon defaults: the shared + // hardening check requires private IPC and cgroup namespaces, and a + // daemon configured with host-mode defaults would otherwise create a + // container its own acceptance check then rejects. + "--ipc", + "private", + "--cgroupns", + "private", "--cap-drop", "ALL", "--cap-add", @@ -708,9 +759,11 @@ export async function containerComputerAction( return containerComputerStatus(runner, platform); } -type ScreenshotCheck = { ok: boolean; mime: "image/png" | "image/jpeg" }; +export type ScreenshotCheck = { ok: boolean; mime: "image/png" | "image/jpeg" }; -function wholeScreenshot(bytes: Buffer): ScreenshotCheck { +/** Shared with the BYO-VPS backend: a truncated base64 transfer must never + * become a "successful" preview frame on either transport. */ +export function wholeScreenshot(bytes: Buffer): ScreenshotCheck { if (bytes.length < 512) return { ok: false, mime: "image/png" }; const png = bytes[0] === 0x89 && bytes[1] === 0x50 && bytes[2] === 0x4e && bytes[3] === 0x47; if (png) { diff --git a/server/container-mcp.ts b/server/container-mcp.ts index 4f2be2f733..4c0562788b 100644 --- a/server/container-mcp.ts +++ b/server/container-mcp.ts @@ -1,8 +1,9 @@ // Transparent stdio bridge into Cua Driver's official MCP server inside the -// Local VM. This process defines no tools and parses no MCP messages. -import { spawn } from "node:child_process"; - -import { augmentedPath } from "./env-path.ts"; +// Local VM. This process defines no tools and parses no MCP messages; the +// piping, drain-safe exit, and watchdog live in mcp-bridge.ts, shared with +// the VPS entry point. +import { cuaExecArgs } from "./container-computer.ts"; +import { runMcpBridge } from "./mcp-bridge.ts"; const [runtime, container, socket] = process.argv.slice(2); if (!runtime || !["docker", "podman", "container"].includes(runtime)) { @@ -14,46 +15,10 @@ if (!container || !/^[a-zA-Z0-9_.-]+$/.test(container) || !socket?.startsWith("/ process.exit(2); } -const child = spawn( - runtime, - [ - "exec", - "-i", - "-u", - "cua", - "-e", - "HOME=/home/cua", - "-e", - "DISPLAY=:1", - "-e", - "CUA_DRIVER_INSTALL_CHANNEL=python_package", - "-e", - "CUA_DRIVER_RS_TELEMETRY_ENABLED=0", - container, - "/usr/local/libexec/openmausbot/cua-driver", - "mcp", - "--socket", - socket, - ], - { - env: { ...process.env, PATH: augmentedPath() }, - stdio: ["pipe", "pipe", "pipe"], - }, -); - -process.stdin.pipe(child.stdin); -child.stdout.pipe(process.stdout); -child.stderr.pipe(process.stderr); - -child.on("error", (error) => { - process.stderr.write(`could not connect to Cua Driver: ${error.message}\n`); - process.exit(1); +runMcpBridge({ + command: runtime, + args: cuaExecArgs(["mcp", "--socket", socket], { container, interactive: true }), + label: "Cua Driver", + // No liveness watchdog: the runtime CLI talks to a local daemon and fails + // fast on its own — there is no silent WAN peer to wedge on. }); -child.on("close", (code, signal) => { - if (signal) process.stderr.write(`Cua Driver connection ended with ${signal}\n`); - process.exit(code ?? 1); -}); - -for (const signal of ["SIGTERM", "SIGINT"] as const) { - process.on(signal, () => child.kill(signal)); -} diff --git a/server/contracts.ts b/server/contracts.ts index bde9dcc3fb..6abd22684b 100644 --- a/server/contracts.ts +++ b/server/contracts.ts @@ -9,6 +9,7 @@ export type DriverKind = string; export type InstanceId = string; export type ThreadId = string; export type TurnId = string; +export type CloudBackend = "box" | "vps"; export type ProviderErrorCode = | "missing_cli" @@ -154,7 +155,7 @@ export interface SendTurnInput { composio?: { command: string; args: string[]; env: Record }; /** Cloud computer, reached through OpenMausBot's REST-to-MCP adapter. */ computer?: { kind?: "box"; boxId: string; token: string }; - /** Direct stdio connection to a Cua Driver MCP server (host or sandbox). */ + /** Direct stdio connection to a Cua Driver MCP server (host, sandbox, or VPS). */ localComputer?: { command: string; args: string[]; env: Record }; /** Peer-agent comms: an MCP proxy (list_bots / ask_bot) that routes back * through the harness so this bot can message other bots. The harness diff --git a/server/index.test.ts b/server/index.test.ts index d489cbfe92..7274f70cb9 100644 --- a/server/index.test.ts +++ b/server/index.test.ts @@ -604,6 +604,26 @@ describe("harness HTTP API", () => { expect(nothing.status).toBe(400); }); + it("validates the non-secret VPS alias and keeps old bots on Box by default", async () => { + const before = await api("GET", "/api/bots"); + const bot = before.body.bots[0]; + expect(bot.cloudBackend).toBeUndefined(); + + const bad = await api("PUT", "/api/config", { vps: { sshAlias: "prod; reboot" } }); + expect(bad.status).toBe(400); + + const saved = await api("PUT", "/api/config", { vps: { sshAlias: "production-vps" } }); + expect(saved.status).toBe(200); + expect(saved.body.vps).toEqual({ configured: true, sshAlias: "production-vps" }); + expect(JSON.stringify(saved.body)).not.toContain("privateKey"); + + const patched = await api("PATCH", `/api/bots/${bot.id}`, { cloudBackend: "vps" }); + expect(patched.status).toBe(200); + expect(patched.body.bot.cloudBackend).toBe("vps"); + const invalid = await api("PATCH", `/api/bots/${bot.id}`, { cloudBackend: "daytona" }); + expect(invalid.status).toBe(400); + }); + it("validates a Composio project key, creates a Session, and keeps externally stored secrets off disk", async () => { const oldKey = await api("PUT", "/api/config", { composio: { apiKey: "old_key" } }); expect(oldKey.status).toBe(400); diff --git a/server/index.ts b/server/index.ts index fdda9a5179..4b94dfd351 100644 --- a/server/index.ts +++ b/server/index.ts @@ -13,6 +13,7 @@ import { approvalKey, autoDecision } from "./auto-approve.ts"; import { validateBotCwd } from "./bot-cwd.ts"; import { groupTurnCwd } from "./room-cwd.ts"; import * as box from "./box.ts"; +import { cloudBackendChangeError, vpsAliasChangeError } from "./cloud-backend.ts"; import * as composio from "./composio.ts"; import { chiefOfStaffSystemPrompt } from "./chief-of-staff.ts"; import { @@ -30,6 +31,7 @@ import { parseConfigPatch, saveConfig, withInstanceCli, + vpsSshAlias, EVENTS_DIR, NATIVE_DIR, } from "./config.ts"; @@ -72,6 +74,7 @@ import { readCuaConnection } from "./local-computer.ts"; import { LocalVmIdleTimer } from "./local-vm-idle.ts"; import { LocalVmLease } from "./local-vm-lease.ts"; import { RepeatDetector, callKey } from "./repeat-detector.ts"; +import * as vps from "./vps-computer.ts"; import { RoutineManager, type RoutineRunOn, type RoutineRunTrigger } from "./routines.ts"; import { fetchGithubTeam, fetchLibraryTeam, fetchTeamCatalog } from "./team-library.ts"; import { createTeamManifest, parseTeamManifest } from "./team-manifest.ts"; @@ -492,6 +495,7 @@ const watchdog = new TurnWatchdog({ const currentBot = store.bot(turn.botId); if (currentBot?.busy) { stopScreenPoller(currentBot.id); + if (activeVpsThreads.get(currentBot.id) === turn.threadId) activeVpsThreads.delete(currentBot.id); store.setActivity(currentBot.id, "idle"); } }, 6_000); @@ -551,6 +555,7 @@ const localVmLease = new LocalVmLease(30 * 60_000); const localVmOwnerBusy = (botId: string) => store.bot(botId)?.busy === true; let localVmLifecycleBusy = false; let localVmActiveThread: string | null = null; +const activeVpsThreads = new Map(); const LOCAL_VM_IDLE_MS = 8 * 60 * 60_000; const localVmIdle = new LocalVmIdleTimer( LOCAL_VM_IDLE_MS, @@ -771,6 +776,10 @@ bus.subscribe((event: RuntimeEvent) => { // tally is not the right home for a shared room's spend, so only // 1:1 task turns are tallied for now. if (bot) { + const vpsTurn = activeVpsThreads.get(bot.id) === event.threadId; + const clearVpsTurn = () => { + if (activeVpsThreads.get(bot.id) === event.threadId) activeVpsThreads.delete(bot.id); + }; // bank what this turn spent before the bot broadcast carries the // task list to every window. The driver's own per-turn figure // (turn.completed.usage) is authoritative; a driver that only @@ -795,7 +804,9 @@ bus.subscribe((event: RuntimeEvent) => { if (frame && store.bot(bot.id)) { pushMessage({ role: "bot", kind: "screen", png: frame.png, mime: frame.mime }); } - }); + }).finally(clearVpsTurn); + } else if (vpsTurn) { + clearVpsTurn(); } } const speaker = groupSpeakers.get(event.threadId); @@ -968,7 +979,7 @@ function drainQueuedSends() { ); } -// ── live screen: poll the bot's box while it works ──────────────────── +// ── live screen: poll the bot's computer while it works ─────────────── // Frames stream to clients as SSE {kind:'screen'} (the "Bot's screen" // panel); the final frame is folded into the transcript on turn end. type Frame = { png: string; mime: string }; @@ -997,8 +1008,12 @@ const SCREEN_MIN_GAP_MS = 3000; /** `screenIsTheWork` starts the turn already counting as screen usage: a * boxAgent's whole session runs ON the box, so every tool it calls acts on * that screen even though none of them is named like a computer tool. */ -function startScreenPoller(botId: string, boxId?: string, { screenIsTheWork = false } = {}) { - if (screenPollers.has(botId) || !box.boxConfigured(cfg)) return; +function startScreenPoller( + botId: string, + capture: () => Promise<{ png: string; format: string }>, + { screenIsTheWork = false } = {}, +) { + if (screenPollers.has(botId)) return; // One capture at a time, shared by the interval, the pokes, and the // turn-end grab: awaiting the in-flight promise (rather than dropping the // call) is what lets the final frame be the settled one. The min-gap keeps @@ -1012,9 +1027,7 @@ function startScreenPoller(botId: string, boxId?: string, { screenIsTheWork = fa if (!current && Date.now() - lastAt < SCREEN_MIN_GAP_MS) return Promise.resolve(); current ??= (async () => { try { - // boxId is resolved once per turn — re-resolving per frame cost a - // full LIST of the account's boxes - const { png, format } = await box.screenshotBox(cfg, botId, boxId); + const { png, format } = await capture(); const frame = { png, mime: format === "jpeg" ? "image/jpeg" : "image/png" }; entry.last = frame; broadcast({ kind: "screen", botId, ...frame }); @@ -1234,10 +1247,13 @@ async function startTurn( const dwebUrl = process.env.DWEB_URL?.trim(); if (dwebUrl) integrations.dweb = { url: dwebUrl }; const wants = opts?.runOn === "cloud" ? "cloud" : bot.computer; // cloud routine overrides the MAUS default + // Cloud routines always use Box/BoxAgent. The per-bot backend applies + // only to ordinary turns that mount a computer into the local agent. + const cloudBackend = opts?.runOn === "cloud" || bot.cloudBackend !== "vps" ? "box" : "vps"; const mountsComputerMcp = instance.adapter.capabilities.computerMcp === true; const mountsCloudComputer = mountsComputerMcp || instance.driverKind === "boxAgent"; - let previewBoxId: string | null = null; - let computerKind: "box" | "vm" | "local" | null = null; + let previewCapture: (() => Promise<{ png: string; format: string }>) | null = null; + let computerKind: "box" | "vps" | "vm" | "local" | null = null; // Explicit destinations are strict. In particular, Local VM must never // fall through to host CUA and accidentally click on the user's Mac. @@ -1272,9 +1288,34 @@ async function startTurn( computerKind = "local"; } + // A VPS is a local-agent computer mount, never a remote agent runner. + // Explicit Cloud may prepare/start it; Auto is read-only and can only + // attach to an already-running, verified container. + if ((wants === "cloud" || wants === undefined) && cloudBackend === "vps") { + const unsupported = vps.vpsDriverError(instance.driverKind, mountsComputerMcp); + if (unsupported && wants === "cloud") throw new Error(unsupported); + if (!unsupported) { + activeVpsThreads.set(bot.id, threadId); + const remote = wants === "cloud" + ? await vps.vpsComputerAction("provision", cfg, bot.id) + : await vps.reuseVps(cfg, bot.id); + if (remote?.ready && remote.sshAlias) { + const targetCfg = { ...cfg, vps: { sshAlias: remote.sshAlias } }; + integrations.localComputer = vps.vpsComputerMcp(targetCfg, bot.id, remote.container_id ?? undefined); + computerKind = "vps"; + previewCapture = () => vps.vpsComputerScreenshot(targetCfg, bot.id); + } else { + activeVpsThreads.delete(bot.id); + if (wants === "cloud") { + throw new Error(remote?.problem ?? "the VPS computer could not be created or reached"); + } + } + } + } + // Cloud is also strict when explicitly selected. Auto (unset) reuses an // existing cloud box, then falls back to host CUA without provisioning. - if ((wants === "cloud" || wants === undefined) && box.boxConfigured(cfg)) { + if ((wants === "cloud" || wants === undefined) && cloudBackend === "box" && box.boxConfigured(cfg)) { if (!mountsCloudComputer && wants === "cloud") { throw new Error("this model engine cannot use computer tools — choose Claude, an ACP engine, or the Computer engine"); } @@ -1295,17 +1336,17 @@ async function startTurn( b = (await box.readyBox(cfg, bot.id).catch(() => null)) ?? b; } if (b) { - previewBoxId = b.id; + previewCapture = () => box.screenshotBox(cfg, bot.id, b!.id); if (mountsCloudComputer) { integrations.computer = { kind: "box", boxId: b.id, token: cfg.box!.token! }; computerKind = "box"; } } } - if (wants === "cloud" && !box.boxConfigured(cfg)) { + if (wants === "cloud" && cloudBackend === "box" && !box.boxConfigured(cfg)) { throw new Error("Cloud box is not configured — add a Box API key or choose Local VM"); } - if (wants === "cloud" && !integrations.computer) { + if (wants === "cloud" && cloudBackend === "box" && !integrations.computer) { throw new Error("the cloud computer could not be created or reached"); } @@ -1347,6 +1388,8 @@ async function startTurn( ? "You can work with the user's other bots through the agents tools — list_bots shows who's available, ask_bot sends one of them a message and returns their reply." : ""; + // (activeVpsThreads was already claimed above, before the provision or + // reuse await, so the backend guards saw this turn the whole time.) watchdog.watch(threadId, bot.id); await instance.adapter.sendTurn({ threadId, @@ -1364,6 +1407,8 @@ async function startTurn( ? " You have a shared, isolated Cua sandbox: a Linux desktop in a container on this machine. Only /home/cua/workspace is durable; save downloads, repositories, working files, and browser profiles there because everything else inside the VM is disposable. No other host folder is mounted. Use the computer tools for desktop, accessibility, window, and shell work. Inspect the desktop state before acting, prefer accessibility targets over raw coordinates, and work carefully." : computerKind === "box" && instance.driverKind !== "boxAgent" ? " You have your own cloud computer. In Chrome, prefer browser_snapshot with browser_click/browser_fill for semantic, trusted actions; use screenshot/click/type_text for visual or non-browser UI, open_url for navigation, and computer_exec for Linux tasks. Every action already returns the resulting screen, so don't follow it with screenshot; batch predictable pixel actions with computer_batch." + : computerKind === "vps" + ? " You have your own self-hosted remote Linux computer through the official Cua tools. Its filesystem is disposable: everything on it is wiped whenever its container is recreated, so keep long-lived work somewhere durable — push it to a remote, or hand the results back in chat — instead of leaving it only on that computer. Inspect the desktop state before acting, prefer accessibility targets over raw coordinates, and act carefully." : computerKind === "local" ? " You can act on the user's computer through the computer tools — take a screenshot or read the desktop state first, prefer accessibility actions over raw coordinates, and act carefully." : "") + @@ -1397,12 +1442,13 @@ async function startTurn( // after its own turn.completed would never be torn down — it would // keep polling the box forever, carrying dead per-turn state. busy // is flipped false in the fold, so it is the honest "still running". - if (previewBoxId && store.bot(bot.id)?.busy) { - startScreenPoller(bot.id, previewBoxId, { screenIsTheWork: instance.driverKind === "boxAgent" }); + if (previewCapture && store.bot(bot.id)?.busy) { + startScreenPoller(bot.id, previewCapture, { screenIsTheWork: instance.driverKind === "boxAgent" }); } } catch (e) { localVmLease.release(threadId); if (localVmActiveThread === threadId) localVmActiveThread = null; + if (activeVpsThreads.get(bot.id) === threadId) activeVpsThreads.delete(bot.id); watchdog.settle(threadId); turnUsage.delete(threadId); const message = e instanceof Error ? e.message : String(e); @@ -1919,6 +1965,7 @@ function configStatus() { mode: composio.connectionMode(cfg), }, box: { configured: Boolean(cfg.box?.token) }, + vps: { configured: Boolean(vpsSshAlias(cfg)), sshAlias: vpsSshAlias(cfg) ?? "" }, opencodeGo: { configured: Boolean(cfg.opencodeGo?.apiKey) }, // the chosen voice is a setting, not a secret; the key is reported the // same configured-or-not way as every other credential @@ -1945,6 +1992,7 @@ async function reloadProviders() { if (localVmActiveThread === vmLease.threadId) localVmActiveThread = null; } stopScreenPoller(b.id); + activeVpsThreads.delete(b.id); finalizeDelegationWatch( b.threadId, false, @@ -2819,7 +2867,7 @@ const server = createServer(async (req, res) => { if (field === "name" && !value.trim()) return json(res, 400, { error: "name must not be empty" }); } const patch: Record = {}; - for (const key of ["name", "title", "description", "notifications", "modelSelection", "unread", "computer", "color", "mascotExpression", "pinned", "hidden", "speakReplies", "voice"] as const) { + for (const key of ["name", "title", "description", "notifications", "modelSelection", "unread", "computer", "cloudBackend", "color", "mascotExpression", "pinned", "hidden", "speakReplies", "voice"] as const) { if (body[key] !== undefined) patch[key] = body[key]; } // per-bot gate on the workspace's connected apps (Composio) @@ -2833,9 +2881,16 @@ const server = createServer(async (req, res) => { ) { return json(res, 400, { error: "computer must be cloud, vm, local, or off" }); } + if (body.cloudBackend !== undefined && !["box", "vps"].includes(String(body.cloudBackend))) { + return json(res, 400, { error: "cloudBackend must be box or vps" }); + } if (body.chiefOfStaff !== undefined && typeof body.chiefOfStaff !== "boolean") { return json(res, 400, { error: "chiefOfStaff must be true or false" }); } + if (body.cloudBackend !== undefined) { + const backendError = cloudBackendChangeError(Boolean(existing?.busy), activeVpsThreads.has(m[1])); + if (backendError) return json(res, 409, { error: backendError }); + } if (body.cwd !== undefined) { const checked = validateBotCwd(body.cwd); if (!checked.ok) return json(res, 400, { error: checked.error }); @@ -2883,6 +2938,7 @@ const server = createServer(async (req, res) => { // a running turn dies with its bot await registry.get(bot.modelSelection.instanceId)?.adapter.interruptTurn(bot.threadId).catch(() => {}); stopScreenPoller(bot.id); + activeVpsThreads.delete(bot.id); routines!.disableForBot(bot.id); webhooks.disableForBot(bot.id); lastReply.delete(bot.threadId); @@ -3280,6 +3336,12 @@ const server = createServer(async (req, res) => { const patch = parseConfigPatch(body); if (!Object.keys(patch).length) return json(res, 400, { error: "nothing to save" }); if (providerConfigBusy) return json(res, 409, { error: "provider settings are already being updated" }); + if (patch.vps !== undefined) { + const currentAlias = vpsSshAlias(cfg); + const nextAlias = vpsSshAlias({ ...cfg, vps: patch.vps }); + const aliasError = vpsAliasChangeError(currentAlias, nextAlias, activeVpsThreads.size > 0); + if (aliasError) return json(res, 409, { error: aliasError }); + } providerConfigBusy = true; try { // A project key is useful only if it can create/reuse the Session that @@ -3331,7 +3393,9 @@ const server = createServer(async (req, res) => { // provider keys change the fleet; a profile or voice edit must not // kill in-flight turns with a pointless reload — no driver reads // either, and picking a voice mid-turn should be free - if (Object.keys(patch).some((k) => k !== "profile" && k !== "tts")) await reloadProviders(); + // The VPS alias is consumed by lifecycle commands, not provider + // engines. Saving it must not interrupt an in-flight turn. + if (Object.keys(patch).some((k) => k !== "profile" && k !== "tts" && k !== "vps")) await reloadProviders(); const status = configStatus(); broadcast({ kind: "config", ...status }); return json(res, 200, status); @@ -3452,12 +3516,44 @@ const server = createServer(async (req, res) => { // ── the bot's cloud computer (Box) ── m = path.match(/^\/api\/bots\/([\w-]+)\/computer$/); - if (m && method === "GET") return json(res, 200, await box.boxStatus(cfg, m[1])); - m = path.match(/^\/api\/bots\/([\w-]+)\/computer\/(provision|join|sleep|exec|screenshot)$/); + if (m && method === "GET") { + const bot = store.bot(m[1]); + if (!bot) return json(res, 404, { error: "no such bot" }); + return bot.cloudBackend === "vps" + ? json(res, 200, { backend: "vps", ...(await vps.vpsComputerStatus(cfg, bot.id)) }) + : json(res, 200, { backend: "box", ...(await box.boxStatus(cfg, bot.id)) }); + } + m = path.match(/^\/api\/bots\/([\w-]+)\/computer\/(provision|join|sleep|exec|screenshot|remove)$/); if (m && method === "POST") { const botId = m[1]; const bot = store.bot(botId); if (!bot) return json(res, 404, { error: "no such bot" }); + // Requiring JSON makes every computer mutation a non-simple browser + // request (same reasoning as the Local VM lifecycle routes above): a + // hostile page cannot submit it with a form, and its cross-origin JSON + // request dies in the preflight this server never answers. Applied to + // both backends — the Box branch runs commands too. + if (!String(req.headers["content-type"] ?? "").toLowerCase().startsWith("application/json")) { + return json(res, 415, { error: "content-type must be application/json" }); + } + if (bot.cloudBackend === "vps") { + if (m[2] === "join" || m[2] === "exec") { + return json(res, 409, { error: "interactive VPS desktop access is not supported" }); + } + if (m[2] === "provision" && bot.computer !== "cloud") { + return json(res, 409, { error: "Auto mode will not provision a VPS; choose Cloud for this bot first" }); + } + if ((m[2] === "sleep" || m[2] === "remove") && (bot.busy || activeVpsThreads.has(botId))) { + return json(res, 409, { error: "the VPS computer is being used by this bot — interrupt the turn first" }); + } + if (m[2] === "screenshot") return json(res, 200, await vps.vpsComputerScreenshot(cfg, botId)); + const action = m[2] === "provision" ? "provision" : m[2] === "remove" ? "remove" : "stop"; + return json(res, 200, await vps.vpsComputerAction(action, cfg, botId)); + } + if (m[2] === "remove") { + // Boxes sleep and wake; only the VPS backend has a container to remove. + return json(res, 409, { error: "the cloud Box backend has no container to remove — use sleep instead" }); + } switch (m[2]) { case "provision": return json(res, 200, await box.provisionBox(cfg, botId, bot.name)); diff --git a/server/mcp-bridge.test.ts b/server/mcp-bridge.test.ts new file mode 100644 index 0000000000..33dc8dc5c4 --- /dev/null +++ b/server/mcp-bridge.test.ts @@ -0,0 +1,124 @@ +// The bridge's dead-transport watchdog, pinned at the unit level: the 45s +// e2e wait is too slow for the suite, and the property that matters is not +// the constant but the decision table — silence alone never kills, only +// silence PLUS a failed liveness probe does, and traffic always vetoes. +import { describe, expect, it, vi } from "vitest"; + +import { createInactivityWatchdog, runLivenessProbe } from "./mcp-bridge.ts"; + +/** a probe whose answers the test scripts one call at a time */ +function scriptedProbe(answers: boolean[]) { + const calls: Array<(alive: boolean) => void> = []; + let handed = 0; + return { + calls, + probe: () => + new Promise((resolve) => { + calls.push(resolve); + const next = answers[handed]; + handed += 1; + if (next !== undefined) resolve(next); + }), + }; +} + +describe("createInactivityWatchdog", () => { + it("kills only after silence AND a failed probe, then never re-arms", async () => { + vi.useFakeTimers(); + try { + const onDead = vi.fn(); + const scripted = scriptedProbe([true, false]); + createInactivityWatchdog({ inactivityMs: 1_000, probe: scripted.probe, onDead }); + + // first silence window: the probe answers alive → no kill, re-armed + await vi.advanceTimersByTimeAsync(1_000); + expect(scripted.calls).toHaveLength(1); + expect(onDead).not.toHaveBeenCalled(); + + // second silence window: the probe fails → dead, exactly once + await vi.advanceTimersByTimeAsync(1_000); + expect(scripted.calls).toHaveLength(2); + expect(onDead).toHaveBeenCalledTimes(1); + + // dead is terminal: no timer survives to fire again + await vi.advanceTimersByTimeAsync(10_000); + expect(onDead).toHaveBeenCalledTimes(1); + } finally { + vi.useRealTimers(); + } + }); + + it("treats traffic as proof of life, resetting the window and vetoing an in-flight probe", async () => { + vi.useFakeTimers(); + try { + const onDead = vi.fn(); + let probeResolvers: Array<(alive: boolean) => void> = []; + const watchdog = createInactivityWatchdog({ + inactivityMs: 1_000, + probe: () => new Promise((resolve) => probeResolvers.push(resolve)), + onDead, + }); + + // steady traffic keeps the probe from ever firing + for (let i = 0; i < 5; i += 1) { + await vi.advanceTimersByTimeAsync(900); + watchdog.touch(); + } + expect(probeResolvers).toHaveLength(0); + + // silence fires the probe — but a byte arriving WHILE it runs must + // outrank even a failed answer (a slow screenshot finishing is life) + await vi.advanceTimersByTimeAsync(1_000); + expect(probeResolvers).toHaveLength(1); + watchdog.touch(); + probeResolvers[0]!(false); // SAFETY: length asserted above + await vi.advanceTimersByTimeAsync(0); + expect(onDead).not.toHaveBeenCalled(); + + // the veto re-armed the window; a stopped watchdog stays quiet + watchdog.stop(); + probeResolvers = []; + await vi.advanceTimersByTimeAsync(10_000); + expect(probeResolvers).toHaveLength(0); + expect(onDead).not.toHaveBeenCalled(); + } finally { + vi.useRealTimers(); + } + }); + + it("reads a rejected probe as not alive", async () => { + vi.useFakeTimers(); + try { + const onDead = vi.fn(); + createInactivityWatchdog({ + inactivityMs: 1_000, + probe: () => Promise.reject(new Error("probe spawn failed")), + onDead, + }); + await vi.advanceTimersByTimeAsync(1_000); + expect(onDead).toHaveBeenCalledTimes(1); + } finally { + vi.useRealTimers(); + } + }); +}); + +describe("runLivenessProbe", () => { + it("maps exit status to liveness and treats an unspawnable probe as dead", async () => { + await expect( + runLivenessProbe({ command: process.execPath, args: ["-e", "process.exit(0)"] }), + ).resolves.toBe(true); + await expect( + runLivenessProbe({ command: process.execPath, args: ["-e", "process.exit(3)"] }), + ).resolves.toBe(false); + await expect( + runLivenessProbe({ command: "/nonexistent/openmausbot-probe", args: [] }), + ).resolves.toBe(false); + }); + + it("times out a probe that hangs instead of inheriting the hang", async () => { + await expect( + runLivenessProbe({ command: process.execPath, args: ["-e", "setInterval(() => {}, 1000)"] }, 300), + ).resolves.toBe(false); + }); +}); diff --git a/server/mcp-bridge.ts b/server/mcp-bridge.ts new file mode 100644 index 0000000000..ee3c9a22af --- /dev/null +++ b/server/mcp-bridge.ts @@ -0,0 +1,186 @@ +// The one transparent stdio bridge behind both MCP entry points +// (container-mcp.ts for the Local VM, vps-container-mcp.ts for the BYO VPS). +// It defines no tools and parses no MCP messages: bytes in, bytes out. +// +// Two behaviors live here so neither entry point can drift: +// 1. Exit without truncation. `process.exit()` in a close/error handler +// discards whatever is still buffered on stdout — a final MCP result +// would be cut mid-frame. The bridge sets exitCode and unpipes instead, +// letting stdio drain before the process ends on its own. +// 2. A dead-transport watchdog (opt-in via `liveness`). docker's ssh +// connection helper accepts no ConnectTimeout/ServerAlive options, so a +// VPS dropping mid-turn leaves the exec silently wedged until the OS +// gives up — the harness sees a hung tool call, not an error. +import { spawn } from "node:child_process"; + +import { augmentedPath } from "./env-path.ts"; + +// 45s of TOTAL silence before the bridge even probes. An MCP session is +// legitimately quiet between tool calls and a slow screenshot can take tens +// of seconds, so silence alone never kills anything — it only triggers a +// liveness probe, and only a probe that FAILS ends the bridge. Any byte on +// stdin/stdout/stderr resets the window. +export const BRIDGE_INACTIVITY_MS = 45_000; +const PROBE_TIMEOUT_MS = 10_000; + +export interface BridgeLiveness { + command: string; + args: string[]; +} + +/** Run the liveness command; alive means "exited 0 within the timeout". The + * probe is its own short-lived process, so it cannot inherit the wedged + * connection it is diagnosing. */ +export function runLivenessProbe(probe: BridgeLiveness, timeoutMs = PROBE_TIMEOUT_MS): Promise { + return new Promise((resolve) => { + const child = spawn(probe.command, probe.args, { + shell: false, + env: { ...process.env, PATH: augmentedPath() }, + stdio: ["ignore", "ignore", "ignore"], + }); + const timer = setTimeout(() => { + child.kill("SIGKILL"); + resolve(false); + }, timeoutMs); + timer.unref?.(); + child.on("error", () => { + clearTimeout(timer); + resolve(false); + }); + child.on("close", (code) => { + clearTimeout(timer); + resolve(code === 0); + }); + }); +} + +export interface WatchdogHandle { + /** Any traffic in either direction — resets the inactivity window. */ + touch: () => void; + stop: () => void; +} + +/** Inactivity → probe → (only then) declare dead. Traffic arriving while a + * probe is in flight vetoes even a failed probe: bytes are better evidence + * of life than a health command racing a congested link. */ +export function createInactivityWatchdog(options: { + inactivityMs: number; + probe: () => Promise; + onDead: () => void; +}): WatchdogHandle { + let timer: ReturnType | null = null; + let stopped = false; + let probing = false; + let touchedWhileProbing = false; + + const arm = () => { + if (stopped) return; + timer = setTimeout(fire, options.inactivityMs); + timer.unref?.(); + }; + const settleProbe = (alive: boolean) => { + probing = false; + if (stopped) return; + if (alive || touchedWhileProbing) { + arm(); + return; + } + options.onDead(); + }; + const fire = () => { + probing = true; + touchedWhileProbing = false; + void options.probe().then(settleProbe, () => settleProbe(false)); + }; + + arm(); + return { + touch() { + if (stopped) return; + if (probing) { + touchedWhileProbing = true; + return; + } + if (timer) clearTimeout(timer); + arm(); + }, + stop() { + stopped = true; + if (timer) clearTimeout(timer); + }, + }; +} + +export interface BridgeOptions { + command: string; + args: string[]; + /** Names the far end in stderr messages, e.g. "Cua Driver". */ + label: string; + /** Enables the dead-transport watchdog. Omitted for the Local VM, whose + * runtime CLI talks to a local daemon and fails fast on its own. */ + liveness?: BridgeLiveness; +} + +export function runMcpBridge(options: BridgeOptions): void { + const child = spawn(options.command, options.args, { + shell: false, + env: { ...process.env, PATH: augmentedPath() }, + stdio: ["pipe", "pipe", "pipe"], + }); + + // docker may exit before it drains stdin; pipe() leaves this error unhandled. + child.stdin.on("error", () => {}); + process.stdin.pipe(child.stdin); + child.stdout.pipe(process.stdout); + child.stderr.pipe(process.stderr); + + const detach = () => { + process.stdin.unpipe(child.stdin); + process.stdin.pause(); + }; + + let watchdog: WatchdogHandle | null = null; + if (options.liveness) { + const liveness = options.liveness; + watchdog = createInactivityWatchdog({ + inactivityMs: BRIDGE_INACTIVITY_MS, + probe: () => runLivenessProbe(liveness), + onDead: () => { + process.stderr.write( + `${options.label} transport went silent and stopped answering liveness probes; ending the bridge\n`, + ); + process.exitCode = 1; + detach(); + child.kill("SIGKILL"); + // A docker wedged on a dead ssh connection may never deliver close. + // Nothing can be buffered on stdout after 45 quiet seconds, so this + // hard exit — unlike the close-handler one this file exists to avoid — + // cannot truncate anything. + const failsafe = setTimeout(() => process.exit(1), 2_000); + failsafe.unref?.(); + }, + }); + const touch = () => watchdog?.touch(); + process.stdin.on("data", touch); + child.stdout.on("data", touch); + child.stderr.on("data", touch); + } + + child.on("error", (error) => { + process.stderr.write(`could not connect to ${options.label}: ${error.message}\n`); + process.exitCode = 1; + watchdog?.stop(); + detach(); + }); + child.on("close", (code, signal) => { + if (signal) process.stderr.write(`${options.label} connection ended with ${signal}\n`); + // Let stdout and stderr drain before the bridge exits. + process.exitCode = process.exitCode ?? code ?? 1; + watchdog?.stop(); + detach(); + }); + + for (const signal of ["SIGTERM", "SIGINT"] as const) { + process.on(signal, () => child.kill(signal)); + } +} diff --git a/server/proxy-paths.ts b/server/proxy-paths.ts index 30abf9a433..f6bd30c332 100644 --- a/server/proxy-paths.ts +++ b/server/proxy-paths.ts @@ -37,6 +37,7 @@ export const SPAWNED_PROXIES = { computer: resolveProxy("computer-proxy"), permission: resolveProxy("permission-proxy"), containerMcp: resolveProxy("container-mcp"), + vpsContainerMcp: resolveProxy("vps-container-mcp"), agents: resolveProxy("drivers/agents-proxy"), dweb: resolveProxy("drivers/dweb-proxy"), connectors: resolveProxy("connector-proxy"), diff --git a/server/store.test.ts b/server/store.test.ts index 4ceeca13f0..2baf0c89ca 100644 --- a/server/store.test.ts +++ b/server/store.test.ts @@ -110,6 +110,32 @@ describe("Store", () => { expect(messages.at(-1)).toMatchObject({ role: "user", text: "hi there" }); }); + it("normalizes persisted cloud backends without changing valid or absent values", () => { + const store = new Store(selection); + const box = store.createBot(); + const vps = store.createBot(); + const invalid = store.createBot(); + const absent = store.createBot(); + const raw: BotRecord[] = JSON.parse(readFileSync(join(DATA_DIR, "bots.json"), "utf8")); + raw.find((bot) => bot.id === box.id)!.cloudBackend = "box"; + raw.find((bot) => bot.id === vps.id)!.cloudBackend = "vps"; + (raw.find((bot) => bot.id === invalid.id) as unknown as { cloudBackend: string }).cloudBackend = "daytona"; + delete raw.find((bot) => bot.id === absent.id)!.cloudBackend; + writeFileSync(join(DATA_DIR, "bots.json"), JSON.stringify(raw)); + + const reloaded = new Store(selection); + expect(reloaded.bot(box.id)?.cloudBackend).toBe("box"); + expect(reloaded.bot(vps.id)?.cloudBackend).toBe("vps"); + expect(reloaded.bot(invalid.id)?.cloudBackend).toBeUndefined(); + expect(reloaded.bot(absent.id)?.cloudBackend).toBeUndefined(); + + const saved: BotRecord[] = JSON.parse(readFileSync(join(DATA_DIR, "bots.json"), "utf8")); + expect(saved.find((bot) => bot.id === box.id)?.cloudBackend).toBe("box"); + expect(saved.find((bot) => bot.id === vps.id)?.cloudBackend).toBe("vps"); + expect(saved.find((bot) => bot.id === invalid.id)).not.toHaveProperty("cloudBackend"); + expect(saved.find((bot) => bot.id === absent.id)).not.toHaveProperty("cloudBackend"); + }); + it("migrates unambiguous legacy peer grants without guessing duplicate names", () => { const store = new Store(selection); const requester = store.createBot(); diff --git a/server/store.ts b/server/store.ts index ca2d0f8e95..045e156b76 100644 --- a/server/store.ts +++ b/server/store.ts @@ -10,7 +10,7 @@ import { peerAllowKey, type PeerAction } from "./peer-approval-key.ts"; import { DATA_DIR } from "./config.ts"; import * as mdb from "./message-db.ts"; import { workspaceDir } from "./workspace.ts"; -import { newId, type ModelSelection, type ThreadId } from "./contracts.ts"; +import { newId, type CloudBackend, type ModelSelection, type ThreadId } from "./contracts.ts"; import { pickBotName } from "./names.ts"; import { redactSecretsInText } from "./redact.ts"; @@ -246,6 +246,8 @@ export interface BotRecord { /** which computer the bot acts on: its cloud box, this Mac (local CUA), * or none. Unset = auto (box when it exists, else local when available). */ computer?: "cloud" | "vm" | "local" | "off"; + /** Which cloud computer backs `computer: "cloud"`; absent means Box. */ + cloudBackend?: CloudBackend; /** where NEW tasks run their shell tools; each task pins its own copy * on its first turn (TaskRecord.cwd). Absent = the home folder. */ cwd?: string; @@ -424,6 +426,10 @@ export class Store { if (b.busy || (b.activity !== undefined && b.activity !== "idle")) botsMigrated = true; b.busy = false; b.activity = "idle"; + if (b.cloudBackend !== undefined && b.cloudBackend !== "box" && b.cloudBackend !== "vps") { + delete b.cloudBackend; + botsMigrated = true; + } } for (const b of this.bots) { if (!b.chiefOfStaff) continue; diff --git a/server/testing/setup.ts b/server/testing/setup.ts index e8cc6fe122..197d31d494 100644 --- a/server/testing/setup.ts +++ b/server/testing/setup.ts @@ -12,6 +12,11 @@ import { removeTempDir } from "./cleanup.ts"; const home = mkdtempSync(join(tmpdir(), "omb-test-home-")); process.env.HOME = home; process.env.USERPROFILE = home; +// OMB_DATA_DIR is an intentional production override, but tests must never +// let it escape the throwaway home they are about to delete. +delete process.env.OMB_DATA_DIR; +// Do not let a developer's Hermes global config path leak into per-test homes. +delete process.env.HERMES_HOME; // The companion keeps its paired devices in its own directory, and resolves // it from homedir() the same way — so the redirect above already covers it. // Named explicitly all the same: the device tests delete this directory diff --git a/server/vps-computer.runner.test.ts b/server/vps-computer.runner.test.ts new file mode 100644 index 0000000000..0dc0c6ed03 --- /dev/null +++ b/server/vps-computer.runner.test.ts @@ -0,0 +1,76 @@ +import { EventEmitter } from "node:events"; +import { PassThrough, Writable } from "node:stream"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const { spawnMock } = vi.hoisted(() => ({ spawnMock: vi.fn() })); + +vi.mock("node:child_process", async () => ({ + ...(await vi.importActual("node:child_process")), + spawn: spawnMock, +})); + +import { defaultRunner } from "./vps-computer.ts"; + +type FakeChild = EventEmitter & { + stdin: Writable; + stdout: PassThrough; + stderr: PassThrough; + kill: ReturnType; +}; + +function fakeChild(): FakeChild { + const child = new EventEmitter() as FakeChild; + child.stdin = new Writable({ write: (_chunk, _encoding, callback) => callback() }); + child.stdout = new PassThrough(); + child.stderr = new PassThrough(); + child.kill = vi.fn(() => true); + spawnMock.mockReturnValue(child); + return child; +} + +describe("default VPS command runner", () => { + beforeEach(() => { + spawnMock.mockReset(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it("collects output and resolves after the child closes", async () => { + const child = fakeChild(); + const result = defaultRunner(["info"], { input: "request" }); + + child.stdout.write("out"); + child.stderr.write("err"); + child.emit("close", 0, null); + + await expect(result).resolves.toEqual({ stdout: "out", stderr: "err" }); + expect(spawnMock).toHaveBeenCalledWith("docker", ["info"], expect.objectContaining({ shell: false })); + }); + + it("turns stdin EPIPE into a rejected command instead of an unhandled error", async () => { + const child = fakeChild(); + const result = defaultRunner(["build", "-"], { input: "Dockerfile" }); + + child.stdin.emit("error", new Error("write EPIPE")); + + await expect(result).rejects.toThrow("Docker-over-SSH stdin failed: write EPIPE"); + child.emit("close", 1, null); + }); + + it("escalates a timed-out command from SIGTERM to SIGKILL", async () => { + vi.useFakeTimers(); + const child = fakeChild(); + const result = defaultRunner(["info"], { timeoutMs: 100 }); + const rejection = expect(result).rejects.toThrow("Docker-over-SSH command timed out"); + + await vi.advanceTimersByTimeAsync(100); + expect(child.kill).toHaveBeenNthCalledWith(1, "SIGTERM"); + + // the WAN-sized grace window: ssh + docker get 5s to tear down cleanly + await vi.advanceTimersByTimeAsync(5_000); + await rejection; + expect(child.kill).toHaveBeenNthCalledWith(2, "SIGKILL"); + }); +}); diff --git a/server/vps-computer.test.ts b/server/vps-computer.test.ts new file mode 100644 index 0000000000..d2f0372d0b --- /dev/null +++ b/server/vps-computer.test.ts @@ -0,0 +1,590 @@ +import { describe, expect, it, vi } from "vitest"; + +import { + BASE_IMAGE, + BASE_IMAGE_DIGEST, + BASE_IMAGE_LABEL, + CUA_DRIVER_VERSION, + DRIVER_LABEL, + DISPLAY, + IMAGE_LAYER_LABEL, + IMAGE_LAYER_VERSION, + MANAGED_LABEL, +} from "./container-computer.ts"; +import type { AppConfig } from "./config.ts"; +import { + VPS_CONTAINER_LABEL, + VPS_IMAGE, + VPS_MANAGED_LABEL, + vpsComputerAction, + vpsComputerScreenshot, + vpsComputerStatus, + vpsComputerMcp, + vpsContainerMcpArgs, + vpsContainerName, + vpsContainerRunArgs, + vpsDockerArgs, + vpsDriverError, + reuseVps, + type VpsCommandRunner, +} from "./vps-computer.ts"; + +const BOT_ID = "bot-1234-abcd"; +const CONFIG: AppConfig = { vps: { sshAlias: "production-vps" } }; +const IMAGE_ID = `sha256:${"a".repeat(64)}`; +const CONTAINER_ID = "b".repeat(64); +const screenshot = Buffer.concat([ + Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]), + Buffer.alloc(600), + Buffer.from("IEND", "ascii"), +]); + +function fixture({ + image = true, + container = true, + running = true, + managed = true, + mounts = false, + publicPorts = false, + publishAllPorts = false, + deviceRequests = false, + networkMode = "default", + containerImageId = IMAGE_ID, + inspectedImageId = IMAGE_ID, + rebuiltImageId, + containerId = CONTAINER_ID, + privileged = false, + pidMode = "", + ipcMode, + capAdd = ["CAP_SETUID", "CAP_SETGID"], + screenshotValid = true, + screenshotCaptureFails = false, + desktopProbeFails = false, + securityOpt = [], + memory = 4 * 1024 * 1024 * 1024, + restartPolicyName = "unless-stopped", + cgroupnsMode, + imageLabelsMatch = true, +}: { + image?: boolean; + container?: boolean; + running?: boolean; + managed?: boolean; + mounts?: boolean; + publicPorts?: boolean; + publishAllPorts?: boolean; + deviceRequests?: boolean; + networkMode?: string; + containerImageId?: string; + inspectedImageId?: string; + rebuiltImageId?: string; + containerId?: string; + privileged?: boolean; + pidMode?: string; + ipcMode?: string; + capAdd?: string[]; + screenshotValid?: boolean; + screenshotCaptureFails?: boolean; + desktopProbeFails?: boolean; + securityOpt?: string[]; + memory?: number; + restartPolicyName?: string; + cgroupnsMode?: string; + imageLabelsMatch?: boolean; +} = {}) { + const name = vpsContainerName(BOT_ID); + const provisioningArgs = vpsContainerRunArgs(name); + const argValue = (flag: string) => { + const index = provisioningArgs.indexOf(flag); + if (index >= 0) return provisioningArgs[index + 1] ?? ""; + return provisioningArgs.find((arg) => arg.startsWith(`${flag}=`))?.slice(flag.length + 1) ?? ""; + }; + const calls: Array<{ args: string[]; options?: { input?: string; timeoutMs?: number } }> = []; + const state = { image, container, running, imageLabelsMatch, inspectedImageId, containerImageId }; + const runner: VpsCommandRunner = async (args, options) => { + calls.push({ args, options }); + const command = args[2]; + if (command === "image") { + // The real daemon phrases a clean absence this way; anything else is + // read as a transport failure, exactly like production. + if (!state.image) throw new Error(`Error: No such image: ${VPS_IMAGE}`); + return { + stdout: JSON.stringify([{ + Config: { Labels: state.imageLabelsMatch ? { + [MANAGED_LABEL]: "1", + [DRIVER_LABEL]: CUA_DRIVER_VERSION, + [BASE_IMAGE_LABEL]: BASE_IMAGE_DIGEST, + [IMAGE_LAYER_LABEL]: IMAGE_LAYER_VERSION, + } : { [MANAGED_LABEL]: "0" } }, + Id: state.inspectedImageId, + }]), + stderr: "", + }; + } + if (command === "inspect") { + if (!state.container) throw new Error(`Error: No such object: ${name}`); + return { + stdout: JSON.stringify([{ + Config: { + Image: state.image ? VPS_IMAGE : "old-image", + Labels: { + [VPS_MANAGED_LABEL]: managed ? "1" : "0", + [VPS_CONTAINER_LABEL]: managed ? name : "other-container", + [MANAGED_LABEL]: "1", + [DRIVER_LABEL]: CUA_DRIVER_VERSION, + [BASE_IMAGE_LABEL]: BASE_IMAGE_DIGEST, + [IMAGE_LAYER_LABEL]: IMAGE_LAYER_VERSION, + }, + }, + Id: containerId, + Image: state.image ? state.containerImageId : "old-image-id", + HostConfig: { + Binds: mounts ? ["/host:/container"] : [], + VolumesFrom: [], + NetworkMode: networkMode, + PortBindings: publicPorts ? { "6901/tcp": [{ HostIp: "0.0.0.0" }] } : {}, + PublishAllPorts: publishAllPorts, + Memory: memory, + MemorySwap: 4 * 1024 * 1024 * 1024, + NanoCpus: 2_000_000_000, + PidsLimit: 512, + CapDrop: ["ALL"], + CapAdd: capAdd, + Privileged: privileged, + PidMode: pidMode, + IpcMode: ipcMode ?? argValue("--ipc"), + UTSMode: "", + ShmSize: 512 * 1024 * 1024, + Devices: [], + DeviceRequests: deviceRequests ? [{ Driver: "nvidia" }] : [], + SecurityOpt: securityOpt, + UsernsMode: "", + CgroupnsMode: cgroupnsMode ?? argValue("--cgroupns"), + OomKillDisable: false, + AutoRemove: false, + RestartPolicy: { Name: restartPolicyName, MaximumRetryCount: 0 }, + }, + NetworkSettings: { + Networks: { [networkMode === "default" ? "bridge" : networkMode]: {} }, + }, + Mounts: mounts ? [{ Source: "/host", Destination: "/container" }] : [], + State: { Running: state.running }, + }]), + stderr: "", + }; + } + if (command === "exec") { + // only the pixel-carrying screenshot call fails; the status path's + // plain get_desktop_state readiness probe keeps answering + if (screenshotCaptureFails && args.includes("--screenshot-out-file")) throw new Error("capture failed"); + if (args.includes("base64")) return { stdout: screenshotValid ? screenshot.toString("base64") : "not-an-image", stderr: "" }; + if (args.includes("tail")) { + return { stdout: "X display :1 did not become ready within 45 seconds\n", stderr: "" }; + } + if (args.at(-1) === "--version") { + if (desktopProbeFails) throw new Error("driver unavailable"); + return { stdout: `cua-driver ${CUA_DRIVER_VERSION}\n`, stderr: "" }; + } + if (args.includes("status")) return { stdout: "running\n", stderr: "" }; + if (args.includes("health_report")) { + return { stdout: JSON.stringify({ schema_version: "1", overall: "ok", checks: [] }), stderr: "" }; + } + if (args.includes("get_desktop_state")) return { stdout: "{}\n", stderr: "" }; + return { stdout: "{}\n", stderr: "" }; + } + if (command === "pull") return { stdout: "pulled\n", stderr: "" }; + if (command === "build") { + state.image = true; + state.imageLabelsMatch = true; + state.inspectedImageId = rebuiltImageId ?? state.inspectedImageId; + expect(options?.input).toContain(`FROM ${BASE_IMAGE}`); + return { stdout: "built\n", stderr: "" }; + } + if (command === "run") { + state.container = true; + state.running = true; + // a fresh container is created FROM the image ref in the run argv + state.containerImageId = args.at(-1) ?? state.containerImageId; + return { stdout: `${name}\n`, stderr: "" }; + } + if (command === "start") { + state.running = true; + return { stdout: `${name}\n`, stderr: "" }; + } + if (command === "stop") { + state.running = false; + return { stdout: `${name}\n`, stderr: "" }; + } + if (command === "rm") { + state.container = false; + state.running = false; + return { stdout: `${name}\n`, stderr: "" }; + } + throw new Error(`unexpected Docker command ${command}`); + }; + return { calls, runner, state, name }; +} + +describe("VPS computer", () => { + it("uses a deterministic, bot-id-derived managed container name", () => { + expect(vpsContainerName(BOT_ID)).toBe(vpsContainerName(BOT_ID)); + expect(vpsContainerName(BOT_ID)).not.toBe(vpsContainerName("another-bot")); + expect(vpsContainerName(BOT_ID)).toMatch(/^openmausbot-vps-[a-z0-9-]+$/); + }); + + it("passes the SSH target as one validated Docker argv value", () => { + expect(vpsDockerArgs("production-vps", ["info"])).toEqual(["-H", "ssh://production-vps", "info"]); + for (const alias of ["production-vps;touch", "ssh://production-vps", "-H", "--host=evil", "prod vps", "prod\n-v"] ) { + expect(() => vpsDockerArgs(alias, ["info"])).toThrow(/alias/); + } + expect(() => vpsContainerMcpArgs("production-vps", "not a container")).toThrow(/connection/); + }); + + it("reports a ready container only when image, labels, limits, mounts, network, and Cua pass", async () => { + const fake = fixture(); + const status = await vpsComputerStatus(CONFIG, BOT_ID, fake.runner); + expect(status).toMatchObject({ + configured: true, + daemonUp: true, + image: true, + imageMatches: true, + managed: true, + network: "private", + mounts: "none", + security: "hardened", + desktopReady: true, + ready: true, + problem: null, + }); + // No standalone `docker info` round-trip: the image inspect doubles as + // the daemon probe, so a healthy status costs 2 docker calls + 4 execs. + expect(fake.calls[0]?.args).toEqual(["-H", "ssh://production-vps", "image", "inspect", VPS_IMAGE]); + expect(fake.calls.some(({ args }) => args[2] === "info")).toBe(false); + const probes = fake.calls.filter( + ({ args }) => + args[2] === "exec" && + (args.at(-1) === "--version" || args.includes("status") || args.includes("health_report") || args.includes("get_desktop_state")), + ); + expect(probes).toHaveLength(4); + expect(probes.every(({ args }) => args.includes(CONTAINER_ID))).toBe(true); + expect(probes.every(({ args }) => !args.includes(fake.name))).toBe(true); + expect(probes.every(({ args }) => args.includes(`DISPLAY=${DISPLAY}`))).toBe(true); + expect(probes.every(({ args }) => args.includes("CUA_DRIVER_RS_TELEMETRY_ENABLED=0"))).toBe(true); + // The status poll must never transfer pixels: readiness is the driver + // answering get_desktop_state, and pixel validation belongs to the + // screenshot path alone. + expect(fake.calls.some(({ args }) => args.includes("base64"))).toBe(false); + expect(fake.calls.some(({ args }) => args.includes("--screenshot-out-file"))).toBe(false); + }); + + it("refuses host mounts, public ports, and unowned containers", async () => { + const mounted = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ mounts: true }).runner); + expect(mounted.ready).toBe(false); + expect(mounted.mounts).toBe("unsafe"); + + const publicPorts = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ publicPorts: true }).runner); + expect(publicPorts.ready).toBe(false); + expect(publicPorts.network).toBe("unsafe"); + + const publishedAll = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ publishAllPorts: true }).runner); + expect(publishedAll.ready).toBe(false); + expect(publishedAll.network).toBe("unsafe"); + + const devices = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ deviceRequests: true }).runner); + expect(devices.ready).toBe(false); + expect(devices.security).toBe("unsafe"); + + const sharedNetwork = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ networkMode: "shared-net" }).runner); + expect(sharedNetwork.ready).toBe(false); + expect(sharedNetwork.network).toBe("unsafe"); + + const hostNetwork = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ networkMode: "host" }).runner); + expect(hostNetwork.ready).toBe(false); + expect(hostNetwork.network).toBe("unsafe"); + + const privileged = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ privileged: true }).runner); + expect(privileged.ready).toBe(false); + expect(privileged.security).toBe("unsafe"); + + const hostNamespaces = await vpsComputerStatus( + CONFIG, + BOT_ID, + fixture({ pidMode: "host", ipcMode: "host" }).runner, + ); + expect(hostNamespaces.ready).toBe(false); + expect(hostNamespaces.security).toBe("unsafe"); + + const extraCapability = await vpsComputerStatus( + CONFIG, + BOT_ID, + fixture({ capAdd: ["CAP_SETUID", "CAP_SETGID", "CAP_SYS_ADMIN"] }).runner, + ); + expect(extraCapability.ready).toBe(false); + expect(extraCapability.security).toBe("unsafe"); + + const unsafeProfile = await vpsComputerStatus( + CONFIG, + BOT_ID, + fixture({ securityOpt: ["seccomp=unconfined"], memory: 1024, restartPolicyName: "always", cgroupnsMode: "host" }).runner, + ); + expect(unsafeProfile.ready).toBe(false); + expect(unsafeProfile.security).toBe("unsafe"); + + // the VPS container must survive an unwatched reboot: exactly + // unless-stopped, so a policy-less container is flagged for recreation + const noRestart = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ restartPolicyName: "no" }).runner); + expect(noRestart.ready).toBe(false); + expect(noRestart.security).toBe("unsafe"); + + const wrongImage = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ containerImageId: "c".repeat(64) }).runner); + expect(wrongImage.ready).toBe(false); + expect(wrongImage.imageMatches).toBe(false); + + const malformedImageId = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ inspectedImageId: "--help" }).runner); + expect(malformedImageId.ready).toBe(false); + expect(malformedImageId.image).toBe(false); + + const malformedContainerId = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ containerId: "--help" }).runner); + expect(malformedContainerId.ready).toBe(false); + expect(malformedContainerId.container_id).toBeNull(); + + const unowned = await vpsComputerStatus(CONFIG, BOT_ID, fixture({ managed: false }).runner); + expect(unowned.ready).toBe(false); + expect(unowned.managed).toBe(false); + }); + + it("lets explicit provisioning build and run the pinned container, but Auto only reuses", async () => { + const auto = fixture({ image: false, container: false }); + expect(await reuseVps(CONFIG, BOT_ID, auto.runner)).toBeNull(); + expect(auto.calls.some(({ args }) => ["run", "start", "build", "pull"].includes(args[2]!))).toBe(false); + + const provision = fixture({ image: false, container: false }); + const status = await vpsComputerAction("provision", CONFIG, BOT_ID, provision.runner); + expect(status.ready).toBe(true); + const run = provision.calls.find(({ args }) => args[2] === "run")?.args ?? []; + expect(run).toContain("--memory"); + expect(run).toContain("--pids-limit"); + expect(run).toContain("--ipc"); + expect(run[run.indexOf("--ipc") + 1]).toBe("private"); + expect(run).toContain("--cgroupns"); + expect(run[run.indexOf("--cgroupns") + 1]).toBe("private"); + expect(run.at(-1)).toBe(IMAGE_ID); + expect(run.join(" ")).toContain(`--label ${VPS_MANAGED_LABEL}=1`); + expect(run.join(" ")).toContain(`--label ${IMAGE_LAYER_LABEL}=${IMAGE_LAYER_VERSION}`); + expect(run.join(" ")).toContain("--restart unless-stopped"); + expect(run).not.toContain("--mount"); + expect(run).not.toContain("-p"); + expect(provision.calls.some(({ args }) => args[2] === "build")).toBe(true); + }); + + it("uses the image id produced by a rebuild", async () => { + const staleImageId = `sha256:${"b".repeat(64)}`; + const provision = fixture({ + image: true, + imageLabelsMatch: false, + container: false, + inspectedImageId: staleImageId, + rebuiltImageId: IMAGE_ID, + }); + + const status = await vpsComputerAction("provision", CONFIG, BOT_ID, provision.runner); + expect(status.ready).toBe(true); + const run = provision.calls.find(({ args }) => args[2] === "run")?.args ?? []; + expect(run.at(-1)).toBe(IMAGE_ID); + expect(run.at(-1)).not.toBe(staleImageId); + }); + + it("starts and sleeps only the managed container, never the VPS", async () => { + const start = fixture({ running: false }); + const started = await vpsComputerAction("start", CONFIG, BOT_ID, start.runner); + expect(started.container).toBe("running"); + expect(start.calls.some(({ args }) => args[2] === "start")).toBe(true); + expect(start.calls.some(({ args }) => ["rm", "system", "reboot", "shutdown"].includes(args[2]!))).toBe(false); + + const stop = fixture(); + const stopped = await vpsComputerAction("stop", CONFIG, BOT_ID, stop.runner); + expect(stopped.container).toBe("stopped"); + expect(stop.calls.some(({ args }) => args[2] === "stop" && args[3] === CONTAINER_ID)).toBe(true); + expect(stop.calls.some(({ args }) => ["rm", "system", "reboot", "shutdown"].includes(args[2]!))).toBe(false); + }); + + it("serializes concurrent provisioning for the same bot", async () => { + const fake = fixture({ image: false, container: false }); + const results = await Promise.all([ + vpsComputerAction("provision", CONFIG, BOT_ID, fake.runner), + vpsComputerAction("provision", CONFIG, BOT_ID, fake.runner), + ]); + expect(results.every((status) => status.ready)).toBe(true); + expect(fake.calls.filter(({ args }) => args[2] === "run")).toHaveLength(1); + }); + + it("mounts the official Cua MCP server through the tiny remote exec bridge", () => { + const connection = vpsComputerMcp(CONFIG, BOT_ID); + expect(connection.command).toBe(process.execPath); + expect(connection.args.slice(-2)).toEqual(["production-vps", vpsContainerName(BOT_ID)]); + expect(connection.env).toEqual({ ELECTRON_RUN_AS_NODE: "1" }); + expect(vpsComputerMcp(CONFIG, BOT_ID, CONTAINER_ID).args.slice(-2)).toEqual(["production-vps", CONTAINER_ID]); + expect(vpsContainerMcpArgs("production-vps", vpsContainerName(BOT_ID))).toEqual([ + "-H", + "ssh://production-vps", + "exec", + "-i", + "-u", + "cua", + "-e", + "HOME=/home/cua", + "-e", + "DISPLAY=:1", + "-e", + "CUA_DRIVER_INSTALL_CHANNEL=python_package", + "-e", + "CUA_DRIVER_RS_TELEMETRY_ENABLED=0", + vpsContainerName(BOT_ID), + "/usr/local/libexec/openmausbot/cua-driver", + "mcp", + "--socket", + "/run/user/1000/openmausbot-cua.sock", + ]); + }); + + it("captures screenshots through Cua Driver and validates the returned image", async () => { + const fake = fixture(); + const frame = await vpsComputerScreenshot(CONFIG, BOT_ID, fake.runner); + expect(frame).toEqual({ png: screenshot.toString("base64"), format: "png" }); + expect(fake.calls.some(({ args }) => args.includes("get_desktop_state"))).toBe(true); + expect(fake.calls.some(({ args }) => args.includes("base64") && args.includes("-u") && args.includes("cua"))).toBe(true); + expect(fake.calls.some(({ args }) => args.includes("rm") && args.includes("-f"))).toBe(true); + + await expect(vpsComputerScreenshot(CONFIG, BOT_ID, fixture({ screenshotValid: false }).runner)).rejects.toThrow(/incomplete/); + + const failedCapture = fixture({ screenshotCaptureFails: true }); + await expect(vpsComputerScreenshot(CONFIG, BOT_ID, failedCapture.runner)).rejects.toThrow(/capture failed/); + expect(failedCapture.calls.some(({ args }) => args.includes("rm") && args.includes("-f"))).toBe(true); + }); + + it("fails clearly for BoxAgent and engines without computer MCP", () => { + expect(vpsDriverError("boxAgent", true)).toMatch(/cannot use a self-hosted VPS/); + expect(vpsDriverError("codex", false)).toMatch(/cannot mount/); + expect(vpsDriverError("claudeAgent", true)).toBeNull(); + }); + + it("fails cleanly when no VPS alias is configured", async () => { + await expect(vpsComputerAction("provision", {}, BOT_ID, fixture().runner)).rejects.toThrow(/not configured/); + }); + + it("attributes a transport failure to the link, never to a missing container", async () => { + const fake = fixture(); + const flaky: VpsCommandRunner = async (args, options) => { + if (args[2] === "inspect") throw new Error("ssh: connect to host production-vps port 22: Connection timed out"); + return fake.runner(args, options); + }; + const status = await vpsComputerStatus(CONFIG, BOT_ID, flaky); + expect(status.daemonUp).toBe(false); + expect(status.ready).toBe(false); + expect(status.problem).toMatch(/Docker over SSH failed while checking the VPS/); + + // and provision must refuse to `docker run --name ` into the fog + await expect(vpsComputerAction("provision", CONFIG, BOT_ID, flaky)).rejects.toThrow(/Docker over SSH failed/); + expect(fake.calls.some(({ args }) => args[2] === "run")).toBe(false); + + const imageFlaky: VpsCommandRunner = async (args, options) => { + if (args[2] === "image") throw new Error("kex_exchange_identification: read: Connection reset by peer"); + return fake.runner(args, options); + }; + const imageStatus = await vpsComputerStatus(CONFIG, BOT_ID, imageFlaky); + expect(imageStatus.daemonUp).toBe(false); + expect(imageStatus.problem).toMatch(/Docker over SSH failed while checking the VPS/); + }); + + it("removes a managed container even when its image is incompatible, then provisions fresh", async () => { + // an IMAGE_LAYER_VERSION bump leaves a running container that provision + // refuses to touch — remove is the in-app escape hatch + const stale = fixture({ containerImageId: `sha256:${"c".repeat(64)}` }); + expect((await vpsComputerStatus(CONFIG, BOT_ID, stale.runner)).imageMatches).toBe(false); + await expect(vpsComputerAction("provision", CONFIG, BOT_ID, stale.runner)).rejects.toThrow(/incompatible|unsafe/); + + const removed = await vpsComputerAction("remove", CONFIG, BOT_ID, stale.runner); + expect(removed.container).toBe("missing"); + expect(stale.calls.some(({ args }) => args[2] === "rm" && args[3] === "-f" && args[4] === CONTAINER_ID)).toBe(true); + + const rebuilt = await vpsComputerAction("provision", CONFIG, BOT_ID, stale.runner); + expect(rebuilt.ready).toBe(true); + }); + + it("never removes a container OpenMausBot did not create", async () => { + const unowned = fixture({ managed: false }); + await expect(vpsComputerAction("remove", CONFIG, BOT_ID, unowned.runner)).rejects.toThrow(/did not create/); + expect(unowned.calls.some(({ args }) => args[2] === "rm")).toBe(false); + + const absent = fixture({ container: false }); + const afterMissing = await vpsComputerAction("remove", CONFIG, BOT_ID, absent.runner); + expect(afterMissing.container).toBe("missing"); + expect(absent.calls.some(({ args }) => args[2] === "rm")).toBe(false); + }); + + it("surfaces the supervisor error log when the desktop probe fails", async () => { + const fake = fixture({ desktopProbeFails: true }); + const status = await vpsComputerStatus(CONFIG, BOT_ID, fake.runner); + expect(status.desktopReady).toBe(false); + expect(status.desktop_error).toContain("did not become ready"); + expect(status.problem).toContain("desktop failed to start"); + expect(fake.calls.some(({ args }) => args[2] === "exec" && args.includes("tail"))).toBe(true); + }); + + it("waits for readiness with a cheap driver probe and backoff, not full re-inspections", async () => { + vi.useFakeTimers(); + try { + const fake = fixture({ container: false }); + let driverProbes = 0; + const runner: VpsCommandRunner = async (args, options) => { + if (args[2] === "exec" && args.includes("status")) { + driverProbes += 1; + if (driverProbes < 4) throw new Error("driver not up yet"); + } + return fake.runner(args, options); + }; + const pending = vpsComputerAction("provision", CONFIG, BOT_ID, runner); + await vi.advanceTimersByTimeAsync(10_000); + const status = await pending; + expect(status.ready).toBe(true); + // the expensive end of the pipeline ran exactly once — every retry in + // between was the single `cua-driver status` predicate + const desktopCalls = fake.calls.filter(({ args }) => args.includes("get_desktop_state")); + expect(desktopCalls).toHaveLength(1); + const healthCalls = fake.calls.filter(({ args }) => args.includes("health_report")); + expect(healthCalls).toHaveLength(1); + expect(driverProbes).toBeGreaterThanOrEqual(4); + } finally { + vi.useRealTimers(); + } + }); + + it("fails lifecycle calls fast instead of queueing behind a long provision", async () => { + vi.useFakeTimers(); + try { + let releaseBuild!: () => void; + const buildGate = new Promise((resolve) => { + releaseBuild = resolve; + }); + const fake = fixture({ image: false, container: false }); + const slowRunner: VpsCommandRunner = async (args, options) => { + if (args[2] === "build") await buildGate; + return fake.runner(args, options); + }; + const first = vpsComputerAction("provision", CONFIG, BOT_ID, slowRunner); + // let the first action reach its (gated) docker build + await vi.advanceTimersByTimeAsync(0); + const second = vpsComputerAction("stop", CONFIG, BOT_ID, slowRunner); + const rejection = expect(second).rejects.toThrow(/being prepared/); + await vi.advanceTimersByTimeAsync(5_000); + await rejection; + + releaseBuild(); + await vi.advanceTimersByTimeAsync(10_000); + const status = await first; + expect(status.ready).toBe(true); + } finally { + vi.useRealTimers(); + } + }); +}); diff --git a/server/vps-computer.ts b/server/vps-computer.ts new file mode 100644 index 0000000000..6a141d2c4f --- /dev/null +++ b/server/vps-computer.ts @@ -0,0 +1,808 @@ +// BYO Linux VPS computer. The agent process stays local; Docker's own SSH +// transport reaches the user's daemon and the official Cua MCP server stays +// inside one managed container per bot. +import { createHash } from "node:crypto"; +import { spawn } from "node:child_process"; + +import { + BASE_IMAGE, + CUA_DRIVER_VERSION, + CUA_SOCKET, + IMAGE as CUA_IMAGE, + cuaExecArgs, + dockerSecurityIsHardened, + imageLabelsMatch, + managedImageDockerfile, + wholeScreenshot, + type DockerHardeningConfig, + BASE_IMAGE_DIGEST, + BASE_IMAGE_LABEL, + DRIVER_LABEL, + IMAGE_LAYER_LABEL, + IMAGE_LAYER_VERSION, + MANAGED_LABEL, +} from "./container-computer.ts"; +import { isValidSshAlias, vpsSshAlias, type AppConfig } from "./config.ts"; +import { augmentedPath } from "./env-path.ts"; +import { SPAWNED_PROXIES } from "./proxy-paths.ts"; + +export const VPS_IMAGE = CUA_IMAGE; +export const VPS_MANAGED_LABEL = "com.openmausbot.vps"; +export const VPS_CONTAINER_LABEL = "com.openmausbot.container"; +export const VPS_CONTAINER_PREFIX = "openmausbot-vps"; +// SIGTERM must give ssh + docker time to tear down the remote exec before the +// SIGKILL escalation; 1s was routinely too short over a WAN round-trip, and an +// orphaned remote exec keeps the driver socket busy for the next command. +const COMMAND_TIMEOUT_KILL_GRACE_MS = 5_000; + +const CONTAINER_NAME = /^[a-zA-Z0-9][a-zA-Z0-9_.-]+$/; +const CONTAINER_ID = /^[a-f0-9]{12,64}$/i; +const IMAGE_ID = /^sha256:[a-f0-9]{64}$/i; +const PIDS_LIMIT = 512; +const SCREENSHOT_PATH = "/tmp/openmausbot-vps-preview.png"; +const lifecycleLocks = new Map>(); +// A held lock means a lifecycle mutation (worst case: a 10-minute image +// build) is running. Waiting it out would wedge Sleep and the screenshot +// poll behind it, so acquisition fails fast instead. +const LOCK_ACQUIRE_TIMEOUT_MS = 5_000; +// The panel polls status every 4-6s and the screen poller re-checks it before +// every frame; each full status is several docker-over-SSH processes. Same +// pattern as container-computer's screenshotStatusCache, and the same TTL. +const STATUS_CACHE_TTL_MS = 10_000; +const statusCache = new Map(); + +export interface VpsCommandOptions { + input?: string; + timeoutMs?: number; +} + +export type VpsCommandRunner = ( + args: string[], + options?: VpsCommandOptions, +) => Promise<{ stdout: string; stderr: string }>; + +export type VpsLifecycleAction = "provision" | "start" | "stop" | "remove"; + +export interface VpsComputerStatus { + configured: boolean; + sshAlias: string | null; + daemonUp: boolean; + image: boolean; + imageMatches: boolean; + managed: boolean; + container: "running" | "stopped" | "missing"; + network: "private" | "unsafe" | "unknown"; + mounts: "none" | "unsafe" | "unknown"; + security: "hardened" | "unsafe" | "unknown"; + desktopReady: boolean; + desktop_error: string | null; + ready: boolean; + problem: string | null; + image_ref: string; + base_image_ref: string; + driver_version: string; + container_name: string; + container_id: string | null; + image_id: string | null; +} + +function containerNamePart(botId: string): string { + return botId.toLowerCase().replace(/[^a-z0-9]/g, "").slice(0, 12) || "bot"; +} + +/** Stable across restarts and independent of the bot's editable display name. */ +export function vpsContainerName(botId: string): string { + const hash = createHash("sha256").update(botId).digest("hex").slice(0, 12); + return `${VPS_CONTAINER_PREFIX}-${containerNamePart(botId)}-${hash}`; +} + +export function vpsDockerArgs(alias: string, args: string[]): string[] { + if (!isValidSshAlias(alias)) { + throw new Error("invalid VPS SSH config alias"); + } + return ["-H", `ssh://${alias}`, ...args]; +} + +const STREAM_CAP_CHARS = 16 * 1024 * 1024; + +/** Keeps the LAST 16MB of a stream without rebuilding one giant string per + * chunk (a 10-minute `docker build` stream made that rebuild quadratic). + * Chunks fall off the front as soon as the tail alone covers the cap; the + * final slice preserves the exact cap semantics of the old accumulator. */ +function tailCollector() { + const chunks: string[] = []; + let total = 0; + return { + push(chunk: string) { + chunks.push(chunk); + total += chunk.length; + for (;;) { + const first = chunks[0]; + if (chunks.length < 2 || first === undefined || total - first.length < STREAM_CAP_CHARS) break; + chunks.shift(); + total -= first.length; + } + }, + text(): string { + return chunks.join("").slice(-STREAM_CAP_CHARS); + }, + }; +} + +export function defaultRunner(args: string[], options: VpsCommandOptions = {}): Promise<{ stdout: string; stderr: string }> { + return new Promise((resolve, reject) => { + const child = spawn("docker", args, { + shell: false, + env: { ...process.env, PATH: augmentedPath() }, + stdio: ["pipe", "pipe", "pipe"], + }); + const stdout = tailCollector(); + const stderr = tailCollector(); + let settled = false; + let timedOut = false; + let killTimer: ReturnType | undefined; + let timeout: ReturnType; + const settle = (finish: () => void) => { + if (settled) return; + settled = true; + clearTimeout(timeout); + if (killTimer) clearTimeout(killTimer); + finish(); + }; + timeout = setTimeout(() => { + if (settled) return; + timedOut = true; + killTimer = setTimeout(() => { + if (settled) return; + child.kill("SIGKILL"); + settle(() => reject(new Error("Docker-over-SSH command timed out"))); + }, COMMAND_TIMEOUT_KILL_GRACE_MS); + killTimer.unref?.(); + child.kill("SIGTERM"); + }, options.timeoutMs ?? 120_000); + timeout.unref?.(); + + child.stdout.setEncoding("utf8"); + child.stderr.setEncoding("utf8"); + child.stdout.on("data", (chunk: string) => { + stdout.push(chunk); + }); + child.stderr.on("data", (chunk: string) => { + stderr.push(chunk); + }); + child.stdin.on("error", (error) => { + if (timedOut) return; + settle(() => reject(new Error(`Docker-over-SSH stdin failed: ${error.message}`))); + }); + child.on("error", (error) => { + settle(() => reject(new Error(`Docker-over-SSH could not start: ${error.message}`))); + }); + child.on("close", (code, signal) => { + if (timedOut) { + settle(() => reject(new Error("Docker-over-SSH command timed out"))); + return; + } + settle(() => { + if (code === 0) return resolve({ stdout: stdout.text(), stderr: stderr.text() }); + const detail = stderr.text().trim().slice(-1000); + reject(new Error(detail || `Docker-over-SSH exited ${code ?? signal ?? "without a status"}`)); + }); + }); + try { + child.stdin.end(options.input); + } catch (error) { + settle(() => reject(new Error(`Docker-over-SSH stdin failed: ${error instanceof Error ? error.message : String(error)}`))); + } + }); +} + +function emptyStatus(botId: string, alias: string | null): VpsComputerStatus { + return { + configured: Boolean(alias), + sshAlias: alias, + daemonUp: false, + image: false, + imageMatches: false, + managed: false, + container: "missing", + network: "unknown", + mounts: "unknown", + security: "unknown", + desktopReady: false, + desktop_error: null, + ready: false, + problem: alias ? "Docker over SSH is not reachable" : "Configure a VPS SSH alias in App Settings → Connections", + image_ref: VPS_IMAGE, + base_image_ref: BASE_IMAGE, + driver_version: CUA_DRIVER_VERSION, + container_name: vpsContainerName(botId), + container_id: null, + image_id: null, + }; +} + +/** Docker and Podman both phrase a clean not-found this way ("No such + * object" / "No such image" / "no such container"); anything else out of an + * inspect is a transport or daemon failure and must NOT be read as absence — + * a flaky WAN link that looked like "missing" used to send provision into + * `docker run --name ` and a baffling name-in-use error. */ +function isMissingObjectMessage(message: string): boolean { + return /no such (object|image|container)/i.test(message); +} + +function transportFailure(message: string): string { + return `Docker over SSH failed while checking the VPS: ${message.trim().slice(0, 200) || "unknown transport error"}`; +} + +function hasNoHostMounts(detail: { + Mounts?: unknown; + HostConfig?: { Binds?: string[] | null; VolumesFrom?: string[] | null }; +}): boolean { + return ( + Array.isArray(detail.Mounts) && + detail.Mounts.length === 0 && + (!detail.HostConfig?.Binds || detail.HostConfig.Binds.length === 0) && + (!detail.HostConfig?.VolumesFrom || detail.HostConfig.VolumesFrom.length === 0) + ); +} + +function hasNoPublishedPorts(config: { + NetworkMode?: string; + PortBindings?: Record | null; + PublishAllPorts?: boolean; +} | undefined, networks?: Record | null): boolean { + if (!config) return false; + const networkMode = (config.NetworkMode ?? "").toLowerCase(); + if (!["default", "bridge"].includes(networkMode)) return false; + const attached = Object.keys(networks ?? {}).map((name) => name.toLowerCase()); + if (attached.length !== 1 || !["default", "bridge"].includes(attached[0]!)) return false; + return ( + config.PublishAllPorts !== true && + !Object.values(config.PortBindings ?? {}).some((value) => Array.isArray(value) && value.length > 0) + ); +} + +function statusProblem(status: VpsComputerStatus): string | null { + if (!status.configured) return "Configure a VPS SSH alias in App Settings → Connections"; + if (!status.daemonUp) return "Docker over SSH could not reach the VPS; check the SSH alias and Docker on the VPS"; + if (!status.image) return `Prepare the pinned OpenMausBot Cua image on the VPS (Driver ${CUA_DRIVER_VERSION})`; + if (status.container === "missing") return "No OpenMausBot container exists for this bot on the VPS"; + if (!status.imageMatches) return "The VPS container uses an incompatible or untrusted OpenMausBot image"; + if (!status.managed) return "The VPS container name is occupied by a container OpenMausBot did not create"; + if (status.network === "unsafe") return "The VPS container uses an unapproved network or publishes ports; refusing to use it"; + if (status.mounts === "unsafe") return "The VPS container has host mounts; refusing to use it"; + if (status.security === "unsafe") return "The VPS container is missing OpenMausBot safety limits"; + if (status.container === "stopped") return "The OpenMausBot VPS container is stopped"; + if (status.desktop_error) return `The VPS Cua desktop failed to start: ${status.desktop_error}`; + if (!status.desktopReady) return "The VPS container started, but Cua Driver is not ready yet"; + return null; +} + +/** The uncached inspection. Lifecycle mutations and their readiness waits + * call this directly — they must see and publish the truth, never a poll's + * snapshot; the exported vpsComputerStatus wraps it with the poll cache. + * + * `docker info` is deliberately NOT probed as its own round-trip: every + * docker-over-SSH invocation is a full process + SSH connection, and the + * image inspect right below already proves the daemon answers — even its + * "No such image" failure is a daemon reply. Anything that is neither JSON + * nor a "no such object" reply is attributed to the transport instead. */ +async function computeVpsComputerStatus( + cfg: AppConfig, + botId: string, + runner: VpsCommandRunner, +): Promise { + const alias = vpsSshAlias(cfg); + const status = emptyStatus(botId, alias); + if (!alias) return status; + const run = (args: string[], timeoutMs = 10_000, input?: string) => + runner(vpsDockerArgs(alias, args), { timeoutMs, input }); + + let inspectedImageId: string | null = null; + try { + const inspected = JSON.parse((await run(["image", "inspect", VPS_IMAGE])).stdout) as Array<{ + Id?: string; + id?: string; + Config?: { Labels?: Record }; + config?: { Labels?: Record; labels?: Record }; + }>; + status.daemonUp = true; + const image = inspected[0]; + const labels = image?.Config?.Labels ?? image?.config?.Labels ?? image?.config?.labels; + const imageId = image?.Id ?? image?.id; + inspectedImageId = imageId && IMAGE_ID.test(imageId) ? imageId : null; + status.image_id = inspectedImageId; + status.image = Boolean(inspectedImageId) && imageLabelsMatch(labels); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + if (!isMissingObjectMessage(message)) { + status.problem = transportFailure(message); + return status; + } + // A clean "no such image" is still a daemon answer. + status.daemonUp = true; + status.image = false; + } + + try { + const inspected = JSON.parse((await run(["inspect", status.container_name])).stdout) as Array<{ + Config?: { Image?: string; Labels?: Record }; + HostConfig?: DockerHardeningConfig & { + Binds?: string[] | null; + VolumesFrom?: string[] | null; + NetworkMode?: string; + PortBindings?: Record | null; + PublishAllPorts?: boolean; + }; + Id?: string; + id?: string; + Image?: string; + NetworkSettings?: { Networks?: Record | null }; + Mounts?: unknown; + State?: { Running?: boolean }; + }>; + const detail = inspected[0]; + const labels = detail?.Config?.Labels; + const containerId = detail?.Id ?? detail?.id; + status.container_id = containerId && CONTAINER_ID.test(containerId) ? containerId : null; + status.container = detail?.State?.Running ? "running" : "stopped"; + status.imageMatches = + status.image && + Boolean(status.container_id) && + (detail?.Config?.Image === VPS_IMAGE || detail?.Config?.Image === inspectedImageId) && + Boolean(inspectedImageId) && + detail?.Image === inspectedImageId && + imageLabelsMatch(labels); + status.managed = + labels?.[VPS_MANAGED_LABEL] === "1" && labels?.[VPS_CONTAINER_LABEL] === status.container_name; + status.network = hasNoPublishedPorts(detail?.HostConfig, detail?.NetworkSettings?.Networks) ? "private" : "unsafe"; + status.mounts = hasNoHostMounts(detail ?? {}) ? "none" : "unsafe"; + status.security = dockerSecurityIsHardened(detail?.HostConfig, { restartPolicy: "unless-stopped" }) + ? "hardened" + : "unsafe"; + + const containerRef = status.container_id; + const canProbe = + status.container === "running" && + status.image && + status.imageMatches && + status.managed && + status.network === "private" && + status.mounts === "none" && + status.security === "hardened"; + if (canProbe && containerRef) { + try { + const version = await run(cuaExecArgs(["--version"], { container: containerRef })); + if (version.stdout.trim() !== `cua-driver ${CUA_DRIVER_VERSION}`) throw new Error("unexpected Cua Driver version"); + await run(cuaExecArgs(["status", "--socket", CUA_SOCKET], { container: containerRef })); + const health = await run( + cuaExecArgs(["call", "health_report", "{}", "--socket", CUA_SOCKET], { container: containerRef }), + 15_000, + ); + const report = JSON.parse(health.stdout) as { + schema_version?: string; + overall?: string; + checks?: unknown[]; + }; + if ( + report.schema_version !== "1" || + !Array.isArray(report.checks) || + (report.overall !== "ok" && report.overall !== "degraded") + ) { + throw new Error(`Cua health report is ${report.overall ?? "invalid"}`); + } + // The desktop must ANSWER, not render: get_desktop_state succeeding + // is the readiness proof. The Local VM also pulls a pixel-validated + // readiness screenshot because a local exec is free; over SSH that is + // a full-frame base64 transfer on every status poll, so pixel + // validation lives solely in vpsComputerScreenshot(). + await run( + cuaExecArgs(["call", "get_desktop_state", "{}", "--socket", CUA_SOCKET], { container: containerRef }), + 20_000, + ); + status.desktopReady = true; + } catch (error) { + status.desktopReady = false; + status.desktop_error = error instanceof Error ? error.message.slice(0, 320) : null; + // Mirror the Local VM's probe: when the desktop fails, the + // supervisor's error log says WHY — a bounded tail turns an endless + // "not ready yet" into something the user can act on. + try { + const errorLog = await run( + ["exec", containerRef, "tail", "-n", "4", "/var/log/supervisor/cua-driver.error.log"], + 10_000, + ); + status.desktop_error = + errorLog.stdout.replace(/\s+/g, " ").trim().slice(0, 320) || status.desktop_error; + } catch { + // The log may not exist during the first seconds of container boot. + } + } + } + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + if (isMissingObjectMessage(message)) { + status.container = "missing"; + } else { + status.daemonUp = false; + status.problem = transportFailure(message); + return status; + } + } + + status.problem = statusProblem(status); + status.ready = status.problem === null; + return status; +} + +/** Poll-facing status. Real SSH invocations (the default runner) are served + * from a short TTL cache — the panel and the screen poller each re-check + * every few seconds, and without the cache one healthy poll cycle cost a + * dozen SSH connections. Injected runners (tests, lifecycle internals) + * always recompute. */ +export async function vpsComputerStatus( + cfg: AppConfig, + botId: string, + runner: VpsCommandRunner = defaultRunner, +): Promise { + const key = vpsLockKey(cfg, botId); + const cacheable = runner === defaultRunner && key !== null; + if (cacheable) { + const cached = statusCache.get(key); + if (cached && cached.expiresAt > Date.now()) return cached.status; + } + const status = await computeVpsComputerStatus(cfg, botId, runner); + if (cacheable) statusCache.set(key, { status, expiresAt: Date.now() + STATUS_CACHE_TTL_MS }); + return status; +} + +export function vpsContainerRunArgs(containerName: string, imageRef = VPS_IMAGE): string[] { + if (!CONTAINER_NAME.test(containerName) || (imageRef !== VPS_IMAGE && !IMAGE_ID.test(imageRef))) { + throw new Error("invalid managed VPS container or image reference"); + } + return [ + "run", + "-d", + "--name", + containerName, + "--label", + `${VPS_MANAGED_LABEL}=1`, + "--label", + `${VPS_CONTAINER_LABEL}=${containerName}`, + "--label", + `${MANAGED_LABEL}=1`, + "--label", + `${DRIVER_LABEL}=${CUA_DRIVER_VERSION}`, + "--label", + `${BASE_IMAGE_LABEL}=${BASE_IMAGE_DIGEST}`, + "--label", + `${IMAGE_LAYER_LABEL}=${IMAGE_LAYER_VERSION}`, + "--memory", + "4g", + "--memory-swap", + "4g", + "--cpus", + "2", + "--pids-limit", + String(PIDS_LIMIT), + "--network", + "bridge", + "--ipc", + "private", + "--cgroupns", + "private", + "--cap-drop", + "ALL", + "--cap-add", + "SETUID", + "--cap-add", + "SETGID", + "--shm-size", + "512m", + // A VPS reboots with nobody watching; without a restart policy the + // container stays down afterwards and every turn silently degrades until + // someone opens the panel. unless-stopped survives reboots while still + // honoring an explicit Stop. The shared hardening check accepts exactly + // this policy for the VPS caller (and only "no"/unset for the Local VM, + // whose desktop cannot safely resume). + "--restart", + "unless-stopped", + imageRef, + ]; +} + +function assertUsableContainer(status: VpsComputerStatus) { + if ( + !status.image || + !status.imageMatches || + !status.managed || + status.network !== "private" || + status.mounts !== "none" || + status.security !== "hardened" + ) { + throw Object.assign(new Error(status.problem ?? "The existing VPS container is unsafe or incompatible"), { + status: 409, + }); + } +} + +async function prepareVpsImage(alias: string, runner: VpsCommandRunner) { + await runner(vpsDockerArgs(alias, ["pull", BASE_IMAGE]), { timeoutMs: 10 * 60_000 }); + await runner(vpsDockerArgs(alias, ["build", "-t", VPS_IMAGE, "-"]), { + input: managedImageDockerfile(), + timeoutMs: 10 * 60_000, + }); +} + +/** Waits for the driver inside a verified, running container to come up. + * Backoff doubles 0.5s→4s because each poll is a real SSH connection, and + * between polls only ONE cheap exec (`cua-driver status`) asks whether the + * driver answers yet — the full multi-invocation status runs again only when + * that predicate flips, and once more at the deadline, so the returned state + * is always a complete inspection. */ +async function waitForVpsReady( + cfg: AppConfig, + botId: string, + runner: VpsCommandRunner, + budgetMs = 60_000, +): Promise { + const alias = vpsSshAlias(cfg); + const deadline = Date.now() + budgetMs; + let status = await computeVpsComputerStatus(cfg, botId, runner); + let delayMs = 500; + while (!status.ready && alias && Date.now() < deadline) { + if ( + !status.daemonUp || + !status.image || + !status.imageMatches || + !status.managed || + status.container !== "running" || + status.network !== "private" || + status.mounts !== "none" || + status.security !== "hardened" + ) { + return status; + } + await new Promise((resolve) => { + const sleep = setTimeout(resolve, Math.min(delayMs, Math.max(0, deadline - Date.now()))); + sleep.unref?.(); + }); + delayMs = Math.min(delayMs * 2, 4_000); + const container = status.container_id ?? status.container_name; + const driverAnswers = await runner( + vpsDockerArgs(alias, cuaExecArgs(["status", "--socket", CUA_SOCKET], { container })), + { timeoutMs: 10_000 }, + ).then( + () => true, + () => false, + ); + if (!driverAnswers && Date.now() < deadline) continue; + status = await computeVpsComputerStatus(cfg, botId, runner); + } + return status; +} + +async function withVpsLifecycleLock(key: string, operation: () => Promise): Promise { + const previous = lifecycleLocks.get(key); + let release!: () => void; + const current = new Promise((resolve) => { + release = resolve; + }); + lifecycleLocks.set(key, current); + if (previous) { + let acquireTimer: ReturnType | undefined; + const acquired = await Promise.race([ + previous.then(() => true), + new Promise((resolve) => { + acquireTimer = setTimeout(() => resolve(false), LOCK_ACQUIRE_TIMEOUT_MS); + acquireTimer.unref?.(); + }), + ]); + if (acquireTimer) clearTimeout(acquireTimer); + if (!acquired) { + // Keep the queue serialized: this slot opens only when the holder's + // does, so a later caller can never run beside the long operation the + // timed-out one refused to wait for. + void previous.then(() => { + release(); + if (lifecycleLocks.get(key) === current) lifecycleLocks.delete(key); + }); + throw Object.assign(new Error("the VPS is being prepared — try again shortly"), { status: 409 }); + } + } + try { + return await operation(); + } finally { + release(); + if (lifecycleLocks.get(key) === current) lifecycleLocks.delete(key); + } +} + +function vpsLockKey(cfg: AppConfig, botId: string): string | null { + const alias = vpsSshAlias(cfg); + return alias ? `${alias}:${vpsContainerName(botId)}` : null; +} + +export async function vpsComputerAction( + action: VpsLifecycleAction, + cfg: AppConfig, + botId: string, + runner: VpsCommandRunner = defaultRunner, +): Promise { + const alias = vpsSshAlias(cfg); + if (!alias) throw Object.assign(new Error("VPS is not configured — add an SSH config alias in App Settings → Connections"), { status: 409 }); + const key = `${alias}:${vpsContainerName(botId)}`; + const operation = async () => { + // A mutation invalidates every cached poll answer, before and after: the + // panel must never keep showing the pre-action world for a TTL. + statusCache.delete(key); + try { + const before = await computeVpsComputerStatus(cfg, botId, runner); + if (!before.daemonUp) throw Object.assign(new Error(before.problem ?? "Docker over SSH is not reachable"), { status: 409 }); + const run = (args: string[], timeoutMs = 2 * 60_000) => runner(vpsDockerArgs(alias, args), { timeoutMs }); + + const containerRef = before.container_id ?? before.container_name; + if (action === "provision") { + if (before.container === "missing") { + let imageRef = before.image ? before.image_id : null; + if (!before.image) { + await prepareVpsImage(alias, runner); + imageRef = (await computeVpsComputerStatus(cfg, botId, runner)).image_id; + } + if (!imageRef) throw Object.assign(new Error("The prepared VPS image could not be identified"), { status: 409 }); + await run(vpsContainerRunArgs(before.container_name, imageRef)); + } else { + assertUsableContainer(before); + if (before.container === "stopped") await run(["start", containerRef]); + } + } else if (action === "start") { + if (before.container === "missing") throw Object.assign(new Error("No VPS container exists for this bot"), { status: 409 }); + if (before.container === "running") throw Object.assign(new Error("The VPS container is already running"), { status: 409 }); + assertUsableContainer(before); + await run(["start", containerRef]); + } else if (action === "remove") { + // remove exists to escape an incompatible or unsafe container (an + // IMAGE_LAYER_VERSION bump otherwise bricks the bot: provision 409s + // on assertUsableContainer forever), so it deliberately skips that + // check. The ownership labels from the inspect are the only gate: + // never docker-rm a container OpenMausBot did not create, even one + // squatting on our name. + if (before.container === "missing") return before; + if (!before.managed) { + throw Object.assign( + new Error("The VPS container name is occupied by a container OpenMausBot did not create — remove it on the VPS yourself"), + { status: 409 }, + ); + } + await run(["rm", "-f", containerRef]); + return computeVpsComputerStatus(cfg, botId, runner); + } else { + if (before.container !== "running") throw Object.assign(new Error("The VPS container is not running"), { status: 409 }); + assertUsableContainer(before); + await run(["stop", containerRef]); + } + return action === "stop" ? computeVpsComputerStatus(cfg, botId, runner) : waitForVpsReady(cfg, botId, runner); + } finally { + statusCache.delete(key); + } + }; + return withVpsLifecycleLock(key, operation); +} + +/** Auto is intentionally read-only: it can attach only to an existing ready + * container. It recomputes rather than reading the poll cache — routing a + * turn onto a container that stopped seconds ago is worse than one extra + * inspection at turn start. */ +export async function reuseVps( + cfg: AppConfig, + botId: string, + runner: VpsCommandRunner = defaultRunner, +): Promise { + const key = vpsLockKey(cfg, botId); + const status = await (key + ? withVpsLifecycleLock(key, () => computeVpsComputerStatus(cfg, botId, runner)) + : computeVpsComputerStatus(cfg, botId, runner)); + return status.ready ? status : null; +} + +export function vpsContainerMcpArgs(alias: string, containerName: string): string[] { + if (!isValidSshAlias(alias) || (!CONTAINER_NAME.test(containerName) && !CONTAINER_ID.test(containerName))) { + throw new Error("invalid VPS MCP connection"); + } + return vpsDockerArgs( + alias, + cuaExecArgs(["mcp", "--socket", CUA_SOCKET], { container: containerName, interactive: true }), + ); +} + +export function vpsComputerMcp(cfg: AppConfig, botId: string, containerRef?: string): { + command: string; + args: string[]; + env: Record; +} { + const alias = vpsSshAlias(cfg); + if (!alias) throw new Error("VPS is not configured — add an SSH config alias first"); + return { + command: process.execPath, + args: [SPAWNED_PROXIES.vpsContainerMcp, alias, containerRef ?? vpsContainerName(botId)], + env: { ELECTRON_RUN_AS_NODE: "1" }, + }; +} + +export function vpsDriverError(driverKind: string, computerMcp: boolean): string | null { + if (driverKind === "boxAgent") { + return "The Computer engine runs its agent on Box and cannot use a self-hosted VPS — choose Claude or an ACP engine"; + } + if (!computerMcp) { + return "This model engine cannot mount a self-hosted VPS computer — choose Claude or an ACP engine"; + } + return null; +} + +export async function vpsComputerScreenshot( + cfg: AppConfig, + botId: string, + runner: VpsCommandRunner = defaultRunner, +): Promise<{ png: string; format: "png" | "jpeg" }> { + const alias = vpsSshAlias(cfg); + if (!alias) throw Object.assign(new Error("VPS is not configured"), { status: 409 }); + const key = `${alias}:${vpsContainerName(botId)}`; + const cacheable = runner === defaultRunner; + return withVpsLifecycleLock(key, async () => { + // Same shape as containerComputerScreenshot's screenshotStatusCache: the + // poller runs every few seconds, and re-verifying the whole container + // between frames multiplied every frame's SSH cost. + const cached = cacheable ? statusCache.get(key) : undefined; + const status = + cached && cached.expiresAt > Date.now() + ? cached.status + : await computeVpsComputerStatus(cfg, botId, runner); + if (!status.ready) { + if (cacheable) statusCache.delete(key); + throw Object.assign(new Error(status.problem ?? "The VPS computer is not ready"), { status: 409 }); + } + if (cacheable) statusCache.set(key, { status, expiresAt: Date.now() + STATUS_CACHE_TTL_MS }); + const containerRef = status.container_id ?? status.container_name; + // The ref goes straight into docker argv, and a cached status is one + // more step removed from the inspect that produced it — revalidate the + // exact shapes before spending an exec on it. + if (!CONTAINER_ID.test(containerRef) && !CONTAINER_NAME.test(containerRef)) { + throw Object.assign(new Error("the VPS container reference is malformed"), { status: 409 }); + } + try { + await runner( + vpsDockerArgs( + alias, + cuaExecArgs( + ["call", "get_desktop_state", "{}", "--socket", CUA_SOCKET, "--screenshot-out-file", SCREENSHOT_PATH], + { container: containerRef }, + ), + ), + { timeoutMs: 30_000 }, + ); + const encoded = (await runner(vpsDockerArgs(alias, [ + "exec", + "-u", + "cua", + "-e", + "HOME=/home/cua", + containerRef, + "base64", + "-w0", + SCREENSHOT_PATH, + ]), { timeoutMs: 30_000 })).stdout.trim(); + const checked = wholeScreenshot(Buffer.from(encoded, "base64")); + if (!checked.ok) throw Object.assign(new Error("Cua Driver returned an incomplete VPS screenshot"), { status: 502 }); + return { png: encoded, format: checked.mime === "image/jpeg" ? "jpeg" : "png" }; + } catch (error) { + // The failure may mean the world changed (container stopped, link + // dropped); a cached "ready" would keep the poller failing for a TTL. + if (cacheable) statusCache.delete(key); + throw error; + } finally { + await runner(vpsDockerArgs(alias, ["exec", "-u", "cua", containerRef, "rm", "-f", SCREENSHOT_PATH]), { + timeoutMs: 10_000, + }).catch(() => {}); + } + }); +} diff --git a/server/vps-container-mcp.test.ts b/server/vps-container-mcp.test.ts new file mode 100644 index 0000000000..c7abdf66ba --- /dev/null +++ b/server/vps-container-mcp.test.ts @@ -0,0 +1,72 @@ +import { spawn } from "node:child_process"; +import { chmod, mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { fileURLToPath } from "node:url"; +import { afterEach, describe, expect, it } from "vitest"; + +import { vpsContainerName } from "./vps-computer.ts"; + +const temporary: string[] = []; + +afterEach(async () => { + await Promise.all(temporary.splice(0).map((path) => rm(path, { recursive: true, force: true }))); +}); + +function runBridge(bin: string, input: string) { + return new Promise<{ code: number | null; stdout: string; stderr: string }>((resolve, reject) => { + const child = spawn( + process.execPath, + [ + fileURLToPath(new URL("./vps-container-mcp.ts", import.meta.url)), + "production-vps", + vpsContainerName("bridge-test"), + ], + { + env: { ...process.env, OMB_EXTRA_PATH: bin, NODE_NO_WARNINGS: "1" }, + stdio: ["pipe", "pipe", "pipe"], + }, + ); + let stdout = ""; + let stderr = ""; + child.stdout.on("data", (chunk: Buffer) => (stdout += chunk.toString())); + child.stderr.on("data", (chunk: Buffer) => (stderr += chunk.toString())); + child.on("error", reject); + child.on("close", (code) => resolve({ code, stdout, stderr })); + child.stdin.on("error", () => {}); + child.stdin.end(input); + }); +} + +describe.skipIf(process.platform === "win32")("VPS Cua MCP bridge", () => { + it("passes MCP bytes unchanged to docker exec over the validated SSH target", async () => { + const bin = await mkdtemp(join(tmpdir(), "openmausbot-vps-mcp-")); + temporary.push(bin); + const fakeDocker = join(bin, "docker"); + await writeFile(fakeDocker, "#!/bin/sh\nprintf 'ARGS:%s\\n' \"$*\" >&2\ncat\n", { mode: 0o700 }); + await chmod(fakeDocker, 0o700); + + const input = `{"jsonrpc":"2.0","id":1,"method":"tools/list","data":"${"x".repeat(2 * 1024 * 1024)}"}\n`; + const result = await runBridge(bin, input); + + expect(result.code).toBe(0); + expect(result.stdout).toBe(input); + expect(result.stderr).toContain( + `ARGS:-H ssh://production-vps exec -i -u cua -e HOME=/home/cua -e DISPLAY=:1 ` + + `-e CUA_DRIVER_INSTALL_CHANNEL=python_package -e CUA_DRIVER_RS_TELEMETRY_ENABLED=0 ${vpsContainerName("bridge-test")} ` + + "/usr/local/libexec/openmausbot/cua-driver mcp --socket /run/user/1000/openmausbot-cua.sock", + ); + }); + + it("survives docker closing stdin before consuming the request", async () => { + const bin = await mkdtemp(join(tmpdir(), "openmausbot-vps-mcp-")); + temporary.push(bin); + const fakeDocker = join(bin, "docker"); + await writeFile(fakeDocker, "#!/bin/sh\nexit 0\n", { mode: 0o700 }); + await chmod(fakeDocker, 0o700); + + const result = await runBridge(bin, "x".repeat(4 * 1024 * 1024)); + + expect(result.code).toBe(0); + }); +}); diff --git a/server/vps-container-mcp.ts b/server/vps-container-mcp.ts new file mode 100644 index 0000000000..b7808f325a --- /dev/null +++ b/server/vps-container-mcp.ts @@ -0,0 +1,27 @@ +// Transparent stdio bridge to the official Cua MCP server in a VPS +// container. Docker's SSH transport handles authentication through the +// user's normal SSH config and agent; this process stores no credentials. +// The piping, drain-safe exit, and dead-transport watchdog live in +// mcp-bridge.ts, shared with the Local VM entry point. +import { runMcpBridge } from "./mcp-bridge.ts"; +import { vpsContainerMcpArgs, vpsDockerArgs } from "./vps-computer.ts"; + +const [alias, containerName] = process.argv.slice(2); +const sshAlias = alias ?? ""; +let args: string[]; +try { + args = vpsContainerMcpArgs(sshAlias, containerName ?? ""); +} catch { + process.stderr.write("invalid VPS MCP connection\n"); + process.exit(2); +} + +runMcpBridge({ + command: "docker", + args, + label: "VPS Cua Driver", + // The probe checks the TRANSPORT (SSH + daemon), deliberately not the + // driver: a busy desktop mid-tool-call must never look dead, while an + // unreachable VPS must, and `docker version` distinguishes exactly that. + liveness: { command: "docker", args: vpsDockerArgs(sshAlias, ["version", "--format", "{{.Server.Version}}"]) }, +}); diff --git a/server/vps-routing.test.ts b/server/vps-routing.test.ts new file mode 100644 index 0000000000..aa7610f8d9 --- /dev/null +++ b/server/vps-routing.test.ts @@ -0,0 +1,306 @@ +// Turn routing for the BYO-VPS backend, end to end on the real harness +// server: a bot patched to cloudBackend:"vps" must get the managed container +// mounted as its computer (integrations.localComputer → the "computer" MCP +// server), carry the VPS system-prompt clause, never provision from Auto, +// and hold/clear its activeVpsThreads claim across the turn. +// +// The "injected VpsCommandRunner" is a fake `docker` executable on +// OMB_EXTRA_PATH: the server runs in its own process, so injection happens +// where defaultRunner actually looks — argv in, canned inspect JSON out, +// every invocation appended to a log the assertions read. The agent is the +// fake ACP CLI in echo-gated mode (see steer-queue.test.ts), whose echo +// reply carries the FULL prompt and whose gate file gives a deterministic +// busy window — no sleeps anywhere. +import { spawn, type ChildProcess } from "node:child_process"; +import { chmodSync, mkdirSync, mkdtempSync, readFileSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { dirname, join } from "node:path"; +import { fileURLToPath } from "node:url"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; + +import { + BASE_IMAGE_DIGEST, + BASE_IMAGE_LABEL, + CUA_DRIVER_VERSION, + DRIVER_LABEL, + IMAGE_LAYER_LABEL, + IMAGE_LAYER_VERSION, + MANAGED_LABEL, +} from "./container-computer.ts"; +import { VPS_CONTAINER_LABEL, VPS_IMAGE, VPS_MANAGED_LABEL } from "./vps-computer.ts"; +import { removeTempDir, waitForExit } from "./testing/cleanup.ts"; + +const SERVER_DIR = dirname(fileURLToPath(import.meta.url)); +const FAKE_CLI = join(SERVER_DIR, "testing", "fake-acp-cli.ts"); +const PORT = 18800 + Math.floor(Math.random() * 10_000); +const BASE = `http://127.0.0.1:${PORT}`; +const IMAGE_ID = `sha256:${"a".repeat(64)}`; +const CONTAINER_ID = "b".repeat(64); + +// the fake docker is a POSIX shell script, like every process fixture here +const posixOnly = describe.skipIf(process.platform === "win32"); + +function imageInspectJson(): string { + return JSON.stringify([ + { + Id: IMAGE_ID, + Config: { + Labels: { + [MANAGED_LABEL]: "1", + [DRIVER_LABEL]: CUA_DRIVER_VERSION, + [BASE_IMAGE_LABEL]: BASE_IMAGE_DIGEST, + [IMAGE_LAYER_LABEL]: IMAGE_LAYER_VERSION, + }, + }, + }, + ]); +} + +/** __NAME__ is substituted by the fake docker from the inspect argv, because + * the container name derives from a bot id that only exists at runtime. */ +function containerInspectTemplate(): string { + return JSON.stringify([ + { + Id: CONTAINER_ID, + Image: IMAGE_ID, + Config: { + Image: VPS_IMAGE, + Labels: { + [VPS_MANAGED_LABEL]: "1", + [VPS_CONTAINER_LABEL]: "__NAME__", + [MANAGED_LABEL]: "1", + [DRIVER_LABEL]: CUA_DRIVER_VERSION, + [BASE_IMAGE_LABEL]: BASE_IMAGE_DIGEST, + [IMAGE_LAYER_LABEL]: IMAGE_LAYER_VERSION, + }, + }, + State: { Running: true }, + NetworkSettings: { Networks: { bridge: {} } }, + Mounts: [], + HostConfig: { + Binds: [], + VolumesFrom: [], + NetworkMode: "bridge", + PortBindings: {}, + PublishAllPorts: false, + Memory: 4 * 1024 * 1024 * 1024, + MemorySwap: 4 * 1024 * 1024 * 1024, + NanoCpus: 2_000_000_000, + PidsLimit: 512, + CapDrop: ["ALL"], + CapAdd: ["CAP_SETUID", "CAP_SETGID"], + Privileged: false, + PidMode: "", + IpcMode: "private", + UTSMode: "", + ShmSize: 512 * 1024 * 1024, + Devices: [], + DeviceRequests: [], + SecurityOpt: [], + UsernsMode: "", + CgroupnsMode: "private", + OomKillDisable: false, + AutoRemove: false, + RestartPolicy: { Name: "unless-stopped", MaximumRetryCount: 0 }, + }, + }, + ]); +} + +const FAKE_DOCKER = `#!/bin/sh +printf '%s\\n' "$*" >> "$FAKE_DOCKER_LOG" +case "$*" in + *" image inspect "*) cat "$FAKE_DOCKER_DIR/image.json" ;; + *" exec "*"--version"*) echo "cua-driver ${CUA_DRIVER_VERSION}" ;; + *" exec "*"--screenshot-out-file"*) echo "{}" ;; + *" exec "*"health_report"*) echo '{"schema_version":"1","overall":"ok","checks":[]}' ;; + *" exec "*"get_desktop_state"*) echo "{}" ;; + *" exec "*"base64"*) cat "$FAKE_DOCKER_DIR/screenshot.b64" ;; + *" exec "*"status"*) echo "running" ;; + *" exec "*"rm -f"*) : ;; + *" inspect "*) for arg in "$@"; do name="$arg"; done; sed "s|__NAME__|$name|g" "$FAKE_DOCKER_DIR/container.json.tpl" ;; + *) echo "unexpected docker invocation: $*" >&2; exit 64 ;; +esac +`; + +posixOnly("VPS turn routing e2e (fake ACP fleet + fake docker over SSH)", () => { + let child: ChildProcess; + let home: string; + let stderr = ""; + let gateFile: string; + let acpDump: string; + let dockerLog: string; + + type ApiBody = Record; + + const api = async (method: string, path: string, body?: ApiBody): Promise<{ status: number; body: any }> => { + const res = await fetch(`${BASE}${path}`, { + method, + headers: body ? { "content-type": "application/json" } : undefined, + body: body ? JSON.stringify(body) : undefined, + }); + return { status: res.status, body: await res.json() }; + }; + + const botById = async (id: string) => + (await api("GET", "/api/bots")).body.bots.find((b: any) => b.id === id); + + const until = async (probe: () => Promise, what: string, timeoutMs = 30_000) => { + const deadline = Date.now() + timeoutMs; + for (;;) { + if (await probe()) return; + if (Date.now() > deadline) throw new Error(`${what} never happened. stderr: ${stderr.slice(-2000)}`); + await new Promise((r) => setTimeout(r, 200)); + } + }; + + beforeAll(async () => { + chmodSync(FAKE_CLI, 0o755); + home = mkdtempSync(join(tmpdir(), "omb-vps-routing-")); + mkdirSync(join(home, ".openmausbot"), { recursive: true }); + const fakeBin = join(home, "fakebin"); + mkdirSync(fakeBin, { recursive: true }); + gateFile = join(home, "turn.gate"); + acpDump = join(home, "acp.dump.json"); + dockerLog = join(fakeBin, "docker.log"); + + writeFileSync(join(fakeBin, "docker"), FAKE_DOCKER, { mode: 0o755 }); + chmodSync(join(fakeBin, "docker"), 0o755); + writeFileSync(join(fakeBin, "image.json"), imageInspectJson()); + writeFileSync(join(fakeBin, "container.json.tpl"), containerInspectTemplate()); + const png = Buffer.concat([ + Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]), + Buffer.alloc(600), + Buffer.from("IEND", "ascii"), + ]); + writeFileSync(join(fakeBin, "screenshot.b64"), png.toString("base64")); + writeFileSync(dockerLog, ""); + + writeFileSync( + join(home, ".openmausbot", "config.json"), + JSON.stringify({ + instances: { + vps: { + driver: "grokAgent", + environment: { FAKE_ACP_MODE: "echo-gated", FAKE_ACP_GATE_FILE: gateFile, FAKE_ACP_DUMP: acpDump }, + config: { cli: FAKE_CLI, fullAuto: true }, + }, + }, + }), + ); + + const env: NodeJS.ProcessEnv = { + HOME: home, + USERPROFILE: home, + OMB_PORT: String(PORT), + OMB_EXTRA_PATH: fakeBin, + FAKE_DOCKER_DIR: fakeBin, + FAKE_DOCKER_LOG: dockerLog, + }; + if (process.env.PATH) env.PATH = process.env.PATH; + if (process.env.SystemRoot) env.SystemRoot = process.env.SystemRoot; + child = spawn(process.execPath, [join(SERVER_DIR, "index.ts")], { + cwd: join(SERVER_DIR, ".."), + env, + stdio: ["ignore", "pipe", "pipe"], + }); + child.stderr!.on("data", (c) => (stderr += c)); + + const deadline = Date.now() + 20_000; + for (;;) { + try { + if ((await fetch(`${BASE}/api/health`)).ok) break; + } catch { + /* not up yet */ + } + if (Date.now() > deadline) throw new Error(`server never came up. stderr:\n${stderr}`); + if (child.exitCode !== null) throw new Error(`server exited ${child.exitCode}. stderr:\n${stderr}`); + await new Promise((r) => setTimeout(r, 150)); + } + }, 40_000); + + afterAll(async () => { + await waitForExit(child, { signal: "SIGTERM" }); + await removeTempDir(home); + }); + + it( + "mounts the VPS computer on the turn, tells the model, reuses without provisioning, and clears the claim", + async () => { + expect((await api("PUT", "/api/config", { vps: { sshAlias: "production-vps" } })).status).toBe(200); + + const bot = (await api("POST", "/api/bots")).body.bot; + await api("PATCH", `/api/bots/${bot.id}`, { + name: "Remote hand", + modelSelection: { instanceId: "vps", model: "fake-model" }, + }); + expect((await api("PATCH", `/api/bots/${bot.id}`, { cloudBackend: "vps" })).status).toBe(200); + // bot.computer stays unset — Auto, the mode that must never provision + + const sent = await api("POST", `/api/bots/${bot.id}/messages`, { text: "check the remote desktop" }); + expect(sent.status).toBe(202); + await until(async () => (await botById(bot.id))?.busy === true, "the gated turn"); + + // the turn is claimed: the SSH alias cannot be swapped under it... + const aliasChange = await api("PUT", "/api/config", { vps: { sshAlias: "other-vps" } }); + expect(aliasChange.status).toBe(409); + expect(aliasChange.body.error).toMatch(/active VPS turn/); + // ...and neither can the bot's cloud backend + expect((await api("PATCH", `/api/bots/${bot.id}`, { cloudBackend: "box" })).status).toBe(409); + + // open the gate: the echo settles carrying the FULL prompt + writeFileSync(gateFile, "open"); + let snapshot: any; + await until(async () => { + snapshot = await botById(bot.id); + return ( + snapshot?.busy === false && + snapshot.messages.some((m: any) => m.kind === "text" && m.text?.startsWith("echo: ")) + ); + }, "the echoed turn"); + + const echo = snapshot.messages.find((m: any) => m.kind === "text" && m.text?.startsWith("echo: ")).text; + // the VPS clause, including the disposable-filesystem warning + expect(echo).toContain("self-hosted remote Linux computer"); + expect(echo).toContain("wiped whenever its container is recreated"); + + // the official Cua MCP server was mounted through the VPS bridge + // SAFETY: the fake ACP CLI wrote this dump itself from session/new's + // mcpServers array; the shape is pinned by fake-acp-cli.ts. + const mcpServers = JSON.parse(readFileSync(`${acpDump}.mcp.json`, "utf8")) as Array<{ + name: string; + command: string; + args?: string[]; + }>; + const computer = mcpServers.find((s) => s.name === "computer"); + expect(computer, "no computer MCP server reached the agent").toBeTruthy(); + const bridgeArgs = computer?.args ?? []; + expect(bridgeArgs.some((a) => a.includes("vps-container-mcp"))).toBe(true); + expect(bridgeArgs.includes("production-vps")).toBe(true); + expect(bridgeArgs.includes(CONTAINER_ID)).toBe(true); + + // Auto attached to the existing container and NEVER provisioned: every + // docker-over-SSH invocation is an inspection or an exec. (The Local VM + // boot probe also hits the fake docker without -H; it is not the VPS.) + const invocations = readFileSync(dockerLog, "utf8") + .split("\n") + .filter((line) => line.startsWith("-H ssh://")); + expect(invocations.length).toBeGreaterThan(0); + for (const line of invocations) { + const command = line.split(" ")[2]; + expect(["image", "inspect", "exec", "version"], line).toContain(command); + } + + // the status route reads the same fake daemon and reports ready + const status = await api("GET", `/api/bots/${bot.id}/computer`); + expect(status.status).toBe(200); + expect(status.body).toMatchObject({ backend: "vps", ready: true, container: "running" }); + + // turn settled → the claim is gone: the alias may change again + const released = await api("PUT", "/api/config", { vps: { sshAlias: "other-vps" } }); + expect(released.status).toBe(200); + expect((await api("PUT", "/api/config", { vps: { sshAlias: "production-vps" } })).status).toBe(200); + }, + 60_000, + ); +}); diff --git a/src/components/ApiKeys.tsx b/src/components/ApiKeys.tsx index 303b92359c..facd968d07 100644 --- a/src/components/ApiKeys.tsx +++ b/src/components/ApiKeys.tsx @@ -209,3 +209,83 @@ export function ApiKeyRow({ ); } + +/** Non-secret Docker-over-SSH target. Keys and passwords stay with SSH. */ +export function VpsConnection() { + const { state, dispatch } = useStore(); + const [alias, setAlias] = useState(""); + const [saving, setSaving] = useState(false); + const [error, setError] = useState(null); + const configured = Boolean(state.config?.vps?.configured); + + useEffect(() => { + setAlias(state.config?.vps?.sshAlias ?? ""); + }, [state.config?.vps?.sshAlias]); + + const save = () => { + if (saving || (!alias.trim() && !configured)) return; + setSaving(true); + setError(null); + api("/api/config", { + method: "PUT", + body: JSON.stringify({ vps: { sshAlias: alias.trim() } }), + }) + .then((status: ConfigStatus) => { + dispatch({ type: "configStatus", config: status }); + setAlias(status.vps?.sshAlias ?? ""); + }) + .catch((e) => setError(e.message)) + .finally(() => setSaving(false)); + }; + + return ( +
+
+ + Self-hosted VPS + + Optional + + {configured && Connected} +
+
+ SSH config alias for the Linux VPS. OpenMausBot uses your normal SSH config and agent; it does not store keys or passwords.{" "} + See the{" "} + + setup guide + {" "} + for the required SSH alias shape. +
+
+ setAlias(e.target.value)} + onKeyDown={(e) => e.key === "Enter" && save()} + placeholder="my-vps" + aria-label="Self-hosted VPS SSH config alias" + autoComplete="off" + className="w-full rounded-lg border border-hairline/40 bg-inset px-3 py-2 text-[13px] text-ink placeholder:text-ink-secondary focus:border-hairline focus:outline-none" + /> + +
+ {error &&
{error}
} +
+ ); +} diff --git a/src/components/CloudBackendPicker.tsx b/src/components/CloudBackendPicker.tsx new file mode 100644 index 0000000000..35e099aa9d --- /dev/null +++ b/src/components/CloudBackendPicker.tsx @@ -0,0 +1,48 @@ +// The Box / Self-hosted VPS segmented control shown under the "Runs on" +// picker whenever a bot can end up on a cloud computer. One component, two +// homes (ComputerPanel and SettingsPanel), so the copy and the disabled +// rules can never drift apart. +import type { CloudBackend } from "../../server/contracts.ts"; +import { cn } from "@/lib/cn"; + +export function CloudBackendPicker({ + value, + vpsSupported, + onChange, +}: { + value: CloudBackend; + vpsSupported: boolean; + onChange: (backend: CloudBackend) => void; +}) { + return ( +
+
Cloud backend
+
+ {value === "vps" + ? "Auto only attaches to a VPS container that is already running — a stopped or missing one is never provisioned or started, and the bot quietly works as if no cloud computer existed. Choose Cloud to provision or start it. No interactive desktop tunnel is exposed." + : "Box is the default hosted computer. Choose Self-hosted VPS to use your SSH-configured Linux Docker host."} +
+
+ {(["box", "vps"] as const).map((backend, i) => { + const disabled = backend === "vps" && !vpsSupported; + return ( + + ); + })} +
+
+ ); +} diff --git a/src/components/ComputerPanel.tsx b/src/components/ComputerPanel.tsx index c9a83814ab..2a1b5bafec 100644 --- a/src/components/ComputerPanel.tsx +++ b/src/components/ComputerPanel.tsx @@ -22,6 +22,7 @@ import { useStore, type Bot } from "@/state/store"; import type { Routine } from "@/lib/routines"; import { ApiKeyRow } from "./ApiKeys"; import { cn } from "@/lib/cn"; +import { CloudBackendPicker } from "./CloudBackendPicker"; import { useDesktopCapabilities } from "./DesktopCapabilities"; import { RoutineEditor } from "./RoutinesPage"; import { AndroidDevicePanel, useAndroidUsbDevices } from "./AndroidDevicePanel"; @@ -40,6 +41,8 @@ type Phase = | "ready" | "vm" | "vm-unavailable" + | "vps-unconfigured" + | "vps-stopped" | "local" | "local-unavailable" | "off" @@ -82,7 +85,7 @@ export function ComputerPanel({ bot }: { bot: Bot }) { const [polledFrame, setPolledFrame] = useState<{ png: string; mime: string } | null>(null); const [vmFrame, setVmFrame] = useState(null); const [localFrame, setLocalFrame] = useState(null); - const [pending, setPending] = useState<"join" | "sleep" | null>(null); + const [pending, setPending] = useState<"join" | "sleep" | "provision" | null>(null); const [error, setError] = useState(null); const [creatingRoutine, setCreatingRoutine] = useState(false); const [panelView, setPanelView] = useState<"computer" | "android">("computer"); @@ -103,7 +106,11 @@ export function ComputerPanel({ bot }: { bot: Bot }) { selectedInstance.driverKind !== "boxAgent", ); const computerToolSupported = selectedInstance?.capabilities?.computerMcp === true; - const cloudSupported = computerToolSupported || selectedInstance?.driverKind === "boxAgent"; + const vpsSupported = Boolean(computerToolSupported && selectedInstance?.driverKind !== "boxAgent"); + const cloudBackend = bot.cloudBackend ?? "box"; + const cloudSupported = cloudBackend === "vps" + ? vpsSupported + : computerToolSupported || selectedInstance?.driverKind === "boxAgent"; const botRoutines = state.routines .filter((routine) => routine.botId === bot.id) .sort((a, b) => Number(b.enabled) - Number(a.enabled) || (a.nextRunAt ?? Infinity) - (b.nextRunAt ?? Infinity)); @@ -116,7 +123,7 @@ export function ComputerPanel({ bot }: { bot: Bot }) { ); const computerDestination = bot.computer === "cloud" - ? "this cloud box" + ? cloudBackend === "vps" ? "this self-hosted VPS" : "this cloud box" : bot.computer === "vm" ? "the Local VM" : bot.computer === "local" @@ -124,7 +131,7 @@ export function ComputerPanel({ bot }: { bot: Bot }) { : bot.computer === "off" ? null : phase === "ready" - ? "the cloud box selected by Auto" + ? cloudBackend === "vps" ? "the self-hosted VPS selected by Auto" : "the cloud box selected by Auto" : "this computer selected by Auto"; // resolve the mode on open; box endpoints are only ever hit on the @@ -175,6 +182,61 @@ export function ComputerPanel({ bot }: { bot: Bot }) { return; } if (bot.computer !== "cloud" && !capabilitiesReady) return; + if (cloudBackend === "vps") { + const autoLocal = bot.computer !== "cloud" && capabilitiesReady && localAvailable && computerToolSupported; + if (!vpsSupported) { + if (autoLocal) setPhase("local"); + else { + setError("This model engine cannot use a self-hosted VPS. Choose Claude or an ACP engine, or switch the cloud backend to Box."); + setPhase("error"); + } + return; + } + api(`/api/bots/${bot.id}/computer`) + .then((status) => { + if (!alive) return; + if (!status.configured) { + if (autoLocal) setPhase("local"); + else { + setError("Add the VPS SSH config alias in App Settings → Connections."); + setPhase("vps-unconfigured"); + } + return; + } + if (status.ready) { + setBoxState(status.container ?? null); + setPhase("ready"); + return; + } + if (bot.computer === "cloud") { + setPhase("starting"); + return api(`/api/bots/${bot.id}/computer/provision`, { method: "POST" }).then((result) => { + if (!alive) return; + setBoxState(result.container ?? null); + if (result.ready) setPhase("ready"); + else { + setError(result.problem ?? "The VPS Cua desktop is not ready yet"); + setPhase("error"); + } + }); + } + if (autoLocal) { + setPhase("local"); + return; + } + setBoxState(status.container ?? null); + setError(`${status.problem ?? "No ready VPS container"}. Auto will not create or start it; choose Cloud to provision it.`); + setPhase(status.container === "stopped" ? "vps-stopped" : "vps-unconfigured"); + }) + .catch((e) => { + if (!alive) return; + setError(e.message); + setPhase("error"); + }); + return () => { + alive = false; + }; + } // cloud, or auto (cloud box wins when one exists, else local in-app) api(`/api/bots/${bot.id}/computer`) .then((status) => { @@ -203,7 +265,7 @@ export function ComputerPanel({ bot }: { bot: Bot }) { return () => { alive = false; }; - }, [bot.id, bot.computer, retry, capabilitiesReady, localAvailable, vmSupported, computerToolSupported, cloudSupported]); + }, [bot.id, bot.computer, retry, capabilitiesReady, localAvailable, vmSupported, computerToolSupported, cloudSupported, cloudBackend, vpsSupported, state.config?.vps?.sshAlias]); // cloud preview: SSE frames win while the bot works; otherwise poll const live = state.screens[bot.id]; @@ -298,14 +360,25 @@ export function ComputerPanel({ bot }: { bot: Bot }) { ? cloudFrame && `data:${cloudFrame.mime};base64,${cloudFrame.png}` : null; - const run = (kind: "join" | "sleep") => { + const run = (kind: "join" | "sleep" | "provision") => { setPending(kind); setError(null); api(`/api/bots/${bot.id}/computer/${kind}`, { method: "POST" }) .then((result) => { // the join URL's stream token rotates — always freshly minted, never cached - if (kind === "join" && result.joinUrl) window.open(result.joinUrl); - if (kind === "sleep") setBoxState("archived"); + if (kind === "join" && result.joinUrl) window.open(result.joinUrl, "_blank", "noopener"); + if (kind === "provision") { + setBoxState(result.container ?? null); + if (result.ready) setPhase("ready"); + else { + setError(result.problem ?? "The VPS Cua desktop is not ready yet"); + setPhase("error"); + } + } + if (kind === "sleep") { + setBoxState(cloudBackend === "vps" ? "stopped" : "archived"); + if (cloudBackend === "vps") setPhase("vps-stopped"); + } }) .catch((e) => setError(e.message)) .finally(() => setPending(null)); @@ -316,10 +389,16 @@ export function ComputerPanel({ bot }: { bot: Bot }) { dispatch({ type: "toggleAppSettings", open: true }); }; + const openConnectionSettings = () => { + dispatch({ type: "toggleAppSettings", open: true, section: "connections" }); + }; + const emptyState = { checking: "Checking…", starting: "Starting your bot's computer…", unconfigured: "No cloud computer configured", + "vps-unconfigured": "No managed VPS computer is configured for this bot", + "vps-stopped": "The managed VPS computer is stopped", "local-unavailable": capabilities.host.platform === "linux" ? "Local computer control isn't available on Linux yet. Use a cloud box instead." @@ -387,6 +466,7 @@ export function ComputerPanel({ bot }: { bot: Bot }) { {bot.name}'s screen {phase === "local" && this computer} {phase === "vm" && Local VM} + {cloudBackend === "vps" && (phase === "ready" || phase === "starting") && self-hosted VPS}
{frameSrc ? ( @@ -427,6 +507,24 @@ export function ComputerPanel({ bot }: { bot: Bot }) { Open Local VM setup )} + {(phase === "vps-unconfigured" || phase === "vps-stopped") && ( + + )} + {phase === "vps-stopped" && bot.computer === "cloud" && ( + + )}
)} @@ -447,23 +545,38 @@ export function ComputerPanel({ bot }: { bot: Bot }) { /> )} + {phase === "vps-unconfigured" && ( +
+
+ Configure the VPS SSH alias in App Settings → Connections. Auto only reuses an existing ready container. +
+ +
+ )} {/* Cloud-only actions */} {phase === "ready" && (
- - {boxState !== "archived" && ( + {cloudBackend === "box" && ( + + )} + {(cloudBackend === "vps" || boxState !== "archived") && (
) : null} +
Self-host connected apps diff --git a/src/components/SettingsPanel.tsx b/src/components/SettingsPanel.tsx index 1c6418bf88..578f8f769e 100644 --- a/src/components/SettingsPanel.tsx +++ b/src/components/SettingsPanel.tsx @@ -8,6 +8,7 @@ import { MAUS_COLORS, MAUS_COLOR_NAMES, } from "@/lib/mascot"; +import { CloudBackendPicker } from "./CloudBackendPicker"; import { ModelPicker } from "./ModelPicker"; import { useDesktopCapabilities } from "./DesktopCapabilities"; import { cn } from "@/lib/cn"; @@ -327,6 +328,7 @@ export function SettingsPanel({ bot }: { bot: Bot }) { | "description" | "notifications" | "computer" + | "cloudBackend" | "color" | "mascotExpression" | "autoApprove" @@ -344,6 +346,7 @@ export function SettingsPanel({ bot }: { bot: Bot }) { const engine = state.instances.find((instance) => instance.instanceId === bot.modelSelection.instanceId); const canCoordinate = engine?.capabilities?.agentsMcp === true; const canUseConnectedApps = engine?.capabilities?.composioMcp === true; + const canUseVps = engine?.capabilities?.computerMcp === true && engine.driverKind !== "boxAgent"; const connectedAppsConfigured = state.config?.composio?.configured === true; const connectedAppsEnabled = bot.composio !== false; const currentChief = state.bots.find((candidate) => candidate.chiefOfStaff); @@ -645,7 +648,12 @@ export function SettingsPanel({ bot }: { bot: Bot }) { Where this bot's computer runs{bot.computer ? "" : " (currently: auto)"}
- {(["cloud", "local", "off"] as const).map((mode, i) => ( + {([ + ["cloud", "Cloud"], + ["vm", "Local VM"], + ["local", "This computer"], + ["off", "Off"], + ] as const).map(([mode, label], i) => ( ))}
+ {(!bot.computer || bot.computer === "cloud") && ( + patch({ cloudBackend: backend })} + /> + )} diff --git a/src/state/store.tsx b/src/state/store.tsx index fd2b724cd3..75ced8bf8c 100644 --- a/src/state/store.tsx +++ b/src/state/store.tsx @@ -13,7 +13,7 @@ import { useState, type ReactNode, } from "react"; -import type { EffortLevel } from "../../server/contracts.ts"; +import type { CloudBackend, EffortLevel } from "../../server/contracts.ts"; import type { MausColor, MausMotion } from "@/lib/mascot"; import type { Routine, RoutineInput, RoutineRun } from "@/lib/routines"; import type { WebhookAttempt, WebhookIngressStatus, WebhookTrigger } from "@/lib/webhooks"; @@ -153,6 +153,8 @@ export interface Bot { modelSelection: ModelSelection; /** Where this bot's computer runs; unset = auto (cloud box if one exists, else local). */ computer?: "cloud" | "vm" | "local" | "off"; + /** Which cloud computer backs `computer: "cloud"`; absent means Box. */ + cloudBackend?: CloudBackend; /** where new tasks run their shell tools; absent = the private bot workspace */ cwd?: string; /** auto mode: the bot approves its own tool permissions */ @@ -210,6 +212,7 @@ export interface ConfigStatus { xai?: { configured: boolean }; composio: { configured: boolean; mode?: "managed" | "self-hosted" | "unavailable" }; box: { configured: boolean }; + vps: { configured: boolean; sshAlias: string }; opencodeGo?: { configured: boolean }; /** Voice (ElevenLabs). `configured` = a key is saved; `ready` = a key AND * a voice, which is what it takes to actually speak. The key itself is @@ -394,6 +397,7 @@ export type Action = | "description" | "notifications" | "computer" + | "cloudBackend" | "color" | "mascotExpression" | "autoApprove" @@ -1099,6 +1103,7 @@ export function StoreProvider({ children }: { children: ReactNode }) { notifications: source.notifications, modelSelection: source.modelSelection, ...(source.computer ? { computer: source.computer } : {}), + ...(source.cloudBackend ? { cloudBackend: source.cloudBackend } : {}), }), }).then(({ bot: patched }) => rawDispatch({ type: "botAdded", bot: { ...bot, ...patched, messages: bot.messages } }), @@ -1409,6 +1414,7 @@ export function StoreProvider({ children }: { children: ReactNode }) { xai: frame.xai, composio: frame.composio, box: frame.box, + vps: frame.vps, tts: frame.tts, profile: frame.profile, },