Compare commits

..

2 Commits

Author SHA1 Message Date
kartik-mem0 782fa392a5 fix(deepseek-plugin): repair recall and ship 0.3.1 bundle 2026-09-11 16:13:01 +05:30
kartik-mem0 9d60fd9f4e test(deepseek-plugin): reproduce recall and activation regressions 2026-09-11 16:11:53 +05:30
128 changed files with 860 additions and 7061 deletions
+1 -2
View File
@@ -31,8 +31,7 @@ export class PlatformBackend implements Backend {
this.headers = {
Authorization: `Token ${config.apiKey}`,
"Content-Type": "application/json",
"X-Mem0-Source": "CLI",
"X-Mem0-Client": `mem0-cli-node/${CLI_VERSION}`,
"X-Mem0-Source": "cli",
"X-Mem0-Client-Language": "node",
"X-Mem0-Client-Version": CLI_VERSION,
};
+1 -2
View File
@@ -27,8 +27,7 @@ class PlatformBackend(Backend):
headers={
"Authorization": f"Token {config.api_key}",
"Content-Type": "application/json",
"X-Mem0-Source": "CLI",
"X-Mem0-Client": f"mem0-cli-python/{__version__}",
"X-Mem0-Source": "cli",
"X-Mem0-Client-Language": "python",
"X-Mem0-Client-Version": __version__,
},
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "API Reference Overview"
sidebarTitle: "Overview"
title: "Overview"
seo:
title: "API Reference Overview - Mem0"
icon: "terminal"
iconType: "solid"
description: "REST APIs for memory management, search, and entity operations"
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Delete Memory API Endpoint"
sidebarTitle: "Delete Memory"
title: 'Delete Memory'
seo:
title: "Delete Memory API Endpoint - Mem0"
description: "Delete a single memory by its unique memory ID from the Mem0 platform using the DELETE endpoint."
openapi: delete /v1/memories/{memory_id}/
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Update Memory API Endpoint"
sidebarTitle: "Update Memory"
title: 'Update Memory'
seo:
title: "Update Memory API Endpoint - Mem0"
description: "Update the content, metadata, timestamp, or expiration date of a single memory by its unique ID using the PUT endpoint."
openapi: put /v1/memories/{memory_id}/
---
@@ -1,6 +1,7 @@
---
title: "Add Organization Member API Endpoint"
sidebarTitle: "Add Member"
title: 'Add Member'
seo:
title: "Add Organization Member API Endpoint - Mem0"
description: "Add a new member to an organization with a specified role such as READER or OWNER access level."
openapi: post /api/v1/orgs/organizations/{org_id}/members/
---
@@ -1,6 +1,7 @@
---
title: "Get Organization Members API Endpoint"
sidebarTitle: "Get Members"
title: 'Get Members'
seo:
title: "Get Organization Members API Endpoint - Mem0"
description: "Retrieve a list of all members belonging to a specific organization on the Mem0 platform."
openapi: get /api/v1/orgs/organizations/{org_id}/members/
---
@@ -1,6 +1,7 @@
---
title: "Add Project Member API Endpoint"
sidebarTitle: "Add Member"
title: 'Add Member'
seo:
title: "Add Project Member API Endpoint - Mem0"
description: "Add a new member to a project with a specified role such as READER or OWNER access level."
openapi: post /api/v1/orgs/organizations/{org_id}/projects/{project_id}/members/
---
@@ -1,6 +1,7 @@
---
title: "Get Project Members API Endpoint"
sidebarTitle: "Get Members"
title: 'Get Members'
seo:
title: "Get Project Members API Endpoint - Mem0"
description: "Retrieve a list of all members belonging to a specific project on the Mem0 platform."
openapi: get /api/v1/orgs/organizations/{org_id}/projects/{project_id}/members/
---
+17
View File
@@ -3106,6 +3106,23 @@ Removes Sidekick and its start/stop hooks. Memory capture, recall, and six skill
<Tab title="DeepSeek Harness">
<Update label="2026-09-11" description="deepseek-plugin v0.3.1">
**Fixed:**
- Automatic recall uses the claimed human inbox message before Harness appends it to session history, so the first request in a fresh session can recall memories and later turns use the current prompt.
- Native `dsh.bundle` activation loads Mem0 when the package is installed into a Harness profile. `pnpm pack` builds the plugin and includes its activation patch.
- Automatic recall guidance points to the available memory search tool for deeper searches.
**Added:**
- Optional `memoryScope: workspace` applies the same canonical workspace filter to automatic recall, capture, and both memory tools. Missing workspace paths skip automatic operations and reject explicit tools without falling back to shared memory.
**Upgrade notes:**
- `memoryScope: user` remains the default and shares memories across workspaces. Existing memories are not migrated into workspace scopes; separate clones and worktrees have separate scopes when workspace mode is enabled.
- Export `MEM0_USER_ID` alongside `MEM0_API_KEY`. Remove the old manual `insert: mem0` overlay; customize the bundle's existing row with an `id: mem0` override. See [DeepSeek Harness setup](/integrations/deepseek-plugin).
- Capture extraction remains asynchronous. This patch does not change extraction policy or guarantee that every preference or numeric constraint becomes a memory.
</Update>
<Update label="2026-09-08" description="deepseek-plugin v0.3.0">
**Added:**
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Embedder Configuration Reference"
sidebarTitle: "Configurations"
title: Configurations
seo:
title: "Embedder Configuration Reference - Mem0"
description: "Reference for embedder configuration options in Mem0, including provider selection and model settings."
---
@@ -1,6 +1,7 @@
---
title: "AWS Bedrock as Embedding Provider"
sidebarTitle: "AWS Bedrock"
title: AWS Bedrock
seo:
title: "AWS Bedrock as Embedding Provider - Mem0"
description: "Configure AWS Bedrock as an embedding provider in Mem0 with IAM credentials and boto3 authentication."
---
@@ -1,6 +1,7 @@
---
title: "Azure OpenAI as Embedding Provider"
sidebarTitle: "Azure OpenAI"
title: Azure OpenAI
seo:
title: "Azure OpenAI as Embedding Provider - Mem0"
description: "Configure Azure OpenAI as an embedding provider in Mem0 with API key, deployment, and endpoint settings."
---
@@ -1,6 +1,7 @@
---
title: "Google AI as Embedding Provider"
sidebarTitle: "Google AI"
title: Google AI
seo:
title: "Google AI as Embedding Provider - Mem0"
description: "Configure Google AI as an embedding provider in Mem0 using Gemini models and the GOOGLE_API_KEY variable."
---
@@ -1,6 +1,7 @@
---
title: "LangChain as Embedding Provider"
sidebarTitle: "LangChain"
title: LangChain
seo:
title: "LangChain as Embedding Provider - Mem0"
description: "Use LangChain as an embedding provider in Mem0 to access a wide range of models through a unified interface."
---
@@ -1,6 +1,7 @@
---
title: "LM Studio as Embedding Provider"
sidebarTitle: "LM Studio"
title: "LM Studio"
seo:
title: "LM Studio as Embedding Provider - Mem0"
description: "Configure LM Studio as an embedding provider in Mem0 for local embedding generation with models like nomic-embed-text."
---
You can use embedding models from LM Studio to run Mem0 locally.
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Ollama as Embedding Provider"
sidebarTitle: "Ollama"
title: "Ollama"
seo:
title: "Ollama as Embedding Provider - Mem0"
description: "Configure Ollama as an embedding provider in Mem0 to generate embeddings locally using open-source models."
---
You can use embedding models from Ollama to run Mem0 locally.
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "OpenAI as Embedding Provider"
sidebarTitle: "OpenAI"
title: OpenAI
seo:
title: "OpenAI as Embedding Provider - Mem0"
description: "Configure OpenAI as an embedding provider in Mem0 using models like text-embedding-3-large for vector generation."
---
@@ -1,6 +1,7 @@
---
title: "Together AI as Embedding Provider"
sidebarTitle: "Together"
title: Together
seo:
title: "Together AI as Embedding Provider - Mem0"
description: "Configure Together AI as an embedding provider in Mem0 with support for 1024-dimensional embedding models."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Embedding Providers Overview"
sidebarTitle: Overview
title: Overview
seo:
title: "Embedding Providers Overview - Mem0"
description: "Overview of all supported embedding model providers in Mem0, including OpenAI, Azure, Ollama, and more."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "LLM Configuration Reference"
sidebarTitle: "Configurations"
title: Configurations
seo:
title: "LLM Configuration Reference - Mem0"
description: "Reference for LLM configuration options in Mem0 for Python and TypeScript, including value precedence rules."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "AWS Bedrock as LLM Provider"
sidebarTitle: "AWS Bedrock"
title: AWS Bedrock
seo:
title: "AWS Bedrock as LLM Provider - Mem0"
description: "Configure AWS Bedrock as an LLM provider in Mem0 with IAM authentication and Claude model support."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Azure OpenAI as LLM Provider"
sidebarTitle: "Azure OpenAI"
title: Azure OpenAI
seo:
title: "Azure OpenAI as LLM Provider - Mem0"
description: "Configure Azure OpenAI as an LLM provider in Mem0 with Azure Identity authentication and deployment settings."
---
+1 -1
View File
@@ -3,7 +3,7 @@ title: DeepSeek
description: "Configure DeepSeek as an LLM provider in Mem0 with API key setup and optional custom endpoint configuration."
---
To use DeepSeek LLM models, you have to set the `DEEPSEEK_API_KEY` environment variable. You can also optionally set `DEEPSEEK_API_BASE` if you need to use a different API endpoint (defaults to `https://api.deepseek.com`).
To use DeepSeek LLM models, you have to set the `DEEPSEEK_API_KEY` environment variable. You can also optionally set `DEEPSEEK_API_BASE` if you need to use a different API endpoint (defaults to "https://api.deepseek.com").
## Usage
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Google AI as LLM Provider"
sidebarTitle: "Google AI"
title: Google AI
seo:
title: "Google AI as LLM Provider - Mem0"
description: "Configure Google Gemini as an LLM provider in Mem0 using the google.genai SDK and GOOGLE_API_KEY variable."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "LangChain as LLM Provider"
sidebarTitle: "LangChain"
title: LangChain
seo:
title: "LangChain as LLM Provider - Mem0"
description: "Use LangChain as an LLM provider in Mem0 to integrate with various chat models through a unified interface."
---
+4 -3
View File
@@ -1,6 +1,7 @@
---
title: "LM Studio as LLM Provider"
sidebarTitle: "LM Studio"
title: LM Studio
seo:
title: "LM Studio as LLM Provider - Mem0"
description: "Configure LM Studio as an LLM provider in Mem0 for running local language models via an OpenAI-compatible API."
---
@@ -77,7 +78,7 @@ m.add(messages, user_id="alice123", metadata={"category": "movies"})
To use LM Studio, you need to:
1. Download and install [LM Studio](https://lmstudio.ai/)
2. Start a local server from the "Server" tab
3. Set the appropriate `lmstudio_base_url` in your configuration (default is usually `http://localhost:1234/v1`)
3. Set the appropriate `lmstudio_base_url` in your configuration (default is usually http://localhost:1234/v1)
</Note>
## Config
+1 -1
View File
@@ -3,7 +3,7 @@ title: MiniMax
description: "Configure MiniMax as an LLM provider in Mem0 with API key setup and optional custom endpoint configuration."
---
To use MiniMax LLM models, you have to set the `MINIMAX_API_KEY` environment variable. You can also optionally set `MINIMAX_API_BASE` if you need to use a different API endpoint (defaults to `https://api.minimax.io/v1`).
To use MiniMax LLM models, you have to set the `MINIMAX_API_KEY` environment variable. You can also optionally set `MINIMAX_API_BASE` if you need to use a different API endpoint (defaults to "https://api.minimax.io/v1").
## Usage
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Ollama as LLM Provider"
sidebarTitle: "Ollama"
title: Ollama
seo:
title: "Ollama as LLM Provider - Mem0"
description: "Configure Ollama as an LLM provider in Mem0 for running local language models with tool-calling support."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "OpenAI as LLM Provider"
sidebarTitle: "OpenAI"
title: OpenAI
seo:
title: "OpenAI as LLM Provider - Mem0"
description: "Configure OpenAI as an LLM provider in Mem0 with support for GPT models and Openrouter compatibility."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Together AI as LLM Provider"
sidebarTitle: "Together"
title: Together
seo:
title: "Together AI as LLM Provider - Mem0"
description: "Configure Together AI as an LLM provider in Mem0 with API key setup and optional custom endpoint configuration."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "xAI Grok as LLM Provider"
sidebarTitle: "xAI"
title: xAI
seo:
title: "xAI Grok as LLM Provider - Mem0"
description: "Configure xAI Grok models as an LLM provider in Mem0 with API key setup and usage examples."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "LLM Providers Overview"
sidebarTitle: Overview
title: Overview
seo:
title: "LLM Providers Overview - Mem0"
description: "Overview of all supported LLM providers in Mem0, including OpenAI, Anthropic, Groq, Ollama, and more."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Reranker Providers Overview"
sidebarTitle: "Overview"
title: Overview
seo:
title: "Reranker Providers Overview - Mem0"
description: 'Pick the right reranker path to boost Mem0 search relevance.'
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Vector Store Configuration Reference"
sidebarTitle: "Configurations"
title: Configurations
seo:
title: "Vector Store Configuration Reference - Mem0"
description: "Reference for vector database configuration options in Mem0, including provider selection and connection settings."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "LangChain as Vector Store Provider"
sidebarTitle: "LangChain"
title: LangChain
seo:
title: "LangChain as Vector Store Provider - Mem0"
description: "Use LangChain as a unified vector store provider in Mem0 to access multiple vector databases through one interface."
---
@@ -102,7 +102,7 @@ Here are the parameters available for configuring Upstash Vector:
| `url` | URL for the Upstash Vector index | `None` |
| `token` | Token for the Upstash Vector index | `None` |
| `client` | An `upstash_vector.Index` instance | `None` |
| `collection_name` | The default namespace used | `"mem0"` |
| `collection_name` | The default namespace used | `""` |
| `enable_embeddings` | Whether to use Upstash embeddings | `False` |
<Note>
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Vector Store Providers Overview"
sidebarTitle: "Overview"
title: Overview
seo:
title: "Vector Store Providers Overview - Mem0"
description: "Overview of all supported vector databases in Mem0, including Qdrant, Chroma, PGVector, Pinecone, Oracle, and more."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Cookbooks and Tutorials"
sidebarTitle: "Overview"
title: Overview
seo:
title: "Cookbooks and Tutorials - Mem0"
description: "Browse cookbook examples and tutorials for building AI applications with Mem0, from companion chatbots to AI agents."
---
@@ -1,6 +1,7 @@
---
title: "Delete Memory Operation"
sidebarTitle: "Delete Memory"
title: Delete Memory
seo:
title: "Delete Memory Operation - Mem0"
description: Remove memories from Mem0 either individually, in bulk, or via filters.
icon: "trash"
iconType: "solid"
@@ -1,6 +1,7 @@
---
title: "Update Memory Operation"
sidebarTitle: "Update Memory"
title: Update Memory
seo:
title: "Update Memory Operation - Mem0"
description: Modify an existing memory by updating its content or metadata.
icon: "pen-to-square"
iconType: "solid"
-1
View File
@@ -41,7 +41,6 @@
"pages": [
"platform/quickstart",
"platform/overview",
"platform/copilot",
"platform/agent-signup",
"vibecoding",
"platform/cli",
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Integrations Overview"
sidebarTitle: "Overview"
title: Overview
seo:
title: "Integrations Overview - Mem0"
description: "Overview of Mem0 integrations with popular AI frameworks and tools for persistent memory and context management."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "AWS Bedrock Integration"
sidebarTitle: "AWS Bedrock"
title: AWS Bedrock
seo:
title: "AWS Bedrock Integration with Mem0"
description: "Use Mem0 with AWS Bedrock and OpenSearch Service for cloud-native persistent semantic memory storage."
---
-30
View File
@@ -172,36 +172,6 @@ claude plugin update mem0@mem0-plugins --scope user
| Sidekick won't start | Must be in a Git repo. Check that your Claude Code version supports plugin agents and worktrees. |
| Remove the plugin | `claude plugin uninstall mem0@mem0-plugins` |
## Telemetry
The plugin sends usage events (which hook ran, timing, result counts, failure
types) so Mem0 can see what's used and what's breaking.
These events are **not anonymous**. When an API key is configured, which
installing the plugin requires, they are sent under your Mem0 account email,
the same way the Python SDK and the CLI attribute theirs. Without a key they
are sent under a random per-machine id.
Each event carries the event name, the plugin version, the harness it ran in,
your OS and Python version, and per-event properties describing what happened:
timings, counts, coarse outcome and failure labels, and which model was
configured. Repository and session identifiers are hashed with a random salt
generated on your machine, so they cannot be linked back to a repository name
or path.
The exact set is enforced in code rather than by this list: every property is
filtered through a denylist of sensitive keys and credential-shaped values are
redacted before anything is sent.
Prompts, memory text, queries, file paths, repository names, and API keys are
never sent.
Turn it off:
```bash
export MEM0_TELEMETRY=false
```
<CardGroup cols={2}>
<Card title="Mem0 MCP Setup" icon="puzzle-piece" href="/platform/mem0-mcp">
Detailed MCP configuration for all clients
+1 -1
View File
@@ -3,7 +3,7 @@ title: Codex
description: "Add persistent memory to OpenAI Codex with automatic capture, automatic recall, a search tool, and six memory skills."
---
Add persistent memory to [**OpenAI Codex**](https://openai.com/codex/) with the Mem0 plugin. Codex forgets everything between tasks. This plugin fixes that by connecting to Mem0's cloud memory layer via MCP, automatically capturing learnings at key lifecycle points, and retrieving relevant context on the first prompt of a session. Codex can use the search tool for recall later in the session.
Add persistent memory to [**OpenAI Codex**](https://openai.com/index/codex/) with the Mem0 plugin. Codex forgets everything between tasks. This plugin fixes that by connecting to Mem0's cloud memory layer via MCP, automatically capturing learnings at key lifecycle points, and retrieving relevant context on the first prompt of a session. Codex can use the search tool for recall later in the session.
<Info>Current plugin version: `0.3.1`.</Info>
+40 -24
View File
@@ -5,7 +5,7 @@ description: "Add persistent memory to DeepSeek Harness with automatic recall, a
Add persistent memory to the [**DeepSeek Harness**](https://github.com/deepseek-ai/deepseek-harness) with `@mem0/deepseek-plugin`. The plugin recalls relevant context before a model request, captures completed turns, and provides explicit Mem0 tools when the agent needs them.
<Info>Current package version: `0.3.0`.</Info>
<Info>Current package version: `0.3.1`.</Info>
Sidekick is available only in the [Claude Code plugin](/integrations/claude-code#sidekick-agent).
@@ -32,7 +32,7 @@ A Cordis plugin is a module exporting `apply(ctx, config)`. This plugin waits fo
Cordis removes the listeners and tools when the plugin unmounts. Memory failures are fail-open, so a Mem0 outage does not stop the agent from completing its normal work.
DeepSeek Harness provides subagents through separate host-composition packages. This Mem0 package does not register a named Sidekick or claim child filesystem isolation. A Harness child uses Mem0 only when its own agent preset includes the Mem0 plugin.
DeepSeek Harness provides subagents through separate host-composition packages. This Mem0 package does not register a named Sidekick or claim child filesystem isolation. Memory availability in children depends on the Harness composition and service scope.
## Prerequisites
@@ -40,25 +40,27 @@ DeepSeek Harness provides subagents through separate host-composition packages.
- <a href="https://app.mem0.ai?utm_source=oss&utm_medium=integration-deepseek-plugin">Sign up at app.mem0.ai</a>
- <a href="https://app.mem0.ai/dashboard/api-keys?utm_source=oss&utm_medium=integration-deepseek-plugin">Get your API key</a> (starts with `m0-`)
2. The DeepSeek Harness installed.
2. The DeepSeek Harness installed and a funded model-provider credential (`DEEPSEEK_API_KEY` for the official DeepSeek provider). This is separate from your Mem0 key.
3. Your API key exported in your shell:
3. Your Mem0 key and a stable user identity exported in the shell that launches Harness:
<CodeGroup>
```bash zsh
echo 'export MEM0_API_KEY="m0-your-api-key"' >> ~/.zshrc
echo 'export MEM0_USER_ID="your-user-id"' >> ~/.zshrc
source ~/.zshrc
```
```bash bash
echo 'export MEM0_API_KEY="m0-your-api-key"' >> ~/.bashrc
echo 'export MEM0_USER_ID="your-user-id"' >> ~/.bashrc
source ~/.bashrc
```
</CodeGroup>
## Try it locally
1. Build and pack the plugin:
1. Build and pack the plugin from the full repository (the source imports the sibling shared core):
```sh
cd integrations/deepseek-plugin
pnpm install --frozen-lockfile
@@ -70,33 +72,32 @@ source ~/.bashrc
2. Install it into a disposable Harness profile so Harness supplies its peer dependencies:
```sh
DSH_HOME=/tmp/mem0-dsh-dev pnpm dlx @deepseek-ai/dsh@0.1.1-rc.2 \
plugin --profile headless add /tmp/mem0-deepseek-plugin/mem0-deepseek-plugin-0.3.0.tgz
plugin --profile web add /tmp/mem0-deepseek-plugin/mem0-deepseek-plugin-0.3.1.tgz
```
3. Copy `cordis.example.yml`, set its installed package path and your `userId`, then load it with the same profile:
3. Launch the same profile. The package's `dsh.bundle` activates Mem0 automatically:
```sh
DSH_HOME=/tmp/mem0-dsh-dev pnpm dlx @deepseek-ai/dsh@0.1.1-rc.2 \
web --patch ./integrations/deepseek-plugin/cordis.example.yml
web
```
4. Open the web UI and ask the agent to remember something, then recall it in a later turn.
4. Open the web UI, select a workspace, and state a synthetic preference without mentioning Mem0. Allow extraction to finish, start a fresh session, and ask for that preference without tools. Inspect the `mem0:recall` context to confirm automatic recall.
The `cordis.yml` entry looks like this:
`pnpm pack` builds the package automatically and ships its activation patch. To customize it, copy `cordis.example.yml` and launch with `web --patch /absolute/path/to/cordis.yml`. The override updates the existing bundle row:
```yaml
- name: "@deepseek-ai/dsh-system-prompt"
- name: "@deepseek-ai/dsh-tools"
- insert:
- id: mem0
name: "/tmp/mem0-dsh-dev/profiles/headless/node_modules/@mem0/deepseek-plugin/dist/index.js"
config:
# apiKey is read from MEM0_API_KEY when omitted here.
userId: "your-user-id"
autoRecall: true
autoCapture: true
# host: "https://your-onprem.mem0.ai" # optional: Platform on-prem / dedicated base URL
- id: mem0
config:
# A config override replaces the whole config, so restate userId.
userId: !!js process.env.MEM0_USER_ID
memoryScope: workspace
autoRecall: true
autoCapture: true
# host: "https://your-onprem.mem0.ai" # optional: Platform dedicated base URL
```
When upgrading from a manual installation, remove the old patch that inserts a `mem0` row; the bundle now inserts it. Keep custom settings as an override by `id`.
For a Mem0 Platform on-prem or dedicated deployment, point `config.host` at that base URL (defaults to `api.mem0.ai`). `host` is a Platform base-URL override, not a switch to self-hosted Mem0 OSS.
## Configuration
@@ -105,16 +106,31 @@ For a Mem0 Platform on-prem or dedicated deployment, point `config.host` at that
|---|---|---|---|
| `apiKey` | no | `$MEM0_API_KEY` | Mem0 platform API key |
| `userId` | yes | | Default entity that owns the memories |
| `memoryScope` | no | `user` | `user` shares memory across workspaces; `workspace` isolates automatic and explicit memory operations |
| `allowUserOverride` | no | `false` | Permit model-selected access to a different user only in a trusted multi-user deployment |
| `host` | no | `api.mem0.ai` | Platform base URL (on-prem / dedicated) |
| `autoRecall` | no | `true` | Recall relevant memory before model requests |
| `autoCapture` | no | `true` | Store completed human and assistant turns |
Both tools also accept optional per-call `userId`, `agentId`, and `runId` params so a single install can partition memory by entity, agent, or session; `userId` defaults to the configured user, while `agentId` and `runId` are omitted unless supplied. Automatic recall and capture use the configured user without an agent, repository, or session filter. Automatic capture preserves full redacted user and assistant message text without a per-message character cutoff.
Both tools also accept optional per-call `userId`, `agentId`, and `runId` params so a single install can partition memory by entity, agent, or session; `userId` defaults to the configured user, while `agentId` and `runId` are omitted unless supplied. Automatic capture preserves full redacted user and assistant message text without a per-message character cutoff.
## Workspace scope
By default, automatic capture and recall use the configured user across all workspaces. This is intentional compatibility behavior. To isolate workspaces, set `memoryScope: workspace` as shown above. All writes then attach an `appId` derived from the canonical absolute session workspace path; all searches require its matching `app_id`. The model cannot override that workspace filter.
Symlink aliases share a scope. Different directories, clones, or worktrees have different scopes; moving a directory changes its scope. Missing or invalid workspace paths skip automatic memory operations and reject explicit tools without falling back to user-wide access. Existing user-only memories are not migrated into workspace scopes.
User scope searches all memories belonging to the user, including workspace-tagged memories. Workspace isolation applies only when workspace scope is configured; it is application filtering, not a separate Mem0 account or credential.
## Extraction and recall limits
Capture sends the full redacted turn to Mem0 for asynchronous extraction. A queued write, or even an initially nonempty memory list, does not mean every extracted fact is ready. Allow processing to finish and inspect stored memories before diagnosing lost preferences or missing filenames and numeric constraints. Exact extraction remains model-dependent; a one-off request should not automatically become a standing preference.
Automatic recall is a shallow first pass: up to five results, 4,000 context characters, and a two-second wait. The explicit `search_memory` tool can use a more focused query and returns up to ten results by default.
## Telemetry
Writes are tagged `source="DEEPSEEK_HARNESS"` so Mem0 can attribute usage to this integration. Usage events include operation names, durations, result counts, and coarse failure kinds. They are **not anonymous**: when an API key is configured they are sent under your Mem0 account email, the same way the SDK attributes its own. Queries, memory text, entity IDs, and API keys are never included. Set `MEM0_TELEMETRY=false` to opt out.
Writes are tagged `source="DEEPSEEK_HARNESS"` so Mem0 can attribute usage to this integration. Anonymous usage events include operation names, durations, result counts, and coarse failure kinds. Queries, memory text, entity IDs, and API keys are never included. Set `MEM0_TELEMETRY=false` to opt out.
<Note>
This plugin is a developer preview and tracks the evolving DeepSeek Harness plugin API.
@@ -124,6 +140,6 @@ Writes are tagged `source="DEEPSEEK_HARNESS"` so Mem0 can attribute usage to thi
- **`MISSING_CREDENTIAL` for `deepseek-official`**: Configure `DEEPSEEK_API_KEY` through Harness's Models page or export it in the shell that launches Harness.
- **`EMFILE: too many open files, watch` on macOS**: Launch Harness with `CHOKIDAR_USEPOLLING=1`.
- **Mem0 tools do not appear**: Run Harness with `--dump-config` and confirm the final composition contains the `mem0` row and the installed `dist/index.js` path.
- **Mem0 tools do not appear**: Run Harness with `--dump-config` for the same profile used at installation and confirm it contains a `mem0` row naming `@mem0/deepseek-plugin`. Export `MEM0_USER_ID` and `MEM0_API_KEY` before launching.
Per-call `userId` overrides are rejected unless the operator enables `allowUserOverride: true`. Automatic recall and capture always use the configured user.
+1 -1
View File
@@ -23,7 +23,7 @@ npm install -g flowise
npx flowise start
```
2. Access to the Flowise UI at `http://localhost:3000`
2. Access to the Flowise UI at http://localhost:3000
3. Basic familiarity with [Flowise's LLM orchestration](https://flowiseai.com/#features) concepts
## Setup and Configuration
+26 -61
View File
@@ -1,35 +1,39 @@
---
title: Hermes Agent
description: "Add long-term memory to Hermes agents using Mem0 Platform, a self-hosted server, or local OSS mode with background fact extraction."
description: "Add long-term memory to Hermes agents with Mem0, on managed Mem0 Cloud or fully self-hosted (OSS), with automatic background sync and zero-latency prefetch."
---
Add long-term memory to [Hermes Agent](https://github.com/NousResearch/hermes-agent), a self-improving AI agent CLI by Nous Research. Hermes has a pluggable memory system, and Mem0 is one of the supported providers. Once enabled, Mem0 learns facts from your conversations and surfaces relevant ones for the current question, without slowing down the chat.
Add long-term memory to [Hermes Agent](https://github.com/NousResearch/hermes-agent), a self-improving AI agent CLI by Nous Research. Hermes has a pluggable memory system, and Mem0 is one of the supported providers. Once enabled, Mem0 learns facts from your conversations and surfaces relevant ones before each turn, without slowing down the chat.
You can run Mem0 in three ways:
You can run Mem0 in two ways:
- **Platform mode** (default): managed Mem0 Cloud. Add your API key and you are ready.
- **Self-hosted server mode**: point the plugin at a Mem0 server you run yourself (the Docker-shipped server). The plugin only talks HTTP to your server.
- **OSS mode**: run Mem0 in-process with your own LLM, embedder, and vector store. No Mem0 server required.
- **OSS mode**: fully self-hosted with your own LLM, embedder, and vector store. No data leaves your machine.
## How It Works
Hermes runs a built-in memory system (file-based `MEMORY.md` and `USER.md`) alongside one external provider. When Mem0 is active, it works additively with the built-in system at two points in every conversation turn.
Hermes runs a built-in memory system (file-based `MEMORY.md` and `USER.md`) alongside one external provider. When Mem0 is active, it works additively with the built-in system at three points in every conversation turn.
### 1. Current-turn recall (bounded wait)
### 1. Before the agent responds (prefetch)
When you send a message, Hermes searches your stored memories for the current question and waits up to 3 seconds for results. If they arrive in time, they are injected into the system prompt so the model can see them. If the backend is slower, Hermes skips the injection and the model can still call `mem0_search` itself — so a slow backend never blocks a turn.
When you send a message, Hermes checks for cached Mem0 search results from the previous turn. If they exist, those memories are injected into the system prompt so the model can see them. This is zero-latency, with no waiting on an API call.
### 2. Background fact extraction (sync)
### 2. After the agent responds (sync)
Once the model finishes, Hermes sends the `(user message, assistant response)` pair to Mem0 in a background thread. Mem0 extracts facts automatically (for example, "user prefers Python" or "user works at Acme Corp"), so you never have to tell it what to remember. Each write is tagged with the gateway channel it came from.
### 3. Background prefetch for the next turn
At the same time, Hermes runs a background search to pre-load relevant memories for your next message. By the time you type, the results are already cached.
## Agent Tools
When Mem0 is active, the model gets four tools it can call during a conversation:
When Mem0 is active, the model gets five tools it can call during a conversation:
| Tool | Description | Parameters |
|------|-------------|------------|
| `mem0_search` | Semantic search by meaning, ranked by relevance | `query` (required), `top_k` (default 10, max 50), `rerank` (default `false`, Platform mode only) |
| `mem0_list` | List all stored memories, for a full overview | `page`, `page_size` (default 100, max 200) |
| `mem0_search` | Semantic search by meaning, ranked by relevance | `query` (required), `top_k` (default 10, max 50), `rerank` (default `true`, Platform mode only) |
| `mem0_add` | Store a fact verbatim, with no LLM extraction | `content` (required) |
| `mem0_update` | Update a memory's text by ID | `memory_id`, `text` (both required) |
| `mem0_delete` | Delete a memory by ID | `memory_id` (required) |
@@ -75,36 +79,6 @@ memory:
That's it. Mem0 runs automatically from here.
## Self-Hosted Server Setup
Run the [Mem0 server](https://github.com/mem0ai/mem0/tree/main/server) (FastAPI + pgvector) from its Docker image and point the plugin at it. Unlike OSS mode, the plugin just talks HTTP to your server.
### Interactive
```bash
hermes memory setup
# Select "mem0", then "Self-hosted server", and enter the server URL
```
### With flags
```bash
hermes memory setup mem0 --mode selfhosted \
--host http://localhost:8888 \
--api-key your-admin-api-key
```
### With environment variables
```bash
echo "MEM0_HOST=http://localhost:8888" >> ~/.hermes/.env
echo "MEM0_API_KEY=your-admin-api-key" >> ~/.hermes/.env
```
Then start a fresh Hermes session and call `mem0_search` — it connects to your server. The plugin authenticates with `X-API-Key` and uses the server's `/search` and `/memories` routes. The API key is optional only for servers running with `AUTH_DISABLED`.
<Note>Setting `host` routes to the self-hosted server automatically. Don't combine it with `mode: oss` — OSS takes precedence and ignores `host`.</Note>
## OSS (Self-Hosted) Setup
OSS mode runs Mem0 entirely on your own infrastructure: your LLM, your embedder, and your vector store. No data is sent to Mem0 Cloud, and no Mem0 API key is required.
@@ -137,20 +111,15 @@ hermes memory setup mem0 --mode oss \
| Flag | Description |
|------|-------------|
| `--mode` | `platform`, `selfhosted`, or `oss` |
| `--api-key` | Platform API key, or the admin key of a self-hosted server |
| `--host` | Self-hosted server URL (with `--mode selfhosted`) |
| `--mode` | `platform` or `oss` |
| `--oss-llm` | LLM provider (`openai` or `ollama`, default `openai`) |
| `--oss-llm-key` | LLM API key (for `openai`) |
| `--oss-llm-model` | Override the LLM model |
| `--oss-llm-url` | LLM base URL (for `ollama` or a custom endpoint) |
| `--oss-embedder` | Embedder provider (default `openai`) |
| `--oss-embedder-key` | Embedder API key |
| `--oss-embedder-model` | Override the embedder model |
| `--oss-embedder-url` | Embedder base URL (for `ollama` or a custom endpoint) |
| `--oss-vector` | Vector store (`qdrant` or `pgvector`, default `qdrant`) |
| `--oss-vector-path` | Local Qdrant storage path |
| `--oss-vector-url` | Qdrant server URL |
| `--oss-vector-host`, `--oss-vector-port` | PGVector or remote Qdrant host and port |
| `--oss-vector-user`, `--oss-vector-password`, `--oss-vector-dbname` | PGVector connection details |
| `--user-id` | Canonical user identifier |
@@ -158,7 +127,7 @@ hermes memory setup mem0 --mode oss \
## Switching Modes
You can move between the three modes at any time. Run the setup command again, or edit `~/.hermes/mem0.json` directly.
You can move between Platform and OSS at any time. Run the setup command again, or edit `~/.hermes/mem0.json` directly.
```bash
# Platform to OSS
@@ -167,9 +136,6 @@ hermes memory setup mem0 --mode oss --oss-llm-key sk-...
# OSS to Platform
hermes memory setup mem0 --mode platform --api-key sk-...
# Platform to a self-hosted server
hermes memory setup mem0 --mode selfhosted --host http://localhost:8888
# Preview without writing anything
hermes memory setup mem0 --mode oss --oss-llm-key sk-... --dry-run
```
@@ -180,7 +146,7 @@ A self-hosted `~/.hermes/mem0.json` looks like this:
{
"mode": "oss",
"oss": {
"llm": {"provider": "openai", "config": {"model": "gpt-5-mini", "is_reasoning_model": true}},
"llm": {"provider": "openai", "config": {"model": "gpt-5-mini"}},
"embedder": {"provider": "openai", "config": {"model": "text-embedding-3-small"}},
"vector_store": {"provider": "qdrant", "config": {"path": "~/.hermes/mem0_qdrant"}}
}
@@ -193,12 +159,11 @@ Behavioral settings live in `~/.hermes/mem0.json` and are written for you by `he
| Key | Default | Description |
|-----|---------|-------------|
| `mode` | `platform` | `platform` (Mem0 Cloud) or `oss` (self-managed, in-process). Self-hosted server routing is set via `host` |
| `host` | none | Self-hosted Mem0 server URL. When set, the plugin talks HTTP to your server instead of the cloud |
| `api_key` | none | Mem0 Platform API key, or the admin key of a self-hosted server. Stored in `.env` as `MEM0_API_KEY` |
| `mode` | `platform` | `platform` (Mem0 Cloud) or `oss` (self-hosted) |
| `api_key` | none | Mem0 Platform API key, required in Platform mode. Stored in `.env` as `MEM0_API_KEY` |
| `user_id` | `hermes-user` | Identifier that scopes memories. See cross-channel behavior below |
| `agent_id` | `hermes` | Agent identifier attached to writes |
| `rerank` | `false` | Rerank search results for relevance (Platform mode only) |
| `rerank` | `true` | Rerank search results for relevance (Platform mode only) |
### Cross-channel memories
@@ -209,11 +174,12 @@ Hermes can run from the CLI and from gateways like Telegram, Slack, and Discord.
Either way, every write is tagged with `metadata.channel` (for example `telegram` or `cli`), so per-channel views are still possible at query time.
## Reliability
- **Circuit breaker**: if Mem0 fails five times in a row, Hermes pauses calls for two minutes, then retries. The agent keeps working without memory during that window. Expected client errors, like a 404 on a missing memory id, do not count toward tripping the breaker.
- **Non-blocking**: fact extraction runs in a background daemon thread, and current-turn recall waits at most 3 seconds, so a slow or failed call never blocks your conversation.
- **Thread-safe**: the client uses lazy initialization with locking, and the background sync and recall threads are guarded so concurrent gateway messages cannot produce duplicate memories.
- **Non-blocking**: every Mem0 call runs in a background daemon thread, so a slow or failed call never blocks your conversation.
- **Thread-safe**: the client uses lazy initialization with locking, and the background sync and prefetch threads are guarded so concurrent gateway messages cannot produce duplicate memories.
## Troubleshooting
@@ -222,7 +188,6 @@ Either way, every write is tagged with `metadata.channel` (for example `telegram
The circuit breaker tripped after five consecutive failures and resets after two minutes.
- **Platform mode**: check your API key and internet connection.
- **Self-hosted server mode**: check that the server is running and reachable at the configured `host` URL.
- **OSS mode**: make sure your vector store (Qdrant or PGVector) is running and reachable.
### OSS: vector store connection refused
@@ -252,8 +217,8 @@ curl http://localhost:11434/api/tags
## Key Features
1. **Three ways to run**: managed Platform, a self-hosted server, or fully local OSS, switchable at any time.
2. **Current-turn recall**: memories for the current question are injected within a 3-second window, with `mem0_search` as the model's own backstop.
1. **Two ways to run**: managed Platform or fully self-hosted OSS, switchable at any time.
2. **Zero-latency recall**: memories are prefetched in the background and cached before you type.
3. **Automatic extraction**: Mem0 extracts and deduplicates facts from each exchange for you.
4. **Non-blocking and fault tolerant**: background threads plus a circuit breaker keep the agent responsive even when Mem0 is unreachable.
5. **Additive memory**: works alongside Hermes' built-in file memory (`MEMORY.md`, `USER.md`).
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "LangChain Integration"
sidebarTitle: "Langchain"
title: Langchain
seo:
title: "LangChain Integration with Mem0"
description: "Build personalized AI agents using LangChain for conversation flow and Mem0 for long-term memory retention."
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "n8n Integration"
sidebarTitle: "n8n"
title: n8n
seo:
title: "n8n Integration with Mem0"
description: "Add long-term memory to n8n workflows and AI Agents with the Mem0 community node, no code required."
---
+1 -3
View File
@@ -481,9 +481,7 @@ Plugin config is stored in `~/.openclaw/openclaw.json` with file permissions `0o
### Telemetry
Usage telemetry (PostHog) is enabled by default to help improve the plugin. No conversation content or memory values are included, only event counts (recall, capture, tool usage, CLI commands).
These events are **not anonymous**. OpenClaw does not send your account email the way the SDK does, but it does send an unsalted SHA-256 hash of it, falling back to a hash of the API key and then to a random per-machine id. Mem0 holds the email the hash is derived from, so the hash identifies your account rather than concealing it. The first run that resolves an account also emits a PostHog `$identify`, which permanently merges any earlier random id into that identity.
Anonymous usage telemetry (PostHog) is enabled by default to help improve the plugin. No conversation content or memory values are included, only event counts (recall, capture, tool usage, CLI commands).
To opt out, set the environment variable:
-1
View File
@@ -171,7 +171,6 @@ If the user is on a pre-current major (Python < 2, TS < 3, or a Platform call st
- [Introduction](https://docs.mem0.ai/introduction) [Both]: Use when the user wants a one-page overview of how memory fits between the LLM and the app.
- [Vibe Code with Mem0](https://docs.mem0.ai/vibecoding) [Both]: Use when the user is in Claude Code, Cursor, or Windsurf and wants memory wired into their editor.
- [Platform Overview](https://docs.mem0.ai/platform/overview) [Platform]: Use when the user picks the managed product - 4-line integration, hosted API, dashboard.
- [Mem0 Copilot](https://docs.mem0.ai/platform/copilot) [Platform]: Use when inspecting project memories, reviewing configuration changes, or testing extraction in the dashboard.
- [Sign up as an agent](https://docs.mem0.ai/platform/agent-signup) [Platform]: Use when an AI agent needs to mint a Mem0 API key autonomously - four commands, no email or dashboard, human claims ownership later.
- [Platform vs Open Source](https://docs.mem0.ai/platform/platform-vs-oss) [Both]: Use when the user is deciding between managed and self-hosted.
- [Platform Quickstart](https://docs.mem0.ai/platform/quickstart) [Platform]: Use for the first Platform integration - API key plus `MemoryClient.add/search`.
@@ -1,6 +1,7 @@
---
title: "Open Source Custom Instructions"
sidebarTitle: "Custom Instructions"
title: Custom Instructions
seo:
title: "Open Source Custom Instructions - Mem0"
description: Tailor fact extraction so Mem0 stores only the details you care about.
icon: "wand-magic-sparkles"
---
@@ -1,6 +1,7 @@
---
title: "Open Source Multimodal Support"
sidebarTitle: "Multimodal Support"
title: Multimodal Support
seo:
title: "Open Source Multimodal Support - Mem0"
description: Capture and recall memories from both text and images.
icon: "image"
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Open Source Features Overview"
sidebarTitle: "Overview"
title: "Overview"
seo:
title: "Open Source Features Overview - Mem0"
description: "Self-hosting features that extend Mem0 beyond basic memory storage"
icon: "list"
---
+1 -1
View File
@@ -56,7 +56,7 @@ The Mem0 REST API server exposes every OSS memory operation over HTTP. Run it al
make bootstrap # starts Compose, creates an admin, issues the first API key
```
Or to start the stack only and finish setup via the browser wizard at `http://localhost:3000`:
Or to start the stack only and finish setup via the browser wizard at http://localhost:3000:
```bash
cd server
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Open Source Overview"
sidebarTitle: "Overview"
title: "Overview"
seo:
title: "Mem0 Open Source Overview"
description: "Self-host Mem0 with full control over your infrastructure and data"
icon: "house"
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "CLI for Terminal Memory Management"
sidebarTitle: "CLI"
title: CLI
seo:
title: "Mem0 CLI for Terminal Memory Management"
description: "Manage memories from your terminal, for both humans and AI agents."
icon: "terminal"
iconType: "solid"
-194
View File
@@ -1,194 +0,0 @@
---
title: "Mem0 Copilot"
description: "Inspect project memories, review configuration changes, and test extraction from the Mem0 dashboard."
icon: "message"
---
Copilot is an AI assistant in the [Mem0 dashboard](https://app.mem0.ai/dashboard/copilot). Ask it to inspect stored memories, suggest extraction rules and categories, change project settings, or explain the platform SDKs and APIs.
You describe the task in chat. Copilot uses tools to read project data or propose changes. In **Review changes** mode, you approve each change before it runs.
## Start with your project
1. Sign in to the dashboard and select the organization and project you want to work on.
2. Open **Copilot** in the sidebar.
3. Select **Review changes** below the message box.
4. Ask: “Show my current project settings and explain what each one does.”
5. Expand an activity card, such as **Read project config**, to see its **Input** and **Result**. Check these results when reviewing an answer or confirming a change.
You can ask about settings or SDK usage with an empty project. Suggestions based on stored data need enough project history to analyze.
## Inspect stored memories
Start with a project overview, then narrow the question:
| Prompt | What to inspect |
| --- | --- |
| “Summarize what this project has stored and how it is categorized.” | The memories and category patterns Copilot reads. Check which records support its summary. |
| “Analyze the memories for user alice.” | Facts stored for `alice`, their categories, and unwanted information. To investigate a missing fact, also provide the original input. |
| “Search alice's memories for dietary preferences.” | Memories relevant to the query within that user's scope. |
Copilot can list memories across the project. Searching by meaning needs a specific User, Agent, App, or Run ID. Select one in **Scope** or name it in your prompt.
Memory lists are paginated, and suggestions use samples. Check the activity results to see which records were read. Ask for more pages when you need to inspect the rest.
### Choose the memory scope
Open **Scope** below the message box. Select an existing ID or type one and choose it.
| Control | Identifies |
| --- | --- |
| **User** | The end user whose memories you want to inspect, such as `alice`. |
| **Agent** | An AI agent associated with the memories. |
| **App** | An application associated with the memories. |
| **Run** | A particular execution or session associated with the memories. |
A field set to **any** adds no filter for that entity type. **Clear scope** removes the selections. Scope gives Copilot default IDs to use; you can request a different entity in a message. Check the activity's **Input** to confirm which IDs it used. Adding a test memory requires a **User** ID, even when another entity is selected.
<Note>
Scope selects memories, not separate settings. Categories, extraction instructions, memory depth, and multilingual settings apply to the selected project. Category and extraction suggestions also analyze project data, regardless of the selected entity.
</Note>
See [Entity-scoped memory](/platform/features/entity-scoped-memory) for how these IDs organize memories in your application.
## Review and apply changes
The mode control below the message box determines when changes run:
- **Review changes**: Copilot pauses before changing project configuration or adding a test memory. Read the proposal and choose whether to apply it.
- **Auto-apply**: Copilot can change settings and add test memories without asking for approval. Check the selected project and scope before using it.
Ask Copilot to show the current settings before requesting an update. In the approval card, review every listed field. Approving the card applies the whole proposal, including fields you did not edit.
| Action | Effect |
| --- | --- |
| **Apply change** | Approve the proposed write. Inspect the subsequent result to confirm it succeeded. |
| **Edit** | When offered, edit extraction instructions or category names and descriptions. Choose **Apply with edits** to submit your revision. |
| **Discard edits** | Return to the original proposal without applying it. |
| **Decline** | Reject this proposal without applying it. |
| **Ask for something else** | Enter feedback and choose **Send**. This declines the current proposal and sends your feedback as a new message. Review the next proposal before applying it. |
Not every field has an inline editor. If a proposal includes both extraction instructions and categories, only the instructions get an editor. The editor cannot apply blank instructions or an empty category list. Use **Ask for something else** to change other fields or request separate proposals.
<Warning>
A generated prompt profile contains separate prompts for extraction and summaries. If your project uses one, saving extraction instructions through Copilot deactivates it. Review the replacement rules and test the affected behavior before using them in production.
</Warning>
## What Copilot cannot do
Copilot cannot perform these actions, even in **Auto-apply** mode:
- Create or delete API keys.
- Directly edit or delete existing memories.
- Export project data.
- Invite or remove organization or project members, or change their roles and permissions.
- Create or delete organizations or projects.
- Change billing details or subscription plans.
Use the dashboard or the relevant platform API for these tasks. Copilot can explain the steps or point you to documentation, but it cannot carry out the actions.
Copilot supports the managed Mem0 platform. It does not help configure or use the [self-hosted open-source library](/open-source/overview).
## Example 1: Stop storing small talk, then test extraction
This walkthrough uses a project where you have noticed greetings or small talk in stored memories. Use a development project when trying configuration changes, since a test user does not isolate project settings.
<Steps>
<Step title="Inspect the problem">
Select the project, set **User** to `alice`, and ask:
> Analyze the memories for user alice.
Inspect the returned memories. Identify examples of small talk you want to exclude and useful facts you still want to keep.
</Step>
<Step title="Request a targeted change">
With **Review changes** selected, ask:
> Update my extraction instructions to ignore greetings and small talk. Keep the existing rules for durable user preferences.
Check for a **Read project config** activity before reviewing the update. If it is missing, ask Copilot to read the current instructions first. If analysis reports insufficient data, ask it to use the rule you provided. You can also set the rule directly through [Custom instructions](/platform/features/custom-instructions).
</Step>
<Step title="Review the proposal">
Review the full replacement text. Keep existing rules your application needs. For this test, the relevant rules could look like:
```text
Remember the following:
* The user's dietary preferences and food restrictions.
Don't remember the following:
* Greetings and small talk.
```
Use **Edit** to adjust the wording, or **Ask for something else** to request a revision.
Choose **Apply change** or **Apply with edits**, then inspect the result to confirm the instructions were saved.
</Step>
<Step title="Add a test conversation">
Choose a new **User** ID for this test, such as `copilot-extraction-test-01`, and clear any other entity selections. Ask:
> Add this as a user message for copilot-extraction-test-01: “Hi! How's your day? I prefer vegetarian meals and avoid peanuts.”
Review the test-memory proposal and choose **Apply change**.
<Warning>
Adding a test memory writes to the selected project. It is not a dry run. Use a dedicated test user so you can find and remove the test data afterward.
</Warning>
</Step>
<Step title="Wait for extraction, then check the result">
Expand **Add test memory** and inspect **Result**. A `PENDING` status means extraction is still running. Use the returned `event_id` with the [Get Event API](/api-reference/events/get-event) to check progress. Wait for `SUCCEEDED` before evaluating the output. If the event is `FAILED`, inspect its details before retrying.
Then ask:
> List all memories for user copilot-extraction-test-01. Then search that user's memories for dietary preferences.
Check that the memories retain the vegetarian preference and peanut restriction, and exclude the greeting. Inspect the full list as well as search results: a relevant search can hide unwanted small-talk records.
If the output is wrong, describe the mismatch and review another instruction change. Use a new test user for the next attempt so earlier memories do not affect the result. Remove test data afterward through the dashboard or [memory deletion API](/api-reference/memory/delete-memories).
</Step>
</Steps>
Test with new inputs after changing settings. Updating configuration is not a cleanup operation for existing memories. See [Custom instructions](/platform/features/custom-instructions) for guidance on writing extraction rules.
## Example 2: Suggest categories and extraction instructions
These workflows sample the selected project's data. Generating a suggestion does not update settings by itself. Copilot can then propose applying it, which follows the selected review mode. Use **Review changes** to inspect suggestions before they are saved.
| Workflow | Data used | Example prompt |
| --- | --- | --- |
| Custom categories | Stored memories, excluding deleted memories. | “Suggest custom categories from this project's memories.” |
| Extraction instructions | Successful requests to add conversation data, including requests that produced no memories. | “Compare recent add requests with the extracted memories and suggest better extraction instructions.” |
Category suggestions identify recurring themes. Review the names, descriptions, and examples. Applying a category list replaces the project's previous list; it does not add to it or re-tag existing memories. See [Custom categories](/platform/features/custom-categories) for project and per-call behavior.
Extraction suggestions compare conversation inputs with what was extracted. One add request can produce several memories or none, so the number of add requests is different from the number of stored memories. Keep your application's existing requirements when reviewing the proposal.
Both workflows require enough project data. If a suggestion fails because there is too little data, check the activity's **Result** for the current and required counts. You can still inspect memories, ask SDK/API questions, or provide your own rules instead of requesting a data-based suggestion.
## Example 3: Adjust detail and language
Ask Copilot to read the current settings, then request the change you need:
- **Memory depth** controls the level of detail: **Less**, **Medium**, or **More detailed**. Try “Show my current memory depth, then propose More detailed memories.” For deciding *which facts* to retain, use extraction instructions.
- **Multilingual** behavior preserves the user's original language. Try “Enable multilingual behavior so memories preserve the language of the input.” Test it with a new conversation in the language your application uses.
Both are project settings. Review the proposed values and test a fresh input after applying them. See [Organization and project settings](/api-reference/organizations-projects) for configuration through the API.
## Example 4: Ask SDK and API questions
Include your language and the task, for example:
> Show me how to search memories for user alice with the managed Mem0 Python SDK. Link the documentation you used.
Copilot can look up the official platform documentation. Check the **Browse mem0 docs** and **Read docs page** activities and open the cited pages. If an answer has no sources, ask for them before using the example. The [Platform quickstart](/platform/quickstart) covers adding and searching memories in code.
## Resume a conversation and check message limits
Use **History** to reopen your 50 most recently active chats for the selected project. Chats belong to the person who created them; other project members cannot open them. Use **New chat** to start a separate conversation, or **Delete chat** in history to remove one. Deleting a chat does not undo settings changes or remove test memories.
Message allowances depend on the organization's plan and are shared across its projects and members. They reset at the start of each calendar month in UTC.
For plans with a limit, a usage notice appears once 80% of the allowance is used. Below that point, no counter is shown. At the limit, new messages are disabled until the allowance resets or the plan is upgraded. Follow the notice's upgrade link to review plan options.
Sending feedback through **Ask for something else** counts as a new message. Approving or declining a saved proposal without feedback does not use another message. Test additions and searches also use the platform APIs and remain subject to their normal quotas.
@@ -1,6 +1,7 @@
---
title: "Platform Custom Instructions"
sidebarTitle: "Custom Instructions"
title: Custom Instructions
seo:
title: "Platform Custom Instructions - Mem0"
description: 'Control how Mem0 extracts and stores memories using natural language guidelines'
---
@@ -1,6 +1,7 @@
---
title: "Platform Multimodal Support"
sidebarTitle: "Multimodal Support"
title: Multimodal Support
seo:
title: "Platform Multimodal Support - Mem0"
description: Integrate images and documents into your interactions with Mem0
---
+3 -2
View File
@@ -1,6 +1,7 @@
---
title: "Platform Overview"
sidebarTitle: "Overview"
title: "Overview"
seo:
title: "Mem0 Platform Overview"
description: "Managed memory layer for AI agents, production-ready in minutes"
icon: "cloud"
---
+1 -1
View File
@@ -28,7 +28,7 @@
"clsx": "^2.1.1",
"js-cookie": "^3.0.6",
"lucide-react": "^0.477.0",
"next": "15.5.24",
"next": "15.5.21",
"react": "^19.0.0",
"react-dom": "^19.0.0",
"react-markdown": "^10.0.1",
-43
View File
@@ -49,48 +49,6 @@ Run the type check after every TypeScript change: `pnpm run typecheck` or `tsc -
- **`zapier-mem0/`** is a Zapier Platform CLI app: add, search, get, delete. It deploys to Zapier, not npm, so it is **not** in the release router. Deploy it with `gh workflow run zapier-mem0-cd.yml --ref main` (needs the `ZAPIER_DEPLOY_KEY` secret).
- **`mem0-strands/`** is a native Strands `MemoryStore` (Python, published to PyPI as `mem0-strands`). It plugs into the Strands `MemoryManager` for automatic recall and server-side extraction, over the hosted Mem0 platform or self-hosted Mem0 OSS. The package lives under `mem0-strands/python/`.
## Surface attribution
Every integration tells the Mem0 platform which surface it is. Three headers,
and the rules on them are what keep one layer from erasing another:
| Header | Carries | Rule |
|--------|---------|------|
| `X-Mem0-Source` | one canonical source value | **set-once** — write only if absent |
| `X-Application` | the host app it runs inside | **set-once** — write only if absent |
| `X-Mem0-Client` | `name/version`, outermost first | **append-only** — add yourself, never replace |
Set-once means check-then-set, never assignment. An integration that wraps the
SDK is the outermost layer and sets the source; the SDK underneath defers to it.
Assignment is exactly how every agent plugin came to be indistinguishable from
every other one at the platform.
How to declare it from an integration, in order of preference:
1. Send the headers yourself, if you make the HTTP call directly.
2. Pass `source` in the call options, if you go through an SDK.
3. Set `MEM0_SOURCE` / `MEM0_APPLICATION` / `MEM0_CLIENT_STACK` in the
environment before constructing the client. The SDKs read these and defer to
anything already present.
Append-only applies where a stack can actually form: an SDK handed a client that
already carries `X-Mem0-Client` appends itself rather than replacing. An SDK
constructed with no outer context simply reports itself, which is correct — it
is the outermost layer in that process.
The backend recognizes a fixed list of source values and buckets everything else
into `OTHERS`. A new value has to land in the platform's `EventSource` enum, so
do not invent one without that change going in too.
`X-Application` is allowlisted the same way, and this one has a rule of its own:
**omit the header when you do not know the host.** A value outside the allowlist
is discarded server-side, so guessing produces an event that claims an
attribution we do not actually have. The portable bundle is the case that
matters. It runs in whatever editor a user drops it into, so its build leaves
`PLATFORM_APPLICATION` empty and `memory_core` sends no header at all, while the
native bundles each name the host they were generated for. If you add a build
target, decide which of those two it is.
## Adding an integration
1. For a native coding-agent host, add `integrations/<name>-plugin/` with `plugin-build.json`, its manifest, and a thin adapter, then generate its shared runtime. Portable clients use the single `mem0-agent-plugin/` package. Independent TypeScript integrations stay self-contained and import shared lifecycle behavior from `agent-plugin-core/typescript/`.
@@ -101,4 +59,3 @@ target, decide which of those two it is.
5. If it is a Claude Code or editor marketplace plugin, register the generated native bundle path in the applicable marketplace files. Preserve the existing public plugin name.
6. Document it under `docs/integrations/` and add the page to `docs/docs.json` and `docs/llms.txt`.
7. Add rows to the table above and to the CI/CD tables in [`../.github/AGENTS.md`](../.github/AGENTS.md).
8. Send the three headers in [Surface attribution](#surface-attribution), and land the matching `EventSource` value on the platform in the same week. Until it exists, your traffic reports as `OTHERS`.
@@ -81,38 +81,6 @@ def replace_output(staged: Path, output: Path) -> Path:
return output
def _render_harness_id(host: str, *, portable: bool = False) -> str:
"""Emit core/_harness_id.py for one host.
Carries both vocabularies from a single definition: the PostHog `source` tag
and the platform's X-Mem0-Source / X-Application pair. Keeping them together
is what stops the two from drifting into separate vocabularies for the same
thing.
The portable bundle runs in whatever editor a user drops it into, so it does
not know its host and must not guess one. HARNESS_ID stays "coding-agent",
which is true and useful for grouping in PostHog, but PLATFORM_APPLICATION is
left empty: X-Application names a real host app, is checked against an
allowlist server-side, and a value that is always discarded is worse than no
value -- it reads like an attribution we have and do not.
"""
tag = host.upper().replace("-", "_") + "_PLUGIN"
application = "" if portable else host
return (
'"""Generated by integrations/agent-plugin-core/build/build.py. Do not edit."""\n'
"\n"
f'HARNESS_ID = "{host}"\n'
f'SOURCE_TAG = "{tag}"\n'
"\n"
"# Platform-side vocabulary (mem0_event.source + X-Application). The whole\n"
"# plugin family is one source; which editor it runs in is the application.\n"
"# An empty application means the host is unknown, and memory_core omits\n"
"# the header entirely rather than sending a placeholder.\n"
'PLATFORM_SOURCE = "MEM0_PLUGIN"\n'
f'PLATFORM_APPLICATION = "{application}"\n'
)
def _bundle_python(
staged: Path,
host: str,
@@ -128,12 +96,6 @@ def _bundle_python(
continue
shutil.copy2(source, core / source.name)
# Generated per host so identity does not depend on an entrypoint remembering
# to call telemetry.init(). mcp_server.py and the detached telemetry.py sender
# never did, which is how MCP searches reported harness=generic and every
# batch they drained was labelled MEM0_PLUGIN regardless of the real host.
(core / "_harness_id.py").write_text(_render_harness_id(host, portable=portable), encoding="utf-8")
values = {
"PLUGIN_ROOT": plugin_root,
"PLUGIN_DATA": "${PLUGIN_DATA}",
@@ -290,11 +290,6 @@ def run(
if args.plugin_data_dir:
os.environ[data_dir_env] = args.plugin_data_dir
# Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key
# writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking
# after them always saw content and every fresh install reported an upgrade.
data_dir_was_empty = telemetry.data_dir_was_empty()
cache_plugin_api_key()
if args.action == "session-start":
clear_stale_api_key_cache()
@@ -310,19 +305,8 @@ def run(
return 0
if args.action == "session-start":
# Claims the marker atomically and says which event to record, so a
# second session starting alongside this one cannot record it too.
first_event = telemetry.claim_install(was_empty=data_dir_was_empty)
if first_event == "install":
if telemetry.is_first_run():
telemetry.record("install")
elif first_event == "upgrade":
# First run after a build that never wrote the marker; the
# predecessor version was never recorded anywhere.
telemetry.record("upgrade", from_version="pre-0.3")
else:
previous = telemetry.claim_version_change()
if previous:
telemetry.record("upgrade", from_version=previous)
recovered = recover_pending_handoffs()
record_session_start(store, hook_input)
if recovered:
@@ -1800,34 +1800,6 @@ def extraction_message_batches(
return batches
# Platform surface attribution. Read from the generated per-host module so a new
# entrypoint is correct without remembering to configure anything.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
def platform_headers(key: str) -> dict[str, str]:
"""Auth plus the three surface-identity headers.
X-Mem0-Source and X-Application are set-once by contract: this is the
outermost layer, so it sets them, and nothing below may overwrite them.
X-Mem0-Client is append-only — anything downstream adds itself to the tail.
"""
headers = {
"Authorization": f"Token {key}",
"Content-Type": "application/json",
"X-Mem0-Source": _PLATFORM_SOURCE,
"X-Mem0-Client": f"mem0-plugin/{PLUGIN_VERSION}",
}
if _PLATFORM_APPLICATION:
headers["X-Application"] = _PLATFORM_APPLICATION
return headers
def _request_json(
url: str, key: str, payload: dict[str, Any], timeout: float
) -> tuple[dict[str, Any] | list[Any], int, int]:
@@ -1835,7 +1807,7 @@ def _request_json(
request = urllib.request.Request(
url,
data=raw,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -1862,7 +1834,7 @@ def _get_json(
) -> tuple[dict[str, Any] | list[Any], int]:
request = urllib.request.Request(
url,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="GET",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -2008,13 +1980,6 @@ def flush_session(
"user_id": write_user,
"app_id": repo.app_id,
"run_id": session_id,
# Top level, not metadata: the backend reads `source` from the body or
# the query string, never from metadata, which is where this used to
# sit. The X-Mem0-Source header is also read, but only from the
# platform release that ships alongside this change, so the body value
# is what makes attribution work on both. The harness tag stays in
# metadata as hook provenance.
"source": _PLATFORM_SOURCE,
"metadata": {**metadata, "author": write_user, "dirs": directory_chain(repo)},
"agent_custom_instructions": PROJECT_MEMORY_INSTRUCTIONS,
"custom_instructions": PERSONAL_MEMORY_INSTRUCTIONS,
@@ -2558,7 +2523,7 @@ def _collect_memory_ids(
def _delete_memory(api_url: str, key: str, memory_id: str) -> bool:
request = urllib.request.Request(
f"{api_url}/v1/memories/{urllib.parse.quote(memory_id)}/",
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="DELETE",
)
try:
@@ -1,9 +1,5 @@
#!/usr/bin/env python3
"""Usage telemetry for Mem0 agent plugins.
Events are linked to your Mem0 account email when an API key is configured, and
to a random per-machine id otherwise. Not anonymous — the Python SDK and CLI
attribute the same way.
"""Anonymous usage telemetry for Mem0 agent plugins.
Hooks run on a 3-6 second budget and fire on every tool call, so recording never
touches the network: `record` appends one JSON line to a local spool and returns.
@@ -13,8 +9,7 @@ started once per session and again from the flush worker that is already detache
Pure stdlib, matching the rest of the plugin. Opt out with MEM0_TELEMETRY=false.
Never sends prompts, memory text, queries, file paths, repository names, or API
keys: only event names, durations, counts, coarse outcomes, and repo/session
identifiers hashed with a random per-install salt.
keys: only event names, durations, counts, coarse outcomes, and salted hashes.
"""
from __future__ import annotations
@@ -34,24 +29,8 @@ from typing import Any
import memory_core
# Seeded from the per-host module the build generates into core/. Two processes
# in this pipeline never call init() — mcp_server.py, and the detached
# `python3 telemetry.py` sender that spawn_flush() starts — so a module default
# was what every one of their events got labelled with.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import HARNESS_ID as _DEFAULT_HARNESS
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
from _harness_id import SOURCE_TAG as _DEFAULT_SOURCE_TAG
except ImportError:
_DEFAULT_HARNESS = "generic"
_DEFAULT_SOURCE_TAG = "MEM0_PLUGIN"
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_harness: str = "generic"
_source_tag: str = "MEM0_PLUGIN"
_PRIVATE_KEYS = {
"apikey",
"authorization",
@@ -77,19 +56,10 @@ _PRIVATE_KEYS = {
}
def init(harness: str = "", source_tag: str = "") -> None:
"""Override the generated identity. Optional — core/_harness_id.py is the default.
The fallback shape matches memory_core.configure_harness's (``<HOST>_PLUGIN``).
It used to be ``MEM0_<HOST>_PLUGIN`` here and ``<host>_plugin`` there, which
meant one plugin could emit three different source values depending on which
process happened to send the batch.
"""
def init(harness: str = "generic", source_tag: str = "") -> None:
global _harness, _source_tag
_harness = harness or _DEFAULT_HARNESS
_source_tag = source_tag or (
f"{_harness.upper().replace('-', '_')}_PLUGIN" if harness else _DEFAULT_SOURCE_TAG
)
_harness = harness
_source_tag = source_tag or f"MEM0_{harness.upper().replace('-', '_')}_PLUGIN"
POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"
POSTHOG_CAPTURE_URL = "https://us.i.posthog.com/i/v0/e/"
@@ -100,16 +70,6 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -123,126 +83,9 @@ def is_enabled() -> bool:
def _digest(value: str, length: int = 16) -> str:
"""Unsalted digest. Only for values that are already secrets (API keys)."""
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
Published atomically, and there is deliberately no derived fallback. Creating
the file with O_CREAT|O_EXCL and then writing into it leaves a window where
the file exists and is empty, and a concurrent hook that reads it in that
window gets nothing. Falling back to a digest of the path would hand that
process a salt an attacker can compute, memoized for its whole run, which is
the privacy control this function exists to provide silently turning itself
off under load. The salt is written to a private temp file first and linked
into place, so the name either does not exist or already has the full value.
Returns "" when it genuinely cannot persist. Callers omit the hash entirely
rather than emit an unsalted one.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
# Read before writing. Hooks are separate processes firing on every tool
# call, so all but the first find the salt already published; going straight
# to create-fsync-link-unlink meant every one of them paid an fsync to
# discover that, on a path whose whole promise is appending a line and
# returning.
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
if _salt_cache:
return _salt_cache
except OSError:
pass
temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp")
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
stream.flush()
os.fsync(stream.fileno())
try:
# Atomic claim: fails if another process already published one.
# os.link rather than replace, which would clobber theirs.
os.link(temporary, path)
except FileExistsError:
pass
except OSError:
# No hardlinks here (some network mounts, some container volumes).
# Claim the name directly instead. That reopens the empty-file
# window, but the window is now benign: a reader that lands in it
# gets "" and omits the hash for that process rather than caching a
# guessable one. Losing the hashes on every run of an entire
# filesystem is the worse failure.
try:
fallback = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(fallback, "w", encoding="utf-8") as stream:
stream.write(temporary.read_text(encoding="utf-8"))
except OSError:
pass
except OSError:
pass
finally:
try:
temporary.unlink()
except OSError:
pass
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
"""Salted digest for values drawn from a guessable space.
repo.identity is a git remote URL, or ``local:<absolute path>`` when there is
no remote — which normally contains the account username. Sixteen unsalted
hex characters over that input space is enumerable, so this is not a
privacy control without the salt. Salting per install keeps every
within-account join the analytics actually use and gives up only
cross-machine joins on the same repository, which nothing computes.
Returns "" when there is no salt, so record() omits the property. An
unsalted digest over this input space is close to plaintext, and emitting one
under a name that implies it is hashed is worse than sending nothing.
"""
if not value:
return ""
salt = _install_salt()
if not salt:
return ""
return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length]
def _safe_value(value: Any) -> Any:
if isinstance(value, str):
return memory_core.redact(value)
@@ -302,176 +145,9 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str:
return created
def _rotate_anonymous_id(identity: dict[str, str]) -> str:
"""Mint a fresh anonymous id because the account context is gone.
The previous id may already have been merged into a person profile by an
$identify, and that merge is permanent. Reusing it after a logout or a key
change attributes everything that follows to the account that just went
away, which is the same misattribution the key fingerprint exists to stop,
only arriving through the anonymous path instead.
`aliased` is cleared with it: the new id has never been merged, so it is
eligible to be aliased into whatever account comes next.
"""
created = f"code-anon-{uuid.uuid4().hex}"
identity["anonymous_id"] = created
identity.pop("aliased", None)
_write_identity(identity)
return created
def _install_state_path() -> Path:
return memory_core.data_dir() / "install-state.json"
def is_first_run() -> bool:
"""Whether install has never been recorded on this machine.
Deliberately NOT the identity file. That file is only written by a
successful flush, so an offline or firewalled user recorded code.install on
every single session, forever — and every 0.2.x user recorded one on their
first 0.3.x session because 0.2.x never wrote it at all.
"""
return not _install_state_path().exists()
def data_dir_was_empty() -> bool:
"""Whether the data directory is untouched. Call BEFORE anything writes to it.
hook_runner reaches claim_install() only after cache_plugin_api_key() has
written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so
asking at claim time always saw content and every fresh install reported an
upgrade. The caller snapshots this at the top of the run instead.
"""
return not _data_dir_has_content()
def claim_install(was_empty: bool | None = None) -> str | None:
"""Claim the one install/upgrade record for this machine, atomically.
Returns the event to record ("install" or "upgrade"), or None if another
session already claimed it. O_CREAT|O_EXCL so two sessions starting together
cannot both win.
`was_empty` must come from data_dir_was_empty() called before this process
wrote anything. Omitting it falls back to checking now, which is only
correct for a caller that has touched nothing.
"""
if not is_enabled():
# Never consume the one-shot claim while the user is opted out, or they
# would silently lose their install event if they later opt in.
return None
path = _install_state_path()
upgrading = not (data_dir_was_empty() if was_empty is None else was_empty)
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
except FileExistsError:
return None
except OSError:
return None
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
json.dump(
{
"plugin_version": memory_core.PLUGIN_VERSION,
"installed_at": memory_core.utc_now(),
"upgraded": upgrading,
},
stream,
)
# Durable before this returns. The O_EXCL open is what makes the
# claim exclusive, so it cannot be replaced by a temp-and-rename
# without losing that, which leaves the content as the thing to make
# safe. A kill between the open and this fsync used to leave a marker
# that exists but parses to nothing: is_first_run reads it as claimed
# and claim_version_change cannot read a version out of it.
stream.flush()
os.fsync(stream.fileno())
except OSError:
pass
return "upgrade" if upgrading else "install"
def _data_dir_has_content() -> bool:
"""Whether anything predates this session in the plugin data directory."""
try:
for entry in memory_core.data_dir().iterdir():
if entry.name != "install-state.json":
return True
except OSError:
pass
return False
def _repair_install_state(path: Path) -> None:
"""Rewrite an unparseable marker so version tracking can resume."""
try:
temporary = path.with_suffix(f".{os.getpid()}.tmp")
temporary.write_text(
json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}),
encoding="utf-8",
)
temporary.replace(path)
except OSError:
pass
def claim_version_change() -> str | None:
"""Return the previously recorded version if it differs, updating the marker.
Only meaningful once the marker exists — the first transition into 0.3.x has
no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting
the marker means the next session sees no change and records nothing.
"""
path = _install_state_path()
try:
state = json.loads(path.read_text(encoding="utf-8"))
except OSError:
return None
except json.JSONDecodeError:
# A crash between O_EXCL and the write leaves an empty marker. Left
# alone it disables every future upgrade event on this machine, because
# claim_install sees the file and this function cannot parse it.
state = None
if not isinstance(state, dict):
_repair_install_state(path)
return None
previous = str(state.get("plugin_version") or "")
if not previous or previous == memory_core.PLUGIN_VERSION:
return None
# Claim the transition with an exclusive sentinel before rewriting the
# marker. A plain read-modify-write let every concurrently starting session
# observe the old version and each record its own upgrade — and the first
# session after a version bump is exactly when several agent windows restart
# together.
sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}")
try:
os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600))
except FileExistsError:
return None
except OSError:
return None
state["plugin_version"] = memory_core.PLUGIN_VERSION
state["upgraded_at"] = memory_core.utc_now()
temporary = path.with_suffix(f".{os.getpid()}.tmp")
try:
temporary.write_text(json.dumps(state), encoding="utf-8")
temporary.replace(path)
except OSError:
# Release the claim. The marker still records the old version, so
# without this the sentinel makes claim_version_change return early on
# every later run and this version's upgrade is never recorded again.
for leftover in (sentinel, temporary):
try:
leftover.unlink()
except OSError:
pass
return None
return previous
"""Whether this machine has never recorded a plugin event before."""
return not _identity_path().exists()
def record(
@@ -492,32 +168,19 @@ def record(
except OSError:
pass
properties = _safe_value(properties)
# Stamped in the RECORDING process, beside harness. `source` used to be
# read in the sending process from a module global, so whichever process
# drained the spool named every event in it. flush() spreads per-event
# properties last, so this now wins over any sender's default.
properties.update(
harness=_harness,
source=_source_tag,
plugin_version=memory_core.PLUGIN_VERSION,
os=sys.platform,
python_version=platform.python_version(),
)
# Assigned only when the digest is real. _scoped_digest returns "" when
# the salt could not be persisted, and an empty property is worse than an
# absent one: it survives the None filter below and reads as a value.
if repo is not None:
repo_hash = _scoped_digest(getattr(repo, "identity", ""))
if repo_hash:
properties["repo_hash"] = repo_hash
properties["repo_hash"] = _digest(getattr(repo, "identity", ""))
if session_id:
session_hash = _scoped_digest(session_id)
if session_hash:
properties["session_hash"] = session_hash
properties["session_hash"] = _digest(session_id)
line = json.dumps(
{
"event": f"{EVENT_PREFIX}.{event}",
"uuid": str(uuid.uuid4()),
"timestamp": memory_core.utc_now(),
"properties": {
key: value for key, value in properties.items() if value is not None
@@ -576,201 +239,38 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / _claim_name()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
for orphan in sorted(directory.glob("telemetry-*.sending")):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
if age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
pass
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
if age < CLAIM_STALE_SECONDS:
continue
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -800,130 +300,34 @@ def _post(payload: dict[str, Any], url: str) -> bool:
def resolve_distinct_id() -> tuple[str, str]:
"""Return the PostHog distinct id and the anonymous id it replaced, if any.
The second value becomes a PostHog $identify alias. It is ONLY ever an
anonymous id: aliasing one account email to another merges two real person
profiles and cannot be undone, so a key that now belongs to a different
account re-resolves with no alias.
"""
"""Return the PostHog distinct id and the anonymous id it replaced, if any."""
identity = _read_identity()
key = memory_core.api_key()
fingerprint = _digest(key) if key else ""
email = identity.get("email", "")
if email and fingerprint:
recorded = identity.get("key_fingerprint", "")
if recorded == fingerprint:
return email, ""
if not recorded:
# Rows written before fingerprints existed. Verify rather than
# adopt: a key changed before the upgrade would otherwise bind the
# new key to the previous account's email, permanently, and the
# fingerprint would then agree with itself forever after.
verified = _resolve_email(key)
if not verified:
# Offline, firewalled, or the API is down. Keep the previous
# behaviour and retry on the next flush rather than dropping a
# real account attribution. Safe because the same network that
# failed /v1/ping/ is about to fail the PostHog POST, so nothing
# is delivered under the unverified identity in the meantime.
return email, ""
identity["email"] = verified
identity["key_fingerprint"] = fingerprint
_write_identity(identity)
return verified, ""
if email:
return email, ""
key = memory_core.api_key()
if not key:
# No key to verify the account with; do not keep attributing to it.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
return anonymous_id(identity), ""
resolved = _resolve_email(key)
if not resolved:
# The key changed and will not resolve (revoked, offline, API down).
# Reaching here with an email means the recorded fingerprint disagreed,
# so the key really did change. Drop the account and rotate: the stored
# anonymous id may already be merged into that account's person, and
# reusing it would keep the events on the profile we are trying to
# leave.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
email = _resolve_email(key)
if not email:
return anonymous_id(identity), ""
# Alias only when going anonymous -> email for the first time. Once an anon
# id has been merged into an account it must never be offered again: an
# alias naming an already-identified id is what could link two real people.
previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "")
if previous:
identity["aliased"] = True
identity["email"] = resolved
identity["key_fingerprint"] = fingerprint
previous = identity.get("anonymous_id", "")
identity["email"] = email
_write_identity(identity)
return resolved, previous
return email, previous
def flush() -> int:
"""Drain the live spool, then any parked claims, and return events sent."""
"""Drain claimed spools to PostHog and return the number of events sent."""
if not is_enabled():
return 0
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
claim = _claim_spool()
if claim is None:
return 0, True
return 0
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# quarantine rather than retry. flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
return 0
events = []
for line in lines:
try:
@@ -933,18 +337,11 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
claim.unlink()
except OSError:
pass
return 0, True
return 0
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -963,17 +360,12 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
# moved into record() have none of their own.
"source": _source_tag,
"language": "python",
"$process_person_profile": False,
@@ -981,24 +373,16 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
**(event.get("properties") or {}),
},
}
for event in chunk
for event in events[start : start + BATCH_SIZE]
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
def main() -> int:
@@ -7,8 +7,8 @@ disable-model-invocation: true
# Pause memory capture
To pause (hooks stop capturing and sending session content; a minimal
telemetry ping still fires at session start, under your Mem0 account email,
unless `MEM0_TELEMETRY=false`):
anonymous telemetry ping still fires at session start unless
`MEM0_TELEMETRY=false`):
```bash
python3 "{{PLUGIN_ROOT}}/core/memory_cli.py" --harness "{{HARNESS_ID}}" {{PLUGIN_DATA_ARG}} pause
@@ -70,38 +70,6 @@ def test_portable_bundle_is_conformant_and_self_contained(tmp_path: Path) -> Non
assert not any(path.is_symlink() for path in root.rglob("*"))
def _harness_identity(root: Path) -> dict[str, str]:
"""Read the generated core/_harness_id.py without importing it."""
values: dict[str, str] = {}
for line in (root / "core" / "_harness_id.py").read_text(encoding="utf-8").splitlines():
if "=" in line and not line.lstrip().startswith("#"):
name, _, raw = line.partition("=")
values[name.strip()] = raw.strip().strip('"')
return values
def test_the_portable_bundle_declares_no_host_application(tmp_path: Path) -> None:
"""It runs in whatever editor a user drops it into, so it cannot know the host.
X-Application is allowlisted server-side. A guessed value is silently dropped
there, which is the worst outcome: the wire says we know the host and the
stored event says we do not.
"""
identity = _harness_identity(build("mem0-agent-plugin", "portable", tmp_path / "portable"))
assert identity["PLATFORM_APPLICATION"] == ""
# The PostHog-side label is still useful for grouping and stays populated.
assert identity["HARNESS_ID"] == "coding-agent"
assert identity["PLATFORM_SOURCE"] == "MEM0_PLUGIN"
@pytest.mark.parametrize("host", ["claude-code", "cursor", "codex", "kimi", "antigravity"])
def test_a_native_bundle_names_the_host_it_was_built_for(host: str, tmp_path: Path) -> None:
identity = _harness_identity(build(host, "native", tmp_path / host))
assert identity["PLATFORM_APPLICATION"] == host
@pytest.mark.parametrize("host", ["claude-code", "cursor", "codex", "kimi", "antigravity"])
def test_native_bundle_is_self_contained(host: str, tmp_path: Path) -> None:
root = build(host, "native", tmp_path / host)
@@ -1,446 +0,0 @@
"""Delivery semantics of the telemetry spool: no duplicates, no starvation.
These run against a built host's core in-process (not a subprocess) because they
need to inject failures into ``_post``. The identity tests next door cover the
uninitialised-process case that needs a real interpreter.
"""
from __future__ import annotations
import importlib
import json
import os
import sys
import time
from pathlib import Path
import pytest
CORE_ROOT = Path(__file__).resolve().parents[1]
REPOSITORY_ROOT = CORE_ROOT.parents[1]
HOST_CORE = REPOSITORY_ROOT / "integrations" / "claude-code-plugin" / "core"
pytestmark = pytest.mark.skipif(not HOST_CORE.exists(), reason="claude-code-plugin is not built")
@pytest.fixture()
def telemetry(tmp_path, monkeypatch):
# CI runs this directory and claude-code-plugin/tests in ONE pytest process,
# and that suite's conftest sets MEM0_TELEMETRY=false at import, process-wide.
# Without this the whole file silently no-ops: record() returns early and
# every assertion sees an empty spool. Do not rely on ambient env.
monkeypatch.setenv("MEM0_TELEMETRY", "true")
monkeypatch.setenv("MEM0_CODE_DATA_DIR", str(tmp_path / "data"))
monkeypatch.syspath_prepend(str(HOST_CORE))
# Save and RESTORE rather than delete. claude-code-plugin/tests/conftest.py
# imports memory_core once at collection and calls configure_harness() on it;
# dropping the module left a later re-import with default harness config, so
# tests in that suite failed depending on collection order.
names = ("telemetry", "memory_core", "_harness_id")
saved = {name: sys.modules.get(name) for name in names}
for name in names:
sys.modules.pop(name, None)
module = importlib.import_module("telemetry")
monkeypatch.setattr(module, "resolve_distinct_id", lambda: ("tester@example.com", ""))
try:
yield module
finally:
for name in names:
sys.modules.pop(name, None)
if saved[name] is not None:
sys.modules[name] = saved[name]
def _delivered(payloads):
return [event for payload in payloads if "batch" in payload for event in payload["batch"]]
def test_a_partial_failure_does_not_redeliver_what_already_arrived(telemetry):
"""Defect 2a: flush kept the whole claim on failure and retried from the top.
150 events across two batches, the second failing, previously delivered 250.
"""
for index in range(150):
telemetry.record("search", index=index)
sent: list[dict] = []
calls = {"n": 0}
def flaky(payload, url):
calls["n"] += 1
if calls["n"] == 2: # second batch fails
return False
sent.append(payload)
return True
telemetry._post = flaky
telemetry.flush()
telemetry._post = lambda payload, url: sent.append(payload) or True
telemetry.flush()
events = _delivered(sent)
assert len(events) == 150
assert len({event["uuid"] for event in events}) == 150
def test_a_fresh_claim_is_not_immediately_stealable(telemetry):
"""Defect 2b: rename preserves mtime, so a claim inherited the spool's age.
With the last write older than the stale threshold, a claim made now looked
abandoned the instant it existed and a second sender took it over.
"""
telemetry.record("search")
spool = telemetry._spool_path()
old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (old, old))
first = telemetry._claim_spool()
assert first is not None
# A second sender starting right now must find nothing to take.
assert telemetry._claim_parked(first.parent) is None
def test_a_live_final_attempt_is_not_deleted_by_another_sender(telemetry):
"""Review finding: exhaustion was judged before liveness, so owners lost batches.
Claiming a parked file bumps its attempt count and refreshes its mtime. Once
the count reaches the budget, the owner draining it looked exhausted to every
other sender, which unlinked the file out from under it. Everything in that
batch was gone, which is precisely the loss this PR exists to stop.
"""
telemetry.record("search", reason="owned-by-the-first-sender")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
# Walk it to the final attempt, ageing it each round so it can be re-claimed.
# _claim_spool hands back a0 and _release_claim keeps the name, so it takes
# one full round per attempt to reach the budget.
for _ in range(telemetry.MAX_CLAIM_ATTEMPTS):
# Carry the marker through each rewrite so the final assertion proves the
# events survived, not merely that some file with the right name did.
telemetry._release_claim(claim, [{"event": "code.search", "uuid": "owned-by-the-first-sender"}])
parked = sorted(claim.parent.glob("telemetry-*.sending"))
assert parked, "the batch was dropped while still inside its budget"
os.utime(parked[0], (stale, stale))
claim = telemetry._claim_parked(claim.parent)
assert claim is not None
assert telemetry._claim_attempt(claim) >= telemetry.MAX_CLAIM_ATTEMPTS
assert claim.exists()
# The owner is draining it right now: fresh mtime, live lease.
second_sender = telemetry._claim_parked(claim.parent)
assert second_sender is None, "a second sender took a batch under a live lease"
assert claim.exists(), "a second sender deleted a batch its owner was draining"
assert "owned-by-the-first-sender" in claim.read_text(encoding="utf-8")
def test_an_exhausted_batch_is_still_discarded_once_its_lease_lapses(telemetry):
"""The liveness check must defer the cleanup, not cancel it.
Guards the obvious over-correction: skipping live claims is only safe if an
abandoned one at the same attempt count is still reaped on a later run.
"""
telemetry.record("search")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
exhausted = claim.parent / telemetry._claim_name(telemetry.MAX_CLAIM_ATTEMPTS)
claim.replace(exhausted)
os.utime(exhausted, (stale, stale))
assert telemetry._claim_parked(exhausted.parent) is None
assert not exhausted.exists(), "an abandoned exhausted batch was left behind forever"
def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry):
"""Defect 6: parked claims were only reachable when no spool existed.
Because sessions keep recording there usually was one, so a batch parked by
a failed send waited until the 7-day expiry deleted it unsent — even though
its own presence is what starts the sender.
"""
telemetry.record("parked")
telemetry._post = lambda payload, url: False
telemetry.flush()
parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
assert len(parked) == 1
old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(parked[0], (old, old))
telemetry.record("fresh")
sent: list[dict] = []
telemetry._post = lambda payload, url: sent.append(payload) or True
telemetry.flush()
names = {event["event"] for event in _delivered(sent)}
assert names == {"code.parked", "code.fresh"}
def test_a_batch_is_retried_until_the_budget_is_spent_not_discarded(telemetry):
"""Expiry discards what failed repeatedly, not what merely sat for a while.
The budget is the attempt count, because age cannot be one: every re-claim
touches the mtime and every release backdates it, so age never accumulates.
"""
telemetry.record("parked")
telemetry._post = lambda payload, url: False
telemetry.flush()
parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
assert len(parked) == 1
assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS
sent: list[dict] = []
telemetry._post = lambda payload, url: sent.append(payload) or True
telemetry.flush()
assert [event["event"] for event in _delivered(sent)] == ["code.parked"]
def test_progress_is_recorded_after_every_batch(telemetry):
"""A crash repeats at most one batch, not the whole file."""
for index in range(250):
telemetry.record("search", index=index)
calls = {"n": 0}
def die_after_two(payload, url):
calls["n"] += 1
if calls["n"] > 2:
return False
return True
telemetry._post = die_after_two
telemetry.flush()
parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
assert len(parked) == 1
remaining = parked[0].read_text(encoding="utf-8").strip().splitlines()
# Two batches of 100 landed; only the last 50 should still be pending.
assert len(remaining) == 50
assert json.loads(remaining[0])["properties"]["index"] == 200
def test_the_heartbeat_actually_refreshes_the_lease(telemetry):
"""The claim rewrite doubles as the lease heartbeat.
Previously asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS`, which compares
two constants and executes none of the code under test. Drive the real
rewrite and watch the mtime move instead.
"""
for index in range(150):
telemetry.record("search", index=index)
claim = telemetry._claim_spool()
assert claim is not None
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(claim, (stale, stale))
assert time.time() - claim.stat().st_mtime > telemetry.CLAIM_STALE_SECONDS
telemetry._rewrite_claim(claim, [{"event": "code.x", "properties": {}}])
assert time.time() - claim.stat().st_mtime < telemetry.CLAIM_STALE_SECONDS
def test_an_undeliverable_batch_is_eventually_given_up_on(telemetry):
"""Expiry has to be reachable from a state the state machine can produce.
It was not: every re-claim touched the mtime and every release backdated it
by a fixed amount, so age hovered near the stale threshold and the 7-day
expiry never fired. An undeliverable batch lived on disk forever, and
spawn_flush saw it and started a sender on every hook.
"""
telemetry.record("doomed")
telemetry._post = lambda payload, url: False
directory = telemetry.memory_core.data_dir()
for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3):
telemetry.flush()
# Attempts now carry a cooldown, so a released claim is not instantly
# reclaimable. Age it to stand in for the wall time a real retry waits;
# without this the loop spins inside one cooldown and proves nothing.
for parked in directory.glob("telemetry-*.sending"):
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(parked, (stale, stale))
leftover = list(directory.glob("telemetry-*.sending"))
assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}"
def test_a_batch_that_cannot_be_read_is_not_counted_as_delivered(telemetry):
"""Review finding: a read failure reported 'everything delivered'.
Nothing was posted, so calling it delivered lets flush() carry on to other
claims as though this batch had arrived, and hides the failure from the one
signal that says the run went badly. It also must not quarantine: a briefly
unreadable file is retryable, and moving it to .corrupt discards the events
over a transient filesystem error, because nothing ever re-globs .corrupt.
"""
telemetry.record("search")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
original = Path.read_text
def unreadable(self, *args, **kwargs):
if self == claim:
raise OSError(5, "I/O error")
return original(self, *args, **kwargs)
Path.read_text = unreadable
try:
sent, delivered = telemetry._drain(claim)
finally:
Path.read_text = original
assert sent == 0
assert delivered is False, "an unread batch was reported as delivered"
assert claim.exists(), "a transient read error discarded the batch"
assert not list(claim.parent.glob("*.corrupt")), "quarantined over a transient error"
def test_undecodable_content_is_still_quarantined_and_the_run_continues(telemetry):
"""The other half: genuinely unrecoverable content must not block the run.
Guards the over-correction. If every read problem returned undelivered, one
torn file would stop every later claim on every flush, forever.
"""
telemetry.record("search")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
claim.write_bytes(b"\xff\xfe torn \x00 write")
sent, delivered = telemetry._drain(claim)
assert (sent, delivered) == (0, True)
assert not claim.exists()
assert list(claim.parent.glob("*.corrupt")), "unrecoverable content was not quarantined"
def test_retries_are_spread_over_real_time_not_burned_at_once(telemetry):
"""Review finding: releasing straight to reclaimable spent the budget instantly.
Two senders hitting one momentary failure could walk a batch from attempt 0
to the limit within seconds and discard it, when a retry a minute later would
have delivered. Each release now has to age past a cooldown that grows with
the attempts already spent.
"""
telemetry.record("doomed")
telemetry._post = lambda payload, url: False
directory = telemetry.memory_core.data_dir()
telemetry.flush()
parked = list(directory.glob("telemetry-*.sending"))
assert parked, "the batch was discarded on its first failure"
assert telemetry._claim_attempt(parked[0]) == 0
# Second sender, immediately: the cooldown has not elapsed, so it must not
# be able to spend another attempt.
telemetry.flush()
still = list(directory.glob("telemetry-*.sending"))
assert len(still) == 1
assert telemetry._claim_attempt(still[0]) <= 1, "burned attempts without waiting"
def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry):
"""The old shape is telemetry-<pid>-<hex>.sending, and hex can start with 'a'."""
assert telemetry._claim_attempt(Path("telemetry-999-deadbeef.sending")) == 0
assert telemetry._claim_attempt(Path("telemetry-999-a1234567.sending")) == 0
assert telemetry._claim_attempt(Path("telemetry-999-deadbeef-a2.sending")) == 2
def test_a_torn_claim_is_quarantined_not_deleted(telemetry):
"""A non-empty file that parses to nothing is the remainder, not garbage."""
telemetry.record("search")
claim = telemetry._claim_spool()
claim.write_bytes(b"\xff\xfe not utf-8 at all")
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(claim, (stale, stale))
sent = telemetry.flush()
assert sent == 0
assert not claim.exists()
quarantined = list(telemetry.memory_core.data_dir().glob("*.corrupt"))
assert len(quarantined) == 1, "torn claim was destroyed instead of kept"
def test_a_failed_rewrite_stops_instead_of_redelivering(telemetry):
"""Ignoring the rewrite result reintroduced the duplicates this PR fixes."""
for index in range(250):
telemetry.record("search", index=index)
telemetry._rewrite_claim = lambda claim, remaining: False
delivered = []
telemetry._post = lambda payload, url: delivered.extend(payload.get("batch", [])) or True
telemetry.flush()
assert len(delivered) == 100, f"kept going after a failed rewrite: {len(delivered)}"
def test_partial_files_are_swept(telemetry):
"""Nothing else globs *.partial, so a crash mid-rename orphans one forever."""
data_dir = telemetry.memory_core.data_dir()
data_dir.mkdir(parents=True, exist_ok=True)
debris = data_dir / "telemetry-1-abc-a0.1.partial"
debris.write_text("x", encoding="utf-8")
old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(debris, (old, old))
telemetry.flush()
assert not debris.exists()
def test_quarantined_batches_are_eventually_collected(telemetry):
"""Nothing re-globs .corrupt, so without a sweep they live on disk forever.
Kept much longer than .partial debris on purpose: a quarantined batch is the
only remaining evidence of events that could not be delivered.
"""
directory = telemetry.memory_core.data_dir()
directory.mkdir(parents=True, exist_ok=True)
fresh = directory / "telemetry-1-aaaaaaaa-a0.corrupt"
old = directory / "telemetry-2-bbbbbbbb-a0.corrupt"
for path in (fresh, old):
path.write_text("torn", encoding="utf-8")
expired = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60)
os.utime(old, (expired, expired))
telemetry._sweep_debris(directory)
assert fresh.exists(), "a recent quarantine was discarded before anyone could look at it"
assert not old.exists(), "an expired quarantine was left on disk forever"
def test_temp_files_orphaned_by_a_kill_are_collected(telemetry):
"""_write_identity and _install_salt unlink in a finally, which SIGKILL skips."""
directory = telemetry.memory_core.data_dir()
directory.mkdir(parents=True, exist_ok=True)
orphan = directory / "telemetry-salt.999.tmp"
orphan.write_text("abandoned", encoding="utf-8")
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(orphan, (stale, stale))
telemetry._sweep_debris(directory)
assert not orphan.exists(), "a killed process left a temp file on disk forever"
@@ -1,274 +0,0 @@
"""Core telemetry behaviour with NO telemetry.init(), in a real subprocess.
Why this file exists
--------------------
``telemetry.py`` lives in ``agent-plugin-core/python/`` but its only tests lived
under ``claude-code-plugin/tests/``, behind a ``conftest.py`` that calls
``configure_harness()`` and ``telemetry.init()`` at import. Core behaviour was
therefore only ever exercised inside an already-configured module.
Two processes in the real pipeline never call ``init()``:
- ``mcp_server.py``, which records every manual search;
- the detached ``python3 telemetry.py`` sender that ``spawn_flush()`` starts at
session start, after every skill command, and when the MCP server exits.
Both fell back to module defaults, so MCP searches reported ``harness=generic``
and everything that sender delivered was labelled ``MEM0_PLUGIN`` regardless of
which of the six plugins produced it. The suite stayed green throughout.
These tests run in a fresh interpreter with no conftest, against a built host
bundle, which is the only arrangement that can catch that class of bug.
"""
from __future__ import annotations
import json
import subprocess
import sys
import tempfile
from pathlib import Path
import pytest
CORE_ROOT = Path(__file__).resolve().parents[1]
REPOSITORY_ROOT = CORE_ROOT.parents[1]
HOSTS = {
"claude-code": ("claude-code-plugin", "CLAUDE_CODE_PLUGIN"),
"cursor": ("cursor-plugin", "CURSOR_PLUGIN"),
"codex": ("codex-plugin", "CODEX_PLUGIN"),
"kimi": ("kimi-plugin", "KIMI_PLUGIN"),
"antigravity": ("antigravity-plugin", "ANTIGRAVITY_PLUGIN"),
# Portable: no flush_worker and no hook_runner, so its ONLY sender is the
# uninitialised telemetry.py. A native-only test passes here vacuously.
"coding-agent": ("mem0-agent-plugin", "CODING_AGENT_PLUGIN"),
}
def _core_dir(directory: str) -> Path:
return REPOSITORY_ROOT / "integrations" / directory / "core"
def _run(core: Path, data_dir: Path, body: str) -> str:
"""Execute `body` in a fresh interpreter with only the host's core on sys.path."""
script = f"import sys; sys.path.insert(0, {str(core)!r})\n{body}"
result = subprocess.run(
[sys.executable, "-c", script],
capture_output=True,
text=True,
env={
"MEM0_CODE_DATA_DIR": str(data_dir),
"PATH": "/usr/bin:/bin",
"HOME": str(data_dir),
},
)
assert result.returncode == 0, result.stderr
return result.stdout.strip()
@pytest.mark.parametrize("harness,spec", sorted(HOSTS.items()))
def test_identity_resolves_without_init(harness, spec):
"""Every built host knows what it is with no configuration call at all."""
directory, source_tag = spec
core = _core_dir(directory)
if not core.exists():
pytest.skip(f"{directory} is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
out = _run(
core,
Path(tmp),
"import telemetry; print(telemetry._harness, telemetry._source_tag)",
)
assert out == f"{harness} {source_tag}"
def test_mcp_server_records_the_real_harness():
"""mcp_server imports telemetry and never initialises it (server.py has no init).
Its recorded events used to carry harness=generic for every plugin.
"""
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
data_dir = Path(tmp)
_run(
core,
data_dir,
"import mcp_server, telemetry; telemetry.record('search', trigger='mcp-search')",
)
spooled = (data_dir / "telemetry.jsonl").read_text(encoding="utf-8").strip()
event = json.loads(spooled)
assert event["properties"]["harness"] == "claude-code"
assert event["properties"]["source"] == "CLAUDE_CODE_PLUGIN"
def test_the_detached_sender_does_not_relabel_events():
"""`python3 telemetry.py` is the sender spawn_flush() starts, and never inits.
source is stamped at record time now, so which process sends is irrelevant.
"""
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
data_dir = Path(tmp)
_run(core, data_dir, "import telemetry; telemetry.record('search')")
captured = data_dir / "captured.json"
# Drain with a fresh, unconfigured interpreter, capturing the payload
# instead of posting it.
_run(
core,
data_dir,
"import json, telemetry\n"
"sent = []\n"
"telemetry._post = lambda payload, url: sent.append(payload) or True\n"
"telemetry.flush()\n"
f"open({str(captured)!r}, 'w').write(json.dumps(sent))",
)
payloads = json.loads(captured.read_text(encoding="utf-8"))
batches = [p for p in payloads if "batch" in p]
assert batches, "nothing was sent"
properties = batches[0]["batch"][0]["properties"]
assert properties["source"] == "CLAUDE_CODE_PLUGIN"
assert properties["harness"] == "claude-code"
def test_every_event_carries_a_uuid_for_dedupe():
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
data_dir = Path(tmp)
_run(core, data_dir, "import telemetry; telemetry.record('search'); telemetry.record('flush')")
lines = (data_dir / "telemetry.jsonl").read_text(encoding="utf-8").strip().splitlines()
ids = [json.loads(line)["uuid"] for line in lines]
assert len(ids) == 2
assert len(set(ids)) == 2
def test_source_tag_defaults_agree_between_the_two_modules():
"""configure_harness and telemetry.init must derive the same tag.
They disagreed: `<host>_plugin` in one and `MEM0_<HOST>_PLUGIN` in the other,
so one plugin could emit three different source values depending on which
process sent the batch.
"""
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
out = _run(
core,
Path(tmp),
"import memory_core, telemetry\n"
"memory_core.configure_harness('kimi')\n"
"telemetry.init(harness='kimi')\n"
"print(memory_core.harness_config()['source_tag'].upper(), telemetry._source_tag)",
)
left, right = out.split()
assert left == right == "KIMI_PLUGIN"
def test_the_plugin_declares_its_surface_in_the_body_and_the_headers():
"""Body and headers both, because only the body works on every backend."""
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
out = _run(
core,
Path(tmp),
"import json, memory_core\n"
"h = memory_core.platform_headers('k')\n"
"print(json.dumps({'source': h.get('X-Mem0-Source'),"
" 'app': h.get('X-Application'),"
" 'client': h.get('X-Mem0-Client'),"
" 'auth': h.get('Authorization'),"
" 'ctype': h.get('Content-Type')}))",
)
headers = json.loads(out)
assert headers["source"] == "MEM0_PLUGIN"
assert headers["app"] == "claude-code"
assert headers["client"].startswith("mem0-plugin/")
# The transport headers the three call sites relied on must survive.
assert headers["auth"] == "Token k"
assert headers["ctype"] == "application/json"
def _session_start(core: Path, data_dir: Path) -> list[str]:
"""Drive the real hook_runner session-start path and return lifecycle events."""
recorded = "\n".join(
[
"import io, json, sys",
f"sys.path.insert(0, {str(core)!r})",
"import telemetry, hook_runner",
"seen = []",
"telemetry.record = lambda event, **kw: seen.append(event) or None",
"telemetry.spawn_flush = lambda: False",
# run() reads sys.argv through argparse; it takes no positional args.
"sys.argv = ['hook_runner', 'session-start']",
"sys.stdin = io.StringIO('{}')",
"hook_runner.run()",
"print(json.dumps([e for e in seen if e in ('install', 'upgrade')]))",
]
)
import json as _json
return _json.loads(_run(core, data_dir, recorded) or "[]")
def test_a_fresh_install_reports_install_not_upgrade():
"""The decision must survive the writes hook_runner does before asking.
claim_install() is reached only after cache_plugin_api_key() has written
`api-key` and EvidenceStore() has created `evidence.sqlite3`. Asking "is the
data dir empty" at that point always saw content, so code.install could
never fire and every new user was counted as an upgrade.
"""
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
data_dir = Path(tmp) / "data"
assert _session_start(core, data_dir) == ["install"]
def test_the_lifecycle_event_fires_exactly_once():
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
data_dir = Path(tmp) / "data"
first = _session_start(core, data_dir)
second = _session_start(core, data_dir)
third = _session_start(core, data_dir)
assert first == ["install"]
assert second == []
assert third == []
def test_an_existing_data_dir_reports_upgrade():
core = _core_dir("claude-code-plugin")
if not core.exists():
pytest.skip("claude-code-plugin is not built in this tree")
with tempfile.TemporaryDirectory() as tmp:
data_dir = Path(tmp) / "data"
data_dir.mkdir(parents=True)
# A 0.2.x leftover: the data dir survives the upgrade.
(data_dir / "requirements.txt").write_text("mem0ai\n", encoding="utf-8")
assert _session_start(core, data_dir) == ["upgrade"]
@@ -149,7 +149,7 @@ export async function buildRecallContext(
if (!unseen.length) return "";
const prefix =
"<mem0-relevant-memories>\nRetrieved automatically for the current request. This is a shallow first pass — search mem0_memory for more if you need it.\n";
"<mem0-relevant-memories>\nRetrieved automatically for the current request. This is a shallow first pass — use the available memory search tool for more if you need it.\n";
const suffix = "\n</mem0-relevant-memories>";
const maxChars = options.maxChars ?? DEFAULT_MAX_CONTEXT_CHARS;
const lines: string[] = [];
@@ -1,11 +0,0 @@
"""Generated by integrations/agent-plugin-core/build/build.py. Do not edit."""
HARNESS_ID = "antigravity"
SOURCE_TAG = "ANTIGRAVITY_PLUGIN"
# Platform-side vocabulary (mem0_event.source + X-Application). The whole
# plugin family is one source; which editor it runs in is the application.
# An empty application means the host is unknown, and memory_core omits
# the header entirely rather than sending a placeholder.
PLATFORM_SOURCE = "MEM0_PLUGIN"
PLATFORM_APPLICATION = "antigravity"
@@ -290,11 +290,6 @@ def run(
if args.plugin_data_dir:
os.environ[data_dir_env] = args.plugin_data_dir
# Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key
# writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking
# after them always saw content and every fresh install reported an upgrade.
data_dir_was_empty = telemetry.data_dir_was_empty()
cache_plugin_api_key()
if args.action == "session-start":
clear_stale_api_key_cache()
@@ -310,19 +305,8 @@ def run(
return 0
if args.action == "session-start":
# Claims the marker atomically and says which event to record, so a
# second session starting alongside this one cannot record it too.
first_event = telemetry.claim_install(was_empty=data_dir_was_empty)
if first_event == "install":
if telemetry.is_first_run():
telemetry.record("install")
elif first_event == "upgrade":
# First run after a build that never wrote the marker; the
# predecessor version was never recorded anywhere.
telemetry.record("upgrade", from_version="pre-0.3")
else:
previous = telemetry.claim_version_change()
if previous:
telemetry.record("upgrade", from_version=previous)
recovered = recover_pending_handoffs()
record_session_start(store, hook_input)
if recovered:
@@ -1800,34 +1800,6 @@ def extraction_message_batches(
return batches
# Platform surface attribution. Read from the generated per-host module so a new
# entrypoint is correct without remembering to configure anything.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
def platform_headers(key: str) -> dict[str, str]:
"""Auth plus the three surface-identity headers.
X-Mem0-Source and X-Application are set-once by contract: this is the
outermost layer, so it sets them, and nothing below may overwrite them.
X-Mem0-Client is append-only — anything downstream adds itself to the tail.
"""
headers = {
"Authorization": f"Token {key}",
"Content-Type": "application/json",
"X-Mem0-Source": _PLATFORM_SOURCE,
"X-Mem0-Client": f"mem0-plugin/{PLUGIN_VERSION}",
}
if _PLATFORM_APPLICATION:
headers["X-Application"] = _PLATFORM_APPLICATION
return headers
def _request_json(
url: str, key: str, payload: dict[str, Any], timeout: float
) -> tuple[dict[str, Any] | list[Any], int, int]:
@@ -1835,7 +1807,7 @@ def _request_json(
request = urllib.request.Request(
url,
data=raw,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -1862,7 +1834,7 @@ def _get_json(
) -> tuple[dict[str, Any] | list[Any], int]:
request = urllib.request.Request(
url,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="GET",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -2008,13 +1980,6 @@ def flush_session(
"user_id": write_user,
"app_id": repo.app_id,
"run_id": session_id,
# Top level, not metadata: the backend reads `source` from the body or
# the query string, never from metadata, which is where this used to
# sit. The X-Mem0-Source header is also read, but only from the
# platform release that ships alongside this change, so the body value
# is what makes attribution work on both. The harness tag stays in
# metadata as hook provenance.
"source": _PLATFORM_SOURCE,
"metadata": {**metadata, "author": write_user, "dirs": directory_chain(repo)},
"agent_custom_instructions": PROJECT_MEMORY_INSTRUCTIONS,
"custom_instructions": PERSONAL_MEMORY_INSTRUCTIONS,
@@ -2558,7 +2523,7 @@ def _collect_memory_ids(
def _delete_memory(api_url: str, key: str, memory_id: str) -> bool:
request = urllib.request.Request(
f"{api_url}/v1/memories/{urllib.parse.quote(memory_id)}/",
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="DELETE",
)
try:
+39 -655
View File
@@ -1,9 +1,5 @@
#!/usr/bin/env python3
"""Usage telemetry for Mem0 agent plugins.
Events are linked to your Mem0 account email when an API key is configured, and
to a random per-machine id otherwise. Not anonymous — the Python SDK and CLI
attribute the same way.
"""Anonymous usage telemetry for Mem0 agent plugins.
Hooks run on a 3-6 second budget and fire on every tool call, so recording never
touches the network: `record` appends one JSON line to a local spool and returns.
@@ -13,8 +9,7 @@ started once per session and again from the flush worker that is already detache
Pure stdlib, matching the rest of the plugin. Opt out with MEM0_TELEMETRY=false.
Never sends prompts, memory text, queries, file paths, repository names, or API
keys: only event names, durations, counts, coarse outcomes, and repo/session
identifiers hashed with a random per-install salt.
keys: only event names, durations, counts, coarse outcomes, and salted hashes.
"""
from __future__ import annotations
@@ -34,24 +29,8 @@ from typing import Any
import memory_core
# Seeded from the per-host module the build generates into core/. Two processes
# in this pipeline never call init() — mcp_server.py, and the detached
# `python3 telemetry.py` sender that spawn_flush() starts — so a module default
# was what every one of their events got labelled with.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import HARNESS_ID as _DEFAULT_HARNESS
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
from _harness_id import SOURCE_TAG as _DEFAULT_SOURCE_TAG
except ImportError:
_DEFAULT_HARNESS = "generic"
_DEFAULT_SOURCE_TAG = "MEM0_PLUGIN"
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_harness: str = "generic"
_source_tag: str = "MEM0_PLUGIN"
_PRIVATE_KEYS = {
"apikey",
"authorization",
@@ -77,19 +56,10 @@ _PRIVATE_KEYS = {
}
def init(harness: str = "", source_tag: str = "") -> None:
"""Override the generated identity. Optional — core/_harness_id.py is the default.
The fallback shape matches memory_core.configure_harness's (``<HOST>_PLUGIN``).
It used to be ``MEM0_<HOST>_PLUGIN`` here and ``<host>_plugin`` there, which
meant one plugin could emit three different source values depending on which
process happened to send the batch.
"""
def init(harness: str = "generic", source_tag: str = "") -> None:
global _harness, _source_tag
_harness = harness or _DEFAULT_HARNESS
_source_tag = source_tag or (
f"{_harness.upper().replace('-', '_')}_PLUGIN" if harness else _DEFAULT_SOURCE_TAG
)
_harness = harness
_source_tag = source_tag or f"MEM0_{harness.upper().replace('-', '_')}_PLUGIN"
POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"
POSTHOG_CAPTURE_URL = "https://us.i.posthog.com/i/v0/e/"
@@ -100,16 +70,6 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -123,126 +83,9 @@ def is_enabled() -> bool:
def _digest(value: str, length: int = 16) -> str:
"""Unsalted digest. Only for values that are already secrets (API keys)."""
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
Published atomically, and there is deliberately no derived fallback. Creating
the file with O_CREAT|O_EXCL and then writing into it leaves a window where
the file exists and is empty, and a concurrent hook that reads it in that
window gets nothing. Falling back to a digest of the path would hand that
process a salt an attacker can compute, memoized for its whole run, which is
the privacy control this function exists to provide silently turning itself
off under load. The salt is written to a private temp file first and linked
into place, so the name either does not exist or already has the full value.
Returns "" when it genuinely cannot persist. Callers omit the hash entirely
rather than emit an unsalted one.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
# Read before writing. Hooks are separate processes firing on every tool
# call, so all but the first find the salt already published; going straight
# to create-fsync-link-unlink meant every one of them paid an fsync to
# discover that, on a path whose whole promise is appending a line and
# returning.
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
if _salt_cache:
return _salt_cache
except OSError:
pass
temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp")
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
stream.flush()
os.fsync(stream.fileno())
try:
# Atomic claim: fails if another process already published one.
# os.link rather than replace, which would clobber theirs.
os.link(temporary, path)
except FileExistsError:
pass
except OSError:
# No hardlinks here (some network mounts, some container volumes).
# Claim the name directly instead. That reopens the empty-file
# window, but the window is now benign: a reader that lands in it
# gets "" and omits the hash for that process rather than caching a
# guessable one. Losing the hashes on every run of an entire
# filesystem is the worse failure.
try:
fallback = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(fallback, "w", encoding="utf-8") as stream:
stream.write(temporary.read_text(encoding="utf-8"))
except OSError:
pass
except OSError:
pass
finally:
try:
temporary.unlink()
except OSError:
pass
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
"""Salted digest for values drawn from a guessable space.
repo.identity is a git remote URL, or ``local:<absolute path>`` when there is
no remote — which normally contains the account username. Sixteen unsalted
hex characters over that input space is enumerable, so this is not a
privacy control without the salt. Salting per install keeps every
within-account join the analytics actually use and gives up only
cross-machine joins on the same repository, which nothing computes.
Returns "" when there is no salt, so record() omits the property. An
unsalted digest over this input space is close to plaintext, and emitting one
under a name that implies it is hashed is worse than sending nothing.
"""
if not value:
return ""
salt = _install_salt()
if not salt:
return ""
return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length]
def _safe_value(value: Any) -> Any:
if isinstance(value, str):
return memory_core.redact(value)
@@ -302,176 +145,9 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str:
return created
def _rotate_anonymous_id(identity: dict[str, str]) -> str:
"""Mint a fresh anonymous id because the account context is gone.
The previous id may already have been merged into a person profile by an
$identify, and that merge is permanent. Reusing it after a logout or a key
change attributes everything that follows to the account that just went
away, which is the same misattribution the key fingerprint exists to stop,
only arriving through the anonymous path instead.
`aliased` is cleared with it: the new id has never been merged, so it is
eligible to be aliased into whatever account comes next.
"""
created = f"code-anon-{uuid.uuid4().hex}"
identity["anonymous_id"] = created
identity.pop("aliased", None)
_write_identity(identity)
return created
def _install_state_path() -> Path:
return memory_core.data_dir() / "install-state.json"
def is_first_run() -> bool:
"""Whether install has never been recorded on this machine.
Deliberately NOT the identity file. That file is only written by a
successful flush, so an offline or firewalled user recorded code.install on
every single session, forever — and every 0.2.x user recorded one on their
first 0.3.x session because 0.2.x never wrote it at all.
"""
return not _install_state_path().exists()
def data_dir_was_empty() -> bool:
"""Whether the data directory is untouched. Call BEFORE anything writes to it.
hook_runner reaches claim_install() only after cache_plugin_api_key() has
written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so
asking at claim time always saw content and every fresh install reported an
upgrade. The caller snapshots this at the top of the run instead.
"""
return not _data_dir_has_content()
def claim_install(was_empty: bool | None = None) -> str | None:
"""Claim the one install/upgrade record for this machine, atomically.
Returns the event to record ("install" or "upgrade"), or None if another
session already claimed it. O_CREAT|O_EXCL so two sessions starting together
cannot both win.
`was_empty` must come from data_dir_was_empty() called before this process
wrote anything. Omitting it falls back to checking now, which is only
correct for a caller that has touched nothing.
"""
if not is_enabled():
# Never consume the one-shot claim while the user is opted out, or they
# would silently lose their install event if they later opt in.
return None
path = _install_state_path()
upgrading = not (data_dir_was_empty() if was_empty is None else was_empty)
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
except FileExistsError:
return None
except OSError:
return None
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
json.dump(
{
"plugin_version": memory_core.PLUGIN_VERSION,
"installed_at": memory_core.utc_now(),
"upgraded": upgrading,
},
stream,
)
# Durable before this returns. The O_EXCL open is what makes the
# claim exclusive, so it cannot be replaced by a temp-and-rename
# without losing that, which leaves the content as the thing to make
# safe. A kill between the open and this fsync used to leave a marker
# that exists but parses to nothing: is_first_run reads it as claimed
# and claim_version_change cannot read a version out of it.
stream.flush()
os.fsync(stream.fileno())
except OSError:
pass
return "upgrade" if upgrading else "install"
def _data_dir_has_content() -> bool:
"""Whether anything predates this session in the plugin data directory."""
try:
for entry in memory_core.data_dir().iterdir():
if entry.name != "install-state.json":
return True
except OSError:
pass
return False
def _repair_install_state(path: Path) -> None:
"""Rewrite an unparseable marker so version tracking can resume."""
try:
temporary = path.with_suffix(f".{os.getpid()}.tmp")
temporary.write_text(
json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}),
encoding="utf-8",
)
temporary.replace(path)
except OSError:
pass
def claim_version_change() -> str | None:
"""Return the previously recorded version if it differs, updating the marker.
Only meaningful once the marker exists — the first transition into 0.3.x has
no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting
the marker means the next session sees no change and records nothing.
"""
path = _install_state_path()
try:
state = json.loads(path.read_text(encoding="utf-8"))
except OSError:
return None
except json.JSONDecodeError:
# A crash between O_EXCL and the write leaves an empty marker. Left
# alone it disables every future upgrade event on this machine, because
# claim_install sees the file and this function cannot parse it.
state = None
if not isinstance(state, dict):
_repair_install_state(path)
return None
previous = str(state.get("plugin_version") or "")
if not previous or previous == memory_core.PLUGIN_VERSION:
return None
# Claim the transition with an exclusive sentinel before rewriting the
# marker. A plain read-modify-write let every concurrently starting session
# observe the old version and each record its own upgrade — and the first
# session after a version bump is exactly when several agent windows restart
# together.
sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}")
try:
os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600))
except FileExistsError:
return None
except OSError:
return None
state["plugin_version"] = memory_core.PLUGIN_VERSION
state["upgraded_at"] = memory_core.utc_now()
temporary = path.with_suffix(f".{os.getpid()}.tmp")
try:
temporary.write_text(json.dumps(state), encoding="utf-8")
temporary.replace(path)
except OSError:
# Release the claim. The marker still records the old version, so
# without this the sentinel makes claim_version_change return early on
# every later run and this version's upgrade is never recorded again.
for leftover in (sentinel, temporary):
try:
leftover.unlink()
except OSError:
pass
return None
return previous
"""Whether this machine has never recorded a plugin event before."""
return not _identity_path().exists()
def record(
@@ -492,32 +168,19 @@ def record(
except OSError:
pass
properties = _safe_value(properties)
# Stamped in the RECORDING process, beside harness. `source` used to be
# read in the sending process from a module global, so whichever process
# drained the spool named every event in it. flush() spreads per-event
# properties last, so this now wins over any sender's default.
properties.update(
harness=_harness,
source=_source_tag,
plugin_version=memory_core.PLUGIN_VERSION,
os=sys.platform,
python_version=platform.python_version(),
)
# Assigned only when the digest is real. _scoped_digest returns "" when
# the salt could not be persisted, and an empty property is worse than an
# absent one: it survives the None filter below and reads as a value.
if repo is not None:
repo_hash = _scoped_digest(getattr(repo, "identity", ""))
if repo_hash:
properties["repo_hash"] = repo_hash
properties["repo_hash"] = _digest(getattr(repo, "identity", ""))
if session_id:
session_hash = _scoped_digest(session_id)
if session_hash:
properties["session_hash"] = session_hash
properties["session_hash"] = _digest(session_id)
line = json.dumps(
{
"event": f"{EVENT_PREFIX}.{event}",
"uuid": str(uuid.uuid4()),
"timestamp": memory_core.utc_now(),
"properties": {
key: value for key, value in properties.items() if value is not None
@@ -576,201 +239,38 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / _claim_name()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
for orphan in sorted(directory.glob("telemetry-*.sending")):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
if age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
pass
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
if age < CLAIM_STALE_SECONDS:
continue
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -800,130 +300,34 @@ def _post(payload: dict[str, Any], url: str) -> bool:
def resolve_distinct_id() -> tuple[str, str]:
"""Return the PostHog distinct id and the anonymous id it replaced, if any.
The second value becomes a PostHog $identify alias. It is ONLY ever an
anonymous id: aliasing one account email to another merges two real person
profiles and cannot be undone, so a key that now belongs to a different
account re-resolves with no alias.
"""
"""Return the PostHog distinct id and the anonymous id it replaced, if any."""
identity = _read_identity()
key = memory_core.api_key()
fingerprint = _digest(key) if key else ""
email = identity.get("email", "")
if email and fingerprint:
recorded = identity.get("key_fingerprint", "")
if recorded == fingerprint:
return email, ""
if not recorded:
# Rows written before fingerprints existed. Verify rather than
# adopt: a key changed before the upgrade would otherwise bind the
# new key to the previous account's email, permanently, and the
# fingerprint would then agree with itself forever after.
verified = _resolve_email(key)
if not verified:
# Offline, firewalled, or the API is down. Keep the previous
# behaviour and retry on the next flush rather than dropping a
# real account attribution. Safe because the same network that
# failed /v1/ping/ is about to fail the PostHog POST, so nothing
# is delivered under the unverified identity in the meantime.
return email, ""
identity["email"] = verified
identity["key_fingerprint"] = fingerprint
_write_identity(identity)
return verified, ""
if email:
return email, ""
key = memory_core.api_key()
if not key:
# No key to verify the account with; do not keep attributing to it.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
return anonymous_id(identity), ""
resolved = _resolve_email(key)
if not resolved:
# The key changed and will not resolve (revoked, offline, API down).
# Reaching here with an email means the recorded fingerprint disagreed,
# so the key really did change. Drop the account and rotate: the stored
# anonymous id may already be merged into that account's person, and
# reusing it would keep the events on the profile we are trying to
# leave.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
email = _resolve_email(key)
if not email:
return anonymous_id(identity), ""
# Alias only when going anonymous -> email for the first time. Once an anon
# id has been merged into an account it must never be offered again: an
# alias naming an already-identified id is what could link two real people.
previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "")
if previous:
identity["aliased"] = True
identity["email"] = resolved
identity["key_fingerprint"] = fingerprint
previous = identity.get("anonymous_id", "")
identity["email"] = email
_write_identity(identity)
return resolved, previous
return email, previous
def flush() -> int:
"""Drain the live spool, then any parked claims, and return events sent."""
"""Drain claimed spools to PostHog and return the number of events sent."""
if not is_enabled():
return 0
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
claim = _claim_spool()
if claim is None:
return 0, True
return 0
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# quarantine rather than retry. flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
return 0
events = []
for line in lines:
try:
@@ -933,18 +337,11 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
claim.unlink()
except OSError:
pass
return 0, True
return 0
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -963,17 +360,12 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
# moved into record() have none of their own.
"source": _source_tag,
"language": "python",
"$process_person_profile": False,
@@ -981,24 +373,16 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
**(event.get("properties") or {}),
},
}
for event in chunk
for event in events[start : start + BATCH_SIZE]
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
def main() -> int:
@@ -7,8 +7,8 @@ disable-model-invocation: true
# Pause memory capture
To pause (hooks stop capturing and sending session content; a minimal
telemetry ping still fires at session start, under your Mem0 account email,
unless `MEM0_TELEMETRY=false`):
anonymous telemetry ping still fires at session start unless
`MEM0_TELEMETRY=false`):
```bash
python3 "${ANTIGRAVITY_PLUGIN_ROOT}/core/memory_cli.py" --harness "antigravity" pause
+2 -12
View File
@@ -136,23 +136,13 @@ Local data lives in `${CLAUDE_PLUGIN_DATA}`:
- `pending/`: sessions waiting to be sent to Mem0 (retried after interruption)
- `flush-worker.log`: whether memory creation succeeded
- `plugin-errors.log`: hook errors (no credentials)
- `telemetry.jsonl` / `telemetry-identity.json`: usage events and the id they are sent under
- `telemetry-salt`: random per-install salt for the repo and session hashes
- `install-state.json`: records that install has been counted once on this machine
- `telemetry.jsonl` / `telemetry-identity.json`: anonymous usage events
Mem0 receives captured user messages, Claude's answers, sidekick assignments and completed responses, and changed file paths. When a failed command is recorded, extraction can also include bounded command details and results. Complete files and general tool output stay on your machine. Values that look like credentials are redacted before anything is sent.
## Telemetry
Usage events (which hook ran, timing, result counts, failure types) so Mem0 can identify what's used and what's breaking.
**These events are not anonymous.** When an API key is configured — which installing the plugin requires — events are sent under your Mem0 account email, the same way the Python SDK and the CLI attribute theirs. Without a key they are sent under a random per-machine id.
What each event carries: the event name, the plugin version, the harness it ran in, your OS and Python version, and per-event properties describing what happened — timings, counts, coarse outcome and failure labels, and which model was configured. Repository and session identifiers are hashed with a random salt generated on your machine, so they cannot be linked back to a repository name or path.
Rather than restate a list that drifts, the exact set is enforced in code: `telemetry.record` filters every property through a denylist of sensitive keys and redacts credential-shaped values. See `_PRIVATE_KEYS` in `core/telemetry.py`.
Prompts, memory text, queries, file paths, repository names, and API keys are never sent.
Anonymous usage events (which hook ran, timing, result counts, failure types) so Mem0 can identify what's used and what's breaking. Repo and session IDs are hashed before leaving your machine. Prompts, memory text, file paths, tool output, and API keys are never sent.
Turn it off:
@@ -1,11 +0,0 @@
"""Generated by integrations/agent-plugin-core/build/build.py. Do not edit."""
HARNESS_ID = "claude-code"
SOURCE_TAG = "CLAUDE_CODE_PLUGIN"
# Platform-side vocabulary (mem0_event.source + X-Application). The whole
# plugin family is one source; which editor it runs in is the application.
# An empty application means the host is unknown, and memory_core omits
# the header entirely rather than sending a placeholder.
PLATFORM_SOURCE = "MEM0_PLUGIN"
PLATFORM_APPLICATION = "claude-code"
@@ -290,11 +290,6 @@ def run(
if args.plugin_data_dir:
os.environ[data_dir_env] = args.plugin_data_dir
# Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key
# writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking
# after them always saw content and every fresh install reported an upgrade.
data_dir_was_empty = telemetry.data_dir_was_empty()
cache_plugin_api_key()
if args.action == "session-start":
clear_stale_api_key_cache()
@@ -310,19 +305,8 @@ def run(
return 0
if args.action == "session-start":
# Claims the marker atomically and says which event to record, so a
# second session starting alongside this one cannot record it too.
first_event = telemetry.claim_install(was_empty=data_dir_was_empty)
if first_event == "install":
if telemetry.is_first_run():
telemetry.record("install")
elif first_event == "upgrade":
# First run after a build that never wrote the marker; the
# predecessor version was never recorded anywhere.
telemetry.record("upgrade", from_version="pre-0.3")
else:
previous = telemetry.claim_version_change()
if previous:
telemetry.record("upgrade", from_version=previous)
recovered = recover_pending_handoffs()
record_session_start(store, hook_input)
if recovered:
@@ -1800,34 +1800,6 @@ def extraction_message_batches(
return batches
# Platform surface attribution. Read from the generated per-host module so a new
# entrypoint is correct without remembering to configure anything.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
def platform_headers(key: str) -> dict[str, str]:
"""Auth plus the three surface-identity headers.
X-Mem0-Source and X-Application are set-once by contract: this is the
outermost layer, so it sets them, and nothing below may overwrite them.
X-Mem0-Client is append-only — anything downstream adds itself to the tail.
"""
headers = {
"Authorization": f"Token {key}",
"Content-Type": "application/json",
"X-Mem0-Source": _PLATFORM_SOURCE,
"X-Mem0-Client": f"mem0-plugin/{PLUGIN_VERSION}",
}
if _PLATFORM_APPLICATION:
headers["X-Application"] = _PLATFORM_APPLICATION
return headers
def _request_json(
url: str, key: str, payload: dict[str, Any], timeout: float
) -> tuple[dict[str, Any] | list[Any], int, int]:
@@ -1835,7 +1807,7 @@ def _request_json(
request = urllib.request.Request(
url,
data=raw,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -1862,7 +1834,7 @@ def _get_json(
) -> tuple[dict[str, Any] | list[Any], int]:
request = urllib.request.Request(
url,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="GET",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -2008,13 +1980,6 @@ def flush_session(
"user_id": write_user,
"app_id": repo.app_id,
"run_id": session_id,
# Top level, not metadata: the backend reads `source` from the body or
# the query string, never from metadata, which is where this used to
# sit. The X-Mem0-Source header is also read, but only from the
# platform release that ships alongside this change, so the body value
# is what makes attribution work on both. The harness tag stays in
# metadata as hook provenance.
"source": _PLATFORM_SOURCE,
"metadata": {**metadata, "author": write_user, "dirs": directory_chain(repo)},
"agent_custom_instructions": PROJECT_MEMORY_INSTRUCTIONS,
"custom_instructions": PERSONAL_MEMORY_INSTRUCTIONS,
@@ -2558,7 +2523,7 @@ def _collect_memory_ids(
def _delete_memory(api_url: str, key: str, memory_id: str) -> bool:
request = urllib.request.Request(
f"{api_url}/v1/memories/{urllib.parse.quote(memory_id)}/",
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="DELETE",
)
try:
+39 -655
View File
@@ -1,9 +1,5 @@
#!/usr/bin/env python3
"""Usage telemetry for Mem0 agent plugins.
Events are linked to your Mem0 account email when an API key is configured, and
to a random per-machine id otherwise. Not anonymous — the Python SDK and CLI
attribute the same way.
"""Anonymous usage telemetry for Mem0 agent plugins.
Hooks run on a 3-6 second budget and fire on every tool call, so recording never
touches the network: `record` appends one JSON line to a local spool and returns.
@@ -13,8 +9,7 @@ started once per session and again from the flush worker that is already detache
Pure stdlib, matching the rest of the plugin. Opt out with MEM0_TELEMETRY=false.
Never sends prompts, memory text, queries, file paths, repository names, or API
keys: only event names, durations, counts, coarse outcomes, and repo/session
identifiers hashed with a random per-install salt.
keys: only event names, durations, counts, coarse outcomes, and salted hashes.
"""
from __future__ import annotations
@@ -34,24 +29,8 @@ from typing import Any
import memory_core
# Seeded from the per-host module the build generates into core/. Two processes
# in this pipeline never call init() — mcp_server.py, and the detached
# `python3 telemetry.py` sender that spawn_flush() starts — so a module default
# was what every one of their events got labelled with.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import HARNESS_ID as _DEFAULT_HARNESS
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
from _harness_id import SOURCE_TAG as _DEFAULT_SOURCE_TAG
except ImportError:
_DEFAULT_HARNESS = "generic"
_DEFAULT_SOURCE_TAG = "MEM0_PLUGIN"
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_harness: str = "generic"
_source_tag: str = "MEM0_PLUGIN"
_PRIVATE_KEYS = {
"apikey",
"authorization",
@@ -77,19 +56,10 @@ _PRIVATE_KEYS = {
}
def init(harness: str = "", source_tag: str = "") -> None:
"""Override the generated identity. Optional — core/_harness_id.py is the default.
The fallback shape matches memory_core.configure_harness's (``<HOST>_PLUGIN``).
It used to be ``MEM0_<HOST>_PLUGIN`` here and ``<host>_plugin`` there, which
meant one plugin could emit three different source values depending on which
process happened to send the batch.
"""
def init(harness: str = "generic", source_tag: str = "") -> None:
global _harness, _source_tag
_harness = harness or _DEFAULT_HARNESS
_source_tag = source_tag or (
f"{_harness.upper().replace('-', '_')}_PLUGIN" if harness else _DEFAULT_SOURCE_TAG
)
_harness = harness
_source_tag = source_tag or f"MEM0_{harness.upper().replace('-', '_')}_PLUGIN"
POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"
POSTHOG_CAPTURE_URL = "https://us.i.posthog.com/i/v0/e/"
@@ -100,16 +70,6 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -123,126 +83,9 @@ def is_enabled() -> bool:
def _digest(value: str, length: int = 16) -> str:
"""Unsalted digest. Only for values that are already secrets (API keys)."""
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
Published atomically, and there is deliberately no derived fallback. Creating
the file with O_CREAT|O_EXCL and then writing into it leaves a window where
the file exists and is empty, and a concurrent hook that reads it in that
window gets nothing. Falling back to a digest of the path would hand that
process a salt an attacker can compute, memoized for its whole run, which is
the privacy control this function exists to provide silently turning itself
off under load. The salt is written to a private temp file first and linked
into place, so the name either does not exist or already has the full value.
Returns "" when it genuinely cannot persist. Callers omit the hash entirely
rather than emit an unsalted one.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
# Read before writing. Hooks are separate processes firing on every tool
# call, so all but the first find the salt already published; going straight
# to create-fsync-link-unlink meant every one of them paid an fsync to
# discover that, on a path whose whole promise is appending a line and
# returning.
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
if _salt_cache:
return _salt_cache
except OSError:
pass
temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp")
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
stream.flush()
os.fsync(stream.fileno())
try:
# Atomic claim: fails if another process already published one.
# os.link rather than replace, which would clobber theirs.
os.link(temporary, path)
except FileExistsError:
pass
except OSError:
# No hardlinks here (some network mounts, some container volumes).
# Claim the name directly instead. That reopens the empty-file
# window, but the window is now benign: a reader that lands in it
# gets "" and omits the hash for that process rather than caching a
# guessable one. Losing the hashes on every run of an entire
# filesystem is the worse failure.
try:
fallback = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(fallback, "w", encoding="utf-8") as stream:
stream.write(temporary.read_text(encoding="utf-8"))
except OSError:
pass
except OSError:
pass
finally:
try:
temporary.unlink()
except OSError:
pass
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
"""Salted digest for values drawn from a guessable space.
repo.identity is a git remote URL, or ``local:<absolute path>`` when there is
no remote — which normally contains the account username. Sixteen unsalted
hex characters over that input space is enumerable, so this is not a
privacy control without the salt. Salting per install keeps every
within-account join the analytics actually use and gives up only
cross-machine joins on the same repository, which nothing computes.
Returns "" when there is no salt, so record() omits the property. An
unsalted digest over this input space is close to plaintext, and emitting one
under a name that implies it is hashed is worse than sending nothing.
"""
if not value:
return ""
salt = _install_salt()
if not salt:
return ""
return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length]
def _safe_value(value: Any) -> Any:
if isinstance(value, str):
return memory_core.redact(value)
@@ -302,176 +145,9 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str:
return created
def _rotate_anonymous_id(identity: dict[str, str]) -> str:
"""Mint a fresh anonymous id because the account context is gone.
The previous id may already have been merged into a person profile by an
$identify, and that merge is permanent. Reusing it after a logout or a key
change attributes everything that follows to the account that just went
away, which is the same misattribution the key fingerprint exists to stop,
only arriving through the anonymous path instead.
`aliased` is cleared with it: the new id has never been merged, so it is
eligible to be aliased into whatever account comes next.
"""
created = f"code-anon-{uuid.uuid4().hex}"
identity["anonymous_id"] = created
identity.pop("aliased", None)
_write_identity(identity)
return created
def _install_state_path() -> Path:
return memory_core.data_dir() / "install-state.json"
def is_first_run() -> bool:
"""Whether install has never been recorded on this machine.
Deliberately NOT the identity file. That file is only written by a
successful flush, so an offline or firewalled user recorded code.install on
every single session, forever — and every 0.2.x user recorded one on their
first 0.3.x session because 0.2.x never wrote it at all.
"""
return not _install_state_path().exists()
def data_dir_was_empty() -> bool:
"""Whether the data directory is untouched. Call BEFORE anything writes to it.
hook_runner reaches claim_install() only after cache_plugin_api_key() has
written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so
asking at claim time always saw content and every fresh install reported an
upgrade. The caller snapshots this at the top of the run instead.
"""
return not _data_dir_has_content()
def claim_install(was_empty: bool | None = None) -> str | None:
"""Claim the one install/upgrade record for this machine, atomically.
Returns the event to record ("install" or "upgrade"), or None if another
session already claimed it. O_CREAT|O_EXCL so two sessions starting together
cannot both win.
`was_empty` must come from data_dir_was_empty() called before this process
wrote anything. Omitting it falls back to checking now, which is only
correct for a caller that has touched nothing.
"""
if not is_enabled():
# Never consume the one-shot claim while the user is opted out, or they
# would silently lose their install event if they later opt in.
return None
path = _install_state_path()
upgrading = not (data_dir_was_empty() if was_empty is None else was_empty)
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
except FileExistsError:
return None
except OSError:
return None
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
json.dump(
{
"plugin_version": memory_core.PLUGIN_VERSION,
"installed_at": memory_core.utc_now(),
"upgraded": upgrading,
},
stream,
)
# Durable before this returns. The O_EXCL open is what makes the
# claim exclusive, so it cannot be replaced by a temp-and-rename
# without losing that, which leaves the content as the thing to make
# safe. A kill between the open and this fsync used to leave a marker
# that exists but parses to nothing: is_first_run reads it as claimed
# and claim_version_change cannot read a version out of it.
stream.flush()
os.fsync(stream.fileno())
except OSError:
pass
return "upgrade" if upgrading else "install"
def _data_dir_has_content() -> bool:
"""Whether anything predates this session in the plugin data directory."""
try:
for entry in memory_core.data_dir().iterdir():
if entry.name != "install-state.json":
return True
except OSError:
pass
return False
def _repair_install_state(path: Path) -> None:
"""Rewrite an unparseable marker so version tracking can resume."""
try:
temporary = path.with_suffix(f".{os.getpid()}.tmp")
temporary.write_text(
json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}),
encoding="utf-8",
)
temporary.replace(path)
except OSError:
pass
def claim_version_change() -> str | None:
"""Return the previously recorded version if it differs, updating the marker.
Only meaningful once the marker exists — the first transition into 0.3.x has
no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting
the marker means the next session sees no change and records nothing.
"""
path = _install_state_path()
try:
state = json.loads(path.read_text(encoding="utf-8"))
except OSError:
return None
except json.JSONDecodeError:
# A crash between O_EXCL and the write leaves an empty marker. Left
# alone it disables every future upgrade event on this machine, because
# claim_install sees the file and this function cannot parse it.
state = None
if not isinstance(state, dict):
_repair_install_state(path)
return None
previous = str(state.get("plugin_version") or "")
if not previous or previous == memory_core.PLUGIN_VERSION:
return None
# Claim the transition with an exclusive sentinel before rewriting the
# marker. A plain read-modify-write let every concurrently starting session
# observe the old version and each record its own upgrade — and the first
# session after a version bump is exactly when several agent windows restart
# together.
sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}")
try:
os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600))
except FileExistsError:
return None
except OSError:
return None
state["plugin_version"] = memory_core.PLUGIN_VERSION
state["upgraded_at"] = memory_core.utc_now()
temporary = path.with_suffix(f".{os.getpid()}.tmp")
try:
temporary.write_text(json.dumps(state), encoding="utf-8")
temporary.replace(path)
except OSError:
# Release the claim. The marker still records the old version, so
# without this the sentinel makes claim_version_change return early on
# every later run and this version's upgrade is never recorded again.
for leftover in (sentinel, temporary):
try:
leftover.unlink()
except OSError:
pass
return None
return previous
"""Whether this machine has never recorded a plugin event before."""
return not _identity_path().exists()
def record(
@@ -492,32 +168,19 @@ def record(
except OSError:
pass
properties = _safe_value(properties)
# Stamped in the RECORDING process, beside harness. `source` used to be
# read in the sending process from a module global, so whichever process
# drained the spool named every event in it. flush() spreads per-event
# properties last, so this now wins over any sender's default.
properties.update(
harness=_harness,
source=_source_tag,
plugin_version=memory_core.PLUGIN_VERSION,
os=sys.platform,
python_version=platform.python_version(),
)
# Assigned only when the digest is real. _scoped_digest returns "" when
# the salt could not be persisted, and an empty property is worse than an
# absent one: it survives the None filter below and reads as a value.
if repo is not None:
repo_hash = _scoped_digest(getattr(repo, "identity", ""))
if repo_hash:
properties["repo_hash"] = repo_hash
properties["repo_hash"] = _digest(getattr(repo, "identity", ""))
if session_id:
session_hash = _scoped_digest(session_id)
if session_hash:
properties["session_hash"] = session_hash
properties["session_hash"] = _digest(session_id)
line = json.dumps(
{
"event": f"{EVENT_PREFIX}.{event}",
"uuid": str(uuid.uuid4()),
"timestamp": memory_core.utc_now(),
"properties": {
key: value for key, value in properties.items() if value is not None
@@ -576,201 +239,38 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / _claim_name()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
for orphan in sorted(directory.glob("telemetry-*.sending")):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
if age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
pass
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
if age < CLAIM_STALE_SECONDS:
continue
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -800,130 +300,34 @@ def _post(payload: dict[str, Any], url: str) -> bool:
def resolve_distinct_id() -> tuple[str, str]:
"""Return the PostHog distinct id and the anonymous id it replaced, if any.
The second value becomes a PostHog $identify alias. It is ONLY ever an
anonymous id: aliasing one account email to another merges two real person
profiles and cannot be undone, so a key that now belongs to a different
account re-resolves with no alias.
"""
"""Return the PostHog distinct id and the anonymous id it replaced, if any."""
identity = _read_identity()
key = memory_core.api_key()
fingerprint = _digest(key) if key else ""
email = identity.get("email", "")
if email and fingerprint:
recorded = identity.get("key_fingerprint", "")
if recorded == fingerprint:
return email, ""
if not recorded:
# Rows written before fingerprints existed. Verify rather than
# adopt: a key changed before the upgrade would otherwise bind the
# new key to the previous account's email, permanently, and the
# fingerprint would then agree with itself forever after.
verified = _resolve_email(key)
if not verified:
# Offline, firewalled, or the API is down. Keep the previous
# behaviour and retry on the next flush rather than dropping a
# real account attribution. Safe because the same network that
# failed /v1/ping/ is about to fail the PostHog POST, so nothing
# is delivered under the unverified identity in the meantime.
return email, ""
identity["email"] = verified
identity["key_fingerprint"] = fingerprint
_write_identity(identity)
return verified, ""
if email:
return email, ""
key = memory_core.api_key()
if not key:
# No key to verify the account with; do not keep attributing to it.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
return anonymous_id(identity), ""
resolved = _resolve_email(key)
if not resolved:
# The key changed and will not resolve (revoked, offline, API down).
# Reaching here with an email means the recorded fingerprint disagreed,
# so the key really did change. Drop the account and rotate: the stored
# anonymous id may already be merged into that account's person, and
# reusing it would keep the events on the profile we are trying to
# leave.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
email = _resolve_email(key)
if not email:
return anonymous_id(identity), ""
# Alias only when going anonymous -> email for the first time. Once an anon
# id has been merged into an account it must never be offered again: an
# alias naming an already-identified id is what could link two real people.
previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "")
if previous:
identity["aliased"] = True
identity["email"] = resolved
identity["key_fingerprint"] = fingerprint
previous = identity.get("anonymous_id", "")
identity["email"] = email
_write_identity(identity)
return resolved, previous
return email, previous
def flush() -> int:
"""Drain the live spool, then any parked claims, and return events sent."""
"""Drain claimed spools to PostHog and return the number of events sent."""
if not is_enabled():
return 0
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
claim = _claim_spool()
if claim is None:
return 0, True
return 0
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# quarantine rather than retry. flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
return 0
events = []
for line in lines:
try:
@@ -933,18 +337,11 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
claim.unlink()
except OSError:
pass
return 0, True
return 0
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -963,17 +360,12 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
# moved into record() have none of their own.
"source": _source_tag,
"language": "python",
"$process_person_profile": False,
@@ -981,24 +373,16 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
**(event.get("properties") or {}),
},
}
for event in chunk
for event in events[start : start + BATCH_SIZE]
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
def main() -> int:
@@ -7,8 +7,8 @@ disable-model-invocation: true
# Pause memory capture
To pause (hooks stop capturing and sending session content; a minimal
telemetry ping still fires at session start, under your Mem0 account email,
unless `MEM0_TELEMETRY=false`):
anonymous telemetry ping still fires at session start unless
`MEM0_TELEMETRY=false`):
```bash
python3 "${CLAUDE_PLUGIN_ROOT}/core/memory_cli.py" --harness "claude-code" --plugin-data-dir "${CLAUDE_PLUGIN_DATA}" pause
@@ -1,7 +1,6 @@
from __future__ import annotations
import json
import os
import sys
from pathlib import Path
from unittest.mock import patch
@@ -199,11 +198,9 @@ def test_a_stale_claim_is_reclaimed(isolated_env, monkeypatch):
telemetry.record("search")
orphan = telemetry._claim_spool()
assert orphan is not None
# Frozen rather than re-stat'd per call: flush() drains the live spool and
# then looks for parked claims in the same run, so by the second look this
# file no longer exists.
stale_now = orphan.stat().st_mtime + telemetry.CLAIM_STALE_SECONDS + 1
monkeypatch.setattr(telemetry.time, "time", lambda: stale_now)
monkeypatch.setattr(
telemetry.time, "time", lambda: orphan.stat().st_mtime + telemetry.CLAIM_STALE_SECONDS + 1
)
with patch.object(telemetry, "_post", lambda payload, url: True):
assert telemetry.flush() == 1
@@ -213,13 +210,9 @@ def test_an_expired_claim_is_dropped(isolated_env, monkeypatch):
telemetry.record("search")
orphan = telemetry._claim_spool()
assert orphan is not None
expired_now = orphan.stat().st_mtime + telemetry.CLAIM_EXPIRY_SECONDS + 1
monkeypatch.setattr(telemetry.time, "time", lambda: expired_now)
# Expiry now only discards a batch that was genuinely retried and failed,
# so age alone is not enough — age it past the attempt budget too.
retried = orphan.parent / orphan.name.replace("-a0.", f"-a{telemetry.MAX_CLAIM_ATTEMPTS}.")
orphan.replace(retried)
os.utime(retried, (expired_now, expired_now - telemetry.CLAIM_EXPIRY_SECONDS - 1))
monkeypatch.setattr(
telemetry.time, "time", lambda: orphan.stat().st_mtime + telemetry.CLAIM_EXPIRY_SECONDS + 1
)
assert telemetry._claim_spool() is None
assert not list(memory_core.data_dir().glob("telemetry-*.sending"))
@@ -265,150 +258,12 @@ def test_an_unresolvable_key_falls_back_to_the_anonymous_id(isolated_env, monkey
assert telemetry.resolve_distinct_id()[0].startswith("code-anon-")
def test_logging_out_does_not_leave_events_on_the_previous_account(isolated_env, monkeypatch):
"""Review finding: clearing the email kept an id already merged into a person.
The anonymous id is offered to PostHog as $anon_distinct_id on first sign-in,
and that merge is permanent. Keeping it after the key goes away means every
later anonymous event lands on the account that just left.
"""
# Run anonymously first, which is the only way an id exists to be merged.
merged = telemetry.anonymous_id()
monkeypatch.setenv("MEM0_API_KEY", "key-for-account-a")
with patch.object(telemetry, "_resolve_email", lambda key: "a@example.com"):
identified, alias = telemetry.resolve_distinct_id()
assert identified == "a@example.com"
assert alias == merged, "the anonymous id was merged into this account"
monkeypatch.delenv("MEM0_API_KEY", raising=False)
after_logout, logout_alias = telemetry.resolve_distinct_id()
assert after_logout.startswith("code-anon-")
assert after_logout != merged, "reused an id already merged into the previous account"
assert logout_alias == ""
assert "aliased" not in telemetry._read_identity(), "rotated id must be aliasable again"
def test_a_changed_key_that_will_not_resolve_rotates_the_anonymous_id(isolated_env, monkeypatch):
"""Same leak by the other route: fingerprint disagrees and the lookup fails."""
merged = telemetry.anonymous_id()
monkeypatch.setenv("MEM0_API_KEY", "key-for-account-a")
with patch.object(telemetry, "_resolve_email", lambda key: "a@example.com"):
telemetry.resolve_distinct_id()
monkeypatch.setenv("MEM0_API_KEY", "key-for-account-b")
with patch.object(telemetry, "_resolve_email", lambda key: ""):
after, alias = telemetry.resolve_distinct_id()
assert after.startswith("code-anon-")
assert after != merged
assert alias == ""
assert "email" not in telemetry._read_identity()
def test_a_legacy_cached_email_is_verified_before_the_key_is_bound(isolated_env, monkeypatch):
"""Review finding: a key changed before upgrading bound the wrong account.
Rows written before fingerprints existed carry an email and no fingerprint.
Adopting the current key without checking pinned that key to the previous
account's email, and every run after that agreed with itself.
"""
telemetry._write_identity({"email": "old@example.com", "anonymous_id": "code-anon-seed"})
monkeypatch.setenv("MEM0_API_KEY", "key-for-account-b")
with patch.object(telemetry, "_resolve_email", lambda key: "new@example.com"):
resolved, alias = telemetry.resolve_distinct_id()
assert resolved == "new@example.com"
assert alias == "", "email to email must never alias; it merges two real people"
stored = telemetry._read_identity()
assert stored["email"] == "new@example.com"
assert stored["key_fingerprint"] == telemetry._digest("key-for-account-b")
def test_a_legacy_row_keeps_working_when_the_account_cannot_be_checked(isolated_env, monkeypatch):
"""Firewalled users must not lose attribution, and must not bind unverified.
The same network that fails /v1/ping/ fails the PostHog POST, so nothing is
delivered under the unverified identity while this holds.
"""
telemetry._write_identity({"email": "old@example.com"})
monkeypatch.setenv("MEM0_API_KEY", "key-for-account-b")
with patch.object(telemetry, "_resolve_email", lambda key: ""):
resolved, _ = telemetry.resolve_distinct_id()
assert resolved == "old@example.com"
assert "key_fingerprint" not in telemetry._read_identity(), "bound an unverified key"
def test_a_failed_upgrade_claim_can_be_retried(isolated_env, monkeypatch):
"""Review finding: a failed rewrite left the sentinel and suppressed forever.
claim_version_change returns early on FileExistsError, and the marker still
holds the old version, so the upgrade for that version was never recorded
again on that machine.
"""
telemetry.claim_install()
state_path = memory_core.data_dir() / "install-state.json"
state = json.loads(state_path.read_text())
state["plugin_version"] = "0.0.1-old"
state_path.write_text(json.dumps(state), encoding="utf-8")
real_replace = Path.replace
def failing_replace(self, target):
raise OSError("disk full")
monkeypatch.setattr(Path, "replace", failing_replace)
assert telemetry.claim_version_change() is None
monkeypatch.setattr(Path, "replace", real_replace)
assert telemetry.claim_version_change() == "0.0.1-old", "sentinel suppressed the retry"
def test_first_run_is_not_flipped_by_writing_the_identity_file(isolated_env):
"""The identity file is written by a successful flush, not by recording.
Keying first-run off it meant an offline user recorded code.install on every
session forever, and every 0.2.x user recorded one on their first 0.3.x run.
"""
def test_is_first_run_flips_after_the_first_identity_write(isolated_env):
assert telemetry.is_first_run()
telemetry.anonymous_id()
assert telemetry.is_first_run()
def test_claiming_install_ends_first_run(isolated_env):
assert telemetry.claim_install() == "install"
assert not telemetry.is_first_run()
def test_install_can_only_be_claimed_once(isolated_env):
"""Two sessions starting together must not both record an install."""
assert telemetry.claim_install() == "install"
assert telemetry.claim_install() is None
def test_a_populated_data_dir_reads_as_an_upgrade(isolated_env):
"""A fresh install has an empty data directory; anything else predates it."""
data_dir = memory_core.data_dir()
data_dir.mkdir(parents=True, exist_ok=True)
(data_dir / "requirements.txt").write_text("mem0ai\n", encoding="utf-8")
assert telemetry.claim_install() == "upgrade"
def test_a_version_change_is_claimed_once(isolated_env):
telemetry.claim_install()
state_path = memory_core.data_dir() / "install-state.json"
state = json.loads(state_path.read_text())
state["plugin_version"] = "0.0.1-old"
state_path.write_text(json.dumps(state), encoding="utf-8")
assert telemetry.claim_version_change() == "0.0.1-old"
assert telemetry.claim_version_change() is None
def test_spawn_flush_does_nothing_without_a_spool(isolated_env):
with patch.object(telemetry.subprocess, "Popen") as popen:
assert telemetry.spawn_flush() is False
@@ -418,114 +273,3 @@ def test_spawn_flush_does_nothing_without_a_spool(isolated_env):
with patch.object(telemetry.subprocess, "Popen") as popen:
assert telemetry.spawn_flush() is True
popen.assert_called_once()
def test_salt_is_stable_across_processes(isolated_env):
"""Hooks are separate short-lived processes; one repo must hash one way.
An unlocked read-modify-write let each process mint its own salt, so a
repository hashed several ways in the window before one writer won.
"""
import subprocess as sp
core = str(Path(__file__).resolve().parents[1] / "core")
script = (
f"import sys; sys.path.insert(0, {core!r})\n"
"import telemetry\n"
"print(telemetry._install_salt())"
)
env = {**os.environ, "MEM0_CODE_DATA_DIR": str(memory_core.data_dir())}
salts = {
sp.run([sys.executable, "-c", script], capture_output=True, text=True, env=env).stdout.strip()
for _ in range(4)
}
assert len(salts) == 1, f"one repo hashed {len(salts)} ways: {salts}"
def test_salt_does_not_touch_the_identity_file(isolated_env):
"""The identity file is is_first_run's marker and the sender's email store.
Writing the salt into it would create it from record(), suppressing the
install event, and would race resolve_distinct_id, which holds a stale copy
of that dict across a network call.
"""
telemetry._install_salt()
assert not telemetry._identity_path().exists()
def test_no_salt_means_no_hash_rather_than_an_unsalted_one(isolated_env, monkeypatch):
"""A read-only data dir drops the property; it must not emit a weak digest.
The previous fallback was a digest of the salt file's own path, which an
attacker can compute, memoized for the whole process. A property named
repo_hash carrying an effectively unsalted digest is worse than no property:
it reads as protected and is not.
"""
telemetry._salt_cache = ""
monkeypatch.setattr(telemetry.os, "open", lambda *a, **k: (_ for _ in ()).throw(OSError("read-only")))
assert telemetry._install_salt() == ""
assert telemetry._scoped_digest("git@github.com:acme/secret.git") == ""
def test_a_half_written_salt_is_never_visible_to_another_process(isolated_env, monkeypatch):
"""The window this closes: file created, value not yet written.
O_CREAT|O_EXCL then write leaves the name present and empty in between. A
hook reading it there used to get "", fall back to the path digest and cache
that for its whole run, so the same repo hashed two ways depending on timing.
Publishing by link means the name either does not exist or is complete.
"""
telemetry._salt_cache = ""
salt_path = telemetry._salt_path()
observed = []
real_link = telemetry.os.link
def observing_link(source, target):
# Stand where the racing reader stands: after the temp file is written,
# before the real name exists.
observed.append(salt_path.exists())
return real_link(source, target)
monkeypatch.setattr(telemetry.os, "link", observing_link)
salt = telemetry._install_salt()
assert observed == [False], "the salt name existed before it held a value"
assert len(salt) == 32
assert salt_path.read_text(encoding="utf-8").strip() == salt
def test_a_filesystem_without_hardlinks_still_gets_a_salt(isolated_env, monkeypatch):
"""Publishing by link must not become a silent loss of the hashes.
Some network mounts and container volumes reject os.link. Returning ""
there would drop repo_hash and session_hash on every run for that whole
cohort, which is a bigger loss than the narrow race the link closes.
"""
telemetry._salt_cache = ""
monkeypatch.setattr(
telemetry.os, "link", lambda src, dst: (_ for _ in ()).throw(OSError(38, "not implemented"))
)
salt = telemetry._install_salt()
assert len(salt) == 32, "no salt on a filesystem without hardlinks"
assert telemetry._salt_path().read_text(encoding="utf-8").strip() == salt
assert telemetry._scoped_digest("git@github.com:acme/x.git") != ""
assert not list(telemetry._salt_path().parent.glob("telemetry-salt.*.tmp"))
def test_a_concurrent_writer_does_not_clobber_the_published_salt(isolated_env):
"""Second process to finish must adopt the first one's salt, not replace it.
os.link rather than os.replace is what makes losing the race harmless.
"""
telemetry._salt_cache = ""
first = telemetry._install_salt()
telemetry._salt_cache = ""
second = telemetry._install_salt()
assert second == first
assert not list(telemetry._salt_path().parent.glob("telemetry-salt.*.tmp")), "temp file left behind"
@@ -1,11 +0,0 @@
"""Generated by integrations/agent-plugin-core/build/build.py. Do not edit."""
HARNESS_ID = "codex"
SOURCE_TAG = "CODEX_PLUGIN"
# Platform-side vocabulary (mem0_event.source + X-Application). The whole
# plugin family is one source; which editor it runs in is the application.
# An empty application means the host is unknown, and memory_core omits
# the header entirely rather than sending a placeholder.
PLATFORM_SOURCE = "MEM0_PLUGIN"
PLATFORM_APPLICATION = "codex"
+1 -17
View File
@@ -290,11 +290,6 @@ def run(
if args.plugin_data_dir:
os.environ[data_dir_env] = args.plugin_data_dir
# Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key
# writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking
# after them always saw content and every fresh install reported an upgrade.
data_dir_was_empty = telemetry.data_dir_was_empty()
cache_plugin_api_key()
if args.action == "session-start":
clear_stale_api_key_cache()
@@ -310,19 +305,8 @@ def run(
return 0
if args.action == "session-start":
# Claims the marker atomically and says which event to record, so a
# second session starting alongside this one cannot record it too.
first_event = telemetry.claim_install(was_empty=data_dir_was_empty)
if first_event == "install":
if telemetry.is_first_run():
telemetry.record("install")
elif first_event == "upgrade":
# First run after a build that never wrote the marker; the
# predecessor version was never recorded anywhere.
telemetry.record("upgrade", from_version="pre-0.3")
else:
previous = telemetry.claim_version_change()
if previous:
telemetry.record("upgrade", from_version=previous)
recovered = recover_pending_handoffs()
record_session_start(store, hook_input)
if recovered:
+3 -38
View File
@@ -1800,34 +1800,6 @@ def extraction_message_batches(
return batches
# Platform surface attribution. Read from the generated per-host module so a new
# entrypoint is correct without remembering to configure anything.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
def platform_headers(key: str) -> dict[str, str]:
"""Auth plus the three surface-identity headers.
X-Mem0-Source and X-Application are set-once by contract: this is the
outermost layer, so it sets them, and nothing below may overwrite them.
X-Mem0-Client is append-only — anything downstream adds itself to the tail.
"""
headers = {
"Authorization": f"Token {key}",
"Content-Type": "application/json",
"X-Mem0-Source": _PLATFORM_SOURCE,
"X-Mem0-Client": f"mem0-plugin/{PLUGIN_VERSION}",
}
if _PLATFORM_APPLICATION:
headers["X-Application"] = _PLATFORM_APPLICATION
return headers
def _request_json(
url: str, key: str, payload: dict[str, Any], timeout: float
) -> tuple[dict[str, Any] | list[Any], int, int]:
@@ -1835,7 +1807,7 @@ def _request_json(
request = urllib.request.Request(
url,
data=raw,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -1862,7 +1834,7 @@ def _get_json(
) -> tuple[dict[str, Any] | list[Any], int]:
request = urllib.request.Request(
url,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="GET",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -2008,13 +1980,6 @@ def flush_session(
"user_id": write_user,
"app_id": repo.app_id,
"run_id": session_id,
# Top level, not metadata: the backend reads `source` from the body or
# the query string, never from metadata, which is where this used to
# sit. The X-Mem0-Source header is also read, but only from the
# platform release that ships alongside this change, so the body value
# is what makes attribution work on both. The harness tag stays in
# metadata as hook provenance.
"source": _PLATFORM_SOURCE,
"metadata": {**metadata, "author": write_user, "dirs": directory_chain(repo)},
"agent_custom_instructions": PROJECT_MEMORY_INSTRUCTIONS,
"custom_instructions": PERSONAL_MEMORY_INSTRUCTIONS,
@@ -2558,7 +2523,7 @@ def _collect_memory_ids(
def _delete_memory(api_url: str, key: str, memory_id: str) -> bool:
request = urllib.request.Request(
f"{api_url}/v1/memories/{urllib.parse.quote(memory_id)}/",
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="DELETE",
)
try:
+39 -655
View File
@@ -1,9 +1,5 @@
#!/usr/bin/env python3
"""Usage telemetry for Mem0 agent plugins.
Events are linked to your Mem0 account email when an API key is configured, and
to a random per-machine id otherwise. Not anonymous — the Python SDK and CLI
attribute the same way.
"""Anonymous usage telemetry for Mem0 agent plugins.
Hooks run on a 3-6 second budget and fire on every tool call, so recording never
touches the network: `record` appends one JSON line to a local spool and returns.
@@ -13,8 +9,7 @@ started once per session and again from the flush worker that is already detache
Pure stdlib, matching the rest of the plugin. Opt out with MEM0_TELEMETRY=false.
Never sends prompts, memory text, queries, file paths, repository names, or API
keys: only event names, durations, counts, coarse outcomes, and repo/session
identifiers hashed with a random per-install salt.
keys: only event names, durations, counts, coarse outcomes, and salted hashes.
"""
from __future__ import annotations
@@ -34,24 +29,8 @@ from typing import Any
import memory_core
# Seeded from the per-host module the build generates into core/. Two processes
# in this pipeline never call init() — mcp_server.py, and the detached
# `python3 telemetry.py` sender that spawn_flush() starts — so a module default
# was what every one of their events got labelled with.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import HARNESS_ID as _DEFAULT_HARNESS
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
from _harness_id import SOURCE_TAG as _DEFAULT_SOURCE_TAG
except ImportError:
_DEFAULT_HARNESS = "generic"
_DEFAULT_SOURCE_TAG = "MEM0_PLUGIN"
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_harness: str = "generic"
_source_tag: str = "MEM0_PLUGIN"
_PRIVATE_KEYS = {
"apikey",
"authorization",
@@ -77,19 +56,10 @@ _PRIVATE_KEYS = {
}
def init(harness: str = "", source_tag: str = "") -> None:
"""Override the generated identity. Optional — core/_harness_id.py is the default.
The fallback shape matches memory_core.configure_harness's (``<HOST>_PLUGIN``).
It used to be ``MEM0_<HOST>_PLUGIN`` here and ``<host>_plugin`` there, which
meant one plugin could emit three different source values depending on which
process happened to send the batch.
"""
def init(harness: str = "generic", source_tag: str = "") -> None:
global _harness, _source_tag
_harness = harness or _DEFAULT_HARNESS
_source_tag = source_tag or (
f"{_harness.upper().replace('-', '_')}_PLUGIN" if harness else _DEFAULT_SOURCE_TAG
)
_harness = harness
_source_tag = source_tag or f"MEM0_{harness.upper().replace('-', '_')}_PLUGIN"
POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"
POSTHOG_CAPTURE_URL = "https://us.i.posthog.com/i/v0/e/"
@@ -100,16 +70,6 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -123,126 +83,9 @@ def is_enabled() -> bool:
def _digest(value: str, length: int = 16) -> str:
"""Unsalted digest. Only for values that are already secrets (API keys)."""
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
Published atomically, and there is deliberately no derived fallback. Creating
the file with O_CREAT|O_EXCL and then writing into it leaves a window where
the file exists and is empty, and a concurrent hook that reads it in that
window gets nothing. Falling back to a digest of the path would hand that
process a salt an attacker can compute, memoized for its whole run, which is
the privacy control this function exists to provide silently turning itself
off under load. The salt is written to a private temp file first and linked
into place, so the name either does not exist or already has the full value.
Returns "" when it genuinely cannot persist. Callers omit the hash entirely
rather than emit an unsalted one.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
# Read before writing. Hooks are separate processes firing on every tool
# call, so all but the first find the salt already published; going straight
# to create-fsync-link-unlink meant every one of them paid an fsync to
# discover that, on a path whose whole promise is appending a line and
# returning.
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
if _salt_cache:
return _salt_cache
except OSError:
pass
temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp")
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
stream.flush()
os.fsync(stream.fileno())
try:
# Atomic claim: fails if another process already published one.
# os.link rather than replace, which would clobber theirs.
os.link(temporary, path)
except FileExistsError:
pass
except OSError:
# No hardlinks here (some network mounts, some container volumes).
# Claim the name directly instead. That reopens the empty-file
# window, but the window is now benign: a reader that lands in it
# gets "" and omits the hash for that process rather than caching a
# guessable one. Losing the hashes on every run of an entire
# filesystem is the worse failure.
try:
fallback = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(fallback, "w", encoding="utf-8") as stream:
stream.write(temporary.read_text(encoding="utf-8"))
except OSError:
pass
except OSError:
pass
finally:
try:
temporary.unlink()
except OSError:
pass
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
"""Salted digest for values drawn from a guessable space.
repo.identity is a git remote URL, or ``local:<absolute path>`` when there is
no remote — which normally contains the account username. Sixteen unsalted
hex characters over that input space is enumerable, so this is not a
privacy control without the salt. Salting per install keeps every
within-account join the analytics actually use and gives up only
cross-machine joins on the same repository, which nothing computes.
Returns "" when there is no salt, so record() omits the property. An
unsalted digest over this input space is close to plaintext, and emitting one
under a name that implies it is hashed is worse than sending nothing.
"""
if not value:
return ""
salt = _install_salt()
if not salt:
return ""
return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length]
def _safe_value(value: Any) -> Any:
if isinstance(value, str):
return memory_core.redact(value)
@@ -302,176 +145,9 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str:
return created
def _rotate_anonymous_id(identity: dict[str, str]) -> str:
"""Mint a fresh anonymous id because the account context is gone.
The previous id may already have been merged into a person profile by an
$identify, and that merge is permanent. Reusing it after a logout or a key
change attributes everything that follows to the account that just went
away, which is the same misattribution the key fingerprint exists to stop,
only arriving through the anonymous path instead.
`aliased` is cleared with it: the new id has never been merged, so it is
eligible to be aliased into whatever account comes next.
"""
created = f"code-anon-{uuid.uuid4().hex}"
identity["anonymous_id"] = created
identity.pop("aliased", None)
_write_identity(identity)
return created
def _install_state_path() -> Path:
return memory_core.data_dir() / "install-state.json"
def is_first_run() -> bool:
"""Whether install has never been recorded on this machine.
Deliberately NOT the identity file. That file is only written by a
successful flush, so an offline or firewalled user recorded code.install on
every single session, forever — and every 0.2.x user recorded one on their
first 0.3.x session because 0.2.x never wrote it at all.
"""
return not _install_state_path().exists()
def data_dir_was_empty() -> bool:
"""Whether the data directory is untouched. Call BEFORE anything writes to it.
hook_runner reaches claim_install() only after cache_plugin_api_key() has
written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so
asking at claim time always saw content and every fresh install reported an
upgrade. The caller snapshots this at the top of the run instead.
"""
return not _data_dir_has_content()
def claim_install(was_empty: bool | None = None) -> str | None:
"""Claim the one install/upgrade record for this machine, atomically.
Returns the event to record ("install" or "upgrade"), or None if another
session already claimed it. O_CREAT|O_EXCL so two sessions starting together
cannot both win.
`was_empty` must come from data_dir_was_empty() called before this process
wrote anything. Omitting it falls back to checking now, which is only
correct for a caller that has touched nothing.
"""
if not is_enabled():
# Never consume the one-shot claim while the user is opted out, or they
# would silently lose their install event if they later opt in.
return None
path = _install_state_path()
upgrading = not (data_dir_was_empty() if was_empty is None else was_empty)
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
except FileExistsError:
return None
except OSError:
return None
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
json.dump(
{
"plugin_version": memory_core.PLUGIN_VERSION,
"installed_at": memory_core.utc_now(),
"upgraded": upgrading,
},
stream,
)
# Durable before this returns. The O_EXCL open is what makes the
# claim exclusive, so it cannot be replaced by a temp-and-rename
# without losing that, which leaves the content as the thing to make
# safe. A kill between the open and this fsync used to leave a marker
# that exists but parses to nothing: is_first_run reads it as claimed
# and claim_version_change cannot read a version out of it.
stream.flush()
os.fsync(stream.fileno())
except OSError:
pass
return "upgrade" if upgrading else "install"
def _data_dir_has_content() -> bool:
"""Whether anything predates this session in the plugin data directory."""
try:
for entry in memory_core.data_dir().iterdir():
if entry.name != "install-state.json":
return True
except OSError:
pass
return False
def _repair_install_state(path: Path) -> None:
"""Rewrite an unparseable marker so version tracking can resume."""
try:
temporary = path.with_suffix(f".{os.getpid()}.tmp")
temporary.write_text(
json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}),
encoding="utf-8",
)
temporary.replace(path)
except OSError:
pass
def claim_version_change() -> str | None:
"""Return the previously recorded version if it differs, updating the marker.
Only meaningful once the marker exists — the first transition into 0.3.x has
no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting
the marker means the next session sees no change and records nothing.
"""
path = _install_state_path()
try:
state = json.loads(path.read_text(encoding="utf-8"))
except OSError:
return None
except json.JSONDecodeError:
# A crash between O_EXCL and the write leaves an empty marker. Left
# alone it disables every future upgrade event on this machine, because
# claim_install sees the file and this function cannot parse it.
state = None
if not isinstance(state, dict):
_repair_install_state(path)
return None
previous = str(state.get("plugin_version") or "")
if not previous or previous == memory_core.PLUGIN_VERSION:
return None
# Claim the transition with an exclusive sentinel before rewriting the
# marker. A plain read-modify-write let every concurrently starting session
# observe the old version and each record its own upgrade — and the first
# session after a version bump is exactly when several agent windows restart
# together.
sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}")
try:
os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600))
except FileExistsError:
return None
except OSError:
return None
state["plugin_version"] = memory_core.PLUGIN_VERSION
state["upgraded_at"] = memory_core.utc_now()
temporary = path.with_suffix(f".{os.getpid()}.tmp")
try:
temporary.write_text(json.dumps(state), encoding="utf-8")
temporary.replace(path)
except OSError:
# Release the claim. The marker still records the old version, so
# without this the sentinel makes claim_version_change return early on
# every later run and this version's upgrade is never recorded again.
for leftover in (sentinel, temporary):
try:
leftover.unlink()
except OSError:
pass
return None
return previous
"""Whether this machine has never recorded a plugin event before."""
return not _identity_path().exists()
def record(
@@ -492,32 +168,19 @@ def record(
except OSError:
pass
properties = _safe_value(properties)
# Stamped in the RECORDING process, beside harness. `source` used to be
# read in the sending process from a module global, so whichever process
# drained the spool named every event in it. flush() spreads per-event
# properties last, so this now wins over any sender's default.
properties.update(
harness=_harness,
source=_source_tag,
plugin_version=memory_core.PLUGIN_VERSION,
os=sys.platform,
python_version=platform.python_version(),
)
# Assigned only when the digest is real. _scoped_digest returns "" when
# the salt could not be persisted, and an empty property is worse than an
# absent one: it survives the None filter below and reads as a value.
if repo is not None:
repo_hash = _scoped_digest(getattr(repo, "identity", ""))
if repo_hash:
properties["repo_hash"] = repo_hash
properties["repo_hash"] = _digest(getattr(repo, "identity", ""))
if session_id:
session_hash = _scoped_digest(session_id)
if session_hash:
properties["session_hash"] = session_hash
properties["session_hash"] = _digest(session_id)
line = json.dumps(
{
"event": f"{EVENT_PREFIX}.{event}",
"uuid": str(uuid.uuid4()),
"timestamp": memory_core.utc_now(),
"properties": {
key: value for key, value in properties.items() if value is not None
@@ -576,201 +239,38 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / _claim_name()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
for orphan in sorted(directory.glob("telemetry-*.sending")):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
if age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
pass
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
if age < CLAIM_STALE_SECONDS:
continue
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -800,130 +300,34 @@ def _post(payload: dict[str, Any], url: str) -> bool:
def resolve_distinct_id() -> tuple[str, str]:
"""Return the PostHog distinct id and the anonymous id it replaced, if any.
The second value becomes a PostHog $identify alias. It is ONLY ever an
anonymous id: aliasing one account email to another merges two real person
profiles and cannot be undone, so a key that now belongs to a different
account re-resolves with no alias.
"""
"""Return the PostHog distinct id and the anonymous id it replaced, if any."""
identity = _read_identity()
key = memory_core.api_key()
fingerprint = _digest(key) if key else ""
email = identity.get("email", "")
if email and fingerprint:
recorded = identity.get("key_fingerprint", "")
if recorded == fingerprint:
return email, ""
if not recorded:
# Rows written before fingerprints existed. Verify rather than
# adopt: a key changed before the upgrade would otherwise bind the
# new key to the previous account's email, permanently, and the
# fingerprint would then agree with itself forever after.
verified = _resolve_email(key)
if not verified:
# Offline, firewalled, or the API is down. Keep the previous
# behaviour and retry on the next flush rather than dropping a
# real account attribution. Safe because the same network that
# failed /v1/ping/ is about to fail the PostHog POST, so nothing
# is delivered under the unverified identity in the meantime.
return email, ""
identity["email"] = verified
identity["key_fingerprint"] = fingerprint
_write_identity(identity)
return verified, ""
if email:
return email, ""
key = memory_core.api_key()
if not key:
# No key to verify the account with; do not keep attributing to it.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
return anonymous_id(identity), ""
resolved = _resolve_email(key)
if not resolved:
# The key changed and will not resolve (revoked, offline, API down).
# Reaching here with an email means the recorded fingerprint disagreed,
# so the key really did change. Drop the account and rotate: the stored
# anonymous id may already be merged into that account's person, and
# reusing it would keep the events on the profile we are trying to
# leave.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
email = _resolve_email(key)
if not email:
return anonymous_id(identity), ""
# Alias only when going anonymous -> email for the first time. Once an anon
# id has been merged into an account it must never be offered again: an
# alias naming an already-identified id is what could link two real people.
previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "")
if previous:
identity["aliased"] = True
identity["email"] = resolved
identity["key_fingerprint"] = fingerprint
previous = identity.get("anonymous_id", "")
identity["email"] = email
_write_identity(identity)
return resolved, previous
return email, previous
def flush() -> int:
"""Drain the live spool, then any parked claims, and return events sent."""
"""Drain claimed spools to PostHog and return the number of events sent."""
if not is_enabled():
return 0
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
claim = _claim_spool()
if claim is None:
return 0, True
return 0
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# quarantine rather than retry. flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
return 0
events = []
for line in lines:
try:
@@ -933,18 +337,11 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
claim.unlink()
except OSError:
pass
return 0, True
return 0
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -963,17 +360,12 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
# moved into record() have none of their own.
"source": _source_tag,
"language": "python",
"$process_person_profile": False,
@@ -981,24 +373,16 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
**(event.get("properties") or {}),
},
}
for event in chunk
for event in events[start : start + BATCH_SIZE]
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
def main() -> int:
@@ -7,8 +7,8 @@ disable-model-invocation: true
# Pause memory capture
To pause (hooks stop capturing and sending session content; a minimal
telemetry ping still fires at session start, under your Mem0 account email,
unless `MEM0_TELEMETRY=false`):
anonymous telemetry ping still fires at session start unless
`MEM0_TELEMETRY=false`):
```bash
python3 "${PLUGIN_ROOT}/core/memory_cli.py" --harness "codex" --plugin-data-dir "${PLUGIN_DATA}" pause
@@ -1,11 +0,0 @@
"""Generated by integrations/agent-plugin-core/build/build.py. Do not edit."""
HARNESS_ID = "cursor"
SOURCE_TAG = "CURSOR_PLUGIN"
# Platform-side vocabulary (mem0_event.source + X-Application). The whole
# plugin family is one source; which editor it runs in is the application.
# An empty application means the host is unknown, and memory_core omits
# the header entirely rather than sending a placeholder.
PLATFORM_SOURCE = "MEM0_PLUGIN"
PLATFORM_APPLICATION = "cursor"
+1 -17
View File
@@ -290,11 +290,6 @@ def run(
if args.plugin_data_dir:
os.environ[data_dir_env] = args.plugin_data_dir
# Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key
# writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking
# after them always saw content and every fresh install reported an upgrade.
data_dir_was_empty = telemetry.data_dir_was_empty()
cache_plugin_api_key()
if args.action == "session-start":
clear_stale_api_key_cache()
@@ -310,19 +305,8 @@ def run(
return 0
if args.action == "session-start":
# Claims the marker atomically and says which event to record, so a
# second session starting alongside this one cannot record it too.
first_event = telemetry.claim_install(was_empty=data_dir_was_empty)
if first_event == "install":
if telemetry.is_first_run():
telemetry.record("install")
elif first_event == "upgrade":
# First run after a build that never wrote the marker; the
# predecessor version was never recorded anywhere.
telemetry.record("upgrade", from_version="pre-0.3")
else:
previous = telemetry.claim_version_change()
if previous:
telemetry.record("upgrade", from_version=previous)
recovered = recover_pending_handoffs()
record_session_start(store, hook_input)
if recovered:
+3 -38
View File
@@ -1800,34 +1800,6 @@ def extraction_message_batches(
return batches
# Platform surface attribution. Read from the generated per-host module so a new
# entrypoint is correct without remembering to configure anything.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
def platform_headers(key: str) -> dict[str, str]:
"""Auth plus the three surface-identity headers.
X-Mem0-Source and X-Application are set-once by contract: this is the
outermost layer, so it sets them, and nothing below may overwrite them.
X-Mem0-Client is append-only — anything downstream adds itself to the tail.
"""
headers = {
"Authorization": f"Token {key}",
"Content-Type": "application/json",
"X-Mem0-Source": _PLATFORM_SOURCE,
"X-Mem0-Client": f"mem0-plugin/{PLUGIN_VERSION}",
}
if _PLATFORM_APPLICATION:
headers["X-Application"] = _PLATFORM_APPLICATION
return headers
def _request_json(
url: str, key: str, payload: dict[str, Any], timeout: float
) -> tuple[dict[str, Any] | list[Any], int, int]:
@@ -1835,7 +1807,7 @@ def _request_json(
request = urllib.request.Request(
url,
data=raw,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -1862,7 +1834,7 @@ def _get_json(
) -> tuple[dict[str, Any] | list[Any], int]:
request = urllib.request.Request(
url,
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="GET",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
@@ -2008,13 +1980,6 @@ def flush_session(
"user_id": write_user,
"app_id": repo.app_id,
"run_id": session_id,
# Top level, not metadata: the backend reads `source` from the body or
# the query string, never from metadata, which is where this used to
# sit. The X-Mem0-Source header is also read, but only from the
# platform release that ships alongside this change, so the body value
# is what makes attribution work on both. The harness tag stays in
# metadata as hook provenance.
"source": _PLATFORM_SOURCE,
"metadata": {**metadata, "author": write_user, "dirs": directory_chain(repo)},
"agent_custom_instructions": PROJECT_MEMORY_INSTRUCTIONS,
"custom_instructions": PERSONAL_MEMORY_INSTRUCTIONS,
@@ -2558,7 +2523,7 @@ def _collect_memory_ids(
def _delete_memory(api_url: str, key: str, memory_id: str) -> bool:
request = urllib.request.Request(
f"{api_url}/v1/memories/{urllib.parse.quote(memory_id)}/",
headers=platform_headers(key),
headers={"Authorization": f"Token {key}", "Content-Type": "application/json"},
method="DELETE",
)
try:
+39 -655
View File
@@ -1,9 +1,5 @@
#!/usr/bin/env python3
"""Usage telemetry for Mem0 agent plugins.
Events are linked to your Mem0 account email when an API key is configured, and
to a random per-machine id otherwise. Not anonymous — the Python SDK and CLI
attribute the same way.
"""Anonymous usage telemetry for Mem0 agent plugins.
Hooks run on a 3-6 second budget and fire on every tool call, so recording never
touches the network: `record` appends one JSON line to a local spool and returns.
@@ -13,8 +9,7 @@ started once per session and again from the flush worker that is already detache
Pure stdlib, matching the rest of the plugin. Opt out with MEM0_TELEMETRY=false.
Never sends prompts, memory text, queries, file paths, repository names, or API
keys: only event names, durations, counts, coarse outcomes, and repo/session
identifiers hashed with a random per-install salt.
keys: only event names, durations, counts, coarse outcomes, and salted hashes.
"""
from __future__ import annotations
@@ -34,24 +29,8 @@ from typing import Any
import memory_core
# Seeded from the per-host module the build generates into core/. Two processes
# in this pipeline never call init() — mcp_server.py, and the detached
# `python3 telemetry.py` sender that spawn_flush() starts — so a module default
# was what every one of their events got labelled with.
try: # pragma: no cover - absent only in the un-built shared source tree
from _harness_id import HARNESS_ID as _DEFAULT_HARNESS
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
from _harness_id import SOURCE_TAG as _DEFAULT_SOURCE_TAG
except ImportError:
_DEFAULT_HARNESS = "generic"
_DEFAULT_SOURCE_TAG = "MEM0_PLUGIN"
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_harness: str = "generic"
_source_tag: str = "MEM0_PLUGIN"
_PRIVATE_KEYS = {
"apikey",
"authorization",
@@ -77,19 +56,10 @@ _PRIVATE_KEYS = {
}
def init(harness: str = "", source_tag: str = "") -> None:
"""Override the generated identity. Optional — core/_harness_id.py is the default.
The fallback shape matches memory_core.configure_harness's (``<HOST>_PLUGIN``).
It used to be ``MEM0_<HOST>_PLUGIN`` here and ``<host>_plugin`` there, which
meant one plugin could emit three different source values depending on which
process happened to send the batch.
"""
def init(harness: str = "generic", source_tag: str = "") -> None:
global _harness, _source_tag
_harness = harness or _DEFAULT_HARNESS
_source_tag = source_tag or (
f"{_harness.upper().replace('-', '_')}_PLUGIN" if harness else _DEFAULT_SOURCE_TAG
)
_harness = harness
_source_tag = source_tag or f"MEM0_{harness.upper().replace('-', '_')}_PLUGIN"
POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"
POSTHOG_CAPTURE_URL = "https://us.i.posthog.com/i/v0/e/"
@@ -100,16 +70,6 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -123,126 +83,9 @@ def is_enabled() -> bool:
def _digest(value: str, length: int = 16) -> str:
"""Unsalted digest. Only for values that are already secrets (API keys)."""
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
Published atomically, and there is deliberately no derived fallback. Creating
the file with O_CREAT|O_EXCL and then writing into it leaves a window where
the file exists and is empty, and a concurrent hook that reads it in that
window gets nothing. Falling back to a digest of the path would hand that
process a salt an attacker can compute, memoized for its whole run, which is
the privacy control this function exists to provide silently turning itself
off under load. The salt is written to a private temp file first and linked
into place, so the name either does not exist or already has the full value.
Returns "" when it genuinely cannot persist. Callers omit the hash entirely
rather than emit an unsalted one.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
# Read before writing. Hooks are separate processes firing on every tool
# call, so all but the first find the salt already published; going straight
# to create-fsync-link-unlink meant every one of them paid an fsync to
# discover that, on a path whose whole promise is appending a line and
# returning.
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
if _salt_cache:
return _salt_cache
except OSError:
pass
temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp")
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
stream.flush()
os.fsync(stream.fileno())
try:
# Atomic claim: fails if another process already published one.
# os.link rather than replace, which would clobber theirs.
os.link(temporary, path)
except FileExistsError:
pass
except OSError:
# No hardlinks here (some network mounts, some container volumes).
# Claim the name directly instead. That reopens the empty-file
# window, but the window is now benign: a reader that lands in it
# gets "" and omits the hash for that process rather than caching a
# guessable one. Losing the hashes on every run of an entire
# filesystem is the worse failure.
try:
fallback = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(fallback, "w", encoding="utf-8") as stream:
stream.write(temporary.read_text(encoding="utf-8"))
except OSError:
pass
except OSError:
pass
finally:
try:
temporary.unlink()
except OSError:
pass
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
"""Salted digest for values drawn from a guessable space.
repo.identity is a git remote URL, or ``local:<absolute path>`` when there is
no remote — which normally contains the account username. Sixteen unsalted
hex characters over that input space is enumerable, so this is not a
privacy control without the salt. Salting per install keeps every
within-account join the analytics actually use and gives up only
cross-machine joins on the same repository, which nothing computes.
Returns "" when there is no salt, so record() omits the property. An
unsalted digest over this input space is close to plaintext, and emitting one
under a name that implies it is hashed is worse than sending nothing.
"""
if not value:
return ""
salt = _install_salt()
if not salt:
return ""
return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length]
def _safe_value(value: Any) -> Any:
if isinstance(value, str):
return memory_core.redact(value)
@@ -302,176 +145,9 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str:
return created
def _rotate_anonymous_id(identity: dict[str, str]) -> str:
"""Mint a fresh anonymous id because the account context is gone.
The previous id may already have been merged into a person profile by an
$identify, and that merge is permanent. Reusing it after a logout or a key
change attributes everything that follows to the account that just went
away, which is the same misattribution the key fingerprint exists to stop,
only arriving through the anonymous path instead.
`aliased` is cleared with it: the new id has never been merged, so it is
eligible to be aliased into whatever account comes next.
"""
created = f"code-anon-{uuid.uuid4().hex}"
identity["anonymous_id"] = created
identity.pop("aliased", None)
_write_identity(identity)
return created
def _install_state_path() -> Path:
return memory_core.data_dir() / "install-state.json"
def is_first_run() -> bool:
"""Whether install has never been recorded on this machine.
Deliberately NOT the identity file. That file is only written by a
successful flush, so an offline or firewalled user recorded code.install on
every single session, forever — and every 0.2.x user recorded one on their
first 0.3.x session because 0.2.x never wrote it at all.
"""
return not _install_state_path().exists()
def data_dir_was_empty() -> bool:
"""Whether the data directory is untouched. Call BEFORE anything writes to it.
hook_runner reaches claim_install() only after cache_plugin_api_key() has
written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so
asking at claim time always saw content and every fresh install reported an
upgrade. The caller snapshots this at the top of the run instead.
"""
return not _data_dir_has_content()
def claim_install(was_empty: bool | None = None) -> str | None:
"""Claim the one install/upgrade record for this machine, atomically.
Returns the event to record ("install" or "upgrade"), or None if another
session already claimed it. O_CREAT|O_EXCL so two sessions starting together
cannot both win.
`was_empty` must come from data_dir_was_empty() called before this process
wrote anything. Omitting it falls back to checking now, which is only
correct for a caller that has touched nothing.
"""
if not is_enabled():
# Never consume the one-shot claim while the user is opted out, or they
# would silently lose their install event if they later opt in.
return None
path = _install_state_path()
upgrading = not (data_dir_was_empty() if was_empty is None else was_empty)
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
except FileExistsError:
return None
except OSError:
return None
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
json.dump(
{
"plugin_version": memory_core.PLUGIN_VERSION,
"installed_at": memory_core.utc_now(),
"upgraded": upgrading,
},
stream,
)
# Durable before this returns. The O_EXCL open is what makes the
# claim exclusive, so it cannot be replaced by a temp-and-rename
# without losing that, which leaves the content as the thing to make
# safe. A kill between the open and this fsync used to leave a marker
# that exists but parses to nothing: is_first_run reads it as claimed
# and claim_version_change cannot read a version out of it.
stream.flush()
os.fsync(stream.fileno())
except OSError:
pass
return "upgrade" if upgrading else "install"
def _data_dir_has_content() -> bool:
"""Whether anything predates this session in the plugin data directory."""
try:
for entry in memory_core.data_dir().iterdir():
if entry.name != "install-state.json":
return True
except OSError:
pass
return False
def _repair_install_state(path: Path) -> None:
"""Rewrite an unparseable marker so version tracking can resume."""
try:
temporary = path.with_suffix(f".{os.getpid()}.tmp")
temporary.write_text(
json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}),
encoding="utf-8",
)
temporary.replace(path)
except OSError:
pass
def claim_version_change() -> str | None:
"""Return the previously recorded version if it differs, updating the marker.
Only meaningful once the marker exists — the first transition into 0.3.x has
no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting
the marker means the next session sees no change and records nothing.
"""
path = _install_state_path()
try:
state = json.loads(path.read_text(encoding="utf-8"))
except OSError:
return None
except json.JSONDecodeError:
# A crash between O_EXCL and the write leaves an empty marker. Left
# alone it disables every future upgrade event on this machine, because
# claim_install sees the file and this function cannot parse it.
state = None
if not isinstance(state, dict):
_repair_install_state(path)
return None
previous = str(state.get("plugin_version") or "")
if not previous or previous == memory_core.PLUGIN_VERSION:
return None
# Claim the transition with an exclusive sentinel before rewriting the
# marker. A plain read-modify-write let every concurrently starting session
# observe the old version and each record its own upgrade — and the first
# session after a version bump is exactly when several agent windows restart
# together.
sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}")
try:
os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600))
except FileExistsError:
return None
except OSError:
return None
state["plugin_version"] = memory_core.PLUGIN_VERSION
state["upgraded_at"] = memory_core.utc_now()
temporary = path.with_suffix(f".{os.getpid()}.tmp")
try:
temporary.write_text(json.dumps(state), encoding="utf-8")
temporary.replace(path)
except OSError:
# Release the claim. The marker still records the old version, so
# without this the sentinel makes claim_version_change return early on
# every later run and this version's upgrade is never recorded again.
for leftover in (sentinel, temporary):
try:
leftover.unlink()
except OSError:
pass
return None
return previous
"""Whether this machine has never recorded a plugin event before."""
return not _identity_path().exists()
def record(
@@ -492,32 +168,19 @@ def record(
except OSError:
pass
properties = _safe_value(properties)
# Stamped in the RECORDING process, beside harness. `source` used to be
# read in the sending process from a module global, so whichever process
# drained the spool named every event in it. flush() spreads per-event
# properties last, so this now wins over any sender's default.
properties.update(
harness=_harness,
source=_source_tag,
plugin_version=memory_core.PLUGIN_VERSION,
os=sys.platform,
python_version=platform.python_version(),
)
# Assigned only when the digest is real. _scoped_digest returns "" when
# the salt could not be persisted, and an empty property is worse than an
# absent one: it survives the None filter below and reads as a value.
if repo is not None:
repo_hash = _scoped_digest(getattr(repo, "identity", ""))
if repo_hash:
properties["repo_hash"] = repo_hash
properties["repo_hash"] = _digest(getattr(repo, "identity", ""))
if session_id:
session_hash = _scoped_digest(session_id)
if session_hash:
properties["session_hash"] = session_hash
properties["session_hash"] = _digest(session_id)
line = json.dumps(
{
"event": f"{EVENT_PREFIX}.{event}",
"uuid": str(uuid.uuid4()),
"timestamp": memory_core.utc_now(),
"properties": {
key: value for key, value in properties.items() if value is not None
@@ -576,201 +239,38 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / _claim_name()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
for orphan in sorted(directory.glob("telemetry-*.sending")):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
if age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
pass
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
if age < CLAIM_STALE_SECONDS:
continue
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -800,130 +300,34 @@ def _post(payload: dict[str, Any], url: str) -> bool:
def resolve_distinct_id() -> tuple[str, str]:
"""Return the PostHog distinct id and the anonymous id it replaced, if any.
The second value becomes a PostHog $identify alias. It is ONLY ever an
anonymous id: aliasing one account email to another merges two real person
profiles and cannot be undone, so a key that now belongs to a different
account re-resolves with no alias.
"""
"""Return the PostHog distinct id and the anonymous id it replaced, if any."""
identity = _read_identity()
key = memory_core.api_key()
fingerprint = _digest(key) if key else ""
email = identity.get("email", "")
if email and fingerprint:
recorded = identity.get("key_fingerprint", "")
if recorded == fingerprint:
return email, ""
if not recorded:
# Rows written before fingerprints existed. Verify rather than
# adopt: a key changed before the upgrade would otherwise bind the
# new key to the previous account's email, permanently, and the
# fingerprint would then agree with itself forever after.
verified = _resolve_email(key)
if not verified:
# Offline, firewalled, or the API is down. Keep the previous
# behaviour and retry on the next flush rather than dropping a
# real account attribution. Safe because the same network that
# failed /v1/ping/ is about to fail the PostHog POST, so nothing
# is delivered under the unverified identity in the meantime.
return email, ""
identity["email"] = verified
identity["key_fingerprint"] = fingerprint
_write_identity(identity)
return verified, ""
if email:
return email, ""
key = memory_core.api_key()
if not key:
# No key to verify the account with; do not keep attributing to it.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
return anonymous_id(identity), ""
resolved = _resolve_email(key)
if not resolved:
# The key changed and will not resolve (revoked, offline, API down).
# Reaching here with an email means the recorded fingerprint disagreed,
# so the key really did change. Drop the account and rotate: the stored
# anonymous id may already be merged into that account's person, and
# reusing it would keep the events on the profile we are trying to
# leave.
if email:
identity.pop("email", None)
identity.pop("key_fingerprint", None)
return _rotate_anonymous_id(identity), ""
email = _resolve_email(key)
if not email:
return anonymous_id(identity), ""
# Alias only when going anonymous -> email for the first time. Once an anon
# id has been merged into an account it must never be offered again: an
# alias naming an already-identified id is what could link two real people.
previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "")
if previous:
identity["aliased"] = True
identity["email"] = resolved
identity["key_fingerprint"] = fingerprint
previous = identity.get("anonymous_id", "")
identity["email"] = email
_write_identity(identity)
return resolved, previous
return email, previous
def flush() -> int:
"""Drain the live spool, then any parked claims, and return events sent."""
"""Drain claimed spools to PostHog and return the number of events sent."""
if not is_enabled():
return 0
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
claim = _claim_spool()
if claim is None:
return 0, True
return 0
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# quarantine rather than retry. flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
return 0
events = []
for line in lines:
try:
@@ -933,18 +337,11 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
claim.unlink()
except OSError:
pass
return 0, True
return 0
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -963,17 +360,12 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
# moved into record() have none of their own.
"source": _source_tag,
"language": "python",
"$process_person_profile": False,
@@ -981,24 +373,16 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
**(event.get("properties") or {}),
},
}
for event in chunk
for event in events[start : start + BATCH_SIZE]
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
def main() -> int:
@@ -7,8 +7,8 @@ disable-model-invocation: true
# Pause memory capture
To pause (hooks stop capturing and sending session content; a minimal
telemetry ping still fires at session start, under your Mem0 account email,
unless `MEM0_TELEMETRY=false`):
anonymous telemetry ping still fires at session start unless
`MEM0_TELEMETRY=false`):
```bash
python3 "${CURSOR_PLUGIN_ROOT}/core/memory_cli.py" --harness "cursor" pause
+25 -8
View File
@@ -13,7 +13,7 @@ It gives a Harness agent automatic long-term memory plus two explicit memory too
Unlike the local/file-based memory plugins in the ecosystem, Mem0 is a managed backend: server-side extraction, semantic dedup and conflict resolution, and memories that other agents can retrieve when their user and entity filters match.
Current package version: `0.3.0`.
Current package version: `0.3.1`.
Sidekick is available only in the [Claude Code plugin](../claude-code-plugin/README.md#sonnet-sidekick-agent).
@@ -47,21 +47,27 @@ Cordis owns listener and tool cleanup when the plugin unmounts. Every automatic
mkdir -p /tmp/mem0-deepseek-plugin
pnpm pack --pack-destination /tmp/mem0-deepseek-plugin
```
2. Set your Mem0 key:
2. Set your Mem0 key, a stable user identity, and a funded model API key:
```sh
export MEM0_API_KEY=...
export MEM0_USER_ID=your-user-id
export DEEPSEEK_API_KEY=...
```
3. Install it into a disposable Harness profile:
```sh
DSH_HOME=/tmp/mem0-dsh-dev pnpm dlx @deepseek-ai/dsh@0.1.1-rc.2 \
plugin --profile headless add /tmp/mem0-deepseek-plugin/mem0-deepseek-plugin-0.3.0.tgz
plugin --profile web add /tmp/mem0-deepseek-plugin/mem0-deepseek-plugin-0.3.1.tgz
```
4. Copy `cordis.example.yml`, set its installed package path and your `userId`, then run Harness with the same profile:
4. Run Harness with the same profile. The installed bundle activates Mem0 automatically:
```sh
DSH_HOME=/tmp/mem0-dsh-dev pnpm dlx @deepseek-ai/dsh@0.1.1-rc.2 \
web --patch ./integrations/deepseek-plugin/cordis.example.yml
web
```
5. Open http://127.0.0.1:3080 and ask the agent to remember something, then recall it in a later turn.
5. Open http://127.0.0.1:3080, select a workspace, and state a synthetic preference without mentioning Mem0. Wait for asynchronous extraction to finish, then start a fresh session and ask for that preference without tools. Inspect the `mem0:recall` context to verify automatic recall.
Build from the full repository: the source imports the sibling shared core. `pnpm pack` runs the build and includes the activation patch. An installed tarball is self-contained. For custom settings, copy `cordis.example.yml` and pass its absolute path with `--patch`. When upgrading from the old manual setup, remove the old patch that **inserts** a `mem0` row; the bundle now inserts it.
On macOS, set `CHOKIDAR_USEPOLLING=1` if Harness reports `EMFILE`. The Mem0 key and model-provider key are separate credentials.
For a Mem0 Platform on-prem or dedicated deployment, point `config.host` at that base URL (defaults to `api.mem0.ai`). `host` overrides the Platform base URL. It does not support the self-hosted Mem0 OSS API.
@@ -71,6 +77,7 @@ For a Mem0 Platform on-prem or dedicated deployment, point `config.host` at that
|---|---|---|---|
| `apiKey` | no | `$MEM0_API_KEY` | Mem0 platform API key |
| `userId` | yes | | Entity that owns the memories |
| `memoryScope` | no | `user` | `user` shares memory across workspaces; `workspace` isolates all automatic and explicit operations by the session workspace |
| `allowUserOverride` | no | `false` | Permit model-selected access to a different user only in a trusted multi-user deployment |
| `host` | no | `api.mem0.ai` | Platform base URL (on-prem / dedicated) |
| `autoRecall` | no | `true` | Recall relevant memory before model requests |
@@ -80,15 +87,25 @@ For a Mem0 Platform on-prem or dedicated deployment, point `config.host` at that
Automatic capture and recall use the configured `userId` across sessions. Automatic writes do not attach a repository ID or `runId`.
This cross-workspace sharing is intentional compatibility behavior. Set `memoryScope: workspace` in the profile patch to isolate memory. The plugin then adds an `appId` derived from the canonical absolute session workspace path to every write and the matching `app_id` to every search. Symlink aliases share a scope; different directories (including separate clones or worktrees) do not. Moving a workspace changes its scope. Tools cannot override it. Missing or invalid workspace paths skip automatic memory operations and reject explicit tools, without falling back to user-wide access.
Workspace scope does not migrate existing user-only memories. User scope still searches all memories for that user, including workspace-tagged memories; isolation applies when the plugin is configured with workspace scope. These filters are application-level separation, not separate Mem0 credentials.
Both `search_memory` and `add_memory` accept optional `agentId` and `runId`. On search, these narrow the returned memories; on add, they attach those identities to the stored memory. Pass a known `runId` to search memories explicitly saved with that session ID. This does not include automatically captured user-only memories or identify the session making the request.
Per-call `userId` overrides are rejected unless the operator enables `allowUserOverride: true`. Automatic recall and capture always use the configured user.
## Extraction and recall limits
The plugin sends full redacted completed turns to Mem0; Mem0 extracts the memories asynchronously. A queued write is not proof that every extracted fact is already searchable. Even an initial nonempty memory list can be incomplete. Check the stored memories after processing before diagnosing extraction loss.
Automatic recall is a bounded first pass (up to five results, 4,000 context characters, and a two-second wait). `search_memory` uses a focused query and returns up to ten results by default. Exact preference extraction remains model-dependent; include filenames and numeric constraints in release tests, and distinguish standing preferences from one-off requests.
## Telemetry
Writes are tagged `source="DEEPSEEK_HARNESS"`. That value has to exist in the backend's `EventSource` enum for usage to surface by name; until it does, these writes read as `OTHERS`. It is added by [mem0ai/platform#3602](https://github.com/mem0ai/platform/pull/3602), which has to ship before this claim is true.
Writes are tagged `source="DEEPSEEK_HARNESS"` so Mem0's backend can attribute usage to this integration. For it to surface by name (rather than bucketing into `OTHERS`), `DEEPSEEK_HARNESS` must be present in the backend's `KNOWN_EVENT_SOURCES` allowlist, a one-line platform change matching the existing `ZAPIER` / `STRANDS` sources.
The plugin also sends usage events (which tool ran, duration, result counts, coarse failure kind) so Mem0 can tell how the plugin is used and where it breaks. These are **not anonymous**: when an API key is configured they are sent under your Mem0 account email, the same way the SDK attributes its own. Queries, memory text, and entity ids are never sent. Turn it off with `MEM0_TELEMETRY=false`.
The plugin also sends anonymous usage events (which tool ran, duration, result counts, coarse failure kind) so Mem0 can tell how the plugin is used and where it breaks. Queries, memory text, and entity ids are never sent. Turn it off with `MEM0_TELEMETRY=false`.
## Status
+11 -17
View File
@@ -1,17 +1,11 @@
# Example DeepSeek Harness config that loads deepseek-plugin alongside the built-in
# tools plugin. Load with:
#
# DSH_HOME=/tmp/mem0-dsh-dev pnpm dlx @deepseek-ai/dsh@0.1.1-rc.2 \
# web --patch ./integrations/deepseek-plugin/cordis.example.yml
#
- insert:
- id: mem0
# Install the packed plugin into this disposable profile first; loading
# raw dist bypasses Harness's peer dependency installation.
name: "/tmp/mem0-dsh-dev/profiles/headless/node_modules/@mem0/deepseek-plugin/dist/index.js"
config:
# apiKey is read from the MEM0_API_KEY env var when omitted here.
userId: "your-user-id"
# autoRecall: true # optional; enabled by default
# autoCapture: true # optional; enabled by default
# host: "https://your-onprem.mem0.ai" # optional: Platform on-prem / dedicated base URL
# Optional override for the mem0 row installed by dsh.bundle.
# Copy this file, then use: dsh web --patch /absolute/path/to/cordis.yml
- id: mem0
config:
# Replacing a row's config requires restating its userId.
# apiKey is read from MEM0_API_KEY when omitted here.
userId: !!js process.env.MEM0_USER_ID
memoryScope: workspace # optional; default is user (shared across workspaces)
# autoRecall: true
# autoCapture: true
# host: "https://your-onprem.mem0.ai" # Mem0 Platform dedicated deployment
@@ -0,0 +1,6 @@
- insert:
- id: mem0
name: "@mem0/deepseek-plugin"
config:
# Credentials stay in the process environment, not the bundle.
userId: !!js process.env.MEM0_USER_ID
+9 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@mem0/deepseek-plugin",
"version": "0.3.0",
"version": "0.3.1",
"description": "Mem0 long-term memory as a native DeepSeek Harness (Cordis) plugin.",
"type": "module",
"license": "Apache-2.0",
@@ -29,13 +29,21 @@
"publishConfig": {
"access": "public"
},
"dsh": {
"bundle": {
"patch": "./cordis.patch.yml"
}
},
"files": [
"dist",
"cordis.patch.yml",
"cordis.example.yml",
"README.md",
"LICENSE"
],
"scripts": {
"build": "tsup",
"prepack": "pnpm build",
"test": "vitest run",
"test:watch": "vitest",
"typecheck": "tsc --noEmit"

Some files were not shown because too many files have changed in this diff Show More