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.
This commit is contained in:
can1357
2026-06-15 11:59:44 +02:00
parent 39f866fc73
commit 3ab675a83b
17 changed files with 226 additions and 26 deletions
+2 -2
View File
@@ -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)`). - **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) => ...)`. - **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`. - **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 ```ts
import { workerHostEntry } from "@oh-my-pi/pi-utils"; import { workerHostEntry } from "@oh-my-pi/pi-utils";
const hostEntry = workerHostEntry(); const hostEntry = workerHostEntry();
const worker = hostEntry const worker = hostEntry
? new Worker(hostEntry, { type: "module", argv: ["__omp_<name>_worker"] }) ? new Worker(hostEntry, { type: "module", argv: ["__omp_worker_<name>"] })
: new Worker(new URL("./<worker>.ts", import.meta.url).href, { type: "module" }); : new Worker(new URL("./<worker>.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. 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.
+2
View File
@@ -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. - 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 - 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. - 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
- Fixed ModelRegistry tests making outbound network calls by automatically stubbing fetch during test execution. - 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 ## [15.13.2] - 2026-06-15
+25 -12
View File
@@ -14,6 +14,7 @@ try {
* CLI entry point — registers all commands explicitly and delegates to the * CLI entry point — registers all commands explicitly and delegates to the
* lightweight CLI runner from pi-utils. * lightweight CLI runner from pi-utils.
*/ */
import { parentPort } from "node:worker_threads";
import type { CliConfig } from "@oh-my-pi/pi-utils/cli"; import type { CliConfig } from "@oh-my-pi/pi-utils/cli";
import { import {
APP_NAME, APP_NAME,
@@ -23,7 +24,7 @@ import {
setProfile, setProfile,
VERSION, VERSION,
} from "@oh-my-pi/pi-utils/dirs"; } 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 { installProfileAlias, resolveProfileAliasCommandFromProcess } from "./cli/profile-alias";
import { extractProfileFlags } from "./cli/profile-bootstrap"; import { extractProfileFlags } from "./cli/profile-bootstrap";
@@ -67,6 +68,7 @@ async function runSmokeTest(): Promise<void> {
const { smokeTestTinyTitleWorker } = await import("./tiny/title-client"); const { smokeTestTinyTitleWorker } = await import("./tiny/title-client");
const { smokeTestSttWorker } = await import("./stt/asr-client"); const { smokeTestSttWorker } = await import("./stt/asr-client");
const { smokeTestTtsWorker } = await import("./tts/tts-client"); const { smokeTestTtsWorker } = await import("./tts/tts-client");
const { smokeTestJsEvalWorker } = await import("./eval/js/context-manager");
await smokeTestSyncWorker(); await smokeTestSyncWorker();
const statsServer = await startServer(0); const statsServer = await startServer(0);
@@ -83,18 +85,23 @@ async function runSmokeTest(): Promise<void> {
await smokeTestTinyTitleWorker(); await smokeTestTinyTitleWorker();
await smokeTestSttWorker(); await smokeTestSttWorker();
await smokeTestJsEvalWorker();
await smokeTestTtsWorker(); await smokeTestTtsWorker();
process.stdout.write("smoke-test: ok\n"); process.stdout.write("smoke-test: ok\n");
} }
const TINY_WORKER_ARGS = new Set(["--tiny-worker", "__tiny_worker"]); const TINY_WORKER_ARG = "__omp_worker_tiny_inference";
const STATS_SYNC_WORKER_ARG = "__omp_stats_sync_worker"; const STATS_SYNC_WORKER_ARG = "__omp_worker_stats_sync";
const TAB_WORKER_ARG = "__omp_tab_worker"; const TAB_WORKER_ARG = "__omp_worker_tab";
const JS_EVAL_WORKER_ARG = "__omp_js_eval_worker"; const JS_EVAL_WORKER_ARG = "__omp_worker_js_eval";
const STT_WORKER_ARG = "__omp_stt_worker"; const STT_WORKER_ARG = "__omp_worker_stt";
const TTS_WORKER_ARG = "__omp_tts_worker"; const TTS_WORKER_ARG = "__omp_worker_tts";
async function runWorkerEntrypoint(arg: string | undefined): Promise<boolean> { async function runWorkerEntrypoint(arg: string | undefined): Promise<boolean> {
if (arg === TINY_WORKER_ARG) {
await runTinyWorker();
return true;
}
if (arg === STATS_SYNC_WORKER_ARG) { if (arg === STATS_SYNC_WORKER_ARG) {
// The sync worker handles messages via `self.onmessage`, assigned during // The sync worker handles messages via `self.onmessage`, assigned during
// this *async* dynamic import. Bun flushes the worker's initial message // this *async* dynamic import. Bun flushes the worker's initial message
@@ -117,11 +124,20 @@ async function runWorkerEntrypoint(arg: string | undefined): Promise<boolean> {
} }
return true; 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 (arg === TAB_WORKER_ARG) {
if (parentPort) installWorkerInbox(parentPort);
await import("./tools/browser/tab-worker-entry"); await import("./tools/browser/tab-worker-entry");
return true; return true;
} }
if (arg === JS_EVAL_WORKER_ARG) { if (arg === JS_EVAL_WORKER_ARG) {
if (parentPort) installWorkerInbox(parentPort);
await import("./eval/js/worker-entry"); await import("./eval/js/worker-entry");
return true; return true;
} }
@@ -251,11 +267,8 @@ export async function runCli(argv: string[]): Promise<void> {
// synchronous prefix of `runWorkerEntrypoint`, and Bun flushes the // synchronous prefix of `runWorkerEntrypoint`, and Bun flushes the
// worker's parked initial messages as soon as the entry module's // worker's parked initial messages as soon as the entry module's
// top-level evaluation finishes. // top-level evaluation finishes.
if (TINY_WORKER_ARGS.has(resolvedArgv[0] ?? "")) { if (resolvedArgv[0]?.startsWith("__omp_worker_")) {
await runTinyWorker(); await runWorkerEntrypoint(resolvedArgv[0]);
return;
}
if (await runWorkerEntrypoint(resolvedArgv[0])) {
return; return;
} }
@@ -141,6 +141,27 @@ export async function disposeAllVmContexts(): Promise<void> {
await Promise.all(all.map(session => killSession(session, new ToolError("JS context disposed"), { force: false }))); 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<void> {
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( async function runOnce(
session: JsSession, session: JsSession,
options: { options: {
@@ -447,7 +468,7 @@ function spawnJsWorker(): WorkerHandle {
try { try {
const hostEntry = workerHostEntry(); const hostEntry = workerHostEntry();
const worker = hostEntry 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" }); : new Worker(new URL("./worker-entry.ts", import.meta.url).href, { type: "module" });
return wrapBunWorker(worker); return wrapBunWorker(worker);
} catch (err) { } catch (err) {
@@ -1,13 +1,20 @@
import { parentPort } from "node:worker_threads"; import { parentPort } from "node:worker_threads";
import { consumeWorkerInbox } from "@oh-my-pi/pi-utils/worker-host";
import { WorkerCore } from "./worker-core"; import { WorkerCore } from "./worker-core";
import type { Transport, WorkerInbound, WorkerOutbound } from "./worker-protocol"; import type { Transport, WorkerInbound, WorkerOutbound } from "./worker-protocol";
if (!parentPort) throw new Error("js worker-entry: missing parentPort"); if (!parentPort) throw new Error("js worker-entry: missing parentPort");
const port = 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 = { const transport: Transport = {
send: (msg: WorkerOutbound) => port.postMessage(msg), send: (msg: WorkerOutbound) => port.postMessage(msg),
onMessage: handler => { onMessage: handler => {
if (inbox) return inbox.bind(data => handler(data as WorkerInbound));
const wrap = (data: unknown): void => handler(data as WorkerInbound); const wrap = (data: unknown): void => handler(data as WorkerInbound);
port.on("message", wrap); port.on("message", wrap);
return () => port.off("message", wrap); return () => port.off("message", wrap);
+1 -1
View File
@@ -72,7 +72,7 @@ const SMOKE_TEST_TIMEOUT_MS = 30_000;
* Hidden subcommand on the main CLI that boots the speech-recognition worker in * 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`. * 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 { function readTinyModelSetting(key: "providers.tinyModelDevice" | "providers.tinyModelDtype"): string | undefined {
try { try {
@@ -69,7 +69,7 @@ function normalizeTinyTitleGenerateOptions(
* Hidden subcommand on the main CLI that boots the tiny-model worker in the * 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`. * 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 { function readTinyModelSetting(path: "providers.tinyModelDevice" | "providers.tinyModelDtype"): string | undefined {
try { try {
@@ -685,7 +685,7 @@ async function spawnTabWorker(): Promise<WorkerHandle> {
try { try {
const hostEntry = workerHostEntry(); const hostEntry = workerHostEntry();
const worker = hostEntry 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" }); : new Worker(new URL("./tab-worker-entry.ts", import.meta.url).href, { type: "module" });
return wrapBunWorker(worker); return wrapBunWorker(worker);
} catch (err) { } catch (err) {
@@ -1,20 +1,28 @@
import { parentPort } from "node:worker_threads"; 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 type { Transport, WorkerInbound, WorkerOutbound } from "./tab-protocol";
import { WorkerCore } from "./tab-worker"; import { WorkerCore } from "./tab-worker";
if (!parentPort) throw new Error("tab-worker-entry: missing parentPort"); 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 = { const transport: Transport = {
send(msg, transferList) { send(msg, transferList) {
parentPort!.postMessage(msg, transferList ?? []); port.postMessage(msg, transferList ?? []);
}, },
onMessage(handler) { onMessage(handler) {
if (inbox) return inbox.bind(data => handler(data as WorkerOutbound | WorkerInbound));
const wrap = (message: unknown): void => handler(message as WorkerOutbound | WorkerInbound); const wrap = (message: unknown): void => handler(message as WorkerOutbound | WorkerInbound);
parentPort!.on("message", wrap); port.on("message", wrap);
return () => parentPort!.off("message", wrap); return () => port.off("message", wrap);
}, },
close() { close() {
parentPort!.close(); port.close();
}, },
}; };
+1 -1
View File
@@ -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 * 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). * 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 { function readTinyModelSetting(path: "providers.tinyModelDevice" | "providers.tinyModelDtype"): string | undefined {
try { try {
@@ -27,7 +27,7 @@ describe("issue #1011 — tab worker must re-enter the CLI entrypoint", () => {
const packageDir = path.resolve(import.meta.dir, ".."); const packageDir = path.resolve(import.meta.dir, "..");
const supervisorPath = path.join(packageDir, "src/tools/browser/tab-supervisor.ts"); const supervisorPath = path.join(packageDir, "src/tools/browser/tab-supervisor.ts");
const buildBinaryPath = path.join(packageDir, "scripts/build-binary.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 () => { it("tab-supervisor re-enters the worker-host entry with the argv selector", async () => {
const source = await Bun.file(supervisorPath).text(); const source = await Bun.file(supervisorPath).text();
@@ -8,7 +8,7 @@
* in the parent's address space and crashed the CLI on exit. * 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 * 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 * `runTinyWorker`, and the parent `SIGKILL`s the child on dispose so the
* native finalizer never runs in either address space. These tests pin 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 * three pieces of that contract so a future refactor cannot quietly land
+3
View File
@@ -1,6 +1,9 @@
# Changelog # Changelog
## [Unreleased] ## [Unreleased]
### Changed
- Renamed `__omp_stats_sync_worker` to `__omp_worker_stats_sync`.
## [15.13.1] - 2026-06-15 ## [15.13.1] - 2026-06-15
+1 -1
View File
@@ -93,7 +93,7 @@ interface WorkerHandle {
function createSyncWorker(): Worker { function createSyncWorker(): Worker {
const hostEntry = workerHostEntry(); const hostEntry = workerHostEntry();
if (hostEntry) { 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" }); return new Worker(new URL("./sync-worker.ts", import.meta.url).href, { type: "module" });
} }
+4
View File
@@ -2,6 +2,10 @@
## [Unreleased] ## [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 ## [15.13.1] - 2026-06-15
### Added ### Added
+71
View File
@@ -17,3 +17,74 @@ export function declareWorkerHostEntry(): void {
export function workerHostEntry(): string | null { export function workerHostEntry(): string | null {
return workerHostMain; 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;
}
+71
View File
@@ -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();
});
});