Compare commits

..

3 Commits

Author SHA1 Message Date
Mgeeeek ed5ee7b9ff fix(plugin): deterministic mem0 user_id resolution
Hooks previously fell back to $USER when MEM0_USER_ID wasn't set,
producing a different user_id on every machine for the same person.
Result: a single account with memories scattered across many user
buckets, none of which can see each other.

Resolution priority (same in bash and python):
  1. MEM0_USER_ID env var (explicit override)
  2. ~/.mem0/identity.json cache (pinned to MEM0_API_KEY fingerprint)
  3. Derived: "mem0-" + sha256(MEM0_API_KEY)[:12]
  4. Fallback: $USER, else "default"

Same MEM0_API_KEY across machines now yields the same user_id without
the user having to set MEM0_USER_ID by hand on every laptop.

Resolver shipped as two tiny files instead of a shared module:
  _identity.sh -- sourced by bash hooks, exports MEM0_RESOLVED_USER_ID
  _identity.py -- imported by on_pre_compact.py, exposes resolve_user_id()

Hook integration:
  on_user_prompt.sh -- sources resolver, USER_ID interpolated into rubric
  on_session_start.sh -- emits an "Active user_id: <X>" header before the
    bootstrap text, so the agent's MCP search_memories/add_memory calls
    use the same bucket the hooks write to (closes the agent-side half
    of the symptom)
  on_pre_compact.py -- replaces inline env lookup with resolve_user_id()

Existing memories under previous $USER values are not auto-migrated.
The cache file regenerates on key rotation (fingerprint mismatch).

CI nudge in pyproject.toml because path filters in ci.yml exclude
plugin-only PRs but build_mem0/build_embedchain are required.

