Files
openclaw/extensions/memory-lancedb/embeddings.lifecycle.test.ts
yt2102 6bd8e0387c fix(memory): close previous embedding provider before replacement (#113471)
* fix(memory): close previous embedding provider before replacement

* fix(memory): increase embedding worker close grace period for slow hosts

* fix(memory): clear this.provider after close in resetProviderInitializationForRetry

Co-authored-by: Sanjay Santhanam <notifications@github.com>

* fix(memory): serialize embedding provider replacement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): unify provider transition lifecycle

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): join worker shutdown lifecycle

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): block worker restart during close

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retain failed provider retirements

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): drain retirements on manager close

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): make worker close joinable

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): preserve sync before provider retirement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): separate worker exit from disposal errors

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): avoid shutdown admission gap

Co-authored-by: yt2102 <yt2102@qq.com>

* style(memory): format provider lifecycle fix

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): bound embedding worker termination

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): fail closed on fallback initialization errors

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retry primary after fallback creation failure

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retain disconnected embedding workers until exit

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): isolate shared fallback transition failures

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): serialize scoped manager retirement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retain failed global manager closes

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): serialize manager admission with teardown

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): serialize outer manager lifecycle

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(gateway): close request-scoped embedding providers

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retire qmd managers before replacement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(gateway): support synchronous embedding cleanup

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): isolate manager lifecycle scopes

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(cli): close request-scoped embedding providers

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): canonicalize manager lifecycle ownership

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retry primary after null fallback result

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(gateway): retain failed embedding provider closes

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(gateway): lease closable embedding providers

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(gateway): drain retained embedding providers

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): join lazy fallback teardown

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(gateway): drain embedding providers after HTTP close

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retain failed qmd candidate cleanup

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): aggregate retained qmd teardown failures

Co-authored-by: yt2102 <yt2102@qq.com>

* refactor(memory): keep qmd lifecycle policy unchanged

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retain failed worker construction clients

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(gateway): scope local embedding retirement by provider

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): drain provider generations before retirement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): drain admitted operations before teardown

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): preserve provider identity through vector search

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): lease provider generation through embedding

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): lease provider through index publication

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): keep sync harness hooks optional

Co-authored-by: yt2102 <yt2102@qq.com>

* refactor(memory): own generations in sync lifecycle

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): preserve query runtime across retirement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): pin FTS-only sync generations

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retry failed worker construction cleanup

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): drain admitted searches before closing

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): close manager and gateway admission races

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): serialize qmd wrapper replacement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): serialize failed qmd retirement

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): satisfy qmd lifecycle type and deadcode gates

Co-authored-by: yt2102 <yt2102@qq.com>

* test(memory): type qmd lifecycle doubles

Co-authored-by: yt2102 <yt2102@qq.com>

* style(gateway): clarify created embedding provider

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): drain availability probes before close

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): preserve cleanup ownership without blocking fallback

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retire lancedb embedding providers

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): preserve generic provider cleanup receiver

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): close lancedb CLI embeddings

Co-authored-by: yt2102 <yt2102@qq.com>

* test(memory): normalize abort rejection reason

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): retain lancedb provider retirements

Co-authored-by: yt2102 <yt2102@qq.com>

* fix(memory): drain lancedb embedding uses

Co-authored-by: yt2102 <yt2102@qq.com>

---------

Co-authored-by: Sanjay Santhanam <notifications@github.com>
Co-authored-by: Peter Steinberger <steipete@gmail.com>
2026-07-26 03:00:22 -04:00

180 lines
6.0 KiB
TypeScript

