Files
mem0/integrations/agent-plugin-core/build/build.py
Saket Aryan df70d7833f fix(plugins): stamp surface identity at record time, not send time
Six defects in 0.3.x plugin telemetry. Defects 1, 2 and 6 were not three bugs:
they were one spool protocol getting three properties wrong.

Identity was decided by the wrong process. `harness` was stamped in record(),
correctly, but `source` was read from a module global in flush() — so whichever
process drained the spool named every event in it. Two processes never call
init(): mcp_server.py, and the detached `python3 telemetry.py` sender that
spawn_flush() starts. record() now stamps source beside harness, and the build
generates core/_harness_id.py per host so identity resolves with no init() call
at all. That also unifies two defaults that disagreed (`<host>_plugin` vs
`MEM0_<HOST>_PLUGIN`), which could yield three source values for one plugin.

Ownership was inferred, not held. Path.replace is os.rename, which preserves
mtime, so a claim made after a quiet minute inherited the spool's age and was
stealable the instant it existed. Claims are touched on creation and the
per-batch rewrite doubles as a lease heartbeat.

Progress was not durable. flush() returned on the first failed batch without
truncating, so the retry re-posted from index 0 — 150 events delivered 250
times. It now rewrites the claim with the unsent remainder after every batch,
bounding a crash to one repeated batch, and each event carries a uuid.

Parked batches starved. They were only reachable when no spool existed, and
because sessions keep recording there usually was one, so a batch parked by a
failed send waited until the 7-day expiry deleted it unsent — despite its own
presence being what starts the sender. flush() drains them in the same run, and
expiry now applies only after a genuine retry has failed.

code.install counted upgrades and repeat sessions. is_first_run() read the
identity file, which only a successful flush writes, so an offline user recorded
an install every session forever. A dedicated install-state.json is claimed
atomically at record time; a non-empty data directory reads as an upgrade.

The docs called this anonymous. Every event carries the account email, and the
hashes were unsalted SHA-256 over a git remote URL or an absolute path
containing the username. READMEs, the module docstring and a new docs section
now say what the code does, and repo/session digests are salted per install.

A cached email outlived an API key change. It is now re-resolved when the key's
fingerprint differs, and $identify aliases anonymous->email only — aliasing one
account to another merges person profiles irreversibly.

All six shipped green because the shared core's only tests lived under one host,
behind a conftest that calls init() at import. Core behaviour was never
exercised uninitialised. Adds agent-plugin-core/tests with no init, including
subprocess tests and coverage for the portable plugin, which has no flush worker
and would pass a native-only test vacuously.

Also puts the three surface headers on the SDKs, CLIs and integrations, and
corrects a README claiming ZAPIER/STRANDS were already in the platform allowlist.

Verified: 59 core tests, 203 claude-code, 11 cursor, 5 codex, 2 kimi, 6
antigravity. ruff and compileall clean. --check clean for all six hosts.
TypeScript changes are not typechecked locally (deps not installed).

Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb
2026-09-15 00:14:11 +05:30

281 lines
11 KiB
Python

