diff --git a/mem0-ts/src/client/tests/integration/helpers.ts b/mem0-ts/src/client/tests/integration/helpers.ts index 56a173534..3e3acf7ae 100644 --- a/mem0-ts/src/client/tests/integration/helpers.ts +++ b/mem0-ts/src/client/tests/integration/helpers.ts @@ -52,13 +52,17 @@ export async function withRetry( } /** - * Poll getAll until memories appear for a user. - * The Mem0 API processes memories asynchronously — after add() - * we need to wait for them to be available. + * Poll getAll until a user's memories reach minCount and the ID set + * stops changing between polls. + * + * The Mem0 API processes memories asynchronously, and consolidating a + * new fact replaces a memory (delete + add) instead of editing it in + * place. Returning on the first non-empty read hands callers IDs the + * pipeline then invalidates, so compare IDs rather than length: a + * replacement leaves the count unchanged. * * Polls every 15 seconds with a maximum of 4 retries to avoid - * hitting rate limits. Throws if results aren't available after - * all retries. + * hitting rate limits. Throws if minCount is never reached. */ export async function waitForMemories( client: MemoryClient, @@ -66,18 +70,32 @@ export async function waitForMemories( minCount: number, maxRetries = 4, ): Promise { + let memories: Memory[] = []; + let previousIds = ""; + for (let attempt = 1; attempt <= maxRetries; attempt++) { const response = await withRetry(() => client.getAll({ filters: { user_id: userId } }), ); - const memories = response.results ?? []; - if (memories.length >= minCount) { + memories = response.results ?? []; + + const ids = memories + .map((m) => m.id) + .sort() + .join(","); + if (memories.length >= minCount && ids === previousIds) { return memories; } + previousIds = ids; + if (attempt < maxRetries) { await new Promise((r) => setTimeout(r, 15_000)); } } + + if (memories.length >= minCount) { + return memories; + } throw new Error( `waitForMemories: expected at least ${minCount} memories for user "${userId}" but did not get them after ${maxRetries} attempts`, ); diff --git a/mem0-ts/src/client/tests/waitForMemories.test.ts b/mem0-ts/src/client/tests/waitForMemories.test.ts new file mode 100644 index 000000000..d12520688 --- /dev/null +++ b/mem0-ts/src/client/tests/waitForMemories.test.ts @@ -0,0 +1,69 @@ +import { waitForMemories } from "./integration/helpers"; + +const clientReturning = (polls: Array>) => { + let call = 0; + return { + getAll: jest.fn(async () => ({ + results: polls[Math.min(call++, polls.length - 1)], + })), + } as any; +}; + +const settle = async (pending: Promise): Promise => { + pending.catch(() => {}); + for (let i = 0; i < 4; i++) { + await jest.advanceTimersByTimeAsync(15_000); + } + return pending; +}; + +describe("waitForMemories", () => { + beforeEach(() => jest.useFakeTimers()); + afterEach(() => jest.useRealTimers()); + + test("waits for the id set to stop changing", async () => { + const client = clientReturning([ + [{ id: "a" }], + [{ id: "a" }, { id: "b" }], + [{ id: "a" }, { id: "b" }], + ]); + + const memories = await settle(waitForMemories(client, "user", 1)); + + expect(memories.map((m) => m.id)).toEqual(["a", "b"]); + expect(client.getAll).toHaveBeenCalledTimes(3); + }); + + test("keeps polling when consolidation replaces a memory at a constant count", async () => { + const client = clientReturning([ + [{ id: "a" }, { id: "b" }], + [{ id: "a" }, { id: "c" }], + [{ id: "a" }, { id: "c" }], + ]); + + const memories = await settle(waitForMemories(client, "user", 1)); + + expect(memories.map((m) => m.id)).toEqual(["a", "c"]); + }); + + test("returns the last read when the set never settles but minCount is met", async () => { + const client = clientReturning([ + [{ id: "a" }], + [{ id: "b" }], + [{ id: "c" }], + [{ id: "d" }], + ]); + + const memories = await settle(waitForMemories(client, "user", 1)); + + expect(memories.map((m) => m.id)).toEqual(["d"]); + }); + + test("throws when minCount is never reached", async () => { + const client = clientReturning([[]]); + + await expect(settle(waitForMemories(client, "user", 1))).rejects.toThrow( + /expected at least 1 memories/, + ); + }); +});