From 3ab675a83bca909c827399f876fd779fa35134ab Mon Sep 17 00:00:00 2001 From: can1357 Date: Mon, 15 Jun 2026 11:59:44 +0200 Subject: [PATCH] fix: added buffered worker inboxing and standardized worker selectors - Added `WorkerInbox` and `installWorkerInbox(port)` to queue worker messages before bind. - Added `consumeWorkerInbox()` to replay buffered messages and clear one active inbox. - Added buffered inbox consumption in JS and tab worker transports before direct message handlers. - Normalized worker selector arguments to the `__omp_worker_*` naming across workers and tests. --- AGENTS.md | 4 +- packages/coding-agent/CHANGELOG.md | 2 + packages/coding-agent/src/cli.ts | 37 ++++++---- .../src/eval/js/context-manager.ts | 23 +++++- .../coding-agent/src/eval/js/worker-entry.ts | 7 ++ packages/coding-agent/src/stt/asr-client.ts | 2 +- .../coding-agent/src/tiny/title-client.ts | 2 +- .../src/tools/browser/tab-supervisor.ts | 2 +- .../src/tools/browser/tab-worker-entry.ts | 16 +++-- packages/coding-agent/src/tts/tts-client.ts | 2 +- .../test/issue-1011-repro.test.ts | 2 +- .../test/issue-1606-repro.test.ts | 2 +- packages/stats/CHANGELOG.md | 3 + packages/stats/src/aggregator.ts | 2 +- packages/utils/CHANGELOG.md | 4 ++ packages/utils/src/worker-host.ts | 71 +++++++++++++++++++ packages/utils/test/worker-host.test.ts | 71 +++++++++++++++++++ 17 files changed, 226 insertions(+), 26 deletions(-) create mode 100644 packages/utils/test/worker-host.test.ts diff --git a/AGENTS.md b/AGENTS.md index bb4a24039..24cc2112d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -32,12 +32,12 @@ This repo contains multiple packages, but **`packages/coding-agent/`** is the pr - **Class privacy**: use ES `#private` fields; leave externally accessible members bare. **No `private`/`protected`/`public` keyword on fields or methods**, except on **constructor parameter properties** where TypeScript requires it (e.g. `constructor(private readonly session: ToolSession)`). - **Promises**: use `Promise.withResolvers()` instead of `new Promise((resolve, reject) => ...)`. - **Prompts**: never build prompts in code (no inline strings, template literals, or concatenation). Prompts live in static `.md` files; use Handlebars for dynamic content. Import them via `import content from "./prompt.md" with { type: "text" }` — not `readFile`. -- **Worker scripts**: workers re-enter the CLI entrypoint; never spawn separate worker entry modules. `cli.ts` declares itself as the worker host at startup (`declareWorkerHostEntry()` from `@oh-my-pi/pi-utils/env`) and dispatches hidden argv selectors (`__omp_stats_sync_worker`, `__omp_tab_worker`, `__omp_js_eval_worker`, `--tiny-worker`) before loading the command registry. Spawn sites use: +- **Worker scripts**: workers re-enter the CLI entrypoint; never spawn separate worker entry modules. `cli.ts` declares itself as the worker host at startup (`declareWorkerHostEntry()` from `@oh-my-pi/pi-utils/env`) and dispatches hidden argv selectors (`__omp_worker_stats_sync`, `__omp_worker_tab`, `__omp_worker_js_eval`, `__omp_worker_tiny_inference`) before loading the command registry. Spawn sites use: ```ts import { workerHostEntry } from "@oh-my-pi/pi-utils"; const hostEntry = workerHostEntry(); const worker = hostEntry - ? new Worker(hostEntry, { type: "module", argv: ["__omp__worker"] }) + ? new Worker(hostEntry, { type: "module", argv: ["__omp_worker_"] }) : new Worker(new URL("./.ts", import.meta.url).href, { type: "module" }); ``` When the process was started from the omp CLI — source `cli.ts`, npm-bundle `dist/cli.js`, or compiled binary — `workerHostEntry()` is `Bun.main` and the worker re-enters the single entry module, so no per-worker `--compile` entrypoints or bundle entries exist. Outside a CLI host (`bun test`, SDK embedding, standalone `omp-stats`) it returns `null` and the direct-module fallback loads the worker source. New worker kinds MUST add their selector to the dispatch table in `cli.ts` and keep the fallback branch. diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 40e51f49c..992c6b6ee 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -11,10 +11,12 @@ - Changed the `job` poll to return early when a steering message is queued, draining the steer immediately instead of waiting out the poll window. - Capped unexpected-stop auto-continuation to three retry attempts before giving up on repeated stops - Updated the `edit` tool's hashline prompt, grammar, and docs to recommend the `.=` inclusive range separator (`SWAP 1.=3:`); the legacy `..` form still parses. +- Normalized all internal worker argv selectors under the `__omp_worker_` prefix, skipping the async worker dispatch check during normal CLI startup. ### Fixed - Fixed ModelRegistry tests making outbound network calls by automatically stubbing fetch during test execution. +- Fixed `eval` JS cells (and browser-tab worker startup) always stalling for the full init timeout — typically the cell's whole 30s budget — before silently falling back to the slower inline worker. The self-dispatching CLI host imports the worker module dynamically from its argv dispatch, so the worker's own `parentPort.on("message")` attached only after Bun flushed the messages the parent posted before spawn; the synchronously-posted `init` handshake was dropped and never answered with `ready`. The host now installs a buffering `parentPort` inbox synchronously in the entry's sync prefix (before importing the worker module) and the worker binds it on load, replaying the buffered handshake. `omp --smoke-test` now also spawns the JS eval worker through the host entry and asserts it handshakes on a real worker thread. ## [15.13.2] - 2026-06-15 diff --git a/packages/coding-agent/src/cli.ts b/packages/coding-agent/src/cli.ts index 417f8cd66..9bee88073 100755 --- a/packages/coding-agent/src/cli.ts +++ b/packages/coding-agent/src/cli.ts @@ -14,6 +14,7 @@ try { * CLI entry point — registers all commands explicitly and delegates to the * lightweight CLI runner from pi-utils. */ +import { parentPort } from "node:worker_threads"; import type { CliConfig } from "@oh-my-pi/pi-utils/cli"; import { APP_NAME, @@ -23,7 +24,7 @@ import { setProfile, VERSION, } from "@oh-my-pi/pi-utils/dirs"; -import { declareWorkerHostEntry } from "@oh-my-pi/pi-utils/worker-host"; +import { declareWorkerHostEntry, installWorkerInbox } from "@oh-my-pi/pi-utils/worker-host"; import { installProfileAlias, resolveProfileAliasCommandFromProcess } from "./cli/profile-alias"; import { extractProfileFlags } from "./cli/profile-bootstrap"; @@ -67,6 +68,7 @@ async function runSmokeTest(): Promise { const { smokeTestTinyTitleWorker } = await import("./tiny/title-client"); const { smokeTestSttWorker } = await import("./stt/asr-client"); const { smokeTestTtsWorker } = await import("./tts/tts-client"); + const { smokeTestJsEvalWorker } = await import("./eval/js/context-manager"); await smokeTestSyncWorker(); const statsServer = await startServer(0); @@ -83,18 +85,23 @@ async function runSmokeTest(): Promise { await smokeTestTinyTitleWorker(); await smokeTestSttWorker(); + await smokeTestJsEvalWorker(); await smokeTestTtsWorker(); process.stdout.write("smoke-test: ok\n"); } -const TINY_WORKER_ARGS = new Set(["--tiny-worker", "__tiny_worker"]); -const STATS_SYNC_WORKER_ARG = "__omp_stats_sync_worker"; -const TAB_WORKER_ARG = "__omp_tab_worker"; -const JS_EVAL_WORKER_ARG = "__omp_js_eval_worker"; -const STT_WORKER_ARG = "__omp_stt_worker"; -const TTS_WORKER_ARG = "__omp_tts_worker"; +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 JS_EVAL_WORKER_ARG = "__omp_worker_js_eval"; +const STT_WORKER_ARG = "__omp_worker_stt"; +const TTS_WORKER_ARG = "__omp_worker_tts"; async function runWorkerEntrypoint(arg: string | undefined): Promise { + if (arg === TINY_WORKER_ARG) { + await runTinyWorker(); + return true; + } if (arg === STATS_SYNC_WORKER_ARG) { // The sync worker handles messages via `self.onmessage`, assigned during // this *async* dynamic import. Bun flushes the worker's initial message @@ -117,11 +124,20 @@ async function runWorkerEntrypoint(arg: string | undefined): Promise { } return true; } + // Bun flushes messages the parent posted before spawn once this entry's + // top-level evaluation completes, delivering them only to listeners present + // at that moment. These worker modules are imported dynamically below, so + // their own `parentPort.on("message")` lands after the flush and the parent's + // synchronous `init` is dropped. Install a buffering inbox synchronously here + // (still inside the entry's sync prefix) so the handshake survives; the worker + // module binds the real handler once loaded. if (arg === TAB_WORKER_ARG) { + if (parentPort) installWorkerInbox(parentPort); await import("./tools/browser/tab-worker-entry"); return true; } if (arg === JS_EVAL_WORKER_ARG) { + if (parentPort) installWorkerInbox(parentPort); await import("./eval/js/worker-entry"); return true; } @@ -251,11 +267,8 @@ export async function runCli(argv: string[]): Promise { // synchronous prefix of `runWorkerEntrypoint`, and Bun flushes the // worker's parked initial messages as soon as the entry module's // top-level evaluation finishes. - if (TINY_WORKER_ARGS.has(resolvedArgv[0] ?? "")) { - await runTinyWorker(); - return; - } - if (await runWorkerEntrypoint(resolvedArgv[0])) { + if (resolvedArgv[0]?.startsWith("__omp_worker_")) { + await runWorkerEntrypoint(resolvedArgv[0]); return; } diff --git a/packages/coding-agent/src/eval/js/context-manager.ts b/packages/coding-agent/src/eval/js/context-manager.ts index 8975beb70..0cabe4e54 100644 --- a/packages/coding-agent/src/eval/js/context-manager.ts +++ b/packages/coding-agent/src/eval/js/context-manager.ts @@ -141,6 +141,27 @@ export async function disposeAllVmContexts(): Promise { await Promise.all(all.map(session => killSession(session, new ToolError("JS context disposed"), { force: false }))); } +/** + * Smoke probe: spawn the JS eval worker through the worker-host entry and prove + * it answers the `init` handshake on a real worker thread (not the inline + * fallback). Catches the silent worker-load and init-message-drop regressions + * that otherwise strand every cell on the init timeout in a distribution build — + * the failure mode that motivated `installWorkerInbox`. Wired into + * `omp --smoke-test` so binary / source / tarball installs all exercise it. + */ +export async function smokeTestJsEvalWorker(): Promise { + const worker = spawnJsWorker(); + const session: JsSession = { sessionKey: "smoke", worker, state: "alive", pending: new Map() }; + try { + await initWorker(session, { cwd: process.cwd(), sessionId: "smoke" }, WORKER_INIT_TIMEOUT_MS); + if (worker.mode !== "worker") { + throw new Error("JS eval worker smoke fell back to the inline worker (real worker failed to start)"); + } + } finally { + await worker.terminate().catch(() => undefined); + } +} + async function runOnce( session: JsSession, options: { @@ -447,7 +468,7 @@ function spawnJsWorker(): WorkerHandle { try { const hostEntry = workerHostEntry(); const worker = hostEntry - ? new Worker(hostEntry, { type: "module", argv: ["__omp_js_eval_worker"] }) + ? new Worker(hostEntry, { type: "module", argv: ["__omp_worker_js_eval"] }) : new Worker(new URL("./worker-entry.ts", import.meta.url).href, { type: "module" }); return wrapBunWorker(worker); } catch (err) { diff --git a/packages/coding-agent/src/eval/js/worker-entry.ts b/packages/coding-agent/src/eval/js/worker-entry.ts index 069da30f0..dad0f7d4a 100644 --- a/packages/coding-agent/src/eval/js/worker-entry.ts +++ b/packages/coding-agent/src/eval/js/worker-entry.ts @@ -1,13 +1,20 @@ import { parentPort } from "node:worker_threads"; +import { consumeWorkerInbox } from "@oh-my-pi/pi-utils/worker-host"; import { WorkerCore } from "./worker-core"; import type { Transport, WorkerInbound, WorkerOutbound } from "./worker-protocol"; if (!parentPort) throw new Error("js worker-entry: missing parentPort"); const port = parentPort; +// When the CLI host pre-buffered messages (it imports this module dynamically), +// bind that inbox so the parent's already-delivered `init` is replayed. Loaded +// directly (test/SDK fallback), this module's top-level runs synchronously at +// worker start, so the direct `parentPort.on` below wins the flush on its own. +const inbox = consumeWorkerInbox(); const transport: Transport = { send: (msg: WorkerOutbound) => port.postMessage(msg), onMessage: handler => { + if (inbox) return inbox.bind(data => handler(data as WorkerInbound)); const wrap = (data: unknown): void => handler(data as WorkerInbound); port.on("message", wrap); return () => port.off("message", wrap); diff --git a/packages/coding-agent/src/stt/asr-client.ts b/packages/coding-agent/src/stt/asr-client.ts index 813f81e29..bef4f739e 100644 --- a/packages/coding-agent/src/stt/asr-client.ts +++ b/packages/coding-agent/src/stt/asr-client.ts @@ -72,7 +72,7 @@ const SMOKE_TEST_TIMEOUT_MS = 30_000; * Hidden subcommand on the main CLI that boots the speech-recognition worker in * the spawned subprocess. Kept in sync with the dispatch in `cli.ts`. */ -export const STT_WORKER_ARG = "__omp_stt_worker"; +export const STT_WORKER_ARG = "__omp_worker_stt"; function readTinyModelSetting(key: "providers.tinyModelDevice" | "providers.tinyModelDtype"): string | undefined { try { diff --git a/packages/coding-agent/src/tiny/title-client.ts b/packages/coding-agent/src/tiny/title-client.ts index cd7dfee58..8edf97f70 100644 --- a/packages/coding-agent/src/tiny/title-client.ts +++ b/packages/coding-agent/src/tiny/title-client.ts @@ -69,7 +69,7 @@ function normalizeTinyTitleGenerateOptions( * Hidden subcommand on the main CLI that boots the tiny-model worker in the * spawned subprocess. Kept in sync with the dispatch in `cli.ts`. */ -export const TINY_WORKER_ARG = "--tiny-worker"; +export const TINY_WORKER_ARG = "__omp_worker_tiny_inference"; function readTinyModelSetting(path: "providers.tinyModelDevice" | "providers.tinyModelDtype"): string | undefined { try { diff --git a/packages/coding-agent/src/tools/browser/tab-supervisor.ts b/packages/coding-agent/src/tools/browser/tab-supervisor.ts index e9ac5c4f8..89e06917e 100644 --- a/packages/coding-agent/src/tools/browser/tab-supervisor.ts +++ b/packages/coding-agent/src/tools/browser/tab-supervisor.ts @@ -685,7 +685,7 @@ async function spawnTabWorker(): Promise { try { const hostEntry = workerHostEntry(); const worker = hostEntry - ? new Worker(hostEntry, { type: "module", argv: ["__omp_tab_worker"] }) + ? new Worker(hostEntry, { type: "module", argv: ["__omp_worker_tab"] }) : new Worker(new URL("./tab-worker-entry.ts", import.meta.url).href, { type: "module" }); return wrapBunWorker(worker); } catch (err) { diff --git a/packages/coding-agent/src/tools/browser/tab-worker-entry.ts b/packages/coding-agent/src/tools/browser/tab-worker-entry.ts index b8d2ee303..38c0d1c7d 100644 --- a/packages/coding-agent/src/tools/browser/tab-worker-entry.ts +++ b/packages/coding-agent/src/tools/browser/tab-worker-entry.ts @@ -1,20 +1,28 @@ import { parentPort } from "node:worker_threads"; +import { consumeWorkerInbox } from "@oh-my-pi/pi-utils/worker-host"; import type { Transport, WorkerInbound, WorkerOutbound } from "./tab-protocol"; import { WorkerCore } from "./tab-worker"; if (!parentPort) throw new Error("tab-worker-entry: missing parentPort"); +const port = parentPort; +// When the CLI host pre-buffered messages (it imports this module dynamically), +// bind that inbox so the parent's already-delivered `init` is replayed. Loaded +// directly (test/SDK fallback), this module's top-level runs synchronously at +// worker start, so the direct `parentPort.on` below wins the flush on its own. +const inbox = consumeWorkerInbox(); const transport: Transport = { send(msg, transferList) { - parentPort!.postMessage(msg, transferList ?? []); + port.postMessage(msg, transferList ?? []); }, onMessage(handler) { + if (inbox) return inbox.bind(data => handler(data as WorkerOutbound | WorkerInbound)); const wrap = (message: unknown): void => handler(message as WorkerOutbound | WorkerInbound); - parentPort!.on("message", wrap); - return () => parentPort!.off("message", wrap); + port.on("message", wrap); + return () => port.off("message", wrap); }, close() { - parentPort!.close(); + port.close(); }, }; diff --git a/packages/coding-agent/src/tts/tts-client.ts b/packages/coding-agent/src/tts/tts-client.ts index 7878504fb..909ab96f4 100644 --- a/packages/coding-agent/src/tts/tts-client.ts +++ b/packages/coding-agent/src/tts/tts-client.ts @@ -142,7 +142,7 @@ const SMOKE_TEST_TIMEOUT_MS = 30_000; * Hidden subcommand on the main CLI that boots the TTS worker in the spawned * subprocess. Kept in sync with the dispatch in `cli.ts` (Main-owned). */ -export const TTS_WORKER_ARG = "__omp_tts_worker"; +export const TTS_WORKER_ARG = "__omp_worker_tts"; function readTinyModelSetting(path: "providers.tinyModelDevice" | "providers.tinyModelDtype"): string | undefined { try { diff --git a/packages/coding-agent/test/issue-1011-repro.test.ts b/packages/coding-agent/test/issue-1011-repro.test.ts index 7f11f3ed4..865ffd9f0 100644 --- a/packages/coding-agent/test/issue-1011-repro.test.ts +++ b/packages/coding-agent/test/issue-1011-repro.test.ts @@ -27,7 +27,7 @@ describe("issue #1011 — tab worker must re-enter the CLI entrypoint", () => { const packageDir = path.resolve(import.meta.dir, ".."); const supervisorPath = path.join(packageDir, "src/tools/browser/tab-supervisor.ts"); const buildBinaryPath = path.join(packageDir, "scripts/build-binary.ts"); - const workerArg = "__omp_tab_worker"; + const workerArg = "__omp_worker_tab"; it("tab-supervisor re-enters the worker-host entry with the argv selector", async () => { const source = await Bun.file(supervisorPath).text(); diff --git a/packages/coding-agent/test/issue-1606-repro.test.ts b/packages/coding-agent/test/issue-1606-repro.test.ts index 9a0ea3d00..fa57764b3 100644 --- a/packages/coding-agent/test/issue-1606-repro.test.ts +++ b/packages/coding-agent/test/issue-1606-repro.test.ts @@ -8,7 +8,7 @@ * in the parent's address space and crashed the CLI on exit. * * The fix relocates the worker to a child process: `title-client.ts` spawns - * `process.execPath … --tiny-worker`, `cli.ts` dispatches that flag into + * `process.execPath … __omp_tiny_inference`, `cli.ts` dispatches that flag into * `runTinyWorker`, and the parent `SIGKILL`s the child on dispose so the * native finalizer never runs in either address space. These tests pin the * three pieces of that contract so a future refactor cannot quietly land diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index 586c8045a..c38a82477 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -1,6 +1,9 @@ # Changelog ## [Unreleased] +### Changed + +- Renamed `__omp_stats_sync_worker` to `__omp_worker_stats_sync`. ## [15.13.1] - 2026-06-15 diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index af697ed68..4021238db 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -93,7 +93,7 @@ interface WorkerHandle { function createSyncWorker(): Worker { const hostEntry = workerHostEntry(); if (hostEntry) { - return new Worker(hostEntry, { type: "module", argv: ["__omp_stats_sync_worker"] }); + return new Worker(hostEntry, { type: "module", argv: ["__omp_worker_stats_sync"] }); } return new Worker(new URL("./sync-worker.ts", import.meta.url).href, { type: "module" }); } diff --git a/packages/utils/CHANGELOG.md b/packages/utils/CHANGELOG.md index f4df82d13..d53a9362d 100644 --- a/packages/utils/CHANGELOG.md +++ b/packages/utils/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Added + +- Added `installWorkerInbox(port)` / `consumeWorkerInbox()` to `@oh-my-pi/pi-utils/worker-host`. A self-dispatching CLI host that imports a Bun worker module dynamically attaches the worker's real `message` listener after Bun flushes the messages the parent posted before spawn, dropping a synchronously-posted `init`. The host installs this buffering inbox synchronously in the entry's sync prefix so a listener exists at flush time; the worker module consumes it and binds the real handler, replaying anything buffered. + ## [15.13.1] - 2026-06-15 ### Added diff --git a/packages/utils/src/worker-host.ts b/packages/utils/src/worker-host.ts index 13eb39d1a..50cfec739 100644 --- a/packages/utils/src/worker-host.ts +++ b/packages/utils/src/worker-host.ts @@ -17,3 +17,74 @@ export function declareWorkerHostEntry(): void { export function workerHostEntry(): string | null { return workerHostMain; } + +/** + * Buffers messages a Bun worker thread receives before its real handler is + * attached, then hands them off once it is. + * + * Bun delivers messages the parent posted before the worker spawned exactly + * once — when the worker entry module's top-level evaluation completes — to the + * `message` listeners present at that moment. A worker whose handler attaches + * via a later `await import(...)` therefore misses that flush. The + * self-dispatching CLI host imports each worker module dynamically from inside + * its argv dispatch, so the worker's own `parentPort.on("message")` lands after + * the flush and the parent's synchronously-posted `init` handshake is dropped — + * every run then stalls until the init timeout fires and silently falls back to + * the inline worker (issue: eval cells always taking the full timeout). + * + * The host calls {@link installWorkerInbox} synchronously in the entry's sync + * prefix (before importing the worker module) so a `parentPort` listener exists + * at flush time; the worker module then {@link consumeWorkerInbox}es it and + * binds the real handler, replaying anything buffered. Re-dispatching through + * `parentPort.emit("message", …)` is not an option — Bun's port is an + * `EventTarget` whose `emit` throws — so the inbox calls the handler directly. + */ +export interface WorkerInbox { + /** Route buffered and subsequent messages to `handler`; returns an unbind fn. */ + bind(handler: (message: unknown) => void): () => void; +} + +/** Minimal `parentPort` surface the inbox needs (Node/Bun `MessagePort`). */ +interface MessageListenerPort { + on(event: "message", listener: (value: unknown) => void): unknown; +} + +let pendingInbox: WorkerInbox | null = null; + +/** + * Attach a buffering `message` listener on `port` synchronously and stash the + * resulting inbox for the worker module to {@link consumeWorkerInbox}. MUST be + * called in the entry module's synchronous prefix — before the worker module is + * imported — so the listener exists when Bun flushes pre-spawn messages. + */ +export function installWorkerInbox(port: MessageListenerPort): WorkerInbox { + const queue: unknown[] = []; + let handler: ((message: unknown) => void) | null = null; + port.on("message", (data: unknown) => { + if (handler) handler(data); + else queue.push(data); + }); + const inbox: WorkerInbox = { + bind(next) { + handler = next; + for (const data of queue) next(data); + queue.length = 0; + return () => { + if (handler === next) handler = null; + }; + }, + }; + pendingInbox = inbox; + return inbox; +} + +/** + * Take the inbox installed by {@link installWorkerInbox} for this worker, or + * `null` when the worker module was loaded directly (no host pre-buffering, so + * the module's own synchronous top-level listener already wins the flush). + */ +export function consumeWorkerInbox(): WorkerInbox | null { + const inbox = pendingInbox; + pendingInbox = null; + return inbox; +} diff --git a/packages/utils/test/worker-host.test.ts b/packages/utils/test/worker-host.test.ts new file mode 100644 index 000000000..4a1854dc9 --- /dev/null +++ b/packages/utils/test/worker-host.test.ts @@ -0,0 +1,71 @@ +import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import { EventEmitter } from "node:events"; +import { consumeWorkerInbox, installWorkerInbox } from "../src/worker-host"; + +/** + * Regression for JS/tab eval workers always stalling until the init timeout. + * + * The self-dispatching CLI host imports each worker module dynamically from its + * argv dispatch, so the worker's own `parentPort.on("message")` attaches after + * Bun flushes the messages the parent posted before spawn — the synchronously + * posted `init` handshake was dropped and every run waited out the init timeout + * before silently falling back to the inline worker. `installWorkerInbox` + * attaches a `message` listener synchronously in the entry's sync prefix and + * buffers until the worker module `bind`s the real handler; these tests pin that + * buffer-replay contract. + */ +describe("worker-host inbox", () => { + // State is a module-global stash (one worker per process); drain it around + // each test so nothing leaks into the next. + beforeEach(() => consumeWorkerInbox()); + afterEach(() => consumeWorkerInbox()); + + it("replays messages buffered before bind, then forwards live ones, in order", () => { + const port = new EventEmitter(); + const inbox = installWorkerInbox(port); + + // Parent's pre-bind delivery (Bun's flush) — handler not attached yet. + port.emit("message", { type: "init" }); + port.emit("message", { type: "run", runId: "r1" }); + + const received: unknown[] = []; + inbox.bind(msg => received.push(msg)); + // Buffered messages replay synchronously on bind, in arrival order. + expect(received).toEqual([{ type: "init" }, { type: "run", runId: "r1" }]); + + // Subsequent deliveries reach the bound handler directly (not re-buffered). + port.emit("message", { type: "run", runId: "r2" }); + expect(received).toEqual([{ type: "init" }, { type: "run", runId: "r1" }, { type: "run", runId: "r2" }]); + }); + + it("delivers nothing twice and re-buffers after unbind", () => { + const port = new EventEmitter(); + const inbox = installWorkerInbox(port); + port.emit("message", "early"); + + const received: unknown[] = []; + const unbind = inbox.bind(msg => received.push(msg)); + expect(received).toEqual(["early"]); + + // After unbind the handler must stop receiving; messages re-queue instead + // of throwing or double-dispatching to the stale handler. + unbind(); + port.emit("message", "after-unbind"); + expect(received).toEqual(["early"]); + + // A fresh bind drains what arrived while unbound — exactly once. + const received2: unknown[] = []; + inbox.bind(msg => received2.push(msg)); + expect(received2).toEqual(["after-unbind"]); + }); + + it("hands the installed inbox to a single consumer, then reports none", () => { + const port = new EventEmitter(); + const inbox = installWorkerInbox(port); + + expect(consumeWorkerInbox()).toBe(inbox); + // A second consume (or a worker loaded directly with no host pre-buffering) + // sees no inbox and falls back to its own synchronous listener. + expect(consumeWorkerInbox()).toBeNull(); + }); +});