#!/usr/bin/env python3
from __future__ import annotations
import argparse
import json
import re
import shutil
import tempfile
from collections.abc import Mapping
from pathlib import Path
try:
from .validate import validate_bundle
except ImportError:
from validate import validate_bundle
CORE_ROOT = Path(__file__).resolve().parents[1]
REPOSITORY_ROOT = CORE_ROOT.parents[1]
INTEGRATIONS_ROOT = REPOSITORY_ROOT / "integrations"
SHARED_SKILLS = CORE_ROOT / "skills"
PORTABLE_PLUGIN = "mem0-agent-plugin"
NATIVE_PLUGINS = {
"claude-code": INTEGRATIONS_ROOT / "claude-code-plugin",
"cursor": INTEGRATIONS_ROOT / "cursor-plugin",
"codex": INTEGRATIONS_ROOT / "codex-plugin",
"kimi": INTEGRATIONS_ROOT / "kimi-plugin",
"antigravity": INTEGRATIONS_ROOT / "antigravity-plugin",
}
PROTECTED_OUTPUTS = {
REPOSITORY_ROOT,
INTEGRATIONS_ROOT,
CORE_ROOT,
INTEGRATIONS_ROOT / PORTABLE_PLUGIN,
*NATIVE_PLUGINS.values(),
}
TEMPLATE_TOKEN = re.compile(r"{{([A-Z_]+)}}")
TEMPLATE_TOKENS = {
"PLUGIN_ROOT",
"PLUGIN_DATA",
"PLUGIN_DATA_ARG",
"COMMAND_PREFIX",
"HARNESS_ID",
"HARNESS_NAME",
}
def render_template(source: str, values: Mapping[str, str]) -> str:
def replace(match: re.Match[str]) -> str:
token = match.group(1)
if token not in TEMPLATE_TOKENS or token not in values:
raise ValueError(f"unknown or unresolved template token: {token}")
return values[token]
rendered = TEMPLATE_TOKEN.sub(replace, source)
if "{{" in rendered or "}}" in rendered:
raise ValueError("unresolved template token")
return rendered
def replace_output(staged: Path, output: Path) -> Path:
"""Replace one explicit build output without touching its siblings."""
if not staged.is_dir():
raise ValueError(f"staged directory does not exist: {staged}")
resolved_output = output.resolve()
if resolved_output in {path.resolve() for path in PROTECTED_OUTPUTS}:
raise ValueError(f"refusing protected output path: {output}")
if output.exists() and not output.is_dir():
raise ValueError(f"output path is not a directory: {output}")
output.parent.mkdir(parents=True, exist_ok=True)
temporary = Path(tempfile.mkdtemp(prefix=f".{output.name}-", dir=output.parent))
try:
shutil.copytree(staged, temporary, dirs_exist_ok=True)
if output.exists():
shutil.rmtree(output)
temporary.replace(output)
finally:
if temporary.exists():
shutil.rmtree(temporary)
return output
def _render_harness_id(host: str) -> 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.
"""
tag = host.upper().replace("-", "_") + "_PLUGIN"
return (
'"""Generated by integrations/agent-plugin-core/build/build.py. Do not edit."""\n'
"\n"
f'HARNESS_ID = "{host}"\n'
f'SOURCE_TAG = "{tag}"\n'
"\n"
"# Platform-side vocabulary (mem0_event.source + X-Application). The whole\n"
"# plugin family is one source; which editor it runs in is the application.\n"
'PLATFORM_SOURCE = "MEM0_PLUGIN"\n'
f'PLATFORM_APPLICATION = "{host}"\n'
)
def _bundle_python(
staged: Path,
host: str,
plugin_root: str,
*,
plugin_data: str = "",
portable: bool = False,
) -> None:
core = staged / "core"
core.mkdir()
for source in sorted((CORE_ROOT / "python").glob("*.py")):
if portable and source.name in {"flush_worker.py", "hook_runner.py"}:
continue
shutil.copy2(source, core / source.name)
# Generated per host so identity does not depend on an entrypoint remembering
# to call telemetry.init(). mcp_server.py and the detached telemetry.py sender
# never did, which is how MCP searches reported harness=generic and every
# batch they drained was labelled MEM0_PLUGIN regardless of the real host.
(core / "_harness_id.py").write_text(_render_harness_id(host), encoding="utf-8")
values = {
"PLUGIN_ROOT": plugin_root,
"PLUGIN_DATA": "${PLUGIN_DATA}",
"PLUGIN_DATA_ARG": f'--plugin-data-dir "{plugin_data}"' if plugin_data else "",
"COMMAND_PREFIX": "mem0",
"HARNESS_ID": host,
"HARNESS_NAME": host.replace("-", " ").title(),
}
for source in sorted(SHARED_SKILLS.glob("*/SKILL.md.tmpl")):
target = staged / "skills" / source.parent.name / "SKILL.md"
target.parent.mkdir(parents=True)
rendered = render_template(source.read_text(encoding="utf-8"), values)
if portable:
rendered = "\n".join(
line
for line in rendered.splitlines()
if not line.startswith(("argument-hint:", "disable-model-invocation:"))
) + "\n"
target.write_text(rendered, encoding="utf-8")
def _build_portable(staged: Path) -> None:
source = INTEGRATIONS_ROOT / PORTABLE_PLUGIN
_copy_declared_files(staged, source, {"plugin.json": "plugin.json", "mcp.json": "mcp.json"})
_bundle_python(staged, "coding-agent", "${PLUGIN_ROOT}", portable=True)
def _copy_declared_files(staged: Path, source_root: Path, files: object) -> None:
if not isinstance(files, dict):
raise ValueError("native files must be an object")
source_root = source_root.resolve()
staged_root = staged.resolve()
for source_name, target_name in files.items():
if not isinstance(source_name, str) or not isinstance(target_name, str):
raise ValueError("native file paths must be strings")
source = (source_root / source_name).resolve()
target = (staged / target_name).resolve()
if not source.is_relative_to(source_root) or not target.is_relative_to(staged_root):
raise ValueError("native file paths must stay inside their roots")
if not source.is_file():
raise ValueError(f"native source file does not exist: {source_name}")
target.parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(source, target)
def _build_native(host: str, source_root: Path, staged: Path, descriptor: dict) -> None:
native = descriptor.get("native")
if not isinstance(native, dict) or not isinstance(native.get("pluginRoot"), str):
raise ValueError(f"native build is not declared for {host}")
_bundle_python(
staged,
host,
native["pluginRoot"],
plugin_data=str(native.get("pluginData") or ""),
)
_copy_declared_files(staged, source_root, native.get("files", {}))
def build(host: str, kind: str, output: Path) -> Path:
if kind not in {"portable", "native"}:
raise ValueError(f"unknown bundle kind: {kind}")
if kind == "portable":
if host != PORTABLE_PLUGIN:
raise ValueError(f"the portable bundle is {PORTABLE_PLUGIN}")
source_root = INTEGRATIONS_ROOT / PORTABLE_PLUGIN
descriptor = None
else:
source_root = NATIVE_PLUGINS.get(host)
if source_root is None:
raise ValueError(f"unknown host: {host}")
descriptor_path = source_root / "plugin-build.json"
descriptor = json.loads(descriptor_path.read_text(encoding="utf-8"))
with tempfile.TemporaryDirectory(prefix=f"mem0-{host}-{kind}-") as temporary:
staged = Path(temporary) / "bundle"
staged.mkdir()
if kind == "portable":
_build_portable(staged)
else:
assert descriptor is not None
_build_native(host, source_root, staged, descriptor)
errors = validate_bundle(staged, kind)
if errors:
raise ValueError("invalid bundle:\n" + "\n".join(errors))
return replace_output(staged, output)
def installable_root(host: str, kind: str) -> Path:
if kind == "portable" and host == PORTABLE_PLUGIN:
return INTEGRATIONS_ROOT / PORTABLE_PLUGIN
if kind == "native" and host in NATIVE_PLUGINS:
return NATIVE_PLUGINS[host]
raise ValueError(f"unknown {kind} plugin: {host}")
def bundle_drift(host: str, kind: str) -> list[str]:
target = installable_root(host, kind)
with tempfile.TemporaryDirectory(prefix=f"mem0-check-{host}-") as temporary:
generated = build(host, kind, Path(temporary) / "bundle")
errors: list[str] = []
for directory in ("core", "skills"):
expected = {
path.relative_to(generated)
for path in (generated / directory).rglob("*")
if path.is_file()
}
actual = {
path.relative_to(target)
for path in (target / directory).rglob("*")
if path.is_file() and "__pycache__" not in path.parts and path.suffix != ".pyc"
}
errors.extend(f"missing generated file: {path}" for path in sorted(expected - actual))
errors.extend(f"stale generated file: {path}" for path in sorted(actual - expected))
for source in sorted(path for path in generated.rglob("*") if path.is_file()):
relative = source.relative_to(generated)
installed = target / relative
if not installed.is_file() or source.read_bytes() != installed.read_bytes():
errors.append(f"generated file differs: {relative}")
return errors
def sync_generated(host: str, kind: str) -> Path:
target = installable_root(host, kind)
with tempfile.TemporaryDirectory(prefix=f"mem0-sync-{host}-") as temporary:
generated = build(host, kind, Path(temporary) / "bundle")
for directory in ("core", "skills"):
replace_output(generated / directory, target / directory)
return target
def main() -> int:
parser = argparse.ArgumentParser(description="Build a self-contained Mem0 agent plugin")
parser.add_argument("host")
parser.add_argument("--kind", choices=("portable", "native"), required=True)
action = parser.add_mutually_exclusive_group(required=True)
action.add_argument("--output", type=Path)
action.add_argument("--check", action="store_true")
action.add_argument("--sync", action="store_true")
args = parser.parse_args()
if args.check:
errors = bundle_drift(args.host, args.kind)
if errors:
print("\n".join(errors))
return 1
print(f"Current {args.kind} bundle: {installable_root(args.host, args.kind)}")
elif args.sync:
print(sync_generated(args.host, args.kind))
else:
assert args.output is not None
print(build(args.host, args.kind, args.output))
return 0
if __name__ == "__main__":
raise SystemExit(main())