mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-12 22:17:00 +00:00
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 <steipete@gmail.com>
This commit is contained in:
@@ -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();
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user