Merge remote-tracking branch 'origin/main' into mintlify/af82bda4
This commit is contained in:
@@ -32,6 +32,9 @@ jobs:
|
||||
- name: Type check
|
||||
run: bun run type-check
|
||||
|
||||
- name: Test
|
||||
run: bun test
|
||||
|
||||
- name: Build
|
||||
run: bun run build
|
||||
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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__,
|
||||
},
|
||||
|
||||
@@ -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/<name>-plugin/` with `plugin-build.json`, its manifest, and a thin adapter, then generate its shared runtime. Portable clients use the single `mem0-agent-plugin/` package. Independent TypeScript integrations stay self-contained and import shared lifecycle behavior from `agent-plugin-core/typescript/`.
|
||||
@@ -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`.
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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<string, unknown>[] = [];
|
||||
let timer: ReturnType<typeof setInterval> | 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<string, unknown>[]) => {
|
||||
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<void> {
|
||||
async function flush(force = false): Promise<void> {
|
||||
// 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<string, unknown> = {}): Record<string, unknown> | 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);
|
||||
|
||||
@@ -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<string, unknown>[][] = [];
|
||||
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<string, unknown>[][] = [];
|
||||
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<string, unknown>[][] = [];
|
||||
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();
|
||||
});
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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 = ""
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -148,6 +148,7 @@ const ALLOWED_KEYS = [
|
||||
"mode",
|
||||
"apiKey",
|
||||
"anonymousTelemetryId",
|
||||
"keyFingerprint",
|
||||
"baseUrl",
|
||||
"userId",
|
||||
"userEmail",
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<string> {
|
||||
const { createHash } = await import("node:crypto");
|
||||
return createHash("sha256").update(apiKey).digest("hex").slice(0, 16);
|
||||
}
|
||||
|
||||
describe("telemetry", () => {
|
||||
let fetchSpy: ReturnType<typeof vi.fn>;
|
||||
|
||||
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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).mockImplementationOnce(() => {
|
||||
throw new Error("config read failed");
|
||||
|
||||
@@ -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": ""
|
||||
},
|
||||
|
||||
@@ -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<string, unknown>;
|
||||
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<string, unknown>;
|
||||
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<string, unknown>;
|
||||
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<string, any>;
|
||||
|
||||
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<string, any>;
|
||||
const b = buildEvent("session_start", {}, "m0-account-b", PROJECT) as Record<string, any>;
|
||||
|
||||
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<string, any>;
|
||||
const second = buildEvent("session_end", {}, KEY, PROJECT) as Record<string, any>;
|
||||
|
||||
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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<string, string> {
|
||||
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<string, string> {
|
||||
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<string, unknown> | 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) });
|
||||
}
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { applySurfaceHeaders, PLATFORM_APPLICATION, PLATFORM_SOURCE } from "./attribution.ts";
|
||||
|
||||
function client(headers: Record<string, string> = {}) {
|
||||
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<string, string> }).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<string, string> }).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<string, string> }).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<string, string> }).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<string, string> }).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<string, string> }).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<string, string> }).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(/^[^/]+\/[^/]+$/);
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -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<string, string>;
|
||||
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}`);
|
||||
}
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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<Message>, 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),
|
||||
};
|
||||
|
||||
@@ -12,6 +12,7 @@
|
||||
"noUnusedLocals": false,
|
||||
"noUnusedParameters": false,
|
||||
"preserveWatchOutput": true,
|
||||
"resolveJsonModule": true,
|
||||
"skipLibCheck": true,
|
||||
"strict": true,
|
||||
"types": ["@types/node", "jest"],
|
||||
|
||||
@@ -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),
|
||||
},
|
||||
},
|
||||
])
|
||||
@@ -95,6 +95,66 @@ interface ClientIdentity {
|
||||
const IDENTITY_CACHE_MAX_DEFAULT = 50;
|
||||
const identityByCredentials = new Map<string, Promise<ClientIdentity>>();
|
||||
|
||||
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<string, string> {
|
||||
const env: Record<string, string | undefined> =
|
||||
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<string, string> = {
|
||||
"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({
|
||||
|
||||
@@ -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 {
|
||||
|
||||
+94
-24
@@ -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()
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user