Manual verification: same key on different $USER values resolves to
identical user_id; MEM0_USER_ID override bypasses cache and key
derivation; cache invalidates on key change; bash and python
implementations produce identical output for all four priority levels.
2026-05-08 20:58:08 +05:30
youneshima a623cfaf76 Oss qdrant hosted memories to platform migration (#5080) 2026-05-08 08:04:09 +05:30
Chaithanya Kumar 92491c00c2 docs(memory-decay): use SDK calls in code samples (#5079)
Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-08 01:53:55 +05:30
9 changed files with 2200 additions and 31 deletions
+26 -28
View File
@@ -59,6 +59,14 @@ The toggle lives on the project. You enable decay by patching the project's `dec
The toggle is exposed on the standard project-update endpoint, the same place where `multilingual` and `custom_categories` live.
<CodeGroup>
```python Python
client.project.update(decay=True)
```
```javascript JavaScript
await client.project.update({ decay: true });
```
```bash cURL
curl -X PATCH https://api.mem0.ai/api/v1/orgs/organizations/$ORG_ID/projects/$PROJECT_ID/ \
-H "Authorization: Token $MEM0_API_KEY" \
@@ -66,34 +74,6 @@ curl -X PATCH https://api.mem0.ai/api/v1/orgs/organizations/$ORG_ID/projects/$PR
-d '{"decay": true}'
```
```python Python
import os
import requests
org_id = os.environ["MEM0_ORG_ID"]
project_id = os.environ["MEM0_PROJECT_ID"]
requests.patch(
f"https://api.mem0.ai/api/v1/orgs/organizations/{org_id}/projects/{project_id}/",
headers={"Authorization": f"Token {os.environ['MEM0_API_KEY']}"},
json={"decay": True},
)
```
```javascript Node.js
const res = await fetch(
`https://api.mem0.ai/api/v1/orgs/organizations/${process.env.MEM0_ORG_ID}/projects/${process.env.MEM0_PROJECT_ID}/`,
{
method: "PATCH",
headers: {
Authorization: `Token ${process.env.MEM0_API_KEY}`,
"Content-Type": "application/json",
},
body: JSON.stringify({ decay: true }),
},
);
```
```json Response
{ "message": "Updated decay" }
```
@@ -104,6 +84,16 @@ const res = await fetch(
`decay` is returned on every project read. To fetch only this field, use `?fields=decay`.
<CodeGroup>
```python Python
response = client.project.get(fields=["decay"])
print(response["decay"])
```
```javascript JavaScript
const response = await client.project.get({ fields: ["decay"] });
console.log(response.decay);
```
```bash cURL
curl "https://api.mem0.ai/api/v1/orgs/organizations/$ORG_ID/projects/$PROJECT_ID/?fields=decay" \
-H "Authorization: Token $MEM0_API_KEY"
@@ -119,6 +109,14 @@ curl "https://api.mem0.ai/api/v1/orgs/organizations/$ORG_ID/projects/$PROJECT_ID
The toggle is fully reversible. Setting it to `false` immediately restores the pre-decay ranking; nothing about your stored memories is modified or lost.
<CodeGroup>
```python Python
client.project.update(decay=False)
```
```javascript JavaScript
await client.project.update({ decay: false });
```
```bash cURL
curl -X PATCH https://api.mem0.ai/api/v1/orgs/organizations/$ORG_ID/projects/$PROJECT_ID/ \
-H "Authorization: Token $MEM0_API_KEY" \
+59
View File
@@ -0,0 +1,59 @@
"""Resolve mem0 user_id with deterministic priority.
Resolution priority:
1. MEM0_USER_ID env var (explicit override)
2. ~/.mem0/identity.json cache (pinned to current MEM0_API_KEY fingerprint)
3. Derived: "mem0-" + sha256(MEM0_API_KEY)[:12]
4. Fallback: $USER, else "default"
Same MEM0_API_KEY across machines yields the same user_id, which fixes
the "47 user buckets per account" symptom from running on multiple
laptops with different $USER values.
"""
from __future__ import annotations
import hashlib
import json
import os
from datetime import datetime, timezone
_CACHE_PATH = os.path.expanduser("~/.mem0/identity.json")
def resolve_user_id() -> str:
explicit = os.environ.get("MEM0_USER_ID", "").strip()
if explicit:
return explicit
api_key = os.environ.get("MEM0_API_KEY", "").strip()
if api_key:
digest = hashlib.sha256(api_key.encode("utf-8")).hexdigest()
fingerprint = digest[:8]
try:
with open(_CACHE_PATH, "r") as f:
cached = json.load(f)
if cached.get("api_key_fingerprint") == fingerprint and cached.get("user_id"):
return cached["user_id"]
except (OSError, json.JSONDecodeError):
pass
derived = "mem0-" + digest[:12]
try:
os.makedirs(os.path.dirname(_CACHE_PATH), exist_ok=True)
with open(_CACHE_PATH, "w") as f:
json.dump(
{
"user_id": derived,
"source": "api_key",
"api_key_fingerprint": fingerprint,
"resolved_at": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"),
},
f,
)
except OSError:
pass
return derived
return os.environ.get("USER") or "default"
+57
View File
@@ -0,0 +1,57 @@
# Source this file. Sets MEM0_RESOLVED_USER_ID.
#
# Resolution priority:
# 1. MEM0_USER_ID env var (explicit override)
# 2. ~/.mem0/identity.json cache (pinned to current MEM0_API_KEY fingerprint)
# 3. Derived: "mem0-" + sha256(MEM0_API_KEY)[:12]
# 4. Fallback: $USER, else "default"
#
# Same MEM0_API_KEY across machines yields the same user_id, which fixes
# the "47 user buckets per account" symptom from running on multiple
# laptops with different $USER values.
_mem0_sha256() {
if command -v sha256sum >/dev/null 2>&1; then
sha256sum | cut -d' ' -f1
else
shasum -a 256 | cut -d' ' -f1
fi
}
_mem0_resolve_identity() {
if [ -n "${MEM0_USER_ID:-}" ]; then
printf '%s' "$MEM0_USER_ID"
return
fi
local api_key="${MEM0_API_KEY:-}"
local cache="$HOME/.mem0/identity.json"
if [ -n "$api_key" ]; then
local digest
digest=$(printf '%s' "$api_key" | _mem0_sha256)
local fp="${digest:0:8}"
if [ -f "$cache" ]; then
local cached_fp cached_id
cached_fp=$(jq -r '.api_key_fingerprint // ""' "$cache" 2>/dev/null)
cached_id=$(jq -r '.user_id // ""' "$cache" 2>/dev/null)
if [ "$cached_fp" = "$fp" ] && [ -n "$cached_id" ]; then
printf '%s' "$cached_id"
return
fi
fi
local derived="mem0-${digest:0:12}"
mkdir -p "$HOME/.mem0" 2>/dev/null && \
printf '{"user_id":"%s","source":"api_key","api_key_fingerprint":"%s","resolved_at":"%s"}\n' \
"$derived" "$fp" "$(date -u +%FT%TZ)" > "$cache" 2>/dev/null
printf '%s' "$derived"
return
fi
printf '%s' "${USER:-default}"
}
MEM0_RESOLVED_USER_ID="$(_mem0_resolve_identity)"
export MEM0_RESOLVED_USER_ID
+5 -2
View File
@@ -18,8 +18,11 @@ import json
import logging
import os
import sys
import urllib.request
import urllib.error
import urllib.request
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from _identity import resolve_user_id
log = logging.getLogger("mem0-capture")
log.setLevel(logging.DEBUG)
@@ -207,7 +210,7 @@ def main():
log.debug("No transcript_path provided")
return
user_id = os.environ.get("MEM0_USER_ID", os.environ.get("USER", "default"))
user_id = resolve_user_id()
lines = tail_lines(transcript_path, MAX_TAIL_LINES)
if not lines:
+15
View File
@@ -11,9 +11,24 @@
# even if jq is missing or stdin is malformed.
set -uo pipefail
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)"
# shellcheck source=_identity.sh
. "$SCRIPT_DIR/_identity.sh"
INPUT=$(cat)
SOURCE=$(echo "$INPUT" | jq -r '.source // "startup"' 2>/dev/null || echo "startup")
# Identity line is emitted before every bootstrap variant so the agent
# uses the same user_id the hooks resolved. Without this, the agent's
# search_memories/add_memory MCP calls may bind to a different bucket
# than what the hooks write to.
echo "## Mem0 Identity"
echo ""
echo "Active user_id: \`$MEM0_RESOLVED_USER_ID\`"
echo ""
echo "Always include \`{\"user_id\": \"$MEM0_RESOLVED_USER_ID\"}\` (wrapped in an \`AND\` clause) in every \`search_memories\` filter and as \`user_id\` on every \`add_memory\` call. This keeps memories under one bucket regardless of which machine you're on."
echo ""
if [ "$SOURCE" = "startup" ]; then
cat <<'EOF'
## Mem0 Session Bootstrap
+4 -1
View File
@@ -26,7 +26,10 @@ if [ -z "${MEM0_API_KEY:-}" ]; then
exit 0
fi
USER_ID="${MEM0_USER_ID:-${USER:-default}}"
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)"
# shellcheck source=_identity.sh
. "$SCRIPT_DIR/_identity.sh"
USER_ID="$MEM0_RESOLVED_USER_ID"
cat <<EOF
## Memory check
+1
View File
@@ -154,3 +154,4 @@ known-first-party = ["mem0", "mem0_cli"]
profile = "black"
known_first_party = ["mem0", "mem0_cli"]
# isort scope kept aligned with [tool.ruff.lint.isort] above.
# Plugin-only PRs need a touch here to fire required CI checks (path-filter trap).
+1195
View File
File diff suppressed because it is too large Load Diff
+838
View File
@@ -0,0 +1,838 @@
from __future__ import annotations
import json
import os
import subprocess
import threading
from hashlib import sha256
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from typing import Any
from urllib.parse import parse_qs, urlparse
SCRIPT = Path(__file__).resolve().parents[1] / "scripts" / "oss-to-platform-migrate.sh"
class MigrationHTTPServer:
def __init__(
self,
*,
ping_emails: dict[str, str] | None = None,
verify_api_key: str = "verified-key",
verify_status: int = 200,
qdrant_api_key: str = "qdrant-key",
qdrant_collection: str = "mem0",
qdrant_pages: list[dict[str, Any]] | None = None,
platform_memories: list[dict[str, Any]] | None = None,
) -> None:
self.ping_emails = ping_emails or {}
self.verify_api_key = verify_api_key
self.verify_status = verify_status
self.qdrant_api_key = qdrant_api_key
self.qdrant_collection = qdrant_collection
self.qdrant_pages = qdrant_pages or [{"points": [], "next_page_offset": None}]
self.platform_memories = platform_memories or []
self.requests: list[dict[str, Any]] = []
self._server = ThreadingHTTPServer(("127.0.0.1", 0), self._handler())
self.url = f"http://127.0.0.1:{self._server.server_port}"
self._thread = threading.Thread(target=self._server.serve_forever, daemon=True)
def __enter__(self) -> "MigrationHTTPServer":
self._thread.start()
return self
def __exit__(self, *_exc: object) -> None:
self._server.shutdown()
self._server.server_close()
self._thread.join(timeout=5)
def _handler(self) -> type[BaseHTTPRequestHandler]:
owner = self
class Handler(BaseHTTPRequestHandler):
def log_message(self, _format: str, *_args: object) -> None:
return
def _read_json(self) -> dict[str, Any]:
length = int(self.headers.get("Content-Length", "0"))
raw = self.rfile.read(length) if length else b""
if not raw:
return {}
return json.loads(raw.decode("utf-8"))
def _send_json(self, status: int, payload: dict[str, Any]) -> None:
body = json.dumps(payload).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def _record(self, body: dict[str, Any] | None = None) -> None:
owner.requests.append(
{
"method": self.command,
"path": self.path,
"headers": dict(self.headers),
"body": body or {},
}
)
def do_GET(self) -> None:
self._record()
if self.path == "/v1/ping/":
auth = self.headers.get("Authorization", "")
token = auth.removeprefix("Token ")
email = owner.ping_emails.get(token)
if email:
self._send_json(200, {"user_email": email})
else:
self._send_json(401, {"detail": "Invalid token"})
return
if self.path == f"/collections/{owner.qdrant_collection}":
if self.headers.get("api-key") != owner.qdrant_api_key:
self._send_json(401, {"status": {"error": "unauthorized"}})
else:
self._send_json(200, {"result": {"status": "green"}, "status": "ok"})
return
self._send_json(404, {"detail": "Not found"})
def do_POST(self) -> None:
body = self._read_json()
self._record(body)
parsed = urlparse(self.path)
if self.path == "/posthog":
self._send_json(200, {"ok": True})
return
if self.path == "/api/v1/auth/email_code/":
self._send_json(200, {"sent": True})
return
if self.path == "/api/v1/auth/email_code/verify/":
if owner.verify_status != 200:
self._send_json(owner.verify_status, {"error": "bad verification code"})
else:
self._send_json(200, {"api_key": owner.verify_api_key})
return
if self.path == f"/collections/{owner.qdrant_collection}/points/scroll":
if self.headers.get("api-key") != owner.qdrant_api_key:
self._send_json(401, {"status": {"error": "unauthorized"}})
return
offset = body.get("offset")
page_index = int(offset) if offset is not None else 0
page = owner.qdrant_pages[page_index]
self._send_json(
200,
{
"result": {
"points": page["points"],
"next_page_offset": page.get("next_page_offset"),
},
"status": "ok",
},
)
return
if parsed.path == "/v3/memories/":
query = parse_qs(parsed.query)
page = int(query.get("page", ["1"])[0])
page_size = int(query.get("page_size", ["100"])[0])
filters = body.get("filters") if isinstance(body.get("filters"), dict) else {}
filtered = owner.platform_memories
for key in ("user_id", "agent_id", "run_id"):
if key in filters:
filtered = [memory for memory in filtered if memory.get(key) == filters[key]]
start = (page - 1) * page_size
end = start + page_size
page_results = filtered[start:end]
next_url = (
f"{owner.url}/v3/memories/?page={page + 1}&page_size={page_size}"
if end < len(filtered)
else None
)
self._send_json(
200,
{
"count": len(filtered),
"next": next_url,
"previous": None,
"results": page_results,
},
)
return
if parsed.path == "/v3/memories/add/":
memory_id = f"platform-{len(owner.platform_memories) + 1}"
message = (body.get("messages") or [{}])[0]
memory = {
"id": memory_id,
"memory": message.get("content"),
"metadata": body.get("metadata"),
"user_id": body.get("user_id"),
"agent_id": body.get("agent_id"),
"run_id": body.get("run_id"),
}
owner.platform_memories.append(memory)
self._send_json(
200,
{
"message": "Memories stored successfully",
"status": "SUCCEEDED",
"event_id": "event-1",
"results": [{"id": memory_id, "data": {"memory": message.get("content")}, "event": "ADD"}],
},
)
return
self._send_json(404, {"detail": "Not found"})
return Handler
def alias_marker(anon_id: str, email: str) -> str:
return sha256(f"{anon_id}\0{email}".encode("utf-8")).hexdigest()
def write_config(mem0_dir: Path, data: dict[str, Any]) -> None:
mem0_dir.mkdir(parents=True, exist_ok=True)
(mem0_dir / "config.json").write_text(json.dumps(data), encoding="utf-8")
def read_config(mem0_dir: Path) -> dict[str, Any]:
return json.loads((mem0_dir / "config.json").read_text(encoding="utf-8"))
def run_migration_script(
tmp_path: Path,
server: MigrationHTTPServer,
*args: str,
config: dict[str, Any] | None = None,
raw_config: str | None = None,
) -> tuple[subprocess.CompletedProcess[str], Path]:
mem0_dir = tmp_path / "mem0"
if config is not None:
write_config(mem0_dir, config)
if raw_config is not None:
mem0_dir.mkdir(parents=True, exist_ok=True)
(mem0_dir / "config.json").write_text(raw_config, encoding="utf-8")
env = os.environ.copy()
env.update(
{
"MEM0_DIR": str(mem0_dir),
"MEM0_MIGRATE_TELEMETRY_URL": f"{server.url}/posthog",
}
)
env.pop("MEM0_API_KEY", None)
env.pop("MEM0_BASE_URL", None)
result = subprocess.run(
["bash", str(SCRIPT), "--auth-only", "--base-url", server.url, *args],
capture_output=True,
text=True,
env=env,
start_new_session=True,
timeout=20,
check=False,
)
return result, mem0_dir
def run_export_script(
tmp_path: Path,
server: MigrationHTTPServer,
*args: str,
config: dict[str, Any] | None = None,
qdrant_api_key: str = "qdrant-key",
) -> tuple[subprocess.CompletedProcess[str], Path, Path]:
mem0_dir = tmp_path / "mem0"
output_path = tmp_path / "export.json"
if config is not None:
write_config(mem0_dir, config)
env = os.environ.copy()
env.update(
{
"MEM0_DIR": str(mem0_dir),
"MEM0_MIGRATE_TELEMETRY_URL": f"{server.url}/posthog",
"QDRANT_API_KEY": qdrant_api_key,
}
)
env.pop("MEM0_API_KEY", None)
env.pop("MEM0_BASE_URL", None)
result = subprocess.run(
[
"bash",
str(SCRIPT),
"--export-only",
"--qdrant-url",
server.url,
"--qdrant-collection",
server.qdrant_collection,
"--output",
str(output_path),
*args,
],
capture_output=True,
text=True,
env=env,
start_new_session=True,
timeout=20,
check=False,
)
return result, mem0_dir, output_path
def run_import_script(
tmp_path: Path,
server: MigrationHTTPServer,
input_path: Path,
*args: str,
config: dict[str, Any] | None = None,
api_key: str = "import-key",
) -> tuple[subprocess.CompletedProcess[str], Path]:
mem0_dir = tmp_path / "mem0"
if config is not None:
write_config(mem0_dir, config)
env = os.environ.copy()
env.update(
{
"MEM0_DIR": str(mem0_dir),
"MEM0_MIGRATE_TELEMETRY_URL": f"{server.url}/posthog",
"MEM0_API_KEY": api_key,
}
)
env.pop("MEM0_BASE_URL", None)
result = subprocess.run(
[
"bash",
str(SCRIPT),
"--import-only",
"--base-url",
server.url,
"--input",
str(input_path),
*args,
],
capture_output=True,
text=True,
env=env,
start_new_session=True,
timeout=20,
check=False,
)
return result, mem0_dir
def run_full_script(
tmp_path: Path,
server: MigrationHTTPServer,
*args: str,
config: dict[str, Any] | None = None,
qdrant_api_key: str = "qdrant-key",
) -> tuple[subprocess.CompletedProcess[str], Path, Path]:
mem0_dir = tmp_path / "mem0"
output_path = tmp_path / "full-export.json"
if config is not None:
write_config(mem0_dir, config)
env = os.environ.copy()
env.update(
{
"MEM0_DIR": str(mem0_dir),
"MEM0_MIGRATE_TELEMETRY_URL": f"{server.url}/posthog",
"QDRANT_API_KEY": qdrant_api_key,
}
)
env.pop("MEM0_API_KEY", None)
env.pop("MEM0_BASE_URL", None)
result = subprocess.run(
[
"bash",
str(SCRIPT),
"--base-url",
server.url,
"--qdrant-url",
server.url,
"--qdrant-collection",
server.qdrant_collection,
"--output",
str(output_path),
*args,
],
capture_output=True,
text=True,
env=env,
start_new_session=True,
timeout=20,
check=False,
)
return result, mem0_dir, output_path
def posthog_events(server: MigrationHTTPServer) -> list[dict[str, Any]]:
return [request["body"] for request in server.requests if request["path"] == "/posthog"]
def test_existing_api_key_authenticates_and_stitches_ids(tmp_path: Path) -> None:
config = {
"user_id": "oss-123",
"platform": {"api_key": "stored-key", "base_url": "https://api.mem0.ai"},
"telemetry": {"anonymous_id": "cli-456"},
}
with MigrationHTTPServer(ping_emails={"stored-key": "bob@example.com"}) as server:
result, mem0_dir = run_migration_script(tmp_path, server, "--yes", config=config)
assert result.returncode == 0, result.stderr
assert "Authenticated as bob@example.com" in result.stdout
assert not any(request["path"] == "/api/v1/auth/email_code/verify/" for request in server.requests)
updated = read_config(mem0_dir)
assert updated["platform"] == config["platform"]
assert alias_marker("oss-123", "bob@example.com") in updated["telemetry"]["aliased_pairs"]
assert alias_marker("cli-456", "bob@example.com") in updated["telemetry"]["aliased_pairs"]
events = posthog_events(server)
event_names = [event["event"] for event in events]
assert "oss.migrate.started" in event_names
assert "oss.migrate.authenticated" in event_names
assert event_names.count("$identify") == 2
authenticated = next(event for event in events if event["event"] == "oss.migrate.authenticated")
assert authenticated["distinct_id"] == "bob@example.com"
assert authenticated["properties"]["local_anonymous_id"] == "oss-123"
assert authenticated["properties"]["authenticated_email"] == "bob@example.com"
def test_email_code_authenticates_without_persisting_credentials(tmp_path: Path) -> None:
with MigrationHTTPServer(ping_emails={"verified-key": "alice@example.com"}) as server:
result, mem0_dir = run_migration_script(
tmp_path,
server,
"--email",
"Alice@Example.COM",
"--code",
"123456",
)
assert result.returncode == 0, result.stderr
assert "Authenticated as alice@example.com" in result.stdout
verify_request = next(request for request in server.requests if request["path"] == "/api/v1/auth/email_code/verify/")
assert verify_request["body"] == {"email": "alice@example.com", "code": "123456"}
assert not any(request["path"] == "/api/v1/auth/email_code/" for request in server.requests)
updated = read_config(mem0_dir)
assert "api_key" not in updated.get("platform", {})
assert "user_email" not in updated.get("platform", {})
assert updated["user_id"]
assert alias_marker(updated["user_id"], "alice@example.com") in updated["telemetry"]["aliased_pairs"]
events = posthog_events(server)
assert [event["event"] for event in events].count("$identify") == 1
authenticated = next(event for event in events if event["event"] == "oss.migrate.authenticated")
assert authenticated["properties"]["auth_method"] == "email_code"
def test_invalid_stored_key_falls_back_to_email_code(tmp_path: Path) -> None:
config = {
"user_id": "oss-fallback",
"platform": {"api_key": "bad-key", "base_url": "https://api.mem0.ai"},
}
with MigrationHTTPServer(ping_emails={"verified-key": "new@example.com"}) as server:
result, mem0_dir = run_migration_script(
tmp_path,
server,
"--email",
"new@example.com",
"--code",
"123456",
config=config,
)
assert result.returncode == 0, result.stderr
assert "Stored Mem0 Platform API key is invalid or expired" in result.stdout
assert "Authenticated as new@example.com" in result.stdout
ping_tokens = [
request["headers"]["Authorization"].removeprefix("Token ")
for request in server.requests
if request["path"] == "/v1/ping/"
]
assert ping_tokens == ["bad-key", "verified-key"]
updated = read_config(mem0_dir)
assert updated["platform"] == config["platform"]
assert alias_marker("oss-fallback", "new@example.com") in updated["telemetry"]["aliased_pairs"]
def test_email_code_failure_reports_failed_telemetry(tmp_path: Path) -> None:
with MigrationHTTPServer(verify_status=400) as server:
result, mem0_dir = run_migration_script(
tmp_path,
server,
"--email",
"fail@example.com",
"--code",
"bad",
)
assert result.returncode == 1
assert "Verification failed: bad verification code" in result.stderr
updated = read_config(mem0_dir)
assert "telemetry" not in updated or "aliased_pairs" not in updated["telemetry"]
events = posthog_events(server)
event_names = [event["event"] for event in events]
assert "oss.migrate.started" in event_names
assert "oss.migrate.failed" in event_names
assert "$identify" not in event_names
failed = next(event for event in events if event["event"] == "oss.migrate.failed")
assert "Verification failed" in failed["properties"]["error"]
def test_malformed_config_does_not_crash_and_authenticates(tmp_path: Path) -> None:
with MigrationHTTPServer(ping_emails={"verified-key": "malformed@example.com"}) as server:
result, mem0_dir = run_migration_script(
tmp_path,
server,
"--email",
"malformed@example.com",
"--code",
"123456",
raw_config="{not valid json",
)
assert result.returncode == 0, result.stderr
assert "Authenticated as malformed@example.com" in result.stdout
updated = read_config(mem0_dir)
assert updated["user_id"]
assert alias_marker(updated["user_id"], "malformed@example.com") in updated["telemetry"]["aliased_pairs"]
def test_weird_telemetry_shape_does_not_crash(tmp_path: Path) -> None:
config = {"user_id": "oss-weird-telemetry", "telemetry": "not-an-object"}
with MigrationHTTPServer(ping_emails={"verified-key": "weird@example.com"}) as server:
result, mem0_dir = run_migration_script(
tmp_path,
server,
"--email",
"weird@example.com",
"--code",
"123456",
config=config,
)
assert result.returncode == 0, result.stderr
updated = read_config(mem0_dir)
assert isinstance(updated["telemetry"], dict)
assert alias_marker("oss-weird-telemetry", "weird@example.com") in updated["telemetry"]["aliased_pairs"]
def test_missing_python3_prints_clear_shell_error(tmp_path: Path) -> None:
env = os.environ.copy()
env["PATH"] = str(tmp_path)
result = subprocess.run(
["/bin/bash", str(SCRIPT), "--help"],
capture_output=True,
text=True,
env=env,
timeout=20,
check=False,
)
assert result.returncode == 1
assert "python3 is required to run the Mem0 migration" in result.stderr
def test_curl_piped_help_works() -> None:
result = subprocess.run(
["bash", "-c", f"curl -fsSL file://{SCRIPT} | bash -s -- --help"],
capture_output=True,
text=True,
timeout=20,
check=False,
)
assert result.returncode == 0, result.stderr
assert "Migrate Python OSS hosted-Qdrant memories" in result.stdout
def test_export_qdrant_memories_to_json_without_vectors_or_api_key(tmp_path: Path) -> None:
pages = [
{
"points": [
{
"id": "point-1",
"vector": [0.1, 0.2],
"payload": {
"data": "User likes dark mode",
"hash": "hash-1",
"created_at": "2026-05-01T00:00:00Z",
"updated_at": "2026-05-01T00:00:00Z",
"user_id": "alice",
"agent_id": "agent-1",
"run_id": "run-1",
"actor_id": "actor-1",
"role": "user",
"topic": "preferences",
"text_lemmatized": "user like dark mode",
},
}
],
"next_page_offset": 1,
},
{
"points": [
{
"id": "point-2",
"vector": [0.3, 0.4],
"payload": {
"data": "User prefers concise answers",
"hash": "hash-2",
"user_id": "alice",
"metadata_note": "extra",
},
}
],
"next_page_offset": None,
},
]
with MigrationHTTPServer(qdrant_pages=pages) as server:
result, _mem0_dir, output_path = run_export_script(
tmp_path,
server,
"--user-id",
"alice",
"--qdrant-page-size",
"1",
config={"user_id": "oss-export-user"},
)
assert result.returncode == 0, result.stderr
assert "Exported 2 memories" in result.stdout
artifact = json.loads(output_path.read_text(encoding="utf-8"))
assert artifact["kind"] == "mem0_oss_qdrant_export"
assert artifact["source"]["sdk"] == "python"
assert artifact["source"]["vector_store"] == "qdrant"
assert artifact["source"]["storage"] == "hosted"
assert artifact["source"]["filters"]["user_id"] == "alice"
assert artifact["record_count"] == 2
assert artifact["local_anonymous_id"] == "oss-export-user"
first = artifact["records"][0]
assert first["id"] == "point-1"
assert first["memory"] == "User likes dark mode"
assert first["hash"] == "hash-1"
assert first["user_id"] == "alice"
assert first["agent_id"] == "agent-1"
assert first["run_id"] == "run-1"
assert first["actor_id"] == "actor-1"
assert first["role"] == "user"
assert first["metadata"] == {"topic": "preferences"}
assert all("vector" not in record for record in artifact["records"])
assert "qdrant-key" not in output_path.read_text(encoding="utf-8")
scroll_requests = [request for request in server.requests if request["path"].endswith("/points/scroll")]
assert len(scroll_requests) == 2
assert scroll_requests[0]["body"]["with_vector"] is False
assert scroll_requests[0]["body"]["filter"] == {"must": [{"key": "user_id", "match": {"value": "alice"}}]}
assert scroll_requests[1]["body"]["offset"] == 1
def test_export_requires_scope_or_all(tmp_path: Path) -> None:
with MigrationHTTPServer() as server:
result, _mem0_dir, output_path = run_export_script(tmp_path, server)
assert result.returncode == 1
assert "Export requires --user-id, --agent-id, --run-id, or --all" in result.stderr
assert not output_path.exists()
assert not any(request["path"].endswith("/points/scroll") for request in server.requests)
def test_export_all_uses_no_qdrant_filter(tmp_path: Path) -> None:
with MigrationHTTPServer(qdrant_pages=[{"points": [], "next_page_offset": None}]) as server:
result, _mem0_dir, output_path = run_export_script(tmp_path, server, "--all")
assert result.returncode == 0, result.stderr
artifact = json.loads(output_path.read_text(encoding="utf-8"))
assert artifact["record_count"] == 0
assert artifact["records"] == []
scroll_request = next(request for request in server.requests if request["path"].endswith("/points/scroll"))
assert "filter" not in scroll_request["body"]
def test_export_invalid_qdrant_credentials_fail_clearly(tmp_path: Path) -> None:
with MigrationHTTPServer(qdrant_api_key="correct-key") as server:
result, _mem0_dir, output_path = run_export_script(
tmp_path,
server,
"--user-id",
"alice",
qdrant_api_key="wrong-key",
)
assert result.returncode == 1
assert "Qdrant authentication failed" in result.stderr
assert not output_path.exists()
def test_import_platform_memories_from_export_json(tmp_path: Path) -> None:
input_path = tmp_path / "import.json"
input_path.write_text(
json.dumps(
{
"source": {"sdk": "python", "vector_store": "qdrant", "collection": "mem0_test"},
"records": [
{
"id": "local-1",
"memory": "User likes barbecue",
"hash": "hash-1",
"created_at": "2026-05-07T00:00:00Z",
"user_id": "alice",
"metadata": {"topic": "food"},
}
],
}
),
encoding="utf-8",
)
with MigrationHTTPServer(ping_emails={"import-key": "alice@example.com"}) as server:
result, _mem0_dir = run_import_script(tmp_path, server, input_path)
assert result.returncode == 0, result.stderr
assert "Imported: 1" in result.stdout
assert "Failed: 0" in result.stdout
add_request = next(request for request in server.requests if request["path"] == "/v3/memories/add/")
body = add_request["body"]
assert body["messages"] == [{"role": "user", "content": "User likes barbecue"}]
assert body["user_id"] == "alice"
assert body["infer"] is False
assert body["source"] == "migration"
assert body["timestamp"] == 1778112000
assert body["metadata"]["topic"] == "food"
assert body["metadata"]["mem0_migration_source"] == "python_oss_qdrant"
assert body["metadata"]["mem0_migration_collection"] == "mem0_test"
assert body["metadata"]["mem0_migration_local_id"] == "local-1"
assert body["metadata"]["mem0_migration_local_hash"] == "hash-1"
events = posthog_events(server)
assert "oss.migrate.completed" in [event["event"] for event in events]
def test_import_skips_existing_identical_memory(tmp_path: Path) -> None:
input_path = tmp_path / "import.json"
source = {"sdk": "python", "vector_store": "qdrant", "collection": "mem0_test"}
record = {"id": "local-1", "memory": "User likes barbecue", "hash": "hash-1", "user_id": "alice"}
input_path.write_text(json.dumps({"source": source, "records": [record]}), encoding="utf-8")
import_key = sha256("python:qdrant:mem0_test:local-1".encode("utf-8")).hexdigest()
existing = [
{
"id": "platform-1",
"memory": "User likes barbecue",
"user_id": "alice",
"metadata": {
"mem0_migration_import_key": import_key,
"mem0_migration_local_hash": "hash-1",
},
}
]
with MigrationHTTPServer(ping_emails={"import-key": "alice@example.com"}, platform_memories=existing) as server:
result, _mem0_dir = run_import_script(tmp_path, server, input_path)
assert result.returncode == 0, result.stderr
assert "Imported: 0" in result.stdout
assert "Skipped existing identical: 1" in result.stdout
assert "Changed existing: 0" in result.stdout
assert not any(request["path"] == "/v3/memories/add/" for request in server.requests)
def test_import_reports_changed_existing_without_update_or_add(tmp_path: Path) -> None:
input_path = tmp_path / "import.json"
source = {"sdk": "python", "vector_store": "qdrant", "collection": "mem0_test"}
record = {"id": "local-1", "memory": "User likes brisket", "hash": "hash-new", "user_id": "alice"}
input_path.write_text(json.dumps({"source": source, "records": [record]}), encoding="utf-8")
import_key = sha256("python:qdrant:mem0_test:local-1".encode("utf-8")).hexdigest()
existing = [
{
"id": "platform-1",
"memory": "User likes barbecue",
"user_id": "alice",
"metadata": {
"mem0_migration_import_key": import_key,
"mem0_migration_local_hash": "hash-old",
},
}
]
with MigrationHTTPServer(ping_emails={"import-key": "alice@example.com"}, platform_memories=existing) as server:
result, _mem0_dir = run_import_script(tmp_path, server, input_path)
assert result.returncode == 0, result.stderr
assert "Imported: 0" in result.stdout
assert "Skipped existing identical: 0" in result.stdout
assert "Changed existing: 1" in result.stdout
assert not any(request["path"] == "/v3/memories/add/" for request in server.requests)
review_path_line = next(line for line in result.stdout.splitlines() if line.startswith("Review file: "))
review_path = Path(review_path_line.removeprefix("Review file: "))
review = json.loads(review_path.read_text(encoding="utf-8"))
assert review["records"][0]["status"] == "changed_existing"
assert review["records"][0]["platform_memory_id"] == "platform-1"
def test_full_flow_auth_export_and_imports_memories(tmp_path: Path) -> None:
qdrant_pages = [
{
"points": [
{
"id": "point-1",
"payload": {
"data": "User likes barbecue",
"hash": "hash-1",
"created_at": "2026-05-07T00:00:00Z",
"user_id": "alice",
},
}
],
"next_page_offset": None,
}
]
with MigrationHTTPServer(ping_emails={"verified-key": "alice@example.com"}, qdrant_pages=qdrant_pages) as server:
result, _mem0_dir, output_path = run_full_script(
tmp_path,
server,
"--email",
"alice@example.com",
"--code",
"123456",
"--user-id",
"alice",
)
assert result.returncode == 0, result.stderr
assert "Phase 1/3: Authenticate with Mem0 Platform" in result.stdout
assert "Phase 2/3: Export Python OSS memories from hosted Qdrant" in result.stdout
assert "Phase 3/3: Import memories into Mem0 Platform" in result.stdout
assert "Imported: 1" in result.stdout
assert output_path.exists()
assert any(request["path"] == "/v3/memories/add/" for request in server.requests)
events = posthog_events(server)
event_names = [event["event"] for event in events]
assert "oss.migrate.authenticated" in event_names
assert "oss.migrate.completed" in event_names