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
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -84,6 +84,21 @@ export interface MnemopiSubprocessEmbeddingModel {
|
||||
embed(texts: string[], batchSize?: number): AsyncIterable<number[][]>;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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<string, PendingRequest>();
|
||||
#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<T>(promise: Promise<T>): Promise<T> {
|
||||
const { promise: timedOut, resolve: fire } = Promise.withResolvers<typeof REQUEST_TIMED_OUT>();
|
||||
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,
|
||||
|
||||
@@ -0,0 +1,141 @@
|
||||
/**
|
||||
* Regression for https://github.com/can1357/oh-my-pi/issues/7352
|
||||
*
|
||||
* A headless `omp --mode json --no-session -p @<file>` 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);
|
||||
});
|
||||
Reference in New Issue
Block a user