返回 CodeWhale
transport.mjs
根目录 / crates / tui / plugins / computer-use / src / transport.mjs
1 // Transport: turn a registered computer into an executor.
2 // - local: the Codewhale Computer Use app when it is running or registered
3 // (it owns the OS permissions), otherwise spawn directly
4 // - ssh: run the codewhale-cu remote agent over ssh (args travel as base64 JSON,
5 // so no tool argument can ever become remote shell syntax)
6 // - hdc: HarmonyOS device over `hdc` shell / file push-pull
7 import { run, runOk, runInputLease, ExecError, currentSignal } from "./exec.mjs";
8 import { ensureApp, appSessionRequest } from "./app-socket.mjs";
9 import { sshArgv, sshDestination, sshOptions, validateSshTarget } from "./ssh-args.mjs";
10 import { spawn } from "node:child_process";
11 import crypto from "node:crypto";
12 import fs from "node:fs";
13 import path from "node:path";
14 import os from "node:os";
15 import url from "node:url";
16
17 const __dirname = path.dirname(url.fileURLToPath(import.meta.url));
18 export const PLUGIN_ROOT = path.resolve(__dirname, "..");
19 // A new MCP process always starts a fresh app/input/raster binding, even when
20 // the permission-owning desktop helper remains running across tasks.
21 export const SESSION_ID = crypto.randomUUID();
22 let usedApp = false;
23 let appSessionClosed = false;
24 export function closeAppSession({ releaseOnly = false } = {}) {
25 if (!usedApp || appSessionClosed) return Promise.resolve();
26 return appSessionRequest({ tool: releaseOnly ? "release_session_input" : "close_session", sessionId: SESSION_ID }, { timeoutMs: 2_500, signal: null }).then((reply) => {
27 if (!reply?.ok) throw Object.assign(new ExecError(reply?.error?.message ?? "Computer input cleanup failed"), { code: reply?.error?.code ?? "input_release_failed" });
28 if (!releaseOnly) appSessionClosed = true;
29 });
30 }
31
32 export function b64(obj) {
33 return Buffer.from(JSON.stringify(obj), "utf8").toString("base64");
34 }
35
36 /**
37 * Validate a remote-side filesystem path we construct ourselves.
38 * Blocks shell metacharacters and traversal outside the agent dir.
39 */
40 export function safeRemotePath(p) {
41 if (typeof p !== "string" || !/^[A-Za-z0-9.][A-Za-z0-9/._-]{0,511}$/.test(p) || p.includes("..")) {
42 throw new ExecError(`refusing unsafe remote path: ${JSON.stringify(p)}`);
43 }
44 return p;
45 }
46
47 /**
48 * Local executor bound to a platform backend name.
49 * All backends receive this shape.
50 */
51 export function localExec() {
52 return {
53 kind: "local",
54 run,
55 runOk,
56 runInputLease,
57 async readFile(p) { return fs.promises.readFile(p); },
58 async writeFile(p, data) { return fs.promises.writeFile(p, data); },
59 tmpFile(prefix) {
60 return path.join(fs.mkdtempSync(path.join(os.tmpdir(), prefix)), "out");
61 },
62 };
63 }
64
65 /**
66 * App executor: the local computer driven through the desktop app's socket.
67 * Same `remote()` contract as ssh, but files the app writes are on this disk.
68 */
69 export function appExec(app, sessionId = SESSION_ID) {
70 return {
71 ...localExec(),
72 kind: "app",
73 app,
74 filesLocal: true,
75 remote(request, opts = {}) {
76 if (sessionId === SESSION_ID && appSessionClosed) throw Object.assign(new ExecError("Computer session was closed; start a new MCP session to use the local helper again"), { code: "app_session_closed" });
77 usedApp = true;
78 return appSessionRequest({ ...request, sessionId }, { timeoutMs: opts.timeoutMs ?? 30_000 });
79 },
80 };
81 }
82
83 /**
84 * A persistent ssh agent channel: one `ssh host node agent.mjs --serve`
85 * process carrying base64-JSON request lines in and JSON receipt lines out.
86 * Unlike the one-shot agent it keeps its backend alive between calls, so an
87 * open_application binding survives into later raw-input calls and held
88 * input/recording can be owned by the session. Requests written before the
89 * channel dies may already have run remotely — their failures are marked
90 * requestDispatched so the server reports outcome_unknown instead of
91 * inviting a blind retry.
92 */
93 export function ensureSshChannel(binding, argv) {
94 let ch = binding.sshChannel;
95 if (ch?.alive) return ch;
96 const next = { alive: false, restarted: !!ch, everReplied: false, pending: new Map(), seq: 1, buf: "", proc: null };
97 binding.sshChannel = next;
98 const failAll = (err) => {
99 for (const [, p] of next.pending) { clearTimeout(p.timer); p.reject(err); }
100 next.pending.clear();
101 };
102 let proc;
103 try {
104 proc = spawn(argv[0], argv.slice(1), { stdio: ["pipe", "pipe", "pipe"], windowsHide: true });
105 } catch (err) {
106 next.spawnError = err;
107 return next;
108 }
109 next.proc = proc;
110 next.alive = true;
111 proc.stdin.on("error", () => {});
112 proc.stderr.on("data", () => {}); // drain; stderr is never parsed
113 proc.stdout.setEncoding("utf8");
114 proc.stdout.on("data", (d) => {
115 next.buf += d;
116 let i;
117 while ((i = next.buf.indexOf("\n")) !== -1) {
118 const line = next.buf.slice(0, i).trim();
119 next.buf = next.buf.slice(i + 1);
120 if (!line.startsWith("{")) continue; // MOTD/banner noise
121 let msg;
122 try { msg = JSON.parse(line); } catch { continue; }
123 next.everReplied = true;
124 const p = next.pending.get(msg.id);
125 if (!p) continue; // timed-out or unknown request: drop the late reply
126 next.pending.delete(msg.id);
127 clearTimeout(p.timer);
128 p.resolve(msg);
129 }
130 });
131 const dead = (why) => {
132 if (!next.alive) return;
133 next.alive = false;
134 failAll(Object.assign(new ExecError(`ssh agent channel closed${why ? `: ${why}` : ""}`), { code: "remote_session_lost", requestDispatched: true }));
135 };
136 proc.on("error", (err) => dead(String(err?.message ?? err)));
137 proc.on("close", (code, sig) => dead(code != null ? `exited ${code}` : `signal ${sig}`));
138 return next;
139 }
140
141 export function closeSshChannel(binding) {
142 const ch = binding.sshChannel;
143 if (!ch) return;
144 binding.sshChannel = null;
145 ch.alive = false;
146 try { ch.proc?.stdin.end(); } catch {}
147 try { ch.proc?.kill("SIGTERM"); } catch {}
148 for (const [, p] of ch.pending ?? []) {
149 clearTimeout(p.timer);
150 p.reject(Object.assign(new ExecError("ssh agent channel closed"), { code: "remote_session_lost", requestDispatched: true }));
151 }
152 ch.pending?.clear();
153 }
154
155 export function channelRequest(ch, request, timeoutMs) {
156 return new Promise((resolve, reject) => {
157 if (!ch.alive) {
158 reject(Object.assign(new ExecError("ssh agent channel is closed"), { code: "remote_session_lost" }));
159 return;
160 }
161 const id = ch.seq++;
162 const timer = setTimeout(() => {
163 ch.pending.delete(id);
164 // The request was written; the remote may still be executing it.
165 reject(Object.assign(new ExecError(`ssh agent timed out after ${timeoutMs}ms`), { code: "remote_timeout", requestDispatched: true }));
166 }, timeoutMs);
167 ch.pending.set(id, { resolve, reject, timer });
168 ch.proc.stdin.write(b64({ id, tool: request.tool, args: request.args ?? {} }) + "\n");
169 });
170 }
171
172 /**
173 * Shared persistent-channel front for executors whose requests ride one
174 * long-lived `<argv> --serve` process (ssh agent, docker exec). Read-only and
175 * identity requests may run on a restarted channel; input tools may not —
176 * the fresh remote agent no longer holds this session's open_application
177 * binding.
178 */
179 const SAFE_AFTER_RESTART = new Set([
180 "platform", "probe", "list_displays", "switch_display", "list_apps", "list_windows",
181 "get_app_state", "resolve_element", "screenshot", "zoom", "cursor_position",
182 "read_clipboard", "recordingList", "recordingStatus", "open_application", "preview",
183 ]);
184 function attachPersistentChannel(ex, binding, serveArgv) {
185 if (!binding) return;
186 ex.persistent = (request, opts = {}) => {
187 const ch = ensureSshChannel(binding, serveArgv);
188 if (ch.spawnError) {
189 return Promise.reject(Object.assign(new ExecError(`remote agent channel failed to start: ${ch.spawnError.message}`), { code: "remote_session_lost" }));
190 }
191 if (ch.restarted) {
192 ch.restarted = false;
193 binding.needsObservation = true;
194 if (!SAFE_AFTER_RESTART.has(request.tool)) {
195 return Promise.reject(Object.assign(new ExecError("the remote agent session restarted — rebind with open_application and observe before acting"), { code: "remote_session_restarted" }));
196 }
197 }
198 return channelRequest(ch, request, opts.timeoutMs ?? 25_000);
199 };
200 ex.closeChannel = () => closeSshChannel(binding);
201 }
202
203 /** ssh executor: speaks to the remote agent installed by installRemoteAgent(). */
204 export function sshExec(computer, binding) {
205 const userHost = sshDestination(computer);
206 const remoteAgent = safeRemotePath(computer.agentPath ?? ".codewhale-cu/agent/agent.mjs");
207 const base = sshArgv(computer);
208 const ex = {
209 kind: "ssh",
210 base,
211 userHost,
212 remoteAgent,
213 run(cmd, args = [], opts = {}) {
214 // Local side commands (e.g. ssh itself) run directly.
215 return run(cmd, args, opts);
216 },
217 async remote(request, opts = {}) {
218 const r = await run("ssh", [...base, "node", remoteAgent, b64({ args: request.args ?? {}, tool: request.tool, nonce: crypto.randomBytes(6).toString("hex") })], {
219 timeoutMs: opts.timeoutMs ?? 25_000,
220 });
221 // A spawned one-shot agent may have acted before it was cancelled or
222 // timed out; say so rather than inviting a blind retry.
223 if (r.aborted) throw Object.assign(new ExecError("computer request cancelled", r), { code: "cancelled" }, r.spawned ? { requestDispatched: true } : {});
224 if (r.timedOut) throw Object.assign(new ExecError(`ssh ${userHost}: timed out`, r), { code: "remote_timeout", requestDispatched: true });
225 if (r.code !== 0) throw new ExecError(`ssh ${userHost} exited ${r.code}: ${r.stderr.trim().slice(0, 400)}`, r);
226 // The agent prints exactly one JSON line; anything before it is MOTD noise.
227 const line = r.stdout.trim().split("\n").filter((l) => l.startsWith("{")).pop();
228 const reply = line ? JSON.parse(line) : null;
229 if (!reply) throw new ExecError(`ssh ${userHost}: agent returned no JSON receipt`, r);
230 return reply;
231 },
232 };
233 attachPersistentChannel(ex, binding, ["ssh", ...base, "node", remoteAgent, "--serve"]);
234 return ex;
235 }
236
237 /**
238 * docker executor: a spawned task-owned desktop container. Same agent contract
239 * as ssh, but the channel is `docker exec` — no sshd, no keys, the container
240 * boundary itself is the isolation. Every call goes through
241 * docker/agent-exec.sh, which joins the desktop session env (display + bus)
242 * the container entrypoint recorded before serving requests.
243 */
244 export function dockerExec(computer, binding) {
245 const container = safeRemotePath(computer.container);
246 const remoteAgent = "/app/docker/agent-exec.sh";
247 const ex = {
248 kind: "docker",
249 container,
250 remoteAgent,
251 run(cmd, args = [], opts = {}) {
252 // Local side commands (docker itself) run directly.
253 return run(cmd, args, opts);
254 },
255 async remote(request, opts = {}) {
256 const r = await run("docker", ["exec", container, "/bin/sh", remoteAgent, b64({ args: request.args ?? {}, tool: request.tool, nonce: crypto.randomBytes(6).toString("hex") })], {
257 timeoutMs: opts.timeoutMs ?? 25_000,
258 });
259 if (r.aborted) throw Object.assign(new ExecError("computer request cancelled", r), { code: "cancelled" }, r.spawned ? { requestDispatched: true } : {});
260 if (r.timedOut) throw Object.assign(new ExecError(`docker exec ${container}: timed out`, r), { code: "remote_timeout", requestDispatched: true });
261 if (r.code !== 0) throw new ExecError(`docker exec ${container} exited ${r.code}: ${r.stderr.trim().slice(0, 400)}`, r);
262 const line = r.stdout.trim().split("\n").filter((l) => l.startsWith("{")).pop();
263 const reply = line ? JSON.parse(line) : null;
264 if (!reply) throw new ExecError(`docker exec ${container}: agent returned no JSON receipt`, r);
265 return reply;
266 },
267 };
268 attachPersistentChannel(ex, binding, ["docker", "exec", "-i", container, "/bin/sh", remoteAgent, "--serve"]);
269 return ex;
270 }
271
272 /** Push the self-contained remote agent + src tree to an ssh computer. */
273 export async function installRemoteAgent(computer) {
274 const ex = sshExec(computer);
275 const srcDir = path.join(PLUGIN_ROOT, "src");
276 const rels = ["agent.mjs"];
277 for (const dir of ["", "backends"]) {
278 const full = path.join(srcDir, dir);
279 for (const f of fs.readdirSync(full)) {
280 if (f.endsWith(".mjs") || f.endsWith(".m") || f.endsWith(".h")) rels.push(`src/${dir ? dir + "/" : ""}${f}`);
281 }
282 }
283 const marker = ".codewhale-cu/agent";
284 let r = await run("ssh", [...ex.base, "mkdir", "-p", `${marker}/src/backends`], { timeoutMs: 15_000 });
285 if (r.code !== 0) throw new ExecError(`ssh ${ex.userHost}: mkdir failed: ${r.stderr.trim().slice(0, 300)}`, r);
286 for (const rel of rels) {
287 const localPath = rel === "agent.mjs" ? path.join(PLUGIN_ROOT, "agent.mjs") : path.join(srcDir, rel.slice(4));
288 const dest = safeRemotePath(`${marker}/${rel}`);
289 validateSshTarget(computer);
290 r = await run("scp", [...sshOptions(computer, { portFlag: "-P" }), "--", localPath, `${ex.userHost}:${dest}`], { timeoutMs: 30_000 });
291 if (r.code !== 0) throw new ExecError(`scp ${rel} failed: ${r.stderr.trim().slice(0, 300)}`, r);
292 }
293 // Probe remote platform via the agent itself.
294 const reply = await ex.remote({ tool: "platform" });
295 return { installed: rels.length, remotePlatform: reply.platform, agentPath: `${marker}/agent.mjs` };
296 }
297
298 /** hdc (HarmonyOS) executor. Commands run on-device; files pull to local tmp. */
299 export function hdcExec(computer) {
300 const targetArgs = computer.target ? ["-t", computer.target] : [];
301 const shell = (args, opts = {}) => run("hdc", [...targetArgs, "shell", ...args], opts);
302 return {
303 kind: "hdc",
304 targetArgs,
305 run,
306 runOk,
307 shell,
308 async pullFile(remotePath, localPath, opts = {}) {
309 // HDC device captures use absolute paths; SSH agent paths are relative.
310 // Validate the remaining path with the same traversal/metacharacter guard.
311 safeRemotePath(typeof remotePath === "string" ? remotePath.replace(/^\//, "") : remotePath);
312 const r = await run("hdc", [...targetArgs, "file", "recv", remotePath, localPath], opts);
313 if (r.code !== 0) throw new ExecError(`hdc file recv failed: ${r.stderr.trim().slice(0, 300)}`, r);
314 return localPath;
315 },
316 async readFile(remotePath, opts = {}) {
317 // Containment: pull into a private mkdtemp dir and remove exactly that
318 // dir. Never rm() the parent of a file placed directly in os.tmpdir() —
319 // that recursively deletes the entire user temp directory.
320 const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "cu-hdc-"));
321 try {
322 const tmp = path.join(dir, "out");
323 await this.pullFile(remotePath, tmp, opts);
324 return await fs.promises.readFile(tmp);
325 } finally {
326 // Cleanup must not replace downloaded bytes or the original I/O error.
327 await fs.promises.rm(dir, { recursive: true, force: true }).catch(() => {});
328 }
329 },
330 };
331 }
332
333 export async function executorFor(computer, binding) {
334 if (computer.transport === "local") {
335 // Test hook: exercise the out-of-process wire path (desktop app / ssh
336 // agent) in-process, so wire argument preparation is covered by tests.
337 if (process.env.CODEWHALE_CU_TEST_REMOTE === "1") {
338 const { handle } = await import("./app-handler.mjs");
339 return { ...appExec({ id: "test", name: "test app" }), remote: (request) => handle(request, { sessionId: SESSION_ID, signal: currentSignal() }) };
340 }
341 const status = await ensureApp();
342 if (status.via === "app") {
343 if (status.app.sessionProtocol !== 2) throw Object.assign(new ExecError("The installed Computer Use helper needs an update for isolated sessions and disconnect cleanup. Rebuild/reinstall it, then retry."), { code: "app_upgrade_required" });
344 if (process.platform === "darwin" && status.app.backgroundProtocol !== 1) throw Object.assign(new ExecError("The installed Computer Use app predates background scrolling, scoped observations and foreground preemption. Update and restart the helper before using it."), { code: "app_upgrade_required" });
345 return appExec(status.app);
346 }
347 return { ...localExec(), appReason: status.reason };
348 }
349 if (computer.transport === "ssh") return sshExec(computer, binding);
350 if (computer.transport === "docker") return dockerExec(computer, binding);
351 if (computer.transport === "hdc") return hdcExec(computer);
352 throw new ExecError(`unknown transport ${computer.transport}`);
353 }
354
355 /**
356 * Map a computer to its backend module. Local platform is fixed; ssh
357 * computers may carry platformHint (probed at registration).
358 */
359 function effectivePlatform(computer) {
360 let platform = computer.platform ?? computer.platformHint;
361 if (!platform) {
362 if (computer.transport === "local") platform = process.platform;
363 else if (computer.transport === "hdc") platform = "harmonyos";
364 else if (computer.transport === "docker") platform = "linux"; // spawned containers are always the Linux desktop image
365 else platform = "linux"; // conservative default for ssh; registration probes it
366 }
367 return platform;
368 }
369
370 /** Identity of the effective route, excluding catalog presentation metadata. */
371 export function routeFingerprint(computer) {
372 const route = [computer.transport, effectivePlatform(computer)];
373 if (computer.transport === "docker") route.push(computer.container || null);
374 if (computer.transport === "hdc") route.push(computer.target || null);
375 if (computer.transport === "ssh") route.push(computer.host, computer.user || null,
376 computer.port || null, computer.knownHosts || null, computer.agentPath ?? ".codewhale-cu/agent/agent.mjs");
377 return JSON.stringify(route);
378 }
379
380 export async function backendFor(computer) {
381 const platform = effectivePlatform(computer);
382 // Test hook: inject a fake local backend by absolute path to an .mjs
383 // module exporting `create` (used by tests/, never set in production).
384 const testBackend = computer.transport === "local" && process.env.CODEWHALE_CU_TEST_BACKEND;
385 const mod = await import(testBackend ? url.pathToFileURL(testBackend).href : `./backends/${platform}.mjs`);
386 // Backends always get a direct executor; routing through the app happens
387 // one level up (the server dispatches to `executor.remote` when present).
388 const exec = computer.transport === "local" ? localExec() : await executorFor(computer);
389 return { backend: mod.create({ exec, computer, platform }), platform };
390 }
391
391 lines Plain Text