diff --git a/docs/integrations/eve.mdx b/docs/integrations/eve.mdx index f3b78b5ac..4e8329617 100644 --- a/docs/integrations/eve.mdx +++ b/docs/integrations/eve.mdx @@ -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. diff --git a/integrations/eve/README.md b/integrations/eve/README.md index 37d129e51..87ead5552 100644 --- a/integrations/eve/README.md +++ b/integrations/eve/README.md @@ -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`. diff --git a/integrations/eve/src/capture.ts b/integrations/eve/src/capture.ts index d28ecf73c..2fdfea959 100644 --- a/integrations/eve/src/capture.ts +++ b/integrations/eve/src/capture.ts @@ -10,6 +10,12 @@ export async function captureCompletedTurn(input: { messages: readonly ConversationMessage[]; turnInput: readonly ConversationMessage[]; }): Promise { + // 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 }, diff --git a/integrations/eve/src/messages.ts b/integrations/eve/src/messages.ts index 18d7dd36f..1d8b50ee0 100644 --- a/integrations/eve/src/messages.ts +++ b/integrations/eve/src/messages.ts @@ -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; } diff --git a/integrations/eve/src/store.ts b/integrations/eve/src/store.ts index 9d7d27b42..7973229bd 100644 --- a/integrations/eve/src/store.ts +++ b/integrations/eve/src/store.ts @@ -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; + add( + messages: readonly MemoryMessage[], + input: AddMemoryInput, + ): Promise; search(query: string, input: SearchMemoryInput): Promise<{ results: readonly SearchHit[] }>; get(memoryId: string): Promise; 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; + 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, { diff --git a/integrations/eve/src/tools.ts b/integrations/eve/src/tools.ts index 3d65b1ede..731e08648 100644 --- a/integrations/eve/src/tools.ts +++ b/integrations/eve/src/tools.ts @@ -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({ diff --git a/integrations/eve/tests/capture.test.ts b/integrations/eve/tests/capture.test.ts index d2a653162..a6fb0931b 100644 --- a/integrations/eve/tests/capture.test.ts +++ b/integrations/eve/tests/capture.test.ts @@ -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([ { diff --git a/integrations/eve/tests/helpers.ts b/integrations/eve/tests/helpers.ts index f84555ee7..c277bf2c8 100644 --- a/integrations/eve/tests/helpers.ts +++ b/integrations/eve/tests/helpers.ts @@ -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 }); diff --git a/integrations/eve/tests/messages.test.ts b/integrations/eve/tests/messages.test.ts index 4e8fca5e1..66ab56a0c 100644 --- a/integrations/eve/tests/messages.test.ts +++ b/integrations/eve/tests/messages.test.ts @@ -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" }, + ]); + }); }); diff --git a/integrations/eve/tests/provider.test.ts b/integrations/eve/tests/provider.test.ts index e6a30af9d..33abee57d 100644 --- a/integrations/eve/tests/provider.test.ts +++ b/integrations/eve/tests/provider.test.ts @@ -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" }, diff --git a/integrations/eve/tests/store.test.ts b/integrations/eve/tests/store.test.ts index 1ed3ccd15..3e58c47ed 100644 --- a/integrations/eve/tests/store.test.ts +++ b/integrations/eve/tests/store.test.ts @@ -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",