From d5a31888edeb70e9f24bc4fc9fa7ea027703f66a Mon Sep 17 00:00:00 2001 From: QiuYuang Date: Thu, 9 Jul 2026 16:26:05 +0800 Subject: [PATCH] fix(memory-host): reject queued worker requests on shutdown (#102451) * fix(memory-host): reject queued worker requests on shutdown * fix(memory-host): settle queued requests on shutdown --------- Co-authored-by: Peter Steinberger --- .../src/host/embeddings-worker.ts | 7 +++++++ .../src/host/embeddings.test.ts | 20 +++++++++++++++---- 2 files changed, 23 insertions(+), 4 deletions(-) diff --git a/packages/memory-host-sdk/src/host/embeddings-worker.ts b/packages/memory-host-sdk/src/host/embeddings-worker.ts index 97e5782e78ac..ae305dfaad9c 100644 --- a/packages/memory-host-sdk/src/host/embeddings-worker.ts +++ b/packages/memory-host-sdk/src/host/embeddings-worker.ts @@ -333,6 +333,13 @@ class LocalEmbeddingWorkerClient { if (!child) { return; } + this.rejectPending( + createLocalEmbeddingWorkerFailureError({ + message: "Local embedding worker exited unexpectedly (shutdown)", + code: LOCAL_EMBEDDING_WORKER_ERROR_CODES.exited, + reason: "exit", + }), + ); if (child.connected) { child.disconnect(); } diff --git a/packages/memory-host-sdk/src/host/embeddings.test.ts b/packages/memory-host-sdk/src/host/embeddings.test.ts index 59bde6f90bc6..1de711eed568 100644 --- a/packages/memory-host-sdk/src/host/embeddings.test.ts +++ b/packages/memory-host-sdk/src/host/embeddings.test.ts @@ -441,7 +441,7 @@ process.on("message", (message) => { await expect(provider.close?.()).resolves.toBeUndefined(); }); - it("terminates the worker when close runs behind a pending request", async () => { + it("rejects pending and queued requests when closing a busy worker", async () => { const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-local-embedding-worker-")); const workerScript = path.join(tempDir, "worker.cjs"); const embedStartedPath = path.join(tempDir, "embed-started"); @@ -478,8 +478,7 @@ process.on("message", (message) => { { workerScriptPath: workerScript }, ); - const embedPromise = provider.embedQuery("stuck"); - const embedError = embedPromise.then( + const firstEmbedError = provider.embedQuery("first").then( () => undefined, (err: unknown) => err, ); @@ -494,6 +493,16 @@ process.on("message", (message) => { }) .toBe(true); + const queuedEmbedResult = Promise.race([ + provider.embedQuery("queued").then( + () => "resolved" as const, + (err: unknown) => err, + ), + new Promise<"timeout">((resolve) => { + setTimeout(() => resolve("timeout"), 1_000); + }), + ]); + const closePromise = provider.close?.() ?? Promise.resolve(); const closeResult = await Promise.race([ closePromise.then(() => "closed" as const), @@ -503,7 +512,10 @@ process.on("message", (message) => { ]); expect(closeResult).toBe("closed"); - await expect(embedError).resolves.toMatchObject({ + await expect(firstEmbedError).resolves.toMatchObject({ + code: LOCAL_EMBEDDING_WORKER_ERROR_CODES.exited, + }); + await expect(queuedEmbedResult).resolves.toMatchObject({ code: LOCAL_EMBEDDING_WORKER_ERROR_CODES.exited, }); });