import { expectDefined } from "@openclaw/normalization-core";
import { describe, expect, it, vi } from "vitest";
import type { OpenClawPluginApi } from "./api.js";
import type { MemoryConfig } from "./config.js";
const providerMocks = vi.hoisted(() => ({
getMemoryEmbeddingProvider: vi.fn(),
resolveDefaultAgentId: vi.fn(() => "main"),
}));
vi.mock("openclaw/plugin-sdk/memory-core-host-engine-embeddings", () => ({
getMemoryEmbeddingProvider: providerMocks.getMemoryEmbeddingProvider,
}));
vi.mock("openclaw/plugin-sdk/memory-host-core", () => ({
resolveDefaultAgentId: providerMocks.resolveDefaultAgentId,
}));
import { createEmbeddings } from "./embeddings.js";
function createApi(): OpenClawPluginApi {
const config = {};
return {
config,
runtime: {
config: { current: () => config },
agent: { resolveAgentDir: () => "/tmp/openclaw-agent" },
},
} as unknown as OpenClawPluginApi;
}
const embeddingConfig = {
provider: "openai",
model: "text-embedding-3-small",
} as MemoryConfig["embedding"];
describe("memory-lancedb provider lifecycle", () => {
it("queues replacement behind close intent while provider creation is pending", async () => {
let releaseFirstCreate: () => void = () => {};
const firstCreateGate = new Promise<void>((resolve) => {
releaseFirstCreate = resolve;
});
const closeProvider = vi.fn(async () => {});
const createProvider = vi.fn(async () => {
if (createProvider.mock.calls.length === 1) {
await firstCreateGate;
}
return {
provider: {
id: "openai",
model: "text-embedding-3-small",
embedQuery: vi.fn(async () => [0.1, 0.2, 0.3]),
embedBatch: vi.fn(async () => [[0.1, 0.2, 0.3]]),
close: closeProvider,
},
};
});
providerMocks.getMemoryEmbeddingProvider.mockReturnValue({
id: "openai",
create: createProvider,
});
const first = createEmbeddings(createApi(), { embedding: embeddingConfig } as MemoryConfig);
const firstEmbed = first.embed("first");
await vi.waitFor(() => expect(createProvider).toHaveBeenCalledTimes(1));
const closePromise = first.close?.();
const replacement = createEmbeddings(createApi(), {
embedding: embeddingConfig,
} as MemoryConfig);
const replacementEmbed = replacement.embed("replacement");
await Promise.resolve();
expect(createProvider).toHaveBeenCalledTimes(1);
releaseFirstCreate();
await firstEmbed;
await closePromise;
await replacementEmbed;
expect(closeProvider).toHaveBeenCalledTimes(1);
expect(createProvider).toHaveBeenCalledTimes(2);
expect(
expectDefined(closeProvider.mock.invocationCallOrder[0], "pending provider close order"),
).toBeLessThan(
expectDefined(createProvider.mock.invocationCallOrder[1], "replacement create order"),
);
await replacement.close?.();
});
it("does not re-close a provider retired while an older provider still fails", async () => {
const closeOlder = vi
.fn<() => Promise<void>>()
.mockRejectedValueOnce(new Error("older close failed once"))
.mockRejectedValueOnce(new Error("older close failed twice"))
.mockResolvedValue(undefined);
const closeCurrent = vi.fn(async () => {});
const createProvider = vi
.fn()
.mockResolvedValueOnce({
provider: {
id: "openai",
model: "older",
embedQuery: vi.fn(async () => [0.1]),
embedBatch: vi.fn(async () => [[0.1]]),
close: closeOlder,
},
})
.mockResolvedValueOnce({
provider: {
id: "openai",
model: "current",
embedQuery: vi.fn(async () => [0.2]),
embedBatch: vi.fn(async () => [[0.2]]),
close: closeCurrent,
},
});
providerMocks.getMemoryEmbeddingProvider.mockReturnValue({
id: "openai",
create: createProvider,
});
const older = createEmbeddings(createApi(), { embedding: embeddingConfig } as MemoryConfig);
const current = createEmbeddings(createApi(), { embedding: embeddingConfig } as MemoryConfig);
await older.embed("older");
await current.embed("current");
await expect(older.close?.()).rejects.toThrow("older close failed once");
await expect(current.close?.()).rejects.toThrow("older close failed twice");
expect(closeCurrent).toHaveBeenCalledTimes(1);
await expect(current.close?.()).resolves.toBeUndefined();
expect(closeOlder).toHaveBeenCalledTimes(3);
expect(closeCurrent).toHaveBeenCalledTimes(1);
});
it("drains an admitted embedding before provider close", async () => {
let markEmbedStarted: () => void = () => {};
const embedStarted = new Promise<void>((resolve) => {
markEmbedStarted = resolve;
});
let releaseEmbed: () => void = () => {};
const embedGate = new Promise<void>((resolve) => {
releaseEmbed = resolve;
});
const closeProvider = vi.fn(async () => {});
providerMocks.getMemoryEmbeddingProvider.mockReturnValue({
id: "openai",
create: vi.fn(async () => ({
provider: {
id: "openai",
model: "text-embedding-3-small",
embedQuery: vi.fn(async () => {
markEmbedStarted();
await embedGate;
return [0.1, 0.2, 0.3];
}),
embedBatch: vi.fn(async () => [[0.1, 0.2, 0.3]]),
close: closeProvider,
},
})),
});
const embeddings = createEmbeddings(createApi(), {
embedding: embeddingConfig,
} as MemoryConfig);
const embedPromise = embeddings.embed("active");
await embedStarted;
const closePromise = embeddings.close?.();
await Promise.resolve();
expect(closeProvider).not.toHaveBeenCalled();
await expect(embeddings.embed("late")).rejects.toThrow("memory-lancedb embeddings are closed");
releaseEmbed();
await expect(embedPromise).resolves.toEqual([0.1, 0.2, 0.3]);
await closePromise;
expect(closeProvider).toHaveBeenCalledTimes(1);
});
});