diff --git a/docs/docs.json b/docs/docs.json index 59698569d..535047752 100644 --- a/docs/docs.json +++ b/docs/docs.json @@ -85,7 +85,8 @@ "platform/features/advanced-retrieval", "platform/advanced-memory-operations", "platform/features/custom-instructions", - "platform/features/memory-decay" + "platform/features/memory-decay", + "platform/features/dream" ] }, { diff --git a/docs/llms.txt b/docs/llms.txt index 618c4eef4..13be0c6d5 100644 --- a/docs/llms.txt +++ b/docs/llms.txt @@ -207,6 +207,7 @@ If the user is on a pre-current major (Python < 2, TS < 3, or Platform `output_f - [Temporal Reasoning](https://docs.mem0.ai/platform/features/temporal-reasoning) [Platform]: Use when time-aware searches like last week, upcoming, or right now need better result ordering. - [Custom Instructions](https://docs.mem0.ai/platform/features/custom-instructions) [Platform]: Use when tailoring what Mem0 extracts and stores on Platform. - [Memory Decay](https://docs.mem0.ai/platform/features/memory-decay) [Platform]: Use when search results should boost recently-reinforced memories and dampen stale ones. Opt in per project; applies at search time and never filters candidates out. +- [Dream](https://docs.mem0.ai/platform/features/dream) [Platform]: Use when long-lived user memory should stay coherent on its own - synthesizing recurring patterns, superseding outdated facts, and merging duplicates in the background. Synthesis is opt-in (Pro+); Supersede and Merge are always on. - [Advanced Memory Operations](https://docs.mem0.ai/platform/advanced-memory-operations) [Platform]: Use when basic CRUD is not enough - batch ops, complex filters, workflows. ### Features - Data Management diff --git a/docs/platform/features/dream.mdx b/docs/platform/features/dream.mdx new file mode 100644 index 000000000..05c07711f --- /dev/null +++ b/docs/platform/features/dream.mdx @@ -0,0 +1,164 @@ +--- +title: Dream +description: "Dream keeps a user's memory clean and insightful over time by distilling recurring patterns, retiring outdated facts, and folding away duplicates, automatically and in the background." +--- + +# Dream + +As an application talks to the same user over weeks and months, their memory grows. Some of that growth is signal, meaning new facts worth keeping. A lot of it is noise: the same preference stated three different ways, an old fact that a newer one has quietly replaced, or a set of small observations that only mean something when you look at them together. + +**Dream** is the background layer that keeps a user's memory coherent as it grows. It continuously reviews each user's memories and does three things. It synthesizes higher-order patterns, supersedes outdated facts, and merges duplicates. The result is that what you read back stays sharp instead of drifting into a pile of overlapping, stale entries. + + + **Dream matters when…** + - Your users interact with your product over a long period and accumulate a lot of memories. + - You want retrieval to return the *current* truth about a user, not a mix of old and new contradictory facts. + - You want higher-level insights ("this user consistently prefers X") without writing your own summarization layer. + + +## The three actions of Dream + +Dream is made of three independent actions. Two of them, Supersede and Merge, keep memory clean and run automatically for everyone. The third, Synthesis, produces new insight and is a toggle you turn on per project. + +| Action | What it does | When it runs | Availability | +|---|---|---|---| +| **Synthesis** | Distills a user's memories into higher-order **pattern memories** | On a schedule, in the background | Opt-in (Pro and above) | +| **Supersede** | Marks an older fact as outdated when a newer one contradicts it | As memories are added | Always on, all plans | +| **Merge** | Folds a duplicate into a single canonical memory | As memories are added | Always on, all plans | + +### Synthesis + +Over time a user's memories often *imply* something larger than any single entry. Ten separate notes about early-morning meetings, workout logs, and coffee orders together say "this user is an early riser." **Synthesis** finds those recurring threads and writes them back as new **pattern memories**: concise, higher-order facts that capture what the individual memories only hint at. + +- Pattern memories are added *alongside* your existing memories, never in place of them. The source memories stay exactly where they are. +- Each pattern memory keeps a link back to the specific memories it was distilled from, so an insight is always traceable to its evidence. +- Synthesis is additive and idempotent. Re-running it will not create duplicate patterns for the same underlying evidence. + +**Example.** These four memories accumulate for the same user over time: + +- *"User runs approximately 40 kilometers per week and is currently training for a marathon."* +- *"User lifts weights at the gym three times a week, primarily focusing on legs and back."* +- *"User tracks all workouts using a Garmin Forerunner watch."* +- *"User's goal for 2026 is to run a sub-4-hour marathon."* + +Synthesis distills them into one pattern memory, kept alongside the originals: + +> *"User follows a structured fitness routine that includes weekly long runs (≈40 km), regular leg-and-back strength training, tracks workouts with a Garmin device, and pursues progressive marathon time goals."* + +Each source memory stays exactly where it was, and the new pattern links back to all of them as its evidence. + + + Synthesis only considers memories created **after** you enable it for a project. Turning it on sets a forward boundary, so historical memories aren't reprocessed in bulk on day one. Everything added from that point on is eligible. + + + + Synthesis only looks at memories scoped to a **`user_id` alone**. Memories that also carry another entity (an `agent_id`, `run_id`, or `app_id`) are left out of a user's synthesis run. This keeps each run tied to a single user's own memories and avoids cross-referencing across agents, runs, or apps. So a memory has to be user-scoped, with no other entity attached, to be eligible for a pattern. + + +### Supersede + +When a user tells you something that **contradicts** an earlier memory ("I moved to Berlin" after an earlier "I live in Lisbon"), Dream marks the older memory as **superseded** and links it to the newer fact that replaced it. Superseded memories are not deleted and not hidden by default. A normal `search` or `get` still returns them alongside your active memories, badged as superseded, so you keep the full history. When you want only the current truth, ask for it explicitly with `latest_only=true` (see [How reads change](#how-reads-change-with-dream) below). Supersede runs automatically as part of adding memories, on every plan. + +**Example.** The user has an existing memory *"User drives a 2019 Subaru Outback."* Later they mention selling it, producing a new memory *"User sold their 2019 Subaru Outback and bought a Tesla Model 3."* Dream marks the Subaru memory as superseded and links it to the newer one. A default `search` returns both: the Tesla memory as active, and the Subaru one labelled superseded. Passing `latest_only=true` returns only *"User sold their 2019 Subaru Outback and bought a Tesla Model 3."* + +### Merge + +When a new memory is effectively a **duplicate** of one you already have, Dream keeps a single canonical memory instead of two near-identical copies. When a duplicate is stored as its own memory, Dream marks it **merged** and links it to the canonical one. The merged record is hidden from reads by default (you get the one canonical memory), retained rather than deleted, and surfaced with `include_merged=true`. Merge runs automatically as memories are added, on every plan, and keeps your memory set compact without you deduplicating by hand. + +In practice, most exact or near-duplicate restatements of a fact you already have are recognised and deduplicated as the memory is added. No second copy is created, so you simply keep the one memory. A distinct `merged` record appears when a fuller version of an existing fact arrives and folds the barer one in. + +**Example.** The user has an existing memory *"User has a dog named Rex."* Later they mention *"My dog Rex is a 3-year-old golden retriever."* The richer statement is stored and Dream marks the barer *"User has a dog named Rex"* as **merged** into it. A default read returns the single canonical memory *"User's dog Rex is a 3-year-old golden retriever"*, and `include_merged=true` also returns the merged original. + + + **Nothing Dream does is destructive.** Superseded and merged memories are retained, never erased. Every change is recorded and reviewable, so you always know why a memory was retired or combined. + + +## How reads change with Dream + +Dream doesn't change the shape of your `add`, `search`, or `get` calls, so you don't touch your application code. What it changes is *which* memories a read returns by default, and it gives you two flags to widen or narrow that set: + +| Read mode | Active | Superseded | Merged | +|---|:---:|:---:|:---:| +| **Default** (`search` / `get`) | ✓ | ✓ | ✗ | +| **`latest_only=true`** | ✓ | — | — | +| **`include_merged=true`** | ✓ | ✓ | ✓ | + +- **By default**, a read returns active plus superseded memories (superseded ones are still there, labelled as history) and hides merged duplicates. Synthesized pattern memories are returned alongside these too. +- **`latest_only=true`** narrows the result to active memories only, meaning the current truth, with superseded and merged both excluded. Use this when you want the cleanest possible snapshot of the user right now. +- **`include_merged=true`** returns everything, including the merged duplicates, when you need the complete picture. + +## Enabling Dream + +Supersede and Merge require no setup. They're always on for every project on every plan. + +Synthesis is opt-in per project: + +1. Open the **[Dream settings for your project](https://app.mem0.ai/dashboard/dream)** in the Mem0 dashboard. +2. Toggle **Synthesis** on. + +Synthesis is a per-project setting, so you can enable it for one project and compare against another with it off. You can turn it off at any time. Doing so is fully reversible and leaves every existing memory (including already-synthesized patterns) untouched. + +From the [Dream page](https://app.mem0.ai/dashboard/dream) you can also review what Dream has done: recent synthesis runs, the patterns produced and their source memories, and the memories that were superseded or merged. + +## Plan availability + +| Capability | Free | Starter | Pro | Enterprise | +|---|:---:|:---:|:---:|:---:| +| **Supersede** (outdated facts flagged) | ✓ | ✓ | ✓ | ✓ | +| **Merge** (duplicates folded, hidden by default) | ✓ | ✓ | ✓ | ✓ | +| **Synthesis** (pattern memories written) | — | — | ✓ | ✓ | +| **Dream dashboard** (runs, activity) | — | — | ✓ | ✓ | + +Synthesis requires a **Pro plan or higher**. Enterprise plans also get a faster, configurable schedule (see below). + +## How often Dream runs, and what delay to expect + +Different actions run on different clocks, so the delay you should expect depends on which action. + +### Supersede & Merge, as memories are added + +Supersede and Merge are part of the memory-addition pipeline. They're evaluated when a memory is added, so an outdated fact is superseded or a duplicate is merged as part of that add being processed, on the same timescale as the memory becoming searchable. There's no separate schedule to wait for. + +### Synthesis, on a schedule in the background + +Synthesis runs as a scheduled background job per user, not on every add. Two conditions gate it: + +- **Enough to work with:** a user must have at least **20** memories before Synthesis considers them. Below that threshold there isn't a meaningful pattern to distill yet. +- **Cadence elapsed:** each user is re-synthesized at most once per cadence window. + +| Plan | Synthesis cadence (per user) | +|---|---| +| Pro | Every **7 days** | +| Enterprise | **Daily** (and configurable) | + +Because Synthesis is processed in batches in the background, **expect new pattern memories to appear within roughly 24 hours of a scheduled run**, not instantly. In practice, the end-to-end delay from crossing a cadence window to seeing new patterns is up to about a day. This background design is deliberate: it keeps Synthesis from adding any latency to your live `add` and `search` calls. + + + Synthesis is **not** real-time. If you enable it today, the first pattern memories for an eligible user will appear on the next scheduled run for that user (governed by the cadence above), and can take up to ~24 hours to complete once that run starts. Supersede and Merge, by contrast, keep pace with your adds. + + +## FAQ + +**Does Dream delete any of my memories?** +No. Nothing Dream does is destructive. Superseded memories stay visible in default reads (labelled as history), merged duplicates are hidden by default but retained, and synthesized patterns are added alongside your existing memories, never in place of them. Every change is reviewable from the dashboard. + +**How do I get only the current facts, without the superseded ones?** +Pass `latest_only=true` on your read. Superseded and merged memories are both excluded, leaving only active memories. By default (no flag) superseded memories are included so you keep the full history. + +**Do I need to change my code to use Dream?** +No. Supersede and Merge are always on, and enabling Synthesis is a project setting. Your `add`, `search`, and `get` calls are unchanged. Dream shapes the memory set behind the same API. + +**Will Synthesis reprocess all my old memories when I turn it on?** +No. Enabling Synthesis sets a forward boundary, so only memories created after you turn it on are eligible. This avoids a bulk reprocess of your entire history on day one. + +**Why don't I see pattern memories immediately after enabling Synthesis?** +Synthesis runs on a schedule (every 7 days on Pro, daily on Enterprise) and only for users with at least 20 memories. Patterns appear on the next scheduled run for an eligible user and can take up to ~24 hours to complete once that run starts. + +**Which memories does Synthesis include?** +Only memories scoped to a `user_id` on its own. If a memory also carries an `agent_id`, `run_id`, or `app_id`, it's excluded from that user's synthesis run. This keeps each run confined to a single user's memories and prevents any cross-referencing across agents, runs, or apps. Supersede and Merge are not affected by this and run across your memories as usual. + +**Are synthesized pattern memories traceable?** +Yes. Every pattern memory links back to the specific source memories it was distilled from, so you can always see the evidence behind an insight from the Dream dashboard. + +**Can I turn Synthesis off?** +Yes, at any time, per project. Turning it off is fully reversible and leaves all existing memories, including already-synthesized patterns, untouched. diff --git a/integrations/n8n-nodes-mem0/package.json b/integrations/n8n-nodes-mem0/package.json index 6e8fe280a..04eb220fe 100644 --- a/integrations/n8n-nodes-mem0/package.json +++ b/integrations/n8n-nodes-mem0/package.json @@ -1,6 +1,6 @@ { "name": "@mem0/n8n-nodes-mem0", - "version": "0.1.1", + "version": "0.1.2", "description": "n8n community node for Mem0 — the memory layer for AI agents. Add, search, get, update, and delete long-term memories.", "keywords": [ "n8n-community-node-package", @@ -14,7 +14,7 @@ "homepage": "https://mem0.ai", "author": { "name": "Mem0", - "email": "founders@mem0.ai" + "email": "integrations@mem0.ai" }, "repository": { "type": "git", diff --git a/integrations/zapier-mem0/index.js b/integrations/zapier-mem0/index.js new file mode 100644 index 000000000..f3d35a1d6 --- /dev/null +++ b/integrations/zapier-mem0/index.js @@ -0,0 +1,3 @@ +// Zapier's Lambda wrapper requires `/index.js` and ignores package.json +// "main", so re-export the compiled app from the root. Run `npm run build` first. +module.exports = require('./dist/index.js'); diff --git a/integrations/zapier-mem0/package.json b/integrations/zapier-mem0/package.json index 47b62a109..64a5510d9 100644 --- a/integrations/zapier-mem0/package.json +++ b/integrations/zapier-mem0/package.json @@ -20,7 +20,7 @@ "directory": "integrations/zapier-mem0" }, "license": "Apache-2.0", - "main": "dist/index.js", + "main": "index.js", "scripts": { "build": "tsc", "test": "jest --testTimeout 180000", diff --git a/integrations/zapier-mem0/src/creates/add_memory.ts b/integrations/zapier-mem0/src/creates/add_memory.ts index a9260cd16..d0e27a034 100644 --- a/integrations/zapier-mem0/src/creates/add_memory.ts +++ b/integrations/zapier-mem0/src/creates/add_memory.ts @@ -115,7 +115,7 @@ export default { choices: { user: 'User', assistant: 'Assistant', system: 'System' }, default: 'user', }, - { key: 'user_id', label: 'User ID', type: 'string' }, + { key: 'user_id', label: 'User ID', type: 'string', required: true, helpText: 'Scope this memory to a user. Mem0 requires at least one entity ID.' }, { key: 'agent_id', label: 'Agent ID', type: 'string' }, { key: 'run_id', label: 'Run ID', type: 'string' }, { key: 'metadata', label: 'Metadata (JSON)', type: 'string' }, diff --git a/mem0-ts/src/oss/src/memory/index.ts b/mem0-ts/src/oss/src/memory/index.ts index 8a6c2cf24..a84543e58 100644 --- a/mem0-ts/src/oss/src/memory/index.ts +++ b/mem0-ts/src/oss/src/memory/index.ts @@ -100,11 +100,20 @@ const ENTITY_PARAMS = [ "agentId", "runId", ]; - -// Identity keys stripped from update() metadata: ENTITY_PARAMS covers user_id/agent_id/run_id -// in both casings (the default store promotes camelCase on read); actor_id has no camelCase alias. +// Identity keys stripped from caller metadata in add() and update(): ENTITY_PARAMS covers +// user_id/agent_id/run_id in both casings (the default store promotes camelCase on read); +// actor_id has no camelCase alias. const IDENTITY_KEYS = [...ENTITY_PARAMS, "actor_id"]; +// Caller metadata must not overwrite or inject an identity scope (#6342 / #6367 / #6371). +function stripIdentityKeys( + metadata: Record = {}, +): Record { + return Object.fromEntries( + Object.entries(metadata).filter(([key]) => !IDENTITY_KEYS.includes(key)), + ); +} + // Batch size for deleteAll pagination. Larger than most vector store default // page limits (~100) to minimize roundtrips while bounded to avoid memory pressure. const DELETE_ALL_BATCH_SIZE = 1000; @@ -740,7 +749,8 @@ export class Memory { has_filters: !!config.filters, infer: config.infer, }); - const { metadata = {}, filters = {}, infer = true } = config; + const { filters = {}, infer = true } = config; + const metadata = stripIdentityKeys(config.metadata); // Validate and trim entity IDs const userId = validateAndTrimEntityId(config.userId, "userId"); @@ -751,6 +761,9 @@ export class Memory { if (userId) filters.user_id = metadata.user_id = userId; if (agentId) filters.agent_id = metadata.agent_id = agentId; if (runId) filters.run_id = metadata.run_id = runId; + if (filters.user_id) metadata.user_id = filters.user_id; + if (filters.agent_id) metadata.agent_id = filters.agent_id; + if (filters.run_id) metadata.run_id = filters.run_id; // Normalize expiration date into the stored metadata (round-trips via get()). if (config.expirationDate != null) { @@ -1971,10 +1984,7 @@ export class Memory { existingEmbeddings[newData] || (await this.embedder.embed(newData, "update")); - // Caller metadata must not overwrite or inject an identity scope (#6342 / #6367). - const sanitizedMetadata = Object.fromEntries( - Object.entries(metadata).filter(([k]) => !IDENTITY_KEYS.includes(k)), - ); + const sanitizedMetadata = stripIdentityKeys(metadata); const newMetadata = { ...existingMemory.payload, diff --git a/mem0-ts/src/oss/tests/memory.add.test.ts b/mem0-ts/src/oss/tests/memory.add.test.ts index 4e7ebe773..059a0ef12 100644 --- a/mem0-ts/src/oss/tests/memory.add.test.ts +++ b/mem0-ts/src/oss/tests/memory.add.test.ts @@ -8,55 +8,78 @@ import type { MemoryConfig, MemoryItem, SearchResult } from "../src/types"; jest.setTimeout(15000); -// Mock Google modules to prevent @google/genai crash in CI -jest.mock("../src/embeddings/google", () => ({ - GoogleEmbedder: jest.fn(), -})); -jest.mock("../src/llms/google", () => ({ - GoogleLLM: jest.fn(), -})); +jest.mock("../src/utils/factory", () => { + const { MemoryVectorStore } = jest.requireActual( + "../src/vector_stores/memory", + ); + const { MemoryHistoryManager } = jest.requireActual( + "../src/storage/MemoryHistoryManager", + ); + const testEmbedding = new Array(1536).fill(0.1); -jest.mock("../src/llms/openai", () => ({ - OpenAILLM: jest.fn().mockImplementation(() => ({ - generateResponse: jest - .fn() - .mockImplementation( - (messages: Array<{ role: string; content: string }>) => { - // V3 pipeline: single LLM call with additive extraction prompt. - const userMsg = messages.find((m) => m.role === "user"); - const content = userMsg?.content ?? ""; - const newMsgMatch = content.match( - /## New Messages\n([\s\S]*?)(?=\n##|$)/, + class MockEmbedder { + embeddingDims = 1536; + + async embed(): Promise { + return testEmbedding; + } + + async embedBatch(texts: string[]): Promise { + return texts.map(() => testEmbedding); + } + } + + class MockLLM { + async generateResponse(messages: Array<{ role: string; content: string }>) { + const userMsg = messages.find((m) => m.role === "user"); + const content = userMsg?.content ?? ""; + const newMsgMatch = content.match( + /## New Messages\n([\s\S]*?)(?=\n##|$)/, + ); + const extracted = newMsgMatch + ? newMsgMatch[1].trim() + : "extracted fact from input"; + + return JSON.stringify({ + memory: [ + { + id: "0", + text: extracted, + attributed_to: "user", + }, + ], + }); + } + } + + return { + __esModule: true, + EmbedderFactory: { + create: jest.fn(() => new MockEmbedder()), + }, + LLMFactory: { + create: jest.fn(() => new MockLLM()), + }, + VectorStoreFactory: { + create: jest.fn((provider: string, config: any) => { + if (provider.toLowerCase() !== "memory") { + throw new Error( + `Unsupported vector store provider in test: ${provider}`, ); - const extracted = newMsgMatch - ? newMsgMatch[1].trim() - : "extracted fact from input"; - return JSON.stringify({ - memory: [ - { - id: "0", - text: extracted, - attributed_to: "user", - }, - ], - }); - }, - ), - })), -})); - -const mockEmbedding = new Array(1536).fill(0.1); -jest.mock("../src/embeddings/openai", () => ({ - OpenAIEmbedder: jest.fn().mockImplementation(() => ({ - embed: jest.fn().mockResolvedValue(mockEmbedding), - embedBatch: jest - .fn() - .mockImplementation((texts: string[]) => - Promise.resolve(texts.map(() => mockEmbedding)), - ), - embeddingDims: 1536, - })), -})); + } + return new MemoryVectorStore(config); + }), + }, + HistoryManagerFactory: { + create: jest.fn(() => new MemoryHistoryManager()), + }, + RerankerFactory: { + create: jest.fn(() => { + throw new Error("RerankerFactory is not used in memory.add.test.ts"); + }), + }, + }; +}); function createMemory(overrides: Partial = {}): Memory { return new Memory({ @@ -172,6 +195,138 @@ describe("Memory - add()", () => { ); }); + test("does not allow metadata to set identity scope", async () => { + const result: SearchResult = await memory.add("I am a software engineer", { + userId: "u1", + metadata: { + agent_id: "other", + agentId: "other-camel", + run_id: "other-run", + runId: "other-run-camel", + actor_id: "x", + source: "issue-6371", + nested: { preserved: true }, + }, + }); + const stored: MemoryItem | null = await memory.get(result.results[0].id); + + expect(stored).toEqual( + expect.objectContaining({ + user_id: "u1", + metadata: expect.objectContaining({ + source: "issue-6371", + nested: { preserved: true }, + }), + }), + ); + expect(stored).not.toHaveProperty("agent_id"); + expect(stored).not.toHaveProperty("run_id"); + expect(stored!.metadata).not.toHaveProperty("agent_id"); + expect(stored!.metadata).not.toHaveProperty("agentId"); + expect(stored!.metadata).not.toHaveProperty("run_id"); + expect(stored!.metadata).not.toHaveProperty("runId"); + expect(stored!.metadata).not.toHaveProperty("actor_id"); + }); + + test.each([ + ["userId", { userId: "u1" }, true, "user_id", "u1"], + ["agentId", { agentId: "a1" }, false, "agent_id", "a1"], + ["runId", { runId: "r1" }, true, "run_id", "r1"], + ] as const)( + "preserves typed %s scope while stripping conflicting metadata identities", + async (_mode, scope, infer, canonicalKey, canonicalValue) => { + const result: SearchResult = await memory.add("scoped content", { + ...scope, + infer, + metadata: { + user_id: "metadata-user", + userId: "metadata-user-camel", + agent_id: "metadata-agent", + agentId: "metadata-agent-camel", + run_id: "metadata-run", + runId: "metadata-run-camel", + actor_id: "metadata-actor", + ordinary: "preserved", + }, + }); + const stored: MemoryItem | null = await memory.get(result.results[0].id); + + expect(stored).toHaveProperty(canonicalKey, canonicalValue); + for (const key of ["user_id", "agent_id", "run_id"]) { + if (key !== canonicalKey) { + expect(stored).not.toHaveProperty(key); + } + } + expect(stored!.metadata).toEqual( + expect.objectContaining({ ordinary: "preserved" }), + ); + for (const key of [ + "user_id", + "userId", + "agent_id", + "agentId", + "run_id", + "runId", + "actor_id", + ]) { + if (key !== canonicalKey) { + expect(stored!.metadata).not.toHaveProperty(key); + } + } + }, + ); + + test.each([ + [true, "user_id", "filter-user"], + [false, "user_id", "filter-user"], + [false, "agent_id", "filter-agent"], + [false, "run_id", "filter-run"], + ] as const)( + "preserves infer=%s %s filters scope after sanitization", + async (infer, filterKey, filterValue) => { + const result: SearchResult = await memory.add("filter-scoped content", { + filters: { [filterKey]: filterValue }, + infer, + metadata: { + user_id: "metadata-user", + userId: "metadata-user-camel", + agent_id: "metadata-agent", + agentId: "metadata-agent-camel", + run_id: "metadata-run", + runId: "metadata-run-camel", + actor_id: "metadata-actor", + ordinary: "preserved", + }, + }); + const stored: MemoryItem | null = await memory.get(result.results[0].id); + + expect(stored).toHaveProperty(filterKey, filterValue); + for (const key of ["user_id", "agent_id", "run_id"]) { + if (key !== filterKey) { + expect(stored).not.toHaveProperty(key); + } + } + expect(stored!.metadata).toEqual( + expect.objectContaining({ + [filterKey]: filterValue, + ordinary: "preserved", + }), + ); + for (const key of [ + "user_id", + "userId", + "agent_id", + "agentId", + "run_id", + "runId", + "actor_id", + ]) { + if (key === filterKey) continue; + expect(stored!.metadata).not.toHaveProperty(key); + } + }, + ); + test("with infer=false skips LLM and stores messages directly", async () => { const result: SearchResult = await memory.add("Direct storage content", { userId, diff --git a/mem0/memory/main.py b/mem0/memory/main.py index f4f3c4bf0..1e0fc0c44 100644 --- a/mem0/memory/main.py +++ b/mem0/memory/main.py @@ -135,18 +135,30 @@ _SENSITIVE_SUFFIXES = ( ENTITY_PARAMS = frozenset({"user_id", "agent_id", "run_id"}) DELETE_ALL_BATCH_SIZE = 1000 -# Tenant-scoping fields that update() must never let caller-supplied metadata overwrite (issues #4490, #6277). +# Tenant-scoping fields that caller-supplied metadata must never set, on either the +# creation or the update path (issues #4490, #6277, #6655). _IDENTITY_KEYS = ENTITY_PARAMS | {"actor_id"} -def _strip_identity_keys(metadata: Dict[str, Any], existing_payload: Dict[str, Any]) -> Dict[str, Any]: - """Drop identity keys from caller metadata; they are immutable after creation (issues #4490, #6277).""" +def _strip_identity_keys( + metadata: Dict[str, Any], + existing_payload: Dict[str, Any], + *, + context: str = "update()", +) -> Dict[str, Any]: + """Drop identity keys from caller metadata; scope is set by the entity params, not metadata. + + On the update path `existing_payload` carries the memory's current scope, so + re-sending an identical value is silently accepted; only a changed value warns. + On the creation path there is no prior payload, so pass an empty dict and every + identity key present in `metadata` is dropped with a warning. + """ clean = {} for key, value in metadata.items(): if key not in _IDENTITY_KEYS: clean[key] = value elif value != existing_payload.get(key): - logger.warning(f"update(): ignoring metadata['{key}'] - identity fields are immutable after creation") + logger.warning(f"{context}: ignoring metadata['{key}'] - identity fields cannot be set through metadata") return clean @@ -315,7 +327,9 @@ def _build_filters_and_metadata( for flexible session scoping and optionally narrows queries to a specific `actor_id`. It returns two dicts: 1. `base_metadata_template`: Used as a template for metadata when storing new memories. - It includes all provided session identifier(s) and any `input_metadata`. + It includes all provided session identifier(s) and any `input_metadata`. Identity + scope is set from the entity params only; identity keys in `input_metadata` are + dropped, so freeform metadata cannot place a memory into an unrequested scope. 2. `effective_query_filters`: Used for querying existing memories. It includes all provided session identifier(s), any `input_filters`, and a resolved actor identifier for targeted filtering if specified by any actor-related inputs. @@ -343,7 +357,12 @@ def _build_filters_and_metadata( scoped to the provided session(s) and potentially a resolved actor. """ - base_metadata_template = deepcopy(input_metadata) if input_metadata else {} + # Identity scope is set below from the entity params only. Stripping the keys here + # stops caller metadata from placing a memory into a scope the caller did not pass, + # which the re-pins below cannot prevent for a param that was left unset (issue #6655). + base_metadata_template = ( + _strip_identity_keys(deepcopy(input_metadata), {}, context="add()") if input_metadata else {} + ) effective_query_filters = deepcopy(input_filters) if input_filters else {} # ---------- validate and add all provided session ids ---------- diff --git a/mem0/vector_stores/upstash_vector.py b/mem0/vector_stores/upstash_vector.py index 49c2bda79..62e7f108d 100644 --- a/mem0/vector_stores/upstash_vector.py +++ b/mem0/vector_stores/upstash_vector.py @@ -1,5 +1,6 @@ import logging -from typing import Dict, List, Optional +import re +from typing import Any, Dict, List, Optional from pydantic import BaseModel @@ -13,6 +14,23 @@ except ImportError: logger = logging.getLogger(__name__) +_SAFE_FILTER_KEY = re.compile(r"[a-zA-Z_][a-zA-Z0-9_]*\Z") + + +def _validate_filter(key: str, value: Any) -> None: + if not isinstance(key, str) or not _SAFE_FILTER_KEY.fullmatch(key): + raise ValueError(f"Invalid filter key: {key!r}") + if not isinstance(value, (str, int, float, bool)): + raise ValueError( + f"Filter value for {key!r} must be str, int, float, or bool, " + f"got {type(value).__name__}" + ) + if isinstance(value, str) and ('"' in value or "\\" in value): + raise ValueError( + f"Filter value for {key!r} contains prohibited characters " + f"(double quote or backslash): {value!r}" + ) + class OutputData(BaseModel): id: Optional[str] # memory id @@ -92,7 +110,9 @@ class UpstashVector(VectorStoreBase): ) def _stringify(self, x): - return f'"{x}"' if isinstance(x, str) else x + if isinstance(x, str): + return f'"{x}"' + return x def search( self, @@ -113,6 +133,9 @@ class UpstashVector(VectorStoreBase): List[OutputData]: Search results. """ + if filters: + for k, v in filters.items(): + _validate_filter(k, v) filters_str = " AND ".join([f"{k} = {self._stringify(v)}" for k, v in filters.items()]) if filters else None response = [] @@ -160,13 +183,16 @@ class UpstashVector(VectorStoreBase): Returns: List[OutputData]: Search results, or None if sparse/BM25 search is not supported. """ - try: - filters_str = ( - " AND ".join([f"{k} = {self._stringify(v)}" for k, v in filters.items()]) - if filters - else None - ) + if filters: + for k, v in filters.items(): + _validate_filter(k, v) + filters_str = ( + " AND ".join([f"{k} = {self._stringify(v)}" for k, v in filters.items()]) + if filters + else None + ) + try: response = self.client.query( data=query, top_k=top_k, @@ -252,6 +278,9 @@ class UpstashVector(VectorStoreBase): Returns: List[OutputData]: Search results. """ + if filters: + for k, v in filters.items(): + _validate_filter(k, v) filters_str = " AND ".join([f"{k} = {self._stringify(v)}" for k, v in filters.items()]) if filters else None info = self.client.info() diff --git a/tests/memory/test_main.py b/tests/memory/test_main.py index d00e953b0..cea7e3842 100644 --- a/tests/memory/test_main.py +++ b/tests/memory/test_main.py @@ -453,6 +453,95 @@ def test_update_memory_metadata_cannot_change_identity_fields(mocker, caplog): assert "ignoring metadata['user_id']" in caplog.text +_ATTACKER_ADD_METADATA = { + "agent_id": "victim-agent", + "run_id": "victim-run", + "actor_id": "victim-actor", + "category": "sports", +} + + +def _captured_add_metadata(memory, mocker, **add_kwargs): + """Run add() with the pipeline stubbed and return the metadata template it produced.""" + captured = {} + + def _capture(messages, metadata, filters, infer, **kwargs): + captured.update(metadata) + return [] + + mocker.patch.object(memory, "_add_to_vector_store", side_effect=_capture) + memory.add("I like coffee", infer=False, **add_kwargs) + return captured + + +def test_add_metadata_cannot_set_identity_fields(mocker, caplog): + """Regression (issue #6655): add() metadata must not inject identity scope. + + The caller scopes by user_id only, so the agent_id/run_id re-pins in + _build_filters_and_metadata never fire and cannot defend the payload. + """ + memory = _build_memory_instance(mocker, Memory) + + with caplog.at_level(logging.WARNING, logger="mem0.memory.main"): + metadata = _captured_add_metadata( + memory, mocker, user_id="attacker", metadata=dict(_ATTACKER_ADD_METADATA) + ) + + assert metadata["user_id"] == "attacker" + for key in ("agent_id", "run_id", "actor_id"): + assert key not in metadata, f"{key} was injected through add() metadata" + # Non-identity metadata is untouched. + assert metadata["category"] == "sports" + assert "ignoring metadata['agent_id']" in caplog.text + + +@pytest.mark.asyncio +async def test_async_add_metadata_cannot_set_identity_fields(mocker): + """Async counterpart of test_add_metadata_cannot_set_identity_fields.""" + memory = _build_memory_instance(mocker, AsyncMemory) + captured = {} + + async def _capture(messages, metadata, filters, infer, **kwargs): + captured.update(metadata) + return [] + + mocker.patch.object(memory, "_add_to_vector_store", side_effect=_capture) + await memory.add( + "I like coffee", + user_id="attacker", + metadata=dict(_ATTACKER_ADD_METADATA), + infer=False, + ) + + assert captured["user_id"] == "attacker" + for key in ("agent_id", "run_id", "actor_id"): + assert key not in captured, f"{key} was injected through async add() metadata" + assert captured["category"] == "sports" + + +def test_add_entity_params_still_set_scope(mocker): + """The documented top-level params remain the only way to set scope.""" + memory = _build_memory_instance(mocker, Memory) + + metadata = _captured_add_metadata( + memory, mocker, user_id="u1", agent_id="a1", run_id="r1", metadata={"category": "sports"} + ) + + assert metadata["user_id"] == "u1" + assert metadata["agent_id"] == "a1" + assert metadata["run_id"] == "r1" + assert metadata["category"] == "sports" + + +def test_add_without_metadata_is_unaffected(mocker): + """No metadata argument means no stripping and no behaviour change.""" + memory = _build_memory_instance(mocker, Memory) + + metadata = _captured_add_metadata(memory, mocker, user_id="u1") + + assert metadata == {"user_id": "u1"} + + @pytest.mark.asyncio async def test_async_update_memory_metadata_cannot_change_identity_fields(mocker): """Async counterpart of test_update_memory_metadata_cannot_change_identity_fields.""" diff --git a/tests/vector_stores/test_upstash_vector.py b/tests/vector_stores/test_upstash_vector.py index 5628028c7..b825cf80b 100644 --- a/tests/vector_stores/test_upstash_vector.py +++ b/tests/vector_stores/test_upstash_vector.py @@ -5,7 +5,7 @@ from unittest.mock import MagicMock, call, patch import pytest from mem0.configs.vector_stores.upstash_vector import UpstashVectorConfig -from mem0.vector_stores.upstash_vector import UpstashVector +from mem0.vector_stores.upstash_vector import UpstashVector, _validate_filter @dataclass @@ -228,6 +228,57 @@ def test_update_vector_with_embeddings(upstash_instance_with_embeddings): ) +def test_filter_rejects_dict_value(): + with pytest.raises(ValueError): + _validate_filter("user_id", {"$ne": ""}) + + +def test_filter_rejects_list_value(): + with pytest.raises(ValueError): + _validate_filter("user_id", ["alice", "bob"]) + + +def test_filter_rejects_invalid_key(): + with pytest.raises(ValueError): + _validate_filter("user_id; DROP", "alice") + + +def test_filter_accepts_scalars(): + _validate_filter("user_id", "alice") + _validate_filter("count", 42) + _validate_filter("score", 0.95) + _validate_filter("active", True) + + +def test_filter_rejects_double_quote_in_value(): + with pytest.raises(ValueError, match="prohibited characters"): + _validate_filter("user_id", 'alice" OR 1=1 --') + + +def test_filter_rejects_backslash_in_value(): + with pytest.raises(ValueError, match="prohibited characters"): + _validate_filter("user_id", "alice\\bob") + + +def test_search_rejects_dict_filter(upstash_instance): + with pytest.raises(ValueError): + upstash_instance.search( + query="test", vectors=[[0.1]], filters={"user_id": {"$ne": ""}} + ) + + +def test_keyword_search_raises_on_invalid_filter(upstash_instance): + with pytest.raises(ValueError): + upstash_instance.keyword_search( + query="test", filters={"user_id": 'alice" OR 1=1'} + ) + + +def test_filter_rejects_key_with_trailing_newline(): + with pytest.raises(ValueError, match="Invalid filter key"): + _validate_filter("user_id\n", "alice") + + def test_insert_vectors_with_embeddings_missing_data(upstash_instance_with_embeddings): vectors = [[0.1, 0.2, 0.3]] payloads = [{"name": "vector1"}] # Missing data field