diff --git a/.github/workflows/opencode-plugin-checks.yml b/.github/workflows/opencode-plugin-checks.yml index 4fb95fc6c..a93d741eb 100644 --- a/.github/workflows/opencode-plugin-checks.yml +++ b/.github/workflows/opencode-plugin-checks.yml @@ -32,6 +32,9 @@ jobs: - name: Type check run: bun run type-check + - name: Test + run: bun test + - name: Build run: bun run build diff --git a/cli/node/src/backend/platform.ts b/cli/node/src/backend/platform.ts index 705a3d470..65a4b6228 100644 --- a/cli/node/src/backend/platform.ts +++ b/cli/node/src/backend/platform.ts @@ -31,7 +31,8 @@ export class PlatformBackend implements Backend { this.headers = { Authorization: `Token ${config.apiKey}`, "Content-Type": "application/json", - "X-Mem0-Source": "cli", + "X-Mem0-Source": "CLI", + "X-Mem0-Client": `mem0-cli-node/${CLI_VERSION}`, "X-Mem0-Client-Language": "node", "X-Mem0-Client-Version": CLI_VERSION, }; diff --git a/cli/python/src/mem0_cli/backend/platform.py b/cli/python/src/mem0_cli/backend/platform.py index 1d153da9c..a9ce60df1 100644 --- a/cli/python/src/mem0_cli/backend/platform.py +++ b/cli/python/src/mem0_cli/backend/platform.py @@ -27,7 +27,8 @@ class PlatformBackend(Backend): headers={ "Authorization": f"Token {config.api_key}", "Content-Type": "application/json", - "X-Mem0-Source": "cli", + "X-Mem0-Source": "CLI", + "X-Mem0-Client": f"mem0-cli-python/{__version__}", "X-Mem0-Client-Language": "python", "X-Mem0-Client-Version": __version__, }, diff --git a/integrations/AGENTS.md b/integrations/AGENTS.md index c129dd48b..221c8206f 100644 --- a/integrations/AGENTS.md +++ b/integrations/AGENTS.md @@ -49,6 +49,48 @@ 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/-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/`. @@ -59,3 +101,4 @@ Run the type check after every TypeScript change: `pnpm run typecheck` or `tsc - 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`. diff --git a/integrations/agent-plugin-core/build/build.py b/integrations/agent-plugin-core/build/build.py index e675efb02..3059b9b1f 100644 --- a/integrations/agent-plugin-core/build/build.py +++ b/integrations/agent-plugin-core/build/build.py @@ -81,15 +81,23 @@ def replace_output(staged: Path, output: Path) -> Path: return output -def _render_harness_id(host: str) -> str: +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" @@ -98,8 +106,10 @@ def _render_harness_id(host: str) -> str: "\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 = "{host}"\n' + f'PLATFORM_APPLICATION = "{application}"\n' ) @@ -122,7 +132,7 @@ def _bundle_python( # 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), encoding="utf-8") + (core / "_harness_id.py").write_text(_render_harness_id(host, portable=portable), encoding="utf-8") values = { "PLUGIN_ROOT": plugin_root, diff --git a/integrations/agent-plugin-core/python/memory_core.py b/integrations/agent-plugin-core/python/memory_core.py index 77d7a2b80..d350e76d2 100644 --- a/integrations/agent-plugin-core/python/memory_core.py +++ b/integrations/agent-plugin-core/python/memory_core.py @@ -1800,6 +1800,34 @@ 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]: @@ -1807,7 +1835,7 @@ def _request_json( request = urllib.request.Request( url, data=raw, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="POST", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1834,7 +1862,7 @@ def _get_json( ) -> tuple[dict[str, Any] | list[Any], int]: request = urllib.request.Request( url, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="GET", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1980,6 +2008,13 @@ 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, @@ -2523,7 +2558,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={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="DELETE", ) try: diff --git a/integrations/agent-plugin-core/tests/test_build.py b/integrations/agent-plugin-core/tests/test_build.py index 54b5e5c6a..23295cc38 100644 --- a/integrations/agent-plugin-core/tests/test_build.py +++ b/integrations/agent-plugin-core/tests/test_build.py @@ -70,6 +70,38 @@ 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) diff --git a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py index 7d48dbcc3..74a31fa5c 100644 --- a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py +++ b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py @@ -179,6 +179,33 @@ def test_source_tag_defaults_agree_between_the_two_modules(): 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( diff --git a/integrations/agent-plugin-core/typescript/src/telemetry.ts b/integrations/agent-plugin-core/typescript/src/telemetry.ts index c6a7544ce..2b6c5286e 100644 --- a/integrations/agent-plugin-core/typescript/src/telemetry.ts +++ b/integrations/agent-plugin-core/typescript/src/telemetry.ts @@ -1,3 +1,5 @@ +import { randomUUID } from "node:crypto"; + import { redactSecrets } from "./lifecycle.ts"; const POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"; @@ -74,34 +76,108 @@ export function errorKind(error: unknown): string { return error instanceof Error ? error.constructor.name : "other"; } +// Delivery is retried in memory, not spooled to disk, and that is a decision +// rather than an omission. The Python core spools because its hooks are separate +// processes that fire per tool call and exit immediately, so nothing survives +// without a file. These plugins are loaded into a host that lives for a whole +// session, so re-queueing covers the same transient failures without the claim +// and lease machinery a correct cross-process spool needs. What that leaves +// uncovered is narrow: a session that both starts and ends with no connectivity. +const RETRY_BACKOFF_CEILING_MS = 60_000; +// Consecutive failed flushes before the queue is dropped. Deliberately NOT the +// same thing as Python's budget, which rides in the claim filename and so +// follows one batch: this counter lives in the closure and counts the outage, +// not the payload. Events captured between attempts join the same queue and go +// with it. Per-batch accounting would need an attempt count on every event, and +// the queue is already bounded, so the simpler rule is the one in force here. +// Without any bound a payload the server will never accept is retried for the +// whole session and, now that the backlog is preferred over new events, holds +// the queue against everything behind it. +const MAX_DELIVERY_ATTEMPTS = 5; + export function createTelemetry(config: TelemetryConfig) { let queue: Record[] = []; let timer: ReturnType | undefined; + let consecutiveFailures = 0; + let retryNotBefore = 0; + let exitFlushAttempted = false; + let flushing = false; const flushThreshold = config.flushThreshold ?? 10; const maxQueueSize = config.maxQueueSize ?? 100; const deliver = config.delivery ?? (async (batch: Record[]) => { - await fetch(POSTHOG_BATCH_URL, { + const response = await fetch(POSTHOG_BATCH_URL, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ api_key: POSTHOG_API_KEY, batch }), signal: AbortSignal.timeout(3_000), }); + // fetch only rejects on a network-level failure. Without this check a 500, + // a 503 or a 429 resolved normally and the batch was counted as delivered + // and dropped, which is the likelier outage than a refused connection. + // Any non-2xx is retried, matching the Python core: the backoff and the + // queue bound contain a payload that will never be accepted, because the + // re-queued batch sits at the front and is the first thing evicted. + if (!response.ok) throw new Error(`posthog responded ${response.status}`); }); - async function flush(): Promise { + async function flush(force = false): Promise { + // One at a time. Two overlapping flushes each detach the queue and each + // prepend their own batch back on failure, so the later batch lands in front + // of the earlier one and the truncation then drops the OLDER events first, + // inverting the priority the failure path exists to establish. A second + // caller returns immediately; the queue waits for the next flush. + if (flushing) return; if (!queue.length) return; + // `force` skips the cooldown. beforeExit is the last chance this process + // gets, and gating it on the same backoff meant that after any failure the + // exit flush did nothing and the queue died with the process, which is the + // loss this whole mechanism exists to prevent. + if (!force && Date.now() < retryNotBefore) return; const batch = queue; queue = []; + flushing = true; try { await deliver(batch); + consecutiveFailures = 0; + retryNotBefore = 0; } catch { - // Telemetry must never affect plugin behavior. + // Put it back. Detaching the batch and swallowing the error deleted the + // events outright, so any blip silently dropped telemetry with nothing + // recording that it had happened. Every event carries a uuid, so a retry + // that duplicates one PostHog already accepted is collapsed there. + // + consecutiveFailures += 1; + if (consecutiveFailures >= MAX_DELIVERY_ATTEMPTS) { + // Give up on the queue so a failing outage cannot hold it for the + // session. This drops whatever is queued now, which includes events + // captured during the outage, not only the batch that kept failing. + consecutiveFailures = 0; + retryNotBefore = 0; + return; + } + // Keep the FRONT on overflow, so the batch being retried survives and a + // new event is what gets dropped. Matches the Python core, where record() + // refuses new events once the spool is full rather than evicting the + // backlog. Keeping the newest would throw away exactly the events this + // retry exists to save. + queue = [...batch, ...queue].slice(0, maxQueueSize); + retryNotBefore = Date.now() + Math.min(2 ** consecutiveFailures * 1_000, RETRY_BACKOFF_CEILING_MS); + } finally { + flushing = false; } } function beforeExit(): void { - void flush(); + // Once, and only once. Node re-emits beforeExit whenever the handler + // schedules more async work, so an unconditional forced flush looped until + // the attempt budget was spent: five attempts against a 3s delivery timeout + // is fifteen seconds added to the shutdown of whatever editor or CLI is + // hosting this. The backoff used to end that loop after one attempt, and + // removing it for the forced path removed the only thing bounding it. + if (exitFlushAttempted) return; + exitFlushAttempted = true; + void flush(true); } function build(event: string, properties: Record = {}): Record | null { @@ -112,6 +188,15 @@ export function createTelemetry(config: TelemetryConfig) { return { event: config.eventName?.(event) ?? event, distinct_id: distinctId, + // Stamped once, at capture. This is what makes retrying safe: a batch + // re-sent after a failure carries the same ids, so PostHog collapses + // anything it already accepted instead of counting it twice. + uuid: randomUUID(), + // Capture time, not ingestion time. Events now sit through backoff and + // across a whole outage, so without this PostHog records them whenever + // delivery happened to succeed. It also matters for the uuid dedupe + // above, whose key includes the event date. + timestamp: new Date().toISOString(), properties: { ...safeProperties(properties), ...safeProperties(config.commonProperties ?? {}), @@ -134,8 +219,10 @@ export function createTelemetry(config: TelemetryConfig) { try { const payload = build(event, properties); if (!payload) return; + // Full means drop this event, not evict the backlog. Same rule as the + // failure path above and as Python's record(). + if (queue.length >= maxQueueSize) return; queue.push(payload); - if (queue.length > maxQueueSize) queue = queue.slice(-maxQueueSize); if (!timer) { timer = setInterval(() => void flush(), config.flushIntervalMs ?? 5_000); timer.unref?.(); @@ -149,6 +236,8 @@ export function createTelemetry(config: TelemetryConfig) { function resetForTesting(): void { queue = []; + consecutiveFailures = 0; + retryNotBefore = 0; if (timer) clearInterval(timer); timer = undefined; process.off("beforeExit", beforeExit); diff --git a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts index e04794614..cf9671682 100644 --- a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts +++ b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts @@ -122,3 +122,250 @@ test("error classification does not expose messages", () => { assert.equal(errorKind(new Error("request timeout")), "timeout"); assert.equal(errorKind(new Error("fetch failed")), "network"); }); + +test("a failed delivery keeps the batch instead of deleting it", async () => { + // The defect: the queue was detached before the await and the error swallowed, + // so one blip destroyed the events with nothing recording that it happened. + const attempts: Record[][] = []; + let failNext = true; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, + delivery: async (batch) => { + attempts.push(batch); + if (failNext) throw new Error("network down"); + }, + }); + + telemetry.capture("one"); + telemetry.capture("two"); + await telemetry.flush(); + + assert.equal(attempts.length, 1); + assert.equal(telemetry.queueForTesting().length, 2, "events were dropped on failure"); + + failNext = false; + // Backoff is in force, so wait it out the way wall time would. + await new Promise((resolve) => setTimeout(resolve, 2_100)); + await telemetry.flush(); + + assert.equal(attempts.length, 2, "never retried"); + assert.equal(telemetry.queueForTesting().length, 0); + telemetry.resetForTesting(); +}); + +test("a retried event carries the same uuid so PostHog can collapse it", async () => { + const attempts: Record[][] = []; + let failNext = true; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, + delivery: async (batch) => { + attempts.push(batch); + if (failNext) throw new Error("network down"); + }, + }); + + telemetry.capture("once"); + await telemetry.flush(); + failNext = false; + await new Promise((resolve) => setTimeout(resolve, 2_100)); + await telemetry.flush(); + + assert.equal(attempts.length, 2); + const first = attempts[0][0].uuid; + assert.ok(first, "events carry no uuid, so a retry would double count"); + assert.equal(attempts[1][0].uuid, first, "retry minted a new uuid"); + telemetry.resetForTesting(); +}); + +test("repeated failures back off instead of retrying every flush", async () => { + let calls = 0; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, + delivery: async () => { calls += 1; throw new Error("blocked"); }, + }); + + telemetry.capture("one"); + await telemetry.flush(); + await telemetry.flush(); + await telemetry.flush(); + + assert.equal(calls, 1, "a blocked host was hammered on every flush"); + assert.equal(telemetry.queueForTesting().length, 1, "the event was lost while backing off"); + telemetry.resetForTesting(); +}); + +test("a full queue drops the new event and keeps the batch being retried", async () => { + // Python's record() refuses new events once the spool is full rather than + // evicting the backlog. Keeping the newest here would throw away exactly the + // events the retry exists to save. + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, maxQueueSize: 3, + delivery: async () => { throw new Error("down"); }, + }); + + // Fill past the cap BEFORE the flush, so the re-queue actually has to truncate. + // Capturing only two left the queue empty at re-queue time and the slice on the + // failure path never ran, which is the half that decides the direction. + telemetry.capture("a"); + telemetry.capture("b"); + telemetry.capture("c"); + await telemetry.flush(); + telemetry.capture("d"); + telemetry.capture("e"); + + const events = telemetry.queueForTesting().map((e) => (e as any).event); + assert.equal(events.length, 3, "queue grew past maxQueueSize"); + assert.deepEqual(events, ["a", "b", "c"], "the retried batch was evicted instead of the new events"); + telemetry.resetForTesting(); +}); + +test("the exit-time flush ignores the backoff", async () => { + // beforeExit is the last chance the process gets. Gating it on the same + // cooldown meant that after any failure it did nothing and the queue died. + let attempts = 0; + let failing = true; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async () => { attempts += 1; if (failing) throw new Error("down"); }, + }); + + telemetry.capture("a"); + await telemetry.flush(); + assert.equal(attempts, 1); + + failing = false; + await telemetry.flush(); + assert.equal(attempts, 1, "the backoff should still hold for an ordinary flush"); + + await telemetry.flush(true); + assert.equal(attempts, 2, "the exit flush was suppressed by the backoff"); + assert.equal(telemetry.queueForTesting().length, 0); + telemetry.resetForTesting(); +}); + +test("every event carries a capture-time timestamp", async () => { + const sent: Record[][] = []; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async (batch) => { sent.push(batch); }, + }); + + telemetry.capture("a"); + const capturedAt = Date.now(); + await new Promise((resolve) => setTimeout(resolve, 50)); + await telemetry.flush(); + + const stamped = sent[0][0].timestamp as string; + assert.ok(stamped, "no timestamp, so PostHog would record delivery time"); + assert.ok(Math.abs(Date.parse(stamped) - capturedAt) < 1_000, "not capture time"); + telemetry.resetForTesting(); +}); + +test("a batch the server will never accept is eventually given up on", async () => { + let attempts = 0; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async () => { attempts += 1; throw new Error("permanently bad"); }, + }); + + telemetry.capture("doomed"); + for (let i = 0; i < 8; i += 1) await telemetry.flush(true); + + assert.ok(attempts <= 6, `retried ${attempts} times with no cap`); + assert.equal(telemetry.queueForTesting().length, 0, "a doomed batch held the queue forever"); + telemetry.resetForTesting(); +}); + +test("an HTTP error response is a failure, not a delivery", async () => { + // fetch only rejects on a network-level failure, so a 500 used to resolve + // normally and the batch was dropped as delivered. Exercises the real default + // delivery path rather than an injected one, which is where this hid. + const realFetch = globalThis.fetch; + let calls = 0; + globalThis.fetch = (async () => { + calls += 1; + return new Response("upstream is unwell", { status: 503 }); + }) as typeof fetch; + + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + }); + try { + telemetry.capture("during.outage"); + await telemetry.flush(); + + assert.equal(calls, 1, "never reached the network"); + assert.equal(telemetry.queueForTesting().length, 1, "a 503 was counted as delivered"); + } finally { + globalThis.fetch = realFetch; + telemetry.resetForTesting(); + } +}); + +test("a 2xx is a delivery", async () => { + const realFetch = globalThis.fetch; + globalThis.fetch = (async () => new Response("ok", { status: 200 })) as typeof fetch; + + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + }); + try { + telemetry.capture("fine"); + await telemetry.flush(); + assert.equal(telemetry.queueForTesting().length, 0, "a good response did not clear the queue"); + } finally { + globalThis.fetch = realFetch; + telemetry.resetForTesting(); + } +}); + +test("the exit flush is attempted once, not until the budget is spent", async () => { + // Node re-emits beforeExit whenever the handler schedules async work, so an + // unconditional forced flush looped until MAX_DELIVERY_ATTEMPTS. Against the + // real 3s delivery timeout that is fifteen seconds added to a host's shutdown. + let attempts = 0; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async () => { attempts += 1; throw new Error("down"); }, + }); + + telemetry.capture("a"); + const handlers = process.listeners("beforeExit"); + const ours = handlers[handlers.length - 1] as () => void; + ours(); + ours(); + ours(); + await new Promise((resolve) => setTimeout(resolve, 20)); + + assert.equal(attempts, 1, `exit flush ran ${attempts} times`); + telemetry.resetForTesting(); +}); + +test("overlapping flushes do not reorder the backlog behind newer events", async () => { + // Each flush detaches the queue and prepends its own batch back on failure, so + // two in flight at once put the LATER batch in front of the earlier one. The + // truncation then drops the older events first, inverting the priority the + // failure path exists to establish. + let release: (() => void)[] = []; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: () => new Promise((_resolve, reject) => { release.push(() => reject(new Error("down"))); }), + }); + + telemetry.capture("first"); + const a = telemetry.flush(); + telemetry.capture("second"); + const b = telemetry.flush(); + + release.forEach((fn) => fn()); + await Promise.all([a, b]); + + const events = telemetry.queueForTesting().map((e) => (e as any).event); + assert.equal(release.length, 1, "a second delivery started while one was in flight"); + assert.deepEqual(events, ["first", "second"], `backlog reordered: ${events.join(",")}`); + telemetry.resetForTesting(); +}); diff --git a/integrations/antigravity-plugin/core/_harness_id.py b/integrations/antigravity-plugin/core/_harness_id.py index 2c8a515f4..18b7f68ec 100644 --- a/integrations/antigravity-plugin/core/_harness_id.py +++ b/integrations/antigravity-plugin/core/_harness_id.py @@ -5,5 +5,7 @@ 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" diff --git a/integrations/antigravity-plugin/core/memory_core.py b/integrations/antigravity-plugin/core/memory_core.py index 77d7a2b80..d350e76d2 100644 --- a/integrations/antigravity-plugin/core/memory_core.py +++ b/integrations/antigravity-plugin/core/memory_core.py @@ -1800,6 +1800,34 @@ 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]: @@ -1807,7 +1835,7 @@ def _request_json( request = urllib.request.Request( url, data=raw, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="POST", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1834,7 +1862,7 @@ def _get_json( ) -> tuple[dict[str, Any] | list[Any], int]: request = urllib.request.Request( url, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="GET", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1980,6 +2008,13 @@ 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, @@ -2523,7 +2558,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={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="DELETE", ) try: diff --git a/integrations/claude-code-plugin/core/_harness_id.py b/integrations/claude-code-plugin/core/_harness_id.py index 6c3e1ce15..9c4949408 100644 --- a/integrations/claude-code-plugin/core/_harness_id.py +++ b/integrations/claude-code-plugin/core/_harness_id.py @@ -5,5 +5,7 @@ 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" diff --git a/integrations/claude-code-plugin/core/memory_core.py b/integrations/claude-code-plugin/core/memory_core.py index 77d7a2b80..d350e76d2 100644 --- a/integrations/claude-code-plugin/core/memory_core.py +++ b/integrations/claude-code-plugin/core/memory_core.py @@ -1800,6 +1800,34 @@ 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]: @@ -1807,7 +1835,7 @@ def _request_json( request = urllib.request.Request( url, data=raw, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="POST", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1834,7 +1862,7 @@ def _get_json( ) -> tuple[dict[str, Any] | list[Any], int]: request = urllib.request.Request( url, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="GET", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1980,6 +2008,13 @@ 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, @@ -2523,7 +2558,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={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="DELETE", ) try: diff --git a/integrations/codex-plugin/core/_harness_id.py b/integrations/codex-plugin/core/_harness_id.py index 152313a21..47938ee08 100644 --- a/integrations/codex-plugin/core/_harness_id.py +++ b/integrations/codex-plugin/core/_harness_id.py @@ -5,5 +5,7 @@ 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" diff --git a/integrations/codex-plugin/core/memory_core.py b/integrations/codex-plugin/core/memory_core.py index 77d7a2b80..d350e76d2 100644 --- a/integrations/codex-plugin/core/memory_core.py +++ b/integrations/codex-plugin/core/memory_core.py @@ -1800,6 +1800,34 @@ 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]: @@ -1807,7 +1835,7 @@ def _request_json( request = urllib.request.Request( url, data=raw, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="POST", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1834,7 +1862,7 @@ def _get_json( ) -> tuple[dict[str, Any] | list[Any], int]: request = urllib.request.Request( url, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="GET", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1980,6 +2008,13 @@ 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, @@ -2523,7 +2558,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={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="DELETE", ) try: diff --git a/integrations/cursor-plugin/core/_harness_id.py b/integrations/cursor-plugin/core/_harness_id.py index 40241a7f4..0e20bf1ff 100644 --- a/integrations/cursor-plugin/core/_harness_id.py +++ b/integrations/cursor-plugin/core/_harness_id.py @@ -5,5 +5,7 @@ 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" diff --git a/integrations/cursor-plugin/core/memory_core.py b/integrations/cursor-plugin/core/memory_core.py index 77d7a2b80..d350e76d2 100644 --- a/integrations/cursor-plugin/core/memory_core.py +++ b/integrations/cursor-plugin/core/memory_core.py @@ -1800,6 +1800,34 @@ 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]: @@ -1807,7 +1835,7 @@ def _request_json( request = urllib.request.Request( url, data=raw, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="POST", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1834,7 +1862,7 @@ def _get_json( ) -> tuple[dict[str, Any] | list[Any], int]: request = urllib.request.Request( url, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="GET", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1980,6 +2008,13 @@ 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, @@ -2523,7 +2558,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={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="DELETE", ) try: diff --git a/integrations/deepseek-plugin/src/index.ts b/integrations/deepseek-plugin/src/index.ts index 2d24b498a..3ac86f4b0 100644 --- a/integrations/deepseek-plugin/src/index.ts +++ b/integrations/deepseek-plugin/src/index.ts @@ -26,8 +26,9 @@ export const name = "mem0"; export const inject = ["tools", "systemPrompt"]; // Tags writes so Mem0's backend attributes them to this integration in -// telemetry. The backend's KNOWN_EVENT_SOURCES allowlist recognizes this value; -// anything outside it buckets into "OTHERS". +// telemetry. Values outside the backend's KNOWN_EVENT_SOURCES allowlist bucket +// into "OTHERS"; this one is added by mem0ai/platform#3602 and reads as OTHERS +// until that ships. const SOURCE = "DEEPSEEK_HARNESS"; const DEFAULT_SEARCH_LIMIT = 10; diff --git a/integrations/kimi-plugin/core/_harness_id.py b/integrations/kimi-plugin/core/_harness_id.py index aa17750b1..83ab061ef 100644 --- a/integrations/kimi-plugin/core/_harness_id.py +++ b/integrations/kimi-plugin/core/_harness_id.py @@ -5,5 +5,7 @@ SOURCE_TAG = "KIMI_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 = "kimi" diff --git a/integrations/kimi-plugin/core/memory_core.py b/integrations/kimi-plugin/core/memory_core.py index 77d7a2b80..d350e76d2 100644 --- a/integrations/kimi-plugin/core/memory_core.py +++ b/integrations/kimi-plugin/core/memory_core.py @@ -1800,6 +1800,34 @@ 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]: @@ -1807,7 +1835,7 @@ def _request_json( request = urllib.request.Request( url, data=raw, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="POST", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1834,7 +1862,7 @@ def _get_json( ) -> tuple[dict[str, Any] | list[Any], int]: request = urllib.request.Request( url, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="GET", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1980,6 +2008,13 @@ 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, @@ -2523,7 +2558,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={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="DELETE", ) try: diff --git a/integrations/mem0-agent-plugin/core/_harness_id.py b/integrations/mem0-agent-plugin/core/_harness_id.py index b0836d6de..3efd000f8 100644 --- a/integrations/mem0-agent-plugin/core/_harness_id.py +++ b/integrations/mem0-agent-plugin/core/_harness_id.py @@ -5,5 +5,7 @@ SOURCE_TAG = "CODING_AGENT_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 = "coding-agent" +PLATFORM_APPLICATION = "" diff --git a/integrations/mem0-agent-plugin/core/memory_core.py b/integrations/mem0-agent-plugin/core/memory_core.py index 77d7a2b80..d350e76d2 100644 --- a/integrations/mem0-agent-plugin/core/memory_core.py +++ b/integrations/mem0-agent-plugin/core/memory_core.py @@ -1800,6 +1800,34 @@ 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]: @@ -1807,7 +1835,7 @@ def _request_json( request = urllib.request.Request( url, data=raw, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="POST", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1834,7 +1862,7 @@ def _get_json( ) -> tuple[dict[str, Any] | list[Any], int]: request = urllib.request.Request( url, - headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="GET", ) with urllib.request.urlopen(request, timeout=timeout) as response: @@ -1980,6 +2008,13 @@ 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, @@ -2523,7 +2558,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={"Authorization": f"Token {key}", "Content-Type": "application/json"}, + headers=platform_headers(key), method="DELETE", ) try: diff --git a/integrations/openclaw/cli/config-file.ts b/integrations/openclaw/cli/config-file.ts index 99453cb1b..a34d6c08b 100644 --- a/integrations/openclaw/cli/config-file.ts +++ b/integrations/openclaw/cli/config-file.ts @@ -35,6 +35,8 @@ export interface PluginAuthConfig { autoCapture?: boolean; topK?: number; anonymousTelemetryId?: string; + /** SHA-256 prefix of the API key userEmail was resolved for. */ + keyFingerprint?: string; } // ============================================================================ @@ -135,6 +137,9 @@ export function readPluginAuth(): PluginAuthConfig { autoCapture: cfg.autoCapture as boolean | undefined, topK: cfg.topK as number | undefined, anonymousTelemetryId: cfg.anonymousTelemetryId as string | undefined, + // Without this the reader silently drops it, every fingerprint comparison + // fails against undefined, and the resolved email is never used again. + keyFingerprint: cfg.keyFingerprint as string | undefined, }; } @@ -236,6 +241,21 @@ export function getBaseUrl(): string { return auth.baseUrl || DEFAULT_BASE_URL; } +/** Forget the resolved account, so the next capture re-resolves for the current key. */ +export function clearResolvedAccount(): void { + const full = readFullConfig() as any; + const cfg = full?.plugins?.entries?.[PLUGIN_ID]?.config; + if (!cfg) return; + let changed = false; + for (const key of ["userEmail", "keyFingerprint"]) { + if (key in cfg) { + delete cfg[key]; + changed = true; + } + } + if (changed) writeFullConfig(full); +} + /** Remove anonymousTelemetryId from config (after PostHog aliasing) */ export function clearAnonymousTelemetryId(): void { const full = readFullConfig() as any; diff --git a/integrations/openclaw/config.ts b/integrations/openclaw/config.ts index 4e43c8c44..09dcbc837 100644 --- a/integrations/openclaw/config.ts +++ b/integrations/openclaw/config.ts @@ -148,6 +148,7 @@ const ALLOWED_KEYS = [ "mode", "apiKey", "anonymousTelemetryId", + "keyFingerprint", "baseUrl", "userId", "userEmail", diff --git a/integrations/openclaw/openclaw.plugin.json b/integrations/openclaw/openclaw.plugin.json index 2d9f46518..02d05be4b 100644 --- a/integrations/openclaw/openclaw.plugin.json +++ b/integrations/openclaw/openclaw.plugin.json @@ -209,6 +209,10 @@ "type": "string", "description": "Persistent anonymous telemetry identifier" }, + "keyFingerprint": { + "type": "string", + "description": "Digest of the API key userEmail was resolved for. Set automatically; a mismatch re-resolves the account." + }, "oss": { "type": "object", "properties": { diff --git a/integrations/openclaw/telemetry.ts b/integrations/openclaw/telemetry.ts index fc1c6b99d..bd77e1656 100644 --- a/integrations/openclaw/telemetry.ts +++ b/integrations/openclaw/telemetry.ts @@ -1,14 +1,20 @@ import { createHash, randomUUID } from "node:crypto"; import { createTelemetry } from "../agent-plugin-core/typescript/src/telemetry.ts"; -import { clearAnonymousTelemetryId, getBaseUrl, readPluginAuth, writePluginAuth } from "./cli/config-file.ts"; +import { + clearAnonymousTelemetryId, + clearResolvedAccount, + getBaseUrl, + readPluginAuth, + writePluginAuth, +} from "./cli/config-file.ts"; declare const __OPENCLAW_PLUGIN_VERSION__: string; export const PLUGIN_VERSION: string = __OPENCLAW_PLUGIN_VERSION__; let cachedAnonymousId: string | undefined; let aliasCheckDone = false; -let emailResolutionAttempted = false; +let resolutionAttemptedFor = ""; let currentDistinctId = ""; function enabled(): boolean { @@ -33,10 +39,35 @@ function anonymousId(): string { return (cachedAnonymousId = created); } +/** SHA-256 prefix of the key an account was resolved for. */ +function keyFingerprint(apiKey?: string): string { + return apiKey ? createHash("sha256").update(apiKey).digest("hex").slice(0, 16) : ""; +} + function distinctId(apiKey?: string): string { try { - const email = readPluginAuth().userEmail; - if (email) return createHash("sha256").update(email).digest("hex"); + const auth = readPluginAuth(); + if (auth.userEmail) { + // Only when it belongs to the key in hand. Without this check a cached + // email was used forever: switch to a different account and every event + // kept reporting under the previous one, with nothing to notice it by. + if (auth.keyFingerprint === keyFingerprint(apiKey)) { + return createHash("sha256").update(auth.userEmail).digest("hex"); + } + // Only a REAL key that disagrees means the account changed. Without the + // apiKey guard the comparison is `undefined === ""` for any call that + // simply omits the key, so a capture with no context wiped a perfectly + // good account out of openclaw.json. + // + // A row with an email and NO fingerprint is the legacy shape, from an + // install predating this field. Clearing it here deleted a real account + // before anything had replaced it, and if the re-resolve then failed + // because the user was offline the email was gone from disk for good. The + // Python core refuses the same trade: verify, and keep what you have until + // the verification succeeds. resolveEmail below overwrites both fields + // when it does, so there is nothing to clear first. + if (apiKey && auth.keyFingerprint) clearResolvedAccount(); + } } catch { // Fall through to the API key or anonymous identity. } @@ -66,8 +97,16 @@ function identifyAnonymous(id: string): void { } function resolveEmail(apiKey: string): void { - if (emailResolutionAttempted) return; - emailResolutionAttempted = true; + // Latched per key, not once per process. A single boolean meant a key changed + // mid-session was never looked up, so the fallback identity stuck until restart. + const fingerprint = keyFingerprint(apiKey); + if (resolutionAttemptedFor === fingerprint) return; + resolutionAttemptedFor = fingerprint; + const releaseLatch = () => { + // A failed lookup must not pin the fallback identity for the rest of the + // process. Released so the next capture tries again. + if (resolutionAttemptedFor === fingerprint) resolutionAttemptedFor = ""; + }; fetch(`${getBaseUrl().replace(/\/+$/, "")}/v1/ping/`, { method: "GET", headers: { Authorization: `Token ${apiKey}`, "Content-Type": "application/json" }, @@ -76,7 +115,7 @@ function resolveEmail(apiKey: string): void { .then((response) => response.json()) .then((data: any) => { if (!data?.user_email) return; - writePluginAuth({ userEmail: data.user_email }); + writePluginAuth({ userEmail: data.user_email, keyFingerprint: fingerprint }); const oldId = createHash("sha256").update(apiKey).digest("hex"); const newId = createHash("sha256").update(data.user_email).digest("hex"); for (const event of telemetry.queueForTesting()) { @@ -84,7 +123,8 @@ function resolveEmail(apiKey: string): void { } }) .catch(() => { - // The API-key hash remains a stable fallback. + // The API-key hash remains a stable fallback, and the next capture retries. + releaseLatch(); }); } @@ -96,13 +136,15 @@ export function captureEvent( if (!enabled()) return; try { currentDistinctId = distinctId(context?.apiKey); - let hasEmail = false; + let resolvedForThisKey = false; try { - hasEmail = Boolean(readPluginAuth().userEmail); + const auth = readPluginAuth(); + resolvedForThisKey = + Boolean(auth.userEmail) && auth.keyFingerprint === keyFingerprint(context?.apiKey); } catch { // Resolve it below when possible. } - if (context?.apiKey && !hasEmail) resolveEmail(context.apiKey); + if (context?.apiKey && !resolvedForThisKey) resolveEmail(context.apiKey); identifyAnonymous(currentDistinctId); telemetry.capture(eventName, { mode: context?.mode, diff --git a/integrations/openclaw/tests/config-file.test.ts b/integrations/openclaw/tests/config-file.test.ts index e059ddff3..ba63180cb 100644 --- a/integrations/openclaw/tests/config-file.test.ts +++ b/integrations/openclaw/tests/config-file.test.ts @@ -237,3 +237,43 @@ describe("getBaseUrl", () => { expect(getBaseUrl()).toBe(DEFAULT_BASE_URL); }); }); + +// --------------------------------------------------------------------------- +// keyFingerprint round trip +// --------------------------------------------------------------------------- + +describe("keyFingerprint survives a write and read", () => { + it("readPluginAuth returns a persisted keyFingerprint", () => { + // It did not. readPluginAuth builds its result field by field, and this one + // was missing, so every fingerprint comparison ran against undefined, the + // resolved email was never used again, and telemetry silently fell back to + // the API key hash. The telemetry tests could not see it because they mock + // this module and their mock returned the field the real reader dropped. + setConfigFile({ + plugins: { + entries: { + "openclaw-mem0": { + config: { + apiKey: "m0-test", + userEmail: "person@example.com", + keyFingerprint: "0123456789abcdef", + }, + }, + }, + }, + }); + + expect(readPluginAuth().keyFingerprint).toBe("0123456789abcdef"); + }); + + it("writePluginAuth persists it where readPluginAuth looks", () => { + setConfigFile({ plugins: { entries: { "openclaw-mem0": { config: {} } } } }); + + writePluginAuth({ userEmail: "person@example.com", keyFingerprint: "abc123" }); + + const written = JSON.parse(mockWriteText.mock.calls.at(-1)![1] as string); + const cfg = written.plugins.entries["openclaw-mem0"].config; + expect(cfg.keyFingerprint).toBe("abc123"); + expect(cfg.userEmail).toBe("person@example.com"); + }); +}); diff --git a/integrations/openclaw/tests/config.test.ts b/integrations/openclaw/tests/config.test.ts index 2c6dbac2e..334c064fb 100644 --- a/integrations/openclaw/tests/config.test.ts +++ b/integrations/openclaw/tests/config.test.ts @@ -485,3 +485,18 @@ describe("mem0ConfigSchema.parse() — apiKey edge cases", () => { expect(cfg.needsSetup).toBe(false); }); }); + +describe("telemetry fingerprint round trip", () => { + it("a config carrying keyFingerprint is accepted by the real schema", () => { + // What writePluginAuth persists after a successful lookup. It was not in + // ALLOWED_KEYS, and assertAllowedKeys throws, so the first successful + // resolve wrote a config that broke every subsequent load of the plugin. + const persisted = { + apiKey: "m0-test", + userEmail: "person@example.com", + keyFingerprint: "0123456789abcdef", + }; + + expect(() => mem0ConfigSchema.parse(persisted)).not.toThrow(); + }); +}); diff --git a/integrations/openclaw/tests/telemetry.test.ts b/integrations/openclaw/tests/telemetry.test.ts index 9a7bee716..ad7a9bbe0 100644 --- a/integrations/openclaw/tests/telemetry.test.ts +++ b/integrations/openclaw/tests/telemetry.test.ts @@ -3,15 +3,30 @@ import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; // Mock config-file before importing telemetry vi.mock("../cli/config-file.ts", () => ({ readPluginAuth: vi.fn().mockReturnValue({}), + writePluginAuth: vi.fn(), + clearAnonymousTelemetryId: vi.fn(), + clearResolvedAccount: vi.fn(), + getBaseUrl: vi.fn().mockReturnValue("https://api.mem0.ai"), })); import { captureEvent } from "../telemetry.ts"; -import { readPluginAuth } from "../cli/config-file.ts"; +import { clearResolvedAccount, readPluginAuth } from "../cli/config-file.ts"; + +/** sha256(key).slice(0, 16), the shape telemetry.ts stores. */ +async function fingerprintOf(apiKey: string): Promise { + const { createHash } = await import("node:crypto"); + return createHash("sha256").update(apiKey).digest("hex").slice(0, 16); +} describe("telemetry", () => { let fetchSpy: ReturnType; beforeEach(() => { + // Call history has to be cleared per test, not just restored: the mocks are + // module-level vi.fn()s, so without this one test's calls are visible to the + // next and assertions on "was not called" pass or fail by ordering. + vi.clearAllMocks(); + (readPluginAuth as ReturnType).mockReturnValue({}); // Reset telemetry enabled state (globalThis as any).__mem0_telemetry_override = undefined; fetchSpy = vi.fn().mockResolvedValue({ ok: true }); @@ -40,11 +55,31 @@ describe("telemetry", () => { expect(() => captureEvent("test_event")).not.toThrow(); }); - it("uses userEmail as distinct ID when available", () => { - (readPluginAuth as ReturnType).mockReturnValueOnce({ + it("uses userEmail as distinct ID when it belongs to the current key", async () => { + // Previously asserted only not.toThrow(), which passed whatever the identity + // turned out to be, and under the fingerprint gate the no-context path does + // not use the email at all. Pin the real condition instead. + (readPluginAuth as ReturnType).mockReturnValue({ userEmail: "test@example.com", + keyFingerprint: await fingerprintOf("key-a"), }); - expect(() => captureEvent("test_event")).not.toThrow(); + + captureEvent("test_event", {}, { apiKey: "key-a" }); + + expect(clearResolvedAccount).not.toHaveBeenCalled(); + }); + + it("a capture with no apiKey leaves a resolved account alone", async () => { + // `undefined === ""` made every keyless capture look like a key change, so + // one context-free call wiped a good account out of openclaw.json. + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "test@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event"); + + expect(clearResolvedAccount).not.toHaveBeenCalled(); }); it("falls back to a generated anonymous id when no apiKey", () => { @@ -52,6 +87,43 @@ describe("telemetry", () => { expect(() => captureEvent("test_event", {}, {})).not.toThrow(); }); + it("keeps using a cached email only while it belongs to the current key", async () => { + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "person@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event", {}, { apiKey: "key-a" }); + + expect(clearResolvedAccount).not.toHaveBeenCalled(); + }); + + it("forgets the account when the API key changes", async () => { + // The defect: the cached email was used forever, so events after an account + // switch kept reporting under the previous account. + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "person@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event", {}, { apiKey: "key-b" }); + + expect(clearResolvedAccount).toHaveBeenCalled(); + }); + + it("re-resolves for a key it has not looked up before", async () => { + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "person@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event", {}, { apiKey: "key-c" }); + + // The resolution latch is per key, not once per process, so a key changed + // mid-session is actually looked up instead of sticking to the fallback. + expect(fetchSpy).toHaveBeenCalled(); + }); + it("handles readPluginAuth errors gracefully", () => { (readPluginAuth as ReturnType).mockImplementationOnce(() => { throw new Error("config read failed"); diff --git a/integrations/opencode-plugin/package.json b/integrations/opencode-plugin/package.json index 79f68c4a8..acdbaf469 100644 --- a/integrations/opencode-plugin/package.json +++ b/integrations/opencode-plugin/package.json @@ -38,6 +38,7 @@ "build": "bun build opencode-mem0.ts --outdir dist --target bun --format esm --entry-naming index.[ext]", "dev": "bun build opencode-mem0.ts --outdir dist --target bun --format esm --entry-naming index.[ext] --watch", "type-check": "tsc --noEmit", + "test": "bun test", "prepack": "bun run build", "postpack": "" }, diff --git a/integrations/opencode-plugin/telemetry.test.ts b/integrations/opencode-plugin/telemetry.test.ts index 393150977..79dc3ffdb 100644 --- a/integrations/opencode-plugin/telemetry.test.ts +++ b/integrations/opencode-plugin/telemetry.test.ts @@ -1,3 +1,5 @@ +import { createHash } from "node:crypto"; + import { afterEach, describe, expect, test } from "bun:test"; import { buildEvent, captureEvent, isTelemetryEnabled } from "./telemetry"; @@ -13,7 +15,10 @@ describe("opencode telemetry", () => { expect(payload).not.toBeNull(); const props = payload!.properties as Record; expect(payload!.event).toBe("plugin.session_start"); - expect(props.source).toBe("plugin"); + // Was "plugin", which named no particular plugin and matched no vocabulary. + // Now shaped like every other surface. Any saved PostHog insight filtering + // source = "plugin" needs repointing; historical data is untouched. + expect(props.source).toBe("OPENCODE_PLUGIN"); expect(props.platform).toBe("opencode"); expect(props.memory_count).toBe(5); expect(props.$process_person_profile).toBe(false); @@ -30,7 +35,7 @@ describe("opencode telemetry", () => { const props = buildEvent("x", { platform: "HACK", source: "HACK" }, KEY)! .properties as Record; expect(props.platform).toBe("opencode"); - expect(props.source).toBe("plugin"); + expect(props.source).toBe("OPENCODE_PLUGIN"); }); test("returns null without an API key (no anonymous events)", () => { @@ -56,9 +61,12 @@ describe("opencode telemetry", () => { expect(typeof props.os_version).toBe("string"); }); - test("project_hash is sha256(projectId) when a project id is supplied", async () => { + test("project_hash is a salted digest of the project id", async () => { + // Previously asserted the bare sha256(projectId), which is the defect: that + // digest is reversible by anyone who can guess a project id. Salted with the + // API key, which is already in play here and is high entropy. const { createHash } = await import("node:crypto"); - const expected = createHash("sha256").update("acme-repo").digest("hex"); + const expected = createHash("sha256").update(`${KEY}:acme-repo`).digest("hex"); const props = buildEvent("session_start", {}, KEY, "acme-repo")! .properties as Record; expect(props.project_hash).toBe(expected); @@ -76,3 +84,38 @@ describe("opencode telemetry", () => { } }); }); + +describe("project_hash salting", () => { + const PROJECT = "my-project"; + + test("is not a bare digest of the project id", () => { + // The defect: an unsalted SHA-256 over a guessable identifier is reversible + // by anyone who can enumerate project ids. + const unsalted = createHash("sha256").update(PROJECT).digest("hex"); + const payload = buildEvent("session_start", {}, KEY, PROJECT) as Record; + + expect(payload.properties.project_hash).toBeDefined(); + expect(payload.properties.project_hash).not.toBe(unsalted); + }); + + test("differs per account for the same project", () => { + const a = buildEvent("session_start", {}, "m0-account-a", PROJECT) as Record; + const b = buildEvent("session_start", {}, "m0-account-b", PROJECT) as Record; + + expect(a.properties.project_hash).not.toBe(b.properties.project_hash); + }); + + test("is stable for one account, so joins still work", () => { + const first = buildEvent("session_start", {}, KEY, PROJECT) as Record; + const second = buildEvent("session_end", {}, KEY, PROJECT) as Record; + + expect(first.properties.project_hash).toBe(second.properties.project_hash); + }); + + test("is omitted rather than unsalted when there is no key", () => { + const payload = buildEvent("session_start", {}, undefined, PROJECT); + + // No key means no event at all, so there is no unsalted hash to leak. + expect(payload).toBeNull(); + }); +}); diff --git a/integrations/opencode-plugin/telemetry.ts b/integrations/opencode-plugin/telemetry.ts index 8d0861d97..e2316dd88 100644 --- a/integrations/opencode-plugin/telemetry.ts +++ b/integrations/opencode-plugin/telemetry.ts @@ -25,7 +25,9 @@ const PLUGIN_VERSION = (() => { let currentDistinctId = ""; const telemetry = createTelemetry({ host: "opencode", - source: "plugin", + // Shaped like the platform's EventSource values, as every other surface is. + // "plugin" said nothing about which plugin and matched no vocabulary. + source: "OPENCODE_PLUGIN", version: PLUGIN_VERSION, distinctId: () => currentDistinctId, eventName: (event) => `plugin.${event}`, @@ -40,8 +42,23 @@ function distinctId(apiKey: string): string { return createHash("sha256").update(apiKey).digest("hex").slice(0, 32); } -function projectHash(projectId?: string): Record { - return projectId ? { project_hash: createHash("sha256").update(projectId).digest("hex") } : {}; +/** + * Salted so the hash is not enumerable. + * + * An unsalted SHA-256 of a project id is reversible by anyone who can guess the + * id, which for a project identifier is a small space. The API key is the salt: + * it is already in play here (distinctId is a digest of it), it is high entropy, + * and using it needs no per-install file and so no write race to get wrong. The + * hash is therefore per account rather than per machine, which also keeps joins + * working for one user across machines. It resets when the key rotates, which is + * consistent, because distinctId resets with it. + * + * Both are omitted without a key. An event cannot be built without a distinctId + * anyway, so this costs nothing. + */ +function projectHash(projectId?: string, apiKey?: string): Record { + if (!projectId || !apiKey) return {}; + return { project_hash: createHash("sha256").update(`${apiKey}:${projectId}`).digest("hex") }; } export function buildEvent( @@ -51,7 +68,7 @@ export function buildEvent( projectId?: string, ): Record | null { currentDistinctId = apiKey ? distinctId(apiKey) : ""; - const event = telemetry.build(eventType, { ...properties, ...projectHash(projectId) }); + const event = telemetry.build(eventType, { ...properties, ...projectHash(projectId, apiKey) }); return event ? { api_key: POSTHOG_API_KEY, ...event } : null; } @@ -62,5 +79,5 @@ export function captureEvent( projectId?: string, ): void { currentDistinctId = apiKey ? distinctId(apiKey) : ""; - telemetry.capture(eventType, { ...properties, ...projectHash(projectId) }); + telemetry.capture(eventType, { ...properties, ...projectHash(projectId, apiKey) }); } diff --git a/integrations/pi-agent-plugin/src/attribution.test.ts b/integrations/pi-agent-plugin/src/attribution.test.ts new file mode 100644 index 000000000..f90d28da3 --- /dev/null +++ b/integrations/pi-agent-plugin/src/attribution.test.ts @@ -0,0 +1,80 @@ +import { describe, expect, it } from "vitest"; +import { applySurfaceHeaders, PLATFORM_APPLICATION, PLATFORM_SOURCE } from "./attribution.ts"; + +function client(headers: Record = {}) { + return { headers: { Authorization: "Token k", ...headers } } as never; +} + +describe("applySurfaceHeaders", () => { + it("stamps the shared client so every path is attributed, not just commands", () => { + const mem0 = client(); + applySurfaceHeaders(mem0); + const headers = (mem0 as unknown as { headers: Record }).headers; + + expect(headers["X-Mem0-Source"]).toBe(PLATFORM_SOURCE); + expect(headers["X-Application"]).toBe(PLATFORM_APPLICATION); + expect(headers["X-Mem0-Client"]).toMatch(/^mem0-pi-agent\//); + expect(headers.Authorization).toBe("Token k"); + }); + + it("defers to a surface an outer wrapper already declared", () => { + const mem0 = client({ "X-Mem0-Source": "OPENCLAW", "X-Application": "vscode" }); + applySurfaceHeaders(mem0); + const headers = (mem0 as unknown as { headers: Record }).headers; + + expect(headers["X-Mem0-Source"]).toBe("OPENCLAW"); + expect(headers["X-Application"]).toBe("vscode"); + }); + + it("appends to the client stack rather than replacing it", () => { + const mem0 = client({ "X-Mem0-Client": "openclaw/2.1.0" }); + applySurfaceHeaders(mem0); + const headers = (mem0 as unknown as { headers: Record }).headers; + + expect(headers["X-Mem0-Client"]).toMatch(/^openclaw\/2\.1\.0, mem0-pi-agent\//); + }); + + it("treats a blank header as absent", () => { + const mem0 = client({ "X-Mem0-Source": " " }); + applySurfaceHeaders(mem0); + const headers = (mem0 as unknown as { headers: Record }).headers; + + expect(headers["X-Mem0-Source"]).toBe(PLATFORM_SOURCE); + }); + + it("bounds the stack so a long chain cannot grow the header without limit", () => { + const mem0 = client({ "X-Mem0-Client": "a/1, b/1, c/1, d/1, e/1" }); + applySurfaceHeaders(mem0); + const headers = (mem0 as unknown as { headers: Record }).headers; + + expect(headers["X-Mem0-Client"].split(",").length).toBeLessThanOrEqual(4); + expect(headers["X-Mem0-Client"].length).toBeLessThanOrEqual(200); + }); +}); + +describe("client stack bounding", () => { + it("keeps our own entry when the caller already filled the stack", () => { + // The defect: pushing then trimming to four dropped exactly the entry this + // function exists to add, so we vanished from our own stack. + const mem0 = client({ "X-Mem0-Client": "a/1, b/2, c/3, d/4" }); + applySurfaceHeaders(mem0); + const stack = (mem0 as unknown as { headers: Record }).headers["X-Mem0-Client"]; + + expect(stack).toMatch(/mem0-pi-agent\//); + expect(stack.split(",").length).toBeLessThanOrEqual(4); + }); + + it("drops whole entries at the character cap, never a fragment", () => { + const long = `${"n".repeat(90)}/1.0, ${"m".repeat(90)}/1.0, ${"o".repeat(90)}/1.0`; + const mem0 = client({ "X-Mem0-Client": long }); + applySurfaceHeaders(mem0); + const stack = (mem0 as unknown as { headers: Record }).headers["X-Mem0-Client"]; + + expect(stack.length).toBeLessThanOrEqual(200); + expect(stack.endsWith("/0.0.0") || /mem0-pi-agent\/[\w.\-]+$/.test(stack)).toBe(true); + // Every surviving entry is whole: name/version, no severed tail. + for (const entry of stack.split(",")) { + expect(entry.trim()).toMatch(/^[^/]+\/[^/]+$/); + } + }); +}); diff --git a/integrations/pi-agent-plugin/src/attribution.ts b/integrations/pi-agent-plugin/src/attribution.ts new file mode 100644 index 000000000..2cbc9a73e --- /dev/null +++ b/integrations/pi-agent-plugin/src/attribution.ts @@ -0,0 +1,65 @@ +import type MemoryClient from "mem0ai"; +import * as fs from "node:fs"; + +/** Surface identity for this plugin, as the platform's EventSource knows it. */ +export const PLATFORM_SOURCE = "PI_AGENT"; + +/** Host app the plugin runs inside. Allowlisted server-side. */ +export const PLATFORM_APPLICATION = "pi"; + +const PLUGIN_VERSION = (() => { + try { + return JSON.parse( + fs.readFileSync(new URL("../package.json", import.meta.url), "utf-8"), + ).version as string; + } catch { + return "unknown"; + } +})(); + +const MAX_STACK_ENTRIES = 4; +const MAX_STACK_CHARS = 200; + +/** + * Append our own entry and bound the result, dropping WHOLE entries. + * + * Neither cap cuts characters: slicing the joined string severs an identifier + * and leaves a fragment the platform parses as a real client name. And the + * reserved slot is ours, since it is the only entry this layer can vouch for. + */ +function boundedStack(callerEntries: string[], own: string): string { + const kept: string[] = []; + let budget = MAX_STACK_CHARS - own.length; + for (const entry of callerEntries.slice(0, MAX_STACK_ENTRIES - 1)) { + const cost = entry.length + ", ".length; + if (cost > budget) break; + budget -= cost; + kept.push(entry); + } + return [...kept, own].join(", "); +} + +/** + * Stamp surface identity onto the shared client, once, at construction. + * + * Tagging individual call sites was not enough: automatic recall, capture, the + * memory tools and deletion all go through this same client, so everything + * except the explicit slash commands reached the platform as generic SDK + * traffic. Every request method in the SDK sends `this.headers`, so setting + * them here covers all of them. + * + * X-Mem0-Source and X-Application are set-once, so a wrapper that already named + * a surface keeps it. X-Mem0-Client is append-only, so the platform sees the + * whole chain rather than only the last speaker. + */ +export function applySurfaceHeaders(client: MemoryClient): void { + const headers = client.headers as Record; + if (!headers["X-Mem0-Source"]?.trim()) headers["X-Mem0-Source"] = PLATFORM_SOURCE; + if (!headers["X-Application"]?.trim()) headers["X-Application"] = PLATFORM_APPLICATION; + + const existing = (headers["X-Mem0-Client"] ?? "") + .split(",") + .map((part) => part.trim()) + .filter(Boolean); + headers["X-Mem0-Client"] = boundedStack(existing, `mem0-pi-agent/${PLUGIN_VERSION}`); +} diff --git a/integrations/pi-agent-plugin/src/commands.ts b/integrations/pi-agent-plugin/src/commands.ts index 9379bed9a..e969478a1 100644 --- a/integrations/pi-agent-plugin/src/commands.ts +++ b/integrations/pi-agent-plugin/src/commands.ts @@ -1,3 +1,4 @@ +import type { SearchMemoryOptions } from "mem0ai"; import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import type MemoryClient from "mem0ai"; import type { Mem0Config, ScopeContext, Scope } from "./types.ts"; @@ -5,6 +6,12 @@ import { DEFAULT_CUSTOM_CATEGORIES } from "./types.ts"; import { resolveSearchFilters, resolveAddParams } from "./memory/scoping.ts"; import { formatMemoryList, formatMemoryCompact, groupByCategory } from "./memory/formatting.ts"; import { captureCommandEvent } from "./telemetry.ts"; +import { PLATFORM_SOURCE } from "./attribution.ts"; + +// Wire identity is set once on the shared client in entry.ts, which covers +// every path including recall, capture, tools and deletion. It stays in the +// body of the two calls below as well: body `source` is what the backend reads +// when the header is absent. const SEARCH_TOP_K = 10; @@ -29,7 +36,12 @@ export function registerCommands( threshold: config.searchThreshold, topK: SEARCH_TOP_K, rerank: true, - }); + source: PLATFORM_SOURCE, + // Widened by exactly this one property. `source` reaches the wire via the + // SDK's camelToSnakeKeys spread, but it is absent from SearchMemoryOptions + // in the published mem0ai types. A blanket `as never` would also disable + // checking of filters, threshold, topK and rerank above. + } as SearchMemoryOptions & { source: string }); return result.results ?? []; }; @@ -45,7 +57,7 @@ export function registerCommands( const addParams = resolveAddParams(config.defaultScope, getScopeCtx()); const result = await mem0.add( [{ role: "user", content: text }], - { ...addParams, customCategories: DEFAULT_CUSTOM_CATEGORIES, infer: false }, + { ...addParams, customCategories: DEFAULT_CUSTOM_CATEGORIES, infer: false, source: PLATFORM_SOURCE }, ); captureCommandEvent("mem0-remember", {}, telemetryCtx); diff --git a/integrations/pi-agent-plugin/src/entry.ts b/integrations/pi-agent-plugin/src/entry.ts index d49d78ad0..256c05aa4 100644 --- a/integrations/pi-agent-plugin/src/entry.ts +++ b/integrations/pi-agent-plugin/src/entry.ts @@ -10,6 +10,7 @@ import { captureEvent } from "./telemetry.ts"; import * as os from "node:os"; import type { ScopeContext } from "./types.ts"; import { createMemoryLifecycle } from "../../agent-plugin-core/typescript/src/lifecycle.ts"; +import { applySurfaceHeaders } from "./attribution.ts"; export { buildRecallContext } from "../../agent-plugin-core/typescript/src/lifecycle.ts"; @@ -29,6 +30,11 @@ export default function mem0Extension(pi: ExtensionAPI): void { } const mem0 = new MemoryClient({ apiKey: config.apiKey }); + // Every path below shares this client: automatic recall, capture, the memory + // tools and deletion as well as the slash commands. Attribution belongs here + // rather than on individual calls, or everything except the commands reports + // as generic SDK traffic. + applySurfaceHeaders(mem0); const scopeCtx: ScopeContext = { userId: resolveUserId(config.userId), diff --git a/integrations/vercel-ai-sdk/src/mem0-utils.ts b/integrations/vercel-ai-sdk/src/mem0-utils.ts index 55d485cbc..b7a041939 100644 --- a/integrations/vercel-ai-sdk/src/mem0-utils.ts +++ b/integrations/vercel-ai-sdk/src/mem0-utils.ts @@ -1,3 +1,10 @@ +declare const __MEM0_PROVIDER_VERSION__: string | undefined; + +// Replaced at build time by tsup `define`. The fallback only applies when the +// source is run unbundled, such as in tests. +const PROVIDER_VERSION = + typeof __MEM0_PROVIDER_VERSION__ !== "undefined" ? __MEM0_PROVIDER_VERSION__ : "dev"; + import { LanguageModelV3Prompt } from '@ai-sdk/provider'; import { Mem0ConfigSettings } from './mem0-types'; import { loadApiKey } from '@ai-sdk/provider-utils'; @@ -277,7 +284,11 @@ const searchInternalMemories = async (query: string, config?: Mem0ConfigSettings method: 'POST', headers: { Authorization: `Token ${apiKey}`, - 'Content-Type': 'application/json' + 'Content-Type': 'application/json', + // Surface attribution. Set-once by contract: this wrapper is the + // outermost layer on these raw fetch calls. + 'X-Mem0-Source': 'VERCEL_AI_SDK', + 'X-Mem0-Client': `mem0-vercel-ai-provider/${PROVIDER_VERSION}` }, body: JSON.stringify(body), }; @@ -331,7 +342,11 @@ const updateMemories = async (messages: Array, config?: Mem0ConfigSetti method: 'POST', headers: { Authorization: `Token ${apiKey}`, - 'Content-Type': 'application/json' + 'Content-Type': 'application/json', + // Surface attribution. Set-once by contract: this wrapper is the + // outermost layer on these raw fetch calls. + 'X-Mem0-Source': 'VERCEL_AI_SDK', + 'X-Mem0-Client': `mem0-vercel-ai-provider/${PROVIDER_VERSION}` }, body: JSON.stringify(body), }; diff --git a/integrations/vercel-ai-sdk/tsconfig.json b/integrations/vercel-ai-sdk/tsconfig.json index b05d5db02..e40350730 100644 --- a/integrations/vercel-ai-sdk/tsconfig.json +++ b/integrations/vercel-ai-sdk/tsconfig.json @@ -12,6 +12,7 @@ "noUnusedLocals": false, "noUnusedParameters": false, "preserveWatchOutput": true, + "resolveJsonModule": true, "skipLibCheck": true, "strict": true, "types": ["@types/node", "jest"], diff --git a/integrations/vercel-ai-sdk/tsup.config.ts b/integrations/vercel-ai-sdk/tsup.config.ts index 2c8f74a6d..ff2d298c0 100644 --- a/integrations/vercel-ai-sdk/tsup.config.ts +++ b/integrations/vercel-ai-sdk/tsup.config.ts @@ -1,4 +1,5 @@ import { defineConfig } from 'tsup' +import pkg from './package.json' export default defineConfig([ { @@ -6,5 +7,11 @@ export default defineConfig([ entry: ['src/index.ts'], format: ['cjs', 'esm'], sourcemap: true, + // Injected rather than written in the source. A hardcoded literal matches + // package.json on the day it is written and misreports the client version + // from the next release bump onwards. Same mechanism as mem0-ts. + define: { + __MEM0_PROVIDER_VERSION__: JSON.stringify(pkg.version), + }, }, ]) \ No newline at end of file diff --git a/mem0-ts/src/client/mem0.ts b/mem0-ts/src/client/mem0.ts index ff259f2fe..33e03f0c2 100644 --- a/mem0-ts/src/client/mem0.ts +++ b/mem0-ts/src/client/mem0.ts @@ -95,6 +95,66 @@ interface ClientIdentity { const IDENTITY_CACHE_MAX_DEFAULT = 50; const identityByCredentials = new Map>(); +declare const __MEM0_SDK_VERSION__: string | undefined; + +// Injected by tsup (see mem0-ts/tsup.config.ts `define`), the same mechanism +// telemetry.ts already uses. A hardcoded literal goes stale at the next release +// bump and then misreports the client version forever. +const SDK_VERSION = + typeof __MEM0_SDK_VERSION__ !== "undefined" ? __MEM0_SDK_VERSION__ : "dev"; + +const MAX_STACK_ENTRIES = 4; +const MAX_STACK_CHARS = 200; + +/** + * Append our own entry and bound the result, dropping WHOLE entries. + * + * Neither cap cuts characters: slicing the joined string severs an identifier + * and leaves a fragment the platform parses as a real client name. And the + * reserved slot is ours. Pushing first and then trimming to four dropped exactly + * the entry this exists to add whenever a caller already sent four, so we + * vanished from our own stack while every caller claim survived. + */ +function boundedStack(callerEntries: string[], own: string): string { + const kept: string[] = []; + let budget = MAX_STACK_CHARS - own.length; + for (const entry of callerEntries.slice(0, MAX_STACK_ENTRIES - 1)) { + const cost = entry.length + ", ".length; + if (cost > budget) break; + budget -= cost; + kept.push(entry); + } + return [...kept, own].join(", "); +} + +/** + * Surface-identity headers. + * + * X-Mem0-Source and X-Application are SET-ONCE by contract: whichever layer is + * outermost sets them and nothing below overwrites, so a plugin wrapping this + * SDK keeps its own identity. X-Mem0-Client is APPEND-ONLY - every layer adds + * itself, so the platform sees the whole stack and not just the last speaker. + */ +function surfaceHeaders(): Record { + const env: Record = + typeof process !== "undefined" && process.env ? process.env : {}; + const existing = (env.MEM0_CLIENT_STACK ?? "").trim(); + const entries = existing + ? existing + .split(",") + .map((part) => part.trim()) + .filter(Boolean) + : []; + const headers: Record = { + "X-Mem0-Client": boundedStack(entries, `mem0-js/${SDK_VERSION}`), + }; + const source = (env.MEM0_SOURCE ?? "").trim(); + if (source) headers["X-Mem0-Source"] = source; + const application = (env.MEM0_APPLICATION ?? "").trim(); + if (application) headers["X-Application"] = application; + return headers; +} + export default class MemoryClient { apiKey: string; host: string; @@ -129,6 +189,7 @@ export default class MemoryClient { this.headers = { Authorization: `Token ${this.apiKey}`, "Content-Type": "application/json", + ...surfaceHeaders(), }; this.client = axios.create({ diff --git a/mem0-ts/src/client/mem0.types.ts b/mem0-ts/src/client/mem0.types.ts index c230441e0..cbe541321 100644 --- a/mem0-ts/src/client/mem0.types.ts +++ b/mem0-ts/src/client/mem0.types.ts @@ -30,6 +30,9 @@ export interface SearchMemoryOptions { showExpired?: boolean; referenceDate?: string | number; keywordSearch?: boolean; + /** Surface that produced the call, e.g. "OPENCLAW". Must be a value the + * backend's EventSource enum knows, or it buckets into OTHERS. */ + source?: string; } export interface GetAllMemoryOptions { diff --git a/mem0/client/main.py b/mem0/client/main.py index a52986dcd..3a82d9e5a 100644 --- a/mem0/client/main.py +++ b/mem0/client/main.py @@ -79,6 +79,95 @@ def _maybe_alias_anon_to_email(user_email): logger.debug("Failed to alias anon telemetry to %r: %s", user_email, e) +def _sdk_version() -> str: + """Resolved here rather than imported from the package root, which would cycle.""" + try: + import importlib.metadata + + return importlib.metadata.version("mem0ai") + except Exception: + return "unknown" + + +def _apply_client_headers(client: Any, api_key: str, user_id: str) -> None: + """Merge our headers into a caller-supplied client without erasing theirs. + + A wrapper may hand us a client already carrying its own X-Mem0-Source or a + partial X-Mem0-Client stack. Blanket update() replaced both, which is the + opposite of the set-once / append-only contract: the outermost layer is the + one whose identity should survive. + """ + existing = client.headers + mine = _client_headers(api_key, user_id) + + outer_stack = existing.get("X-Mem0-Client") + if outer_stack: + entries = [part.strip() for part in str(outer_stack).split(",") if part.strip()] + mine["X-Mem0-Client"] = _bounded_stack(entries, f"mem0-python/{_sdk_version()}") + + for name, value in mine.items(): + if name in ("X-Mem0-Source", "X-Application") and existing.get(name): + continue + existing[name] = value + + +MAX_STACK_ENTRIES = 4 +MAX_STACK_CHARS = 200 + + +def _bounded_stack(caller_entries, own: str) -> str: + """Append our own entry and bound the result, dropping WHOLE entries. + + Two rules, and the second is the one that was wrong. Neither cap cuts + characters: a blunt slice severs an identifier and leaves a fragment that + parses as a real client name. And the reserved slot is OURS. Appending first + and then trimming to four dropped exactly the entry this function exists to + add, every time a caller already sent four, so the SDK vanished from its own + stack while the caller's claims all survived. + """ + kept = [] + budget = MAX_STACK_CHARS - len(own) + for entry in list(caller_entries)[: MAX_STACK_ENTRIES - 1]: + cost = len(entry) + len(", ") + if cost > budget: + break + budget -= cost + kept.append(entry) + return ", ".join(kept + [own]) + + +def _client_headers(api_key: str, user_id: str) -> Dict[str, str]: + """Auth plus surface-identity headers. + + X-Mem0-Source and X-Application are SET-ONCE by contract: whichever layer is + outermost sets them, and nothing below overwrites. A plugin or harness that + wraps this SDK therefore keeps its own identity — it declares via MEM0_SOURCE + / MEM0_APPLICATION and the SDK defers. + + X-Mem0-Client is APPEND-ONLY: every layer adds itself, so the platform sees + the whole stack rather than only whoever spoke last. + """ + headers = { + "Authorization": f"Token {api_key}", + "Mem0-User-ID": user_id, + "X-Mem0-Client": _client_stack(), + } + source = os.getenv("MEM0_SOURCE", "").strip() + if source: + headers["X-Mem0-Source"] = source + application = os.getenv("MEM0_APPLICATION", "").strip() + if application: + headers["X-Application"] = application + return headers + + +def _client_stack() -> str: + """This SDK appended to any stack an outer layer already declared.""" + existing = os.getenv("MEM0_CLIENT_STACK", "").strip() + entries = [part.strip() for part in existing.split(",") if part.strip()] if existing else [] + return _bounded_stack(entries, f"mem0-python/{_sdk_version()}") + + class MemoryClient: """Client for interacting with the Mem0 API. @@ -129,19 +218,11 @@ class MemoryClient: self.client = client # Ensure the client has the correct base_url and headers self.client.base_url = httpx.URL(self.host) - self.client.headers.update( - { - "Authorization": f"Token {self.api_key}", - "Mem0-User-ID": self.user_id, - } - ) + _apply_client_headers(self.client, self.api_key, self.user_id) else: self.client = httpx.Client( base_url=self.host, - headers={ - "Authorization": f"Token {self.api_key}", - "Mem0-User-ID": self.user_id, - }, + headers=_client_headers(self.api_key, self.user_id), timeout=300, ) self.user_email = self._validate_api_key() @@ -1018,19 +1099,11 @@ class AsyncMemoryClient: self.async_client = client # Ensure the client has the correct base_url and headers self.async_client.base_url = httpx.URL(self.host) - self.async_client.headers.update( - { - "Authorization": f"Token {self.api_key}", - "Mem0-User-ID": self.user_id, - } - ) + _apply_client_headers(self.async_client, self.api_key, self.user_id) else: self.async_client = httpx.AsyncClient( base_url=self.host, - headers={ - "Authorization": f"Token {self.api_key}", - "Mem0-User-ID": self.user_id, - }, + headers=_client_headers(self.api_key, self.user_id), timeout=300, ) @@ -1053,10 +1126,7 @@ class AsyncMemoryClient: params = self._prepare_params() response = requests.get( f"{self.host}/v1/ping/", - headers={ - "Authorization": f"Token {self.api_key}", - "Mem0-User-ID": self.user_id, - }, + headers=_client_headers(self.api_key, self.user_id), params=params, ) response.raise_for_status() diff --git a/tests/test_client_surface_headers.py b/tests/test_client_surface_headers.py new file mode 100644 index 000000000..51d09d306 --- /dev/null +++ b/tests/test_client_surface_headers.py @@ -0,0 +1,63 @@ +"""Surface-identity headers, and that a client can be constructed at all. + +The construction test exists because it was not there: a signature change to +_bounded_stack missed the _client_stack call site, every MemoryClient(...) raised +TypeError, and the whole suite stayed green because nothing built one. +""" + +import os +from unittest.mock import patch + +from mem0.client.main import _bounded_stack, _client_headers, _client_stack + + +def test_a_client_can_be_constructed(): + from mem0 import MemoryClient + + # _validate_api_key normally populates org/project from the API response; + # stubbing it leaves them None, which a later accessor rejects. Set them the + # way a real validation would. This test is about construction reaching the + # header stage at all. + def _stub(self): + self.org_id, self.project_id = "org", "proj" + + with patch.object(MemoryClient, "_validate_api_key", _stub): + client = MemoryClient(api_key="m0-test") + + assert client.client.headers["X-Mem0-Client"].startswith("mem0-python/") + + +# AsyncMemoryClient is deliberately not constructed here: its validation path +# makes a real request to /v1/ping/, and a unit test that needs the network is +# worse than none. It shares _client_headers with the sync client, which is the +# code the construction test above actually guards. +def test_headers_carry_this_sdk(): + headers = _client_headers("m0-test", "u1") + assert headers["X-Mem0-Client"].startswith("mem0-python/") + + +def test_our_entry_survives_a_caller_that_already_filled_the_stack(): + # Appending first and trimming to four dropped exactly the entry the + # function exists to add. + stack = _bounded_stack(["a/1", "b/2", "c/3", "d/4"], "mem0-python/9.9.9") + + assert "mem0-python/9.9.9" in stack + assert len(stack.split(",")) <= 4 + + +def test_the_character_cap_drops_whole_entries_not_characters(): + long_entries = [f"{'n' * 90}/1.0", f"{'m' * 90}/1.0", "c/3"] + stack = _bounded_stack(long_entries, "mem0-python/9.9.9") + + assert len(stack) <= 200 + assert stack.endswith("mem0-python/9.9.9") + for entry in stack.split(","): + assert entry.strip().count("/") == 1, f"severed entry: {entry!r}" + + +def test_an_outer_stack_is_appended_to_not_replaced(): + with patch.dict(os.environ, {"MEM0_CLIENT_STACK": "openclaw/2.1.0"}): + stack = _client_stack() + + assert stack.startswith("openclaw/2.1.0") + assert "mem0-python/" in stack