fix(eve): address review — async-write status, reasoning filter, turn-bounded capture, best-effort dedup

Review by @kartik-mem0. Root cause across all four: Mem0's hosted add with
infer:true is async and returns a PENDING event before the memory is visible.

- store.ts/tools.ts: add() returns AddResult{status,eventId} via parseAddResult;
  the v3 add path returns {status:"PENDING", event_id}, so remember reports
  {status:"queued", eventId} for async writes and {status:"saved"} only when
  resolved. FAILED/unknown status maps to completed, not queued.
- messages.ts: drop typed non-text parts (reasoning no longer leaks into memory);
  bound the captured assistant reply to the current turn and collect all of its
  text segments (text->tool->text) instead of a stale prior answer.
- capture.ts + README + eve.mdx: document operationId dedup as best-effort
  (lookup+add is not atomic; pending/concurrent windows can double-write);
  removed the "restart replays cannot write twice" claim.

Tests: 68 pass (+15). Added reasoning, tool-only/cross-turn, multi-segment,
pending/concurrent capture replay, and PENDING/FAILED add-result cases.
Typecheck + build clean.
This commit is contained in:
Himanshu-Sangshetti
2026-09-12 00:36:29 +09:00
parent 10d6717068
commit 75e9efa261
11 changed files with 327 additions and 22 deletions
+2 -2
View File
@@ -36,8 +36,8 @@ export default defineMemory({
## How it works
1. **Recall.** Before each turn and after compaction, Mem0 searches memories for the locked Eve scope and returns `{ id, content }` messages. Eve injects those as user-role messages attributed to the slot.
2. **Capture.** After a successful turn, Mem0 adds the user messages from that turn and, when present, the latest assistant reply. By default (`infer: true`) Mem0 extracts durable facts.
3. **Tools.** The model can call `search`, `remember`, and `forget`. Eve qualifies them as `mem0__search`, `mem0__remember`, and `mem0__forget` when the slot file is `mem0.ts`.
2. **Capture.** After a successful turn, Mem0 adds the user messages from that turn and, when present, that turn's assistant reply. By default (`infer: true`) Mem0 extracts durable facts. Capture is deduplicated per `operationId` on a best-effort basis (an in-process gate plus a durable `metadata.operation_id` lookup); because that check and the write are not atomic and platform writes are asynchronous, a restart or concurrent worker can occasionally capture a turn twice.
3. **Tools.** The model can call `search`, `remember`, and `forget`. Eve qualifies them as `mem0__search`, `mem0__remember`, and `mem0__forget` when the slot file is `mem0.ts`. `remember` returns `{ status: "saved" }` on a resolved write or `{ status: "queued" }` when the platform accepts it for asynchronous extraction.
Eve owns namespace, scope, and when recall/capture run. Mem0 owns storage, ranking, and extraction. Search, capture, remember, and forget are partitioned by `memory.scope.key`. Forget deletes by id only after confirming the memory belongs to that key.
+4 -2
View File
@@ -36,12 +36,14 @@ The live MCP path today is `eve add connection/mem0`. That is a different integr
|---|---|
| `recall["turn.started"]` | Semantic search over memories for the locked scope |
| `recall["compaction.completed"]` | Same search after Eve compacting history (`turn` may be null) |
| `capture["turn.completed"]` | Add the completed user turn and, when present, the latest assistant reply |
| `capture["turn.completed"]` | Add the completed user turn and, when present, this turn's assistant reply |
| `tools()` | `search`, `remember`, `forget` |
Search, capture, remember, and forget are partitioned by `memory.scope.key`. Forget loads the memory first and deletes only when `userId` matches that key.
Capture uses `operationId` as an idempotency key: an in-process gate plus a durable `metadata.operation_id` lookup so Eve replays after restart do not write twice.
Capture uses `operationId` for **best-effort** deduplication: an in-process gate plus a durable `metadata.operation_id` lookup. This is not atomic — with `infer: true` a prior write can be a PENDING event whose memory is not yet visible, so a restart during that window (or two concurrent workers) can capture the same turn twice. Eliminating that would require a backend-supported atomic idempotency key.
`remember` returns `{ status: "saved" }` when the write resolves, or `{ status: "queued" }` when the platform accepts it for asynchronous extraction (it becomes searchable a moment later).
If the slot file is named `mem0.ts`, Eve qualifies tools as `mem0__search`, `mem0__remember`, and `mem0__forget`.
+6
View File
@@ -10,6 +10,12 @@ export async function captureCompletedTurn(input: {
messages: readonly ConversationMessage[];
turnInput: readonly ConversationMessage[];
}): Promise<boolean> {
// Best-effort deduplication, not a guarantee. This lookup + add is not atomic,
// and with infer:true the prior write can still be a PENDING event whose
// memory is not yet visible here. So a restart during that window, or two
// workers that both clear this check before either write, can capture the
// same operation twice. The in-process gate in provider.ts narrows the common
// case; eliminating it would need a backend-supported atomic idempotency key.
const existing = await input.store.listByMetadata({
userId: input.scopeKey,
metadata: { operation_id: input.operationId },
+47 -9
View File
@@ -9,9 +9,16 @@ function textFromPart(part: unknown): string {
if (typeof part === "string") {
return part;
}
if (part && typeof part === "object" && "text" in part) {
const text = (part as { text: unknown }).text;
return typeof text === "string" ? text : "";
if (part && typeof part === "object") {
const record = part as { type?: unknown; text?: unknown };
// Typed parts must be explicit text. A `reasoning` part also carries a
// `text` field, so accepting any object with `text` would leak the model's
// private reasoning into long-term memory. Untyped `{ text }` (no `type`
// key) stays supported for callers that pass a bare text object.
if ("type" in record && record.type !== "text") {
return "";
}
return typeof record.text === "string" ? record.text : "";
}
return "";
}
@@ -72,11 +79,42 @@ export function completedTurnMessages(input: {
return [];
}
// Pair this turn's users with the latest assistant in the settled history.
const history = conversationMessages(input.messages);
const lastAssistant = [...history]
.reverse()
.find((message) => message.role === "assistant");
// Pair this turn's users with the assistant reply produced *in this turn*
// only. The current turn begins at the last user message in the settled
// history; any assistant text after it belongs to this turn. Walking the
// whole history instead would attach an older answer to the new user input
// when this turn is tool-only (no assistant text of its own).
let lastUserIndex = -1;
for (let index = input.messages.length - 1; index >= 0; index -= 1) {
if (input.messages[index]?.role === "user") {
lastUserIndex = index;
break;
}
}
return lastAssistant ? [...users, lastAssistant] : users;
// Collect every assistant text segment produced in this turn, in order, so a
// text -> tool-call -> text answer is captured whole rather than losing the
// pre-tool content.
const assistantSegments: string[] = [];
if (lastUserIndex >= 0) {
for (
let index = lastUserIndex + 1;
index < input.messages.length;
index += 1
) {
const message = input.messages[index];
if (message?.role !== "assistant") {
continue;
}
const text = extractText(message.content);
if (text.length > 0) {
assistantSegments.push(text);
}
}
}
const assistantText = assistantSegments.join("\n");
return assistantText
? [...users, { role: "assistant", content: assistantText }]
: users;
}
+44 -2
View File
@@ -27,6 +27,15 @@ export interface AddMemoryInput {
readonly metadata: StoreMetadata;
}
// A hosted-platform write with `infer: true` is asynchronous: Mem0 accepts the
// request and returns an `event_id` with a PENDING status before extraction has
// produced (or rejected) any memory. `queued` reflects that accepted-not-yet-
// stored state; `completed` means the write resolved synchronously.
export interface AddResult {
readonly status: "queued" | "completed";
readonly eventId?: string;
}
export interface SearchMemoryInput {
readonly userId: string;
readonly topK: number;
@@ -35,7 +44,10 @@ export interface SearchMemoryInput {
}
export interface MemoryStore {
add(messages: readonly MemoryMessage[], input: AddMemoryInput): Promise<void>;
add(
messages: readonly MemoryMessage[],
input: AddMemoryInput,
): Promise<AddResult>;
search(query: string, input: SearchMemoryInput): Promise<{ results: readonly SearchHit[] }>;
get(memoryId: string): Promise<MemoryRecord | null>;
listByMetadata(input: {
@@ -70,6 +82,35 @@ export function parseSearchHit(item: unknown): SearchHit | null {
return { id, memory };
}
// Event statuses that mean "accepted, extraction not finished". Only these map
// to `queued`; a terminal status (SUCCEEDED/FAILED) or an unknown/absent status
// falls through to `completed`, so a FAILED write is never reported as still
// pending. (The hosted add path returns PENDING synchronously; FAILED only
// appears later via event polling.)
const PENDING_STATUSES = new Set(["PENDING", "RUNNING", "PROCESSING", "QUEUED"]);
export function parseAddResult(response: unknown): AddResult {
if (response && typeof response === "object" && !Array.isArray(response)) {
const record = response as Record<string, unknown>;
const eventIdValue = record.event_id ?? record.eventId;
const eventId =
typeof eventIdValue === "string" && eventIdValue.length > 0
? eventIdValue
: undefined;
const status =
typeof record.status === "string" ? record.status.toUpperCase() : "";
if (eventId && PENDING_STATUSES.has(status)) {
return { status: "queued", eventId };
}
if (eventId) {
return { status: "completed", eventId };
}
}
// A memory-results array (or any non-event payload) means the write already
// resolved; there is nothing left pending to report.
return { status: "completed" };
}
export function parseMemoryRecord(item: unknown): MemoryRecord | null {
if (!item || typeof item !== "object") {
return null;
@@ -138,11 +179,12 @@ export async function createMem0Store(input: {
const store: MemoryStore = {
async add(messages, options) {
await client.add([...messages], {
const response = await client.add([...messages], {
userId: options.userId,
infer: options.infer,
metadata: { ...options.metadata },
});
return parseAddResult(response);
},
async search(query, options) {
const response = await client.search(query, {
+15 -3
View File
@@ -37,17 +37,29 @@ export function createMem0Tools(
}),
remember: defineTool({
description:
"Save one durable fact or preference about the current caller.",
"Save one durable fact or preference about the current caller. On the " +
"hosted platform the write is queued for extraction and may take a " +
"moment to become searchable; a `queued` status means accepted, not yet " +
"stored. Do not claim the fact is saved when the status is queued.",
inputSchema: z.object({
text: z.string().min(1).max(4000),
}),
async execute({ text }) {
await store.add([{ role: "user", content: text }], {
const result = await store.add([{ role: "user", content: text }], {
userId: scopeKey,
infer: options.infer,
metadata: { source: "eve-tool" },
});
return { saved: true };
// Report the true write state. With infer:true the platform returns a
// PENDING event, so promising "saved" would let the model tell the user
// a fact is stored when only the request was queued. Surface the eventId
// so the write can be correlated/confirmed rather than left a dead end.
return result.status === "queued"
? {
status: "queued" as const,
...(result.eventId ? { eventId: result.eventId } : {}),
}
: { status: "saved" as const };
},
}),
forget: defineTool({
+67
View File
@@ -47,6 +47,73 @@ describe("captureCompletedTurn", () => {
expect(store.added).toEqual([]);
});
it("dedupes a replay once the prior write is visible", async () => {
const store = createFakeStore();
const args = {
store,
scopeKey: "scope_abc",
operationId: "op_dedup",
infer: true,
metadata: { source: "eve", operation_id: "op_dedup" },
turnInput: [{ role: "user", content: "I am vegetarian" }],
messages: [
{ role: "user", content: "I am vegetarian" },
{ role: "assistant", content: "Got it." },
],
};
// First write is immediately visible (completed), so the replay is deduped.
expect(await captureCompletedTurn(args)).toBe(true);
expect(await captureCompletedTurn(args)).toBe(false);
expect(store.added).toHaveLength(1);
});
it("cannot dedupe while the prior write is still pending (best-effort)", async () => {
// Documents the known limitation: with infer:true the first write is a
// PENDING event whose memory is not yet visible, so the metadata lookup
// finds nothing and a restart replay writes the same operation again.
const store = createFakeStore([], { pendingWrites: true });
const args = {
store,
scopeKey: "scope_abc",
operationId: "op_pending",
infer: true,
metadata: { source: "eve", operation_id: "op_pending" },
turnInput: [{ role: "user", content: "I am vegetarian" }],
messages: [
{ role: "user", content: "I am vegetarian" },
{ role: "assistant", content: "Got it." },
],
};
expect(await captureCompletedTurn(args)).toBe(true);
expect(await captureCompletedTurn(args)).toBe(true);
expect(store.added).toHaveLength(2);
});
it("does not dedupe concurrent replays that both pass the lookup (best-effort)", async () => {
// Two workers on separate provider instances both clear the metadata check
// before either write lands, so both capture the same operation.
const store = createFakeStore([], { pendingWrites: true });
const args = {
store,
scopeKey: "scope_abc",
operationId: "op_concurrent",
infer: true,
metadata: { source: "eve", operation_id: "op_concurrent" },
turnInput: [{ role: "user", content: "I am vegetarian" }],
messages: [
{ role: "user", content: "I am vegetarian" },
{ role: "assistant", content: "Got it." },
],
};
const [first, second] = await Promise.all([
captureCompletedTurn(args),
captureCompletedTurn(args),
]);
expect(first).toBe(true);
expect(second).toBe(true);
expect(store.added).toHaveLength(2);
});
it("skips add when the operation id was already captured", async () => {
const store = createFakeStore([
{
+11 -1
View File
@@ -23,7 +23,10 @@ export type FakeSearch = SearchMemoryInput & {
query: string;
};
export function createFakeStore(hits: FakeHit[] = []): MemoryStore & {
export function createFakeStore(
hits: FakeHit[] = [],
options: { pendingWrites?: boolean } = {},
): MemoryStore & {
added: FakeAdd[];
searched: FakeSearch[];
deleted: string[];
@@ -50,11 +53,18 @@ export function createFakeStore(hits: FakeHit[] = []): MemoryStore & {
infer: input.infer,
metadata: input.metadata,
});
// Model Mem0's async extraction: a pending write is accepted but its
// memory is not yet visible to listByMetadata/get, so it cannot dedupe a
// replay. A resolved write becomes immediately visible.
if (options.pendingWrites) {
return { status: "queued" as const, eventId: `evt_${added.length}` };
}
records.push({
id: `mem_added_${added.length}`,
userId: input.userId,
metadata: input.metadata,
});
return { status: "completed" as const };
},
async search(query, input) {
searched.push({ query, ...input });
+63
View File
@@ -27,6 +27,19 @@ describe("extractText", () => {
it("ignores non-text parts", () => {
expect(extractText([{ type: "image", url: "x" }])).toBe("");
});
it("ignores reasoning parts and keeps only the final text", () => {
expect(
extractText([
{ type: "reasoning", text: "speculation" },
{ type: "text", text: "Noted." },
]),
).toBe("Noted.");
});
it("still reads an untyped text object", () => {
expect(extractText({ text: " bare " })).toBe("bare");
});
});
describe("lastUserText", () => {
@@ -87,4 +100,54 @@ describe("completedTurnMessages", () => {
}),
).toEqual([]);
});
it("does not attach a previous turn's answer to a tool-only turn", () => {
// Older Q&A, then a new user message whose only assistant output is a
// tool call (no text). The stale "older reply" must not be captured.
expect(
completedTurnMessages({
turnInput: [{ role: "user", content: "new question" }],
messages: [
{ role: "user", content: "older question" },
{ role: "assistant", content: "older reply" },
{ role: "user", content: "new question" },
{ role: "assistant", content: [{ type: "tool-call", id: "t1" }] },
],
}),
).toEqual([{ role: "user", content: "new question" }]);
});
it("captures all assistant text in a text/tool-call/text turn", () => {
expect(
completedTurnMessages({
turnInput: [{ role: "user", content: "Q" }],
messages: [
{ role: "user", content: "Q" },
{ role: "assistant", content: "part A" },
{ role: "assistant", content: [{ type: "tool-call", id: "t1" }] },
{ role: "assistant", content: "part B" },
],
}),
).toEqual([
{ role: "user", content: "Q" },
{ role: "assistant", content: "part A\npart B" },
]);
});
it("captures this turn's assistant reply, not an earlier one", () => {
expect(
completedTurnMessages({
turnInput: [{ role: "user", content: "new question" }],
messages: [
{ role: "user", content: "older question" },
{ role: "assistant", content: "older reply" },
{ role: "user", content: "new question" },
{ role: "assistant", content: "new reply" },
],
}),
).toEqual([
{ role: "user", content: "new question" },
{ role: "assistant", content: "new reply" },
]);
});
});
+13 -1
View File
@@ -247,7 +247,7 @@ describe("mem0Provider", () => {
expect(search).toEqual({
memories: [{ id: "mem_1", memory: "likes tea" }],
});
expect(remember).toEqual({ saved: true });
expect(remember).toEqual({ status: "saved" });
expect(forget).toEqual({ deleted: true });
expect(store.searched[0]).toMatchObject({
userId: "scope_abc",
@@ -264,6 +264,18 @@ describe("mem0Provider", () => {
expect(store.deleted).toEqual(["mem_1"]);
});
it("reports remember as queued when the write is still pending", async () => {
const store = createFakeStore([], { pendingWrites: true });
const provider = mem0Provider({ store, infer: true });
const tools = await provider.tools!(toolsContext());
const rememberTool = tools?.remember;
if (!rememberTool) {
throw new Error("expected remember tool");
}
const remember = await runTool(rememberTool, { text: "Allergic to peanuts" });
expect(remember).toEqual({ status: "queued", eventId: "evt_1" });
});
it("refuses to forget a memory from another scope", async () => {
const store = createFakeStore([
{ id: "mem_other", memory: "secret", userId: "scope_other" },
+55 -2
View File
@@ -25,11 +25,48 @@ vi.mock("mem0ai", () => ({
import {
createLazyStore,
createMem0Store,
parseAddResult,
parseMemoryRecord,
parseSearchHit,
resolveApiKey,
} from "../src/store.js";
describe("parseAddResult", () => {
it("reports a pending event as queued", () => {
expect(parseAddResult({ event_id: "evt_1", status: "PENDING" })).toEqual({
status: "queued",
eventId: "evt_1",
});
});
it("reports a succeeded event as completed", () => {
expect(parseAddResult({ event_id: "evt_1", status: "SUCCEEDED" })).toEqual({
status: "completed",
eventId: "evt_1",
});
});
it("treats a memory-results array as completed", () => {
expect(parseAddResult([{ id: "m1", memory: "tea" }])).toEqual({
status: "completed",
});
});
it("does not report a failed event as queued", () => {
expect(parseAddResult({ event_id: "evt_1", status: "FAILED" })).toEqual({
status: "completed",
eventId: "evt_1",
});
});
it("does not report an event with no status as queued", () => {
expect(parseAddResult({ event_id: "evt_1" })).toEqual({
status: "completed",
eventId: "evt_1",
});
});
});
describe("parseSearchHit", () => {
it("accepts id, memory_id, nested text, and numeric ids", () => {
expect(parseSearchHit({ id: 12, memory: "tea" })).toEqual({
@@ -122,18 +159,19 @@ describe("createMem0Store", () => {
getAll.mockResolvedValue({
results: [{ id: "m9", userId: "scope_abc", metadata: { operation_id: "op_1" } }],
});
add.mockResolvedValue({ eventId: "evt_1" });
add.mockResolvedValue({ event_id: "evt_1", status: "SUCCEEDED" });
const store = await createMem0Store({
apiKey: "m0-test",
host: "https://api.mem0.ai",
});
await store.add([{ role: "user", content: "I like tea" }], {
const addResult = await store.add([{ role: "user", content: "I like tea" }], {
userId: "scope_abc",
infer: false,
metadata: { source: "eve", operation_id: "op_1" },
});
expect(addResult).toEqual({ status: "completed", eventId: "evt_1" });
const found = await store.search("tea", {
userId: "scope_abc",
topK: 3,
@@ -169,6 +207,21 @@ describe("createMem0Store", () => {
]);
});
it("surfaces a pending platform write as queued", async () => {
const store = await createMem0Store({
apiKey: "m0-test",
host: "https://api.mem0.ai",
});
add.mockResolvedValueOnce({ event_id: "evt_9", status: "PENDING" });
await expect(
store.add([{ role: "user", content: "I like tea" }], {
userId: "scope_abc",
infer: true,
metadata: { source: "eve" },
}),
).resolves.toEqual({ status: "queued", eventId: "evt_9" });
});
it("accepts a bare search array and rejects unknown envelopes", async () => {
const store = await createMem0Store({
apiKey: "m0-test",