From a4987ba51858eb9a950fef0713b2ec7f46de3b7e Mon Sep 17 00:00:00 2001 From: roboomp Date: Sun, 2 Aug 2026 05:33:11 +0000 Subject: [PATCH] fix(mnemopi): bounded embed worker IPC and reap on timeout Embed-worker init/embed requests awaited the reply with no timeout, so a wedged fastembed/onnxruntime runtime blocked the turn memory recall or the shutdown consolidation forever, hanging headless -p/--mode json runs and leaving __omp_worker_mnemopi_embed unreaped. The #5753 fix only bounded the dispose-time consolidate await, not the embed IPC beneath it. Bound every request with an unref-ed timeout; on expiry fail the request and SIGKILL-reap the worker so the next call respawns a fresh child. Fixes #7352 --- packages/coding-agent/CHANGELOG.md | 4 + .../coding-agent/src/mnemopi/embed-client.ts | 51 ++++++- .../test/issue-7352-repro.test.ts | 141 ++++++++++++++++++ 3 files changed, 193 insertions(+), 3 deletions(-) create mode 100644 packages/coding-agent/test/issue-7352-repro.test.ts diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 11d4830cc..e952b3443 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed headless `-p` / `--mode json` runs with `memory.backend: mnemopi` hanging after a completed turn and leaving `__omp_worker_mnemopi_embed` unreaped when the embed worker's fastembed/onnxruntime runtime wedged. Embed-worker IPC requests (`init`/`embed`) were unbounded, so a stuck native runtime blocked the turn's memory recall or the shutdown consolidation forever; requests are now bounded and the wedged worker is reaped on timeout so the next call respawns a fresh child (regression of [#5753](https://github.com/can1357/oh-my-pi/issues/5753); [#7352](https://github.com/can1357/oh-my-pi/issues/7352)). + ## [17.2.4] - 2026-08-01 ### Added diff --git a/packages/coding-agent/src/mnemopi/embed-client.ts b/packages/coding-agent/src/mnemopi/embed-client.ts index be7f95859..a1dd561f1 100644 --- a/packages/coding-agent/src/mnemopi/embed-client.ts +++ b/packages/coding-agent/src/mnemopi/embed-client.ts @@ -84,6 +84,21 @@ export interface MnemopiSubprocessEmbeddingModel { embed(texts: string[], batchSize?: number): AsyncIterable; } +/** + * Upper bound on how long a single embed-worker IPC round-trip (init or embed) + * may block before the worker is treated as wedged. fastembed's steady-state + * embed and even a cold model load settle well within this; a longer stall + * means a hung native runtime (issue #4792) that would otherwise pin whatever + * awaits the embed — a turn's memory recall or the headless shutdown + * consolidation — indefinitely, leaving the process alive with an unreaped + * `__omp_worker_mnemopi_embed` child (issue #7352). On expiry the request fails + * and the worker is SIGKILL-reaped so the next request respawns a fresh one. + */ +const EMBED_REQUEST_TIMEOUT_MS = 120_000; + +/** Race marker for {@link MnemopiEmbedClient.#awaitRequest}. */ +const REQUEST_TIMED_OUT = Symbol("mnemopi.embed.timedOut"); + export class MnemopiEmbedClient { #worker: MnemopiEmbedWorkerHandle | null = null; #unsubscribeMessage: (() => void) | null = null; @@ -91,9 +106,14 @@ export class MnemopiEmbedClient { #pending = new Map(); #nextRequestId = 0; #spawnWorker: () => MnemopiEmbedWorkerHandle; + #requestTimeoutMs: number; - constructor(spawnWorker: () => MnemopiEmbedWorkerHandle = spawnMnemopiEmbedWorker) { + constructor( + spawnWorker: () => MnemopiEmbedWorkerHandle = spawnMnemopiEmbedWorker, + requestTimeoutMs: number = EMBED_REQUEST_TIMEOUT_MS, + ) { this.#spawnWorker = spawnWorker; + this.#requestTimeoutMs = requestTimeoutMs; } /** @@ -115,7 +135,7 @@ export class MnemopiEmbedClient { this.#pending.set(id, { kind: "init", model, resolve }); try { worker.send({ type: "init", id, model, cacheDir }); - const ok = await promise; + const ok = await this.#awaitRequest(promise); if (!ok) return null; } finally { this.#pending.delete(id); @@ -166,7 +186,7 @@ export class MnemopiEmbedClient { // worker's "embed before init" guard. Worker `ensureLoaded` is // idempotent so steady-state embeds pay no extra cost. worker.send({ type: "embed", id, model, cacheDir, texts, batchSize }); - const result = await promise; + const result = await this.#awaitRequest(promise); if (result instanceof Error) throw result; return result; } finally { @@ -174,6 +194,31 @@ export class MnemopiEmbedClient { } } + /** + * Await one embed-worker IPC reply, bounded by {@link EMBED_REQUEST_TIMEOUT_MS}. + * The timeout timer is `unref`'d so a pending request never keeps the parent + * event loop alive on its own (the awaiting caller does). On expiry the + * wedged worker is SIGKILL-reaped via {@link terminate} — faulting any other + * in-flight request and letting the next call respawn a fresh child — before + * the request rejects, so a hung native runtime cannot pin a turn's recall or + * the shutdown consolidation forever (issue #7352). + */ + async #awaitRequest(promise: Promise): Promise { + const { promise: timedOut, resolve: fire } = Promise.withResolvers(); + const timer = setTimeout(() => fire(REQUEST_TIMED_OUT), this.#requestTimeoutMs); + timer.unref(); + try { + const winner = await Promise.race([promise, timedOut]); + if (winner === REQUEST_TIMED_OUT) { + void this.terminate(); + throw new Error("mnemopi embed worker request timed out"); + } + return winner; + } finally { + clearTimeout(timer); + } + } + async *#streamEmbed( model: MnemopiEmbedModelId, cacheDir: string | undefined, diff --git a/packages/coding-agent/test/issue-7352-repro.test.ts b/packages/coding-agent/test/issue-7352-repro.test.ts new file mode 100644 index 000000000..946586f4c --- /dev/null +++ b/packages/coding-agent/test/issue-7352-repro.test.ts @@ -0,0 +1,141 @@ +/** + * Regression for https://github.com/can1357/oh-my-pi/issues/7352 + * + * A headless `omp --mode json --no-session -p @` run with + * `memory.backend: mnemopi` hung after its turn completed and left an + * unreaped `__omp_worker_mnemopi_embed` child. The embed-worker IPC request + * (`init` / `embed`) had no timeout, so a wedged native runtime (fastembed / + * onnxruntime hanging, cf. #4792) blocked whatever awaited the embed — the + * turn's memory recall or the shutdown consolidation — forever. #5753 only + * bounded the dispose-time consolidate *await*; the embed IPC underneath it + * stayed unbounded, so the wedge escaped that budget. + * + * The fix bounds every embed-worker request: on expiry the request fails and + * the wedged worker is SIGKILL-reaped so the next call respawns a fresh child. + * These tests drive the client with fake, deliberately-silent workers so the + * contract is exercised without fastembed/onnxruntime. + */ +import { describe, expect, it } from "bun:test"; +import { MnemopiEmbedClient, type MnemopiEmbedWorkerHandle } from "@oh-my-pi/pi-coding-agent/mnemopi/embed-client"; +import type { + MnemopiEmbedWorkerInbound, + MnemopiEmbedWorkerOutbound, +} from "@oh-my-pi/pi-coding-agent/mnemopi/embed-protocol"; + +/** A fake worker that answers `init` but never answers `embed`. */ +function silentEmbedWorker(state: { spawns: number; terminated: number }): () => MnemopiEmbedWorkerHandle { + return () => { + state.spawns += 1; + let handler: ((message: MnemopiEmbedWorkerOutbound) => void) | undefined; + return { + send(message: MnemopiEmbedWorkerInbound) { + // Reply to init/ping so the model handle resolves, but stay silent + // on `embed` to simulate a wedged native runtime. + queueMicrotask(() => { + if (message.type === "ping") handler?.({ type: "pong", id: message.id }); + else if (message.type === "init") handler?.({ type: "ready", id: message.id }); + }); + }, + onMessage(next) { + handler = next; + return () => { + if (handler === next) handler = undefined; + }; + }, + onError() { + return () => {}; + }, + async terminate() { + state.terminated += 1; + handler = undefined; + }, + }; + }; +} + +describe("issue #7352 — mnemopi embed requests are bounded and reap a wedged worker", () => { + it("fails a wedged embed within the budget instead of hanging forever", async () => { + const state = { spawns: 0, terminated: 0 }; + const client = new MnemopiEmbedClient(silentEmbedWorker(state), 50); + try { + const model = await client.initialize("fast-bge-base-en-v1.5", "/tmp/cache"); + expect(model).not.toBeNull(); + + const start = Date.now(); + let threw = false; + try { + for await (const _ of model!.embed(["hello"])) { + /* drain */ + } + } catch (error) { + threw = true; + expect(String(error)).toMatch(/timed out/i); + } + expect(threw).toBe(true); + // Bounded: nowhere near an indefinite hang. + expect(Date.now() - start).toBeLessThan(5_000); + // The wedged worker was reaped so it cannot linger as an orphan child. + expect(state.terminated).toBeGreaterThanOrEqual(1); + } finally { + await client.terminate(); + } + }, 10_000); + + it("respawns a fresh worker for the next request after reaping a wedged one", async () => { + const state = { spawns: 0, terminated: 0 }; + const client = new MnemopiEmbedClient(silentEmbedWorker(state), 50); + try { + const model = await client.initialize("fast-bge-base-en-v1.5", "/tmp/cache"); + const spawnsAfterInit = state.spawns; + + await expect( + (async () => { + for await (const _ of model!.embed(["a"])) { + /* drain */ + } + })(), + ).rejects.toThrow(/timed out/i); + + // The reap nulled the handle; a second embed must spawn a new child + // rather than reuse the dead one. + await expect( + (async () => { + for await (const _ of model!.embed(["b"])) { + /* drain */ + } + })(), + ).rejects.toThrow(/timed out/i); + + expect(state.spawns).toBeGreaterThan(spawnsAfterInit); + } finally { + await client.terminate(); + } + }, 10_000); + + it("returns null and reaps the worker when init itself wedges", async () => { + const state = { spawns: 0, terminated: 0 }; + // Worker that never answers anything, including init. + const client = new MnemopiEmbedClient(() => { + state.spawns += 1; + return { + send() {}, + onMessage() { + return () => {}; + }, + onError() { + return () => {}; + }, + async terminate() { + state.terminated += 1; + }, + }; + }, 50); + try { + const model = await client.initialize("fast-bge-base-en-v1.5", undefined); + expect(model).toBeNull(); + expect(state.terminated).toBeGreaterThanOrEqual(1); + } finally { + await client.terminate(); + } + }, 10_000); +});