fix(mem0-ts): wait for memory set to settle before returning from waitForMemories
The live-API integration suites failed intermittently on a different test each run: CRUD delete 404ing on a known ID, batch update/delete failing, search returning nothing. All three share one cause. waitForMemories returned as soon as one memory existed, but the v3 pipeline was still consolidating, and consolidation replaces a memory (delete + add) rather than editing it in place. Tests captured memoryIds at that moment and later acted on IDs the platform had already invalidated. Wait for quiescence instead: two consecutive polls returning the same sorted ID set. Compare IDs rather than length, since a replacement leaves the count unchanged. Fixed in the one shared helper every integration suite seeds through.
This commit is contained in:
@@ -52,13 +52,17 @@ export async function withRetry<T>(
|
||||
}
|
||||
|
||||
/**
|
||||
* 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<Memory[]> {
|
||||
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`,
|
||||
);
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
import { waitForMemories } from "./integration/helpers";
|
||||
|
||||
const clientReturning = (polls: Array<Array<{ id: string }>>) => {
|
||||
let call = 0;
|
||||
return {
|
||||
getAll: jest.fn(async () => ({
|
||||
results: polls[Math.min(call++, polls.length - 1)],
|
||||
})),
|
||||
} as any;
|
||||
};
|
||||
|
||||
const settle = async <T>(pending: Promise<T>): Promise<T> => {
|
||||
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/,
|
||||
);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user