fix(computer-use): isolate worker process entry
This commit is contained in:
@@ -386,8 +386,14 @@ function buildParams(
|
||||
if (serializedTools.length > 0) {
|
||||
params.tools = serializedTools;
|
||||
if (options?.toolChoice) {
|
||||
const toolChoice = mapToOpenAIResponsesToolChoice(options.toolChoice);
|
||||
let toolChoice = mapToOpenAIResponsesToolChoice(options.toolChoice);
|
||||
const hasComputerTool = serializedTools.some(tool => tool.type === "computer");
|
||||
if (toolChoice && typeof toolChoice !== "string" && toolChoice.type === "computer" && !hasComputerTool) {
|
||||
const computer = context.tools.find(tool => tool.native?.type === "computer");
|
||||
if (computer && serializedTools.some(tool => tool.type === "function" && tool.name === computer.name)) {
|
||||
toolChoice = { type: "function", name: computer.name };
|
||||
}
|
||||
}
|
||||
if (
|
||||
toolChoice &&
|
||||
(typeof toolChoice === "string" ||
|
||||
|
||||
@@ -208,7 +208,7 @@ describe("azure openai responses streaming", () => {
|
||||
expect(Array.isArray(tools[0].parameters.properties.item.anyOf)).toBe(true);
|
||||
});
|
||||
|
||||
it("serializes computer as a function tool on unsupported models and gates the native forced choice", async () => {
|
||||
it("serializes computer and its forced choice as a function on unsupported models", async () => {
|
||||
const computer: Tool = {
|
||||
name: "computer",
|
||||
description: "Control the desktop",
|
||||
@@ -233,7 +233,7 @@ describe("azure openai responses streaming", () => {
|
||||
expect.objectContaining({ type: "function", name: "read_file" }),
|
||||
]);
|
||||
expect(JSON.stringify(payload.tools)).not.toContain('{"type":"computer"}');
|
||||
expect(payload.tool_choice).toBeUndefined();
|
||||
expect(payload.tool_choice).toEqual({ type: "function", name: "computer" });
|
||||
});
|
||||
|
||||
it("serializes native GA computer and forced choice for a supported GPT-5.4 Azure model", async () => {
|
||||
|
||||
@@ -106,6 +106,7 @@
|
||||
"files": [
|
||||
"src",
|
||||
"dist/cli.js",
|
||||
"dist/computer-worker-process-entry.js",
|
||||
"dist/*.node",
|
||||
"scripts",
|
||||
"examples",
|
||||
|
||||
@@ -96,6 +96,7 @@ async function main(): Promise<void> {
|
||||
await compileCodingAgent({
|
||||
repoRoot,
|
||||
entrypoint: path.join(packageDir, "src", "cli.ts"),
|
||||
workerEntrypoints: [path.join(packageDir, "src", "computer-worker-process-entry.ts")],
|
||||
outfile: outputPath,
|
||||
transformersVersion,
|
||||
target: crossBuild?.target,
|
||||
|
||||
@@ -75,7 +75,13 @@ async function cleanBundleOutputs(): Promise<void> {
|
||||
}
|
||||
await Promise.all(
|
||||
entries
|
||||
.filter(entry => entry === "cli.js" || entry.endsWith(".node") || entry.endsWith(".js.map"))
|
||||
.filter(
|
||||
entry =>
|
||||
entry === "cli.js" ||
|
||||
entry === "computer-worker-process-entry.js" ||
|
||||
entry.endsWith(".node") ||
|
||||
entry.endsWith(".js.map"),
|
||||
)
|
||||
.map(entry => fs.rm(path.join(outDir, entry), { force: true })),
|
||||
);
|
||||
}
|
||||
@@ -92,7 +98,10 @@ async function main(): Promise<void> {
|
||||
// 128KiB per-argv-string cap, so it can never be passed as a CLI
|
||||
// `--define` (posix_spawn fails with E2BIG).
|
||||
const output = await Bun.build({
|
||||
entrypoints: [path.join(packageDir, "src/cli.ts")],
|
||||
entrypoints: [
|
||||
path.join(packageDir, "src/cli.ts"),
|
||||
path.join(packageDir, "src/computer-worker-process-entry.ts"),
|
||||
],
|
||||
outdir: outDir,
|
||||
target: "bun",
|
||||
external: [...ALWAYS_EXTERNAL, ...RUNTIME_EXTERNAL],
|
||||
|
||||
@@ -10,6 +10,8 @@ export interface CodingAgentCompileOptions {
|
||||
readonly repoRoot: string;
|
||||
/** Absolute CLI entrypoint. */
|
||||
readonly entrypoint: string;
|
||||
/** Additional worker modules embedded as independently evaluated process entries. */
|
||||
readonly workerEntrypoints?: readonly string[];
|
||||
/** Absolute standalone executable output path. */
|
||||
readonly outfile: string;
|
||||
/** Concrete Transformers.js version baked into the tiny-model worker. */
|
||||
@@ -33,7 +35,7 @@ export async function compileCodingAgent(options: CodingAgentCompileOptions): Pr
|
||||
}
|
||||
try {
|
||||
const output = await Bun.build({
|
||||
entrypoints: [options.entrypoint],
|
||||
entrypoints: [options.entrypoint, ...(options.workerEntrypoints ?? [])],
|
||||
root: options.repoRoot,
|
||||
external: [...COMPILED_EXTERNAL_DEPENDENCIES],
|
||||
define: {
|
||||
|
||||
@@ -31,6 +31,7 @@ import { extractProfileFlags } from "./cli/profile-bootstrap";
|
||||
import { startJsEvalProcess } from "./eval/js/process-entry";
|
||||
import type { WorkerInbound as JsWorkerInbound, WorkerOutbound as JsWorkerOutbound } from "./eval/js/worker-protocol";
|
||||
import { DAEMON_BROKER_WORKER_ARG } from "./launch/protocol";
|
||||
import { smokeTestComputerWorker } from "./tools/computer/supervisor";
|
||||
|
||||
if (Bun.semver.order(Bun.version, MIN_BUN_VERSION) < 0) {
|
||||
process.stderr.write(
|
||||
@@ -82,8 +83,6 @@ async function runSmokeTest(): Promise<void> {
|
||||
const { smokeTestTtsWorker } = await import("./tts/tts-client");
|
||||
const { smokeTestMnemopiEmbedWorker } = await import("./mnemopi/embed-client");
|
||||
const { smokeTestJsEvalWorker } = await import("./eval/js/context-manager");
|
||||
// Computer modules value-load the native desktop addon; keep them behind the explicit smoke path.
|
||||
const { smokeTestComputerWorker } = await import("./tools/computer/supervisor");
|
||||
// Other smoke dependencies stay lazy so normal CLI startup does not load their worker clients.
|
||||
const { smokeTestDaemonBroker } = await import("./launch/client");
|
||||
await smokeTestSyncWorker();
|
||||
@@ -113,7 +112,6 @@ async function runSmokeTest(): Promise<void> {
|
||||
const TINY_WORKER_ARG = "__omp_worker_tiny_inference";
|
||||
const STATS_SYNC_WORKER_ARG = "__omp_worker_stats_sync";
|
||||
const TAB_WORKER_ARG = "__omp_worker_tab";
|
||||
const COMPUTER_WORKER_ARG = "__omp_worker_computer";
|
||||
const JS_EVAL_WORKER_ARG = "__omp_worker_js_eval";
|
||||
const JS_EVAL_PROCESS_ARG = "__omp_worker_js_eval_process";
|
||||
const STT_WORKER_ARG = "__omp_worker_stt";
|
||||
@@ -157,13 +155,6 @@ async function runWorkerEntrypoint(arg: string | undefined): Promise<boolean> {
|
||||
await import("./tools/browser/tab-worker-entry");
|
||||
return true;
|
||||
}
|
||||
if (arg === COMPUTER_WORKER_ARG) {
|
||||
if (parentPort) installWorkerInbox(parentPort);
|
||||
// This selector is the lazy native-addon boundary; normal CLI startup must not evaluate the worker graph.
|
||||
const { startComputerWorker } = await import("./tools/computer/worker-entry");
|
||||
startComputerWorker();
|
||||
return true;
|
||||
}
|
||||
if (arg === JS_EVAL_WORKER_ARG) {
|
||||
if (parentPort) installWorkerInbox(parentPort);
|
||||
await import("./eval/js/worker-entry");
|
||||
|
||||
@@ -0,0 +1,8 @@
|
||||
import { parentPort } from "node:worker_threads";
|
||||
import { installWorkerInbox } from "@oh-my-pi/pi-utils/worker-host";
|
||||
import { startComputerWorker } from "./tools/computer/worker-entry";
|
||||
|
||||
if (!parentPort) throw new Error("computer-worker-process-entry: missing parentPort");
|
||||
|
||||
installWorkerInbox(parentPort);
|
||||
startComputerWorker();
|
||||
@@ -26,15 +26,32 @@ it("imports the CLI entry graph without loading dotenv before profile bootstrap"
|
||||
expect(stderr).toBe("");
|
||||
});
|
||||
|
||||
it("statically links the JS process selector through the bootstrap-safe rejection seam", async () => {
|
||||
it("starts ordinary CLI paths without evaluating the computer worker entry", async () => {
|
||||
const cliPath = path.resolve(import.meta.dir, "../../cli.ts");
|
||||
const source = await Bun.file(cliPath).text();
|
||||
const imports = new Bun.Transpiler({ loader: "tsx" }).scanImports(source.replace(/^#![^\n]*\n/, ""));
|
||||
|
||||
expect(imports).toContainEqual({
|
||||
kind: "import-statement",
|
||||
path: "@oh-my-pi/pi-utils/postmortem",
|
||||
});
|
||||
expect(imports).toContainEqual({ kind: "import-statement", path: "./eval/js/process-entry" });
|
||||
expect(imports).not.toContainEqual({ kind: "dynamic-import", path: "@oh-my-pi/pi-utils" });
|
||||
for (const args of [
|
||||
["--no-addons", cliPath, "--version"],
|
||||
[cliPath, "--help"],
|
||||
]) {
|
||||
const proc = Bun.spawn([process.execPath, ...args], {
|
||||
stdout: "pipe",
|
||||
stderr: "pipe",
|
||||
});
|
||||
const [exitCode, stderr] = await Promise.all([proc.exited, new Response(proc.stderr).text()]);
|
||||
expect(exitCode, `${args.at(-1)}: ${stderr}`).toBe(0);
|
||||
}
|
||||
});
|
||||
|
||||
it("dispatches the computer worker through its dedicated process entry", async () => {
|
||||
const fixture = path.resolve(import.meta.dir, "../../../test/fixtures/computer-worker-process-entry.ts");
|
||||
const proc = Bun.spawn([process.execPath, fixture], {
|
||||
stdout: "pipe",
|
||||
stderr: "pipe",
|
||||
});
|
||||
const [exitCode, stdout, stderr] = await Promise.all([
|
||||
proc.exited,
|
||||
new Response(proc.stdout).text(),
|
||||
new Response(proc.stderr).text(),
|
||||
]);
|
||||
expect(exitCode, stderr).toBe(0);
|
||||
expect(stdout).toBe('{"type":"pong","id":"computer-process-entry"}\n');
|
||||
});
|
||||
|
||||
@@ -1,7 +1,5 @@
|
||||
import type { DesktopAction, DesktopCapabilities, DesktopCapture, DesktopSessionOptions } from "@oh-my-pi/pi-natives";
|
||||
|
||||
export const COMPUTER_WORKER_ARG = "__omp_worker_computer";
|
||||
|
||||
export type ComputerWorkerInbound =
|
||||
| { type: "ping"; id: string }
|
||||
| { type: "init"; options: DesktopSessionOptions }
|
||||
|
||||
@@ -1,12 +1,8 @@
|
||||
import type { DesktopAction, DesktopCapabilities, DesktopCapture, DesktopSessionOptions } from "@oh-my-pi/pi-natives";
|
||||
import { logger, withTimeout, workerHostEntry } from "@oh-my-pi/pi-utils";
|
||||
import { withTimeout } from "@oh-my-pi/pi-utils/async";
|
||||
import * as logger from "@oh-my-pi/pi-utils/logger";
|
||||
import { ToolAbortError, ToolError } from "../tool-errors";
|
||||
import {
|
||||
COMPUTER_WORKER_ARG,
|
||||
type ComputerWorkerError,
|
||||
type ComputerWorkerInbound,
|
||||
type ComputerWorkerOutbound,
|
||||
} from "./protocol";
|
||||
import type { ComputerWorkerError, ComputerWorkerInbound, ComputerWorkerOutbound } from "./protocol";
|
||||
|
||||
const START_TIMEOUT_MS = 10_000;
|
||||
const CLOSE_TIMEOUT_MS = 1_500;
|
||||
@@ -66,11 +62,11 @@ function wrapWorker(worker: Worker): ComputerWorkerHandle {
|
||||
}
|
||||
|
||||
export function spawnComputerWorker(): ComputerWorkerHandle {
|
||||
const hostEntry = workerHostEntry();
|
||||
const worker = hostEntry
|
||||
? new Worker(hostEntry, { type: "module", argv: [COMPUTER_WORKER_ARG] })
|
||||
: new Worker(new URL("./worker-entry.ts", import.meta.url).href, { type: "module" });
|
||||
return wrapWorker(worker);
|
||||
const processEntry =
|
||||
process.env.PI_BUNDLED === "true" || process.env.PI_COMPILED === "true"
|
||||
? new URL("./computer-worker-process-entry.js", import.meta.url).href
|
||||
: new URL("../../computer-worker-process-entry.ts", import.meta.url).href;
|
||||
return wrapWorker(new Worker(processEntry, { type: "module" }));
|
||||
}
|
||||
|
||||
interface PendingRequest {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { parentPort } from "node:worker_threads";
|
||||
import { consumeWorkerInbox } from "@oh-my-pi/pi-utils/worker-host";
|
||||
import { COMPUTER_WORKER_ARG, type ComputerWorkerInbound, type ComputerWorkerTransport } from "./protocol";
|
||||
import type { ComputerWorkerInbound, ComputerWorkerTransport } from "./protocol";
|
||||
import { ComputerWorkerCore } from "./worker";
|
||||
|
||||
export function startComputerWorker(): void {
|
||||
@@ -25,10 +25,3 @@ export function startComputerWorker(): void {
|
||||
|
||||
new ComputerWorkerCore(transport);
|
||||
}
|
||||
|
||||
// Bun workers report `import.meta.main === false`. The source fallback still
|
||||
// enters this file directly, while the bundled CLI carries the selector and
|
||||
// starts the named entry only after installing its inbox.
|
||||
if (!Bun.isMainThread && !process.argv.includes(COMPUTER_WORKER_ARG) && import.meta.path === Bun.main) {
|
||||
startComputerWorker();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,12 @@
|
||||
const worker = new Worker(new URL("../../src/computer-worker-process-entry.ts", import.meta.url).href, {
|
||||
type: "module",
|
||||
});
|
||||
const response = Promise.withResolvers<unknown>();
|
||||
worker.addEventListener("message", event => response.resolve(event.data));
|
||||
worker.addEventListener("error", event => response.reject(event.error ?? new Error(event.message)));
|
||||
worker.postMessage({ type: "ping", id: "computer-process-entry" });
|
||||
try {
|
||||
process.stdout.write(`${JSON.stringify(await response.promise)}\n`);
|
||||
} finally {
|
||||
worker.terminate();
|
||||
}
|
||||
@@ -41,20 +41,4 @@ describe("computer worker entry", () => {
|
||||
it("is side-effect-free to import and exposes a named start function", () => {
|
||||
expect(computerWorkerEntry.startComputerWorker).toBeFunction();
|
||||
});
|
||||
|
||||
it("dispatches an immediately posted message through the CLI worker selector", async () => {
|
||||
const worker = new Worker(new URL("../src/cli.ts", import.meta.url).href, {
|
||||
type: "module",
|
||||
argv: ["__omp_worker_computer"],
|
||||
});
|
||||
const response = Promise.withResolvers<unknown>();
|
||||
worker.addEventListener("message", event => response.resolve(event.data));
|
||||
|
||||
try {
|
||||
worker.postMessage({ type: "ping", id: "computer-worker-selector" });
|
||||
expect(await response.promise).toEqual({ type: "pong", id: "computer-worker-selector" });
|
||||
} finally {
|
||||
worker.terminate();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
@@ -16,6 +16,7 @@ interface BinaryTarget {
|
||||
const repoRoot = path.join(import.meta.dir, "..");
|
||||
const binariesDir = path.join(repoRoot, "packages", "coding-agent", "binaries");
|
||||
const entrypoint = path.join(repoRoot, "packages", "coding-agent", "src", "cli.ts");
|
||||
const workerEntrypoint = path.join(repoRoot, "packages", "coding-agent", "src", "computer-worker-process-entry.ts");
|
||||
const transformersManifest: unknown = createRequire(import.meta.url)("@huggingface/transformers/package.json");
|
||||
if (
|
||||
typeof transformersManifest !== "object" ||
|
||||
@@ -26,9 +27,8 @@ if (
|
||||
throw new Error("@huggingface/transformers package manifest has no string version");
|
||||
}
|
||||
const transformersVersion = transformersManifest.version;
|
||||
// Worker threads re-enter the binary's CLI entry module. Legacy Pi host
|
||||
// modules are supplied by the in-memory compile plugin, so neither subsystem
|
||||
// needs extra `--compile` entrypoints.
|
||||
// The computer worker is an independent compiled entry so its native graph is
|
||||
// evaluated only when the supervisor selects it. Other workers re-enter the CLI.
|
||||
const isDryRun = process.argv.includes("--dry-run");
|
||||
const targets: BinaryTarget[] = [
|
||||
{
|
||||
@@ -144,6 +144,7 @@ async function buildBinary(target: BinaryTarget): Promise<void> {
|
||||
await compileCodingAgent({
|
||||
repoRoot,
|
||||
entrypoint,
|
||||
workerEntrypoints: [workerEntrypoint],
|
||||
outfile: path.join(repoRoot, target.outfile),
|
||||
transformersVersion,
|
||||
target: target.target,
|
||||
|
||||
Reference in New Issue
Block a user