refactor(integrations): shared agent plugin runtimes and native adapters (#7203)

This commit is contained in:
Kartik
2026-09-08 23:32:25 +05:30
committed by GitHub
parent dae67f74f5
commit 73e7b8763a
369 changed files with 31056 additions and 20464 deletions
@@ -0,0 +1,100 @@
#!/usr/bin/env python3
"""Detached remote checkpoint worker.
Claude Code may cancel SessionEnd hooks as a print-mode process exits. The hook
therefore persists its input first and launches this process in a new session.
"""
from __future__ import annotations
import json
import os
import sys
import time
from pathlib import Path
import telemetry
from memory_core import (
EvidenceStore,
checkpoint_session,
configure_harness,
touch_handoff_heartbeat,
)
def main() -> int:
if len(sys.argv) != 2:
return 2
handoff_path = Path(sys.argv[1])
os.environ["MEM0_CODE_HANDOFF_PATH"] = str(handoff_path)
harness = os.environ.get("MEM0_PLUGIN_HARNESS")
if harness:
source_tag = os.environ.get("MEM0_PLUGIN_SOURCE_TAG", "")
configure_harness(
harness,
env_prefix=os.environ.get("MEM0_PLUGIN_ENV_PREFIX", ""),
data_dir_name=os.environ.get("MEM0_PLUGIN_DATA_DIR_NAME", ""),
source_tag=source_tag,
)
telemetry.init(harness=harness, source_tag=source_tag.upper())
completed = False
try:
payload = json.loads(handoff_path.read_text(encoding="utf-8"))
delay = float(payload.get("delay_seconds") or 0)
if delay > 0:
payload.pop("delay_seconds", None)
temporary = handoff_path.with_suffix(f".{os.getpid()}.tmp")
try:
temporary.write_text(json.dumps(payload), encoding="utf-8")
temporary.replace(handoff_path)
finally:
temporary.unlink(missing_ok=True)
time.sleep(delay)
if not handoff_path.exists():
return 0
hook_input = payload.get("hook_input") or {}
reason = str(payload.get("reason") or "checkpoint")
wait_for_inflight = bool(payload.get("wait_for_inflight"))
store = EvidenceStore()
try:
if wait_for_inflight:
session_id = str(
hook_input.get("session_id") or "unknown-session"
)
repo = store.repo_for_session(session_id, hook_input.get("cwd"))
deadline = time.monotonic() + float(
os.environ.get("MEM0_CODE_EXTRACTION_WAIT_SECONDS", "120")
)
while (
store.has_inflight_flush(repo.identity, session_id)
and time.monotonic() < deadline
):
touch_handoff_heartbeat()
time.sleep(0.25)
# Hooks capture the conversation before handoff; the worker only flushes it.
result = checkpoint_session(store, hook_input, reason)
print(json.dumps(result, sort_keys=True), flush=True)
completed = result.get("status") in {
"semantic-succeeded",
"explicitly-stored",
"nothing-to-flush",
}
finally:
store.close()
return 0
finally:
telemetry.flush()
if completed:
try:
handoff_path.unlink()
except OSError:
pass
elif handoff_path.suffix == ".running":
try:
handoff_path.replace(handoff_path.with_suffix(".json"))
except OSError:
pass
if __name__ == "__main__":
raise SystemExit(main())
@@ -0,0 +1,372 @@
"""Shared hook orchestration for all Mem0 agent plugins."""
from __future__ import annotations
import argparse
import hashlib
import json
import os
import subprocess
import sys
import time
import uuid
from pathlib import Path
import telemetry
from memory_core import (
EvidenceStore,
_session_id,
api_key,
bounded,
cache_plugin_api_key,
checkpoint_session,
clear_stale_api_key_cache,
configure_harness,
data_dir,
detached_process_kwargs,
format_context,
harness_config,
record_session_start,
record_tool,
record_user_prompt,
redact,
search_memories,
)
STALE_RUNNING_SECONDS = 300
PENDING_EXPIRY_SECONDS = 7 * 24 * 60 * 60
PENDING_LAUNCH_LIMIT = 5
DEFAULT_IDLE_FLUSH_SECONDS = 300
_core_dir: Path = Path(__file__).resolve().parent
def read_hook_input() -> dict:
try:
value = json.load(sys.stdin)
return value if isinstance(value, dict) else {}
except (json.JSONDecodeError, OSError):
return {}
def default_record_stop(store: EvidenceStore, hook_input: dict):
"""Record the assistant's response without transcript parsing."""
session_id = _session_id(hook_input)
repo = store.repo_for_session(session_id, hook_input.get("cwd"))
message = redact(hook_input.get("last_assistant_message", "")).strip()
if message:
store.record_assistant_response(repo, session_id, message)
return repo, session_id
def first_prompt_memory_output(store: EvidenceStore, hook_input: dict) -> dict:
"""Search once before the agent handles the first prompt in a session."""
repo, session_id, prompt, is_first_prompt = record_user_prompt(store, hook_input)
if not is_first_prompt:
return {}
try:
minimum_query_chars = int(os.environ.get("MEM0_CODE_MIN_QUERY_CHARS", "20"))
except ValueError:
minimum_query_chars = 20
if len(prompt.strip()) < max(minimum_query_chars, 1):
return {}
result = search_memories(
store, repo, session_id, bounded(prompt, 6000),
top_k=5, operation="first-prompt-search", timeout=2,
)
if not result.memories:
return {}
context = format_context(
result.memories,
"Mem0 found these relevant memories from earlier work in this repository:",
)
telemetry.record(
"context_injected",
repo=repo, session_id=session_id, trigger="first-prompt",
memory_count=len(result.memories), context_chars=len(context),
prompt_chars=len(prompt),
)
return {
"hookSpecificOutput": {
"hookEventName": "UserPromptSubmit",
"additionalContext": context,
},
}
def _launch_handoff(handoff_path: Path) -> bool:
running_path = handoff_path.with_suffix(".running")
try:
handoff_path.replace(running_path)
except OSError:
return False
worker = _core_dir / "flush_worker.py"
log_path = data_dir() / "flush-worker.log"
log_handle = open(log_path, "a", encoding="utf-8")
harness = harness_config()
child_env = os.environ.copy()
child_env.update(
{
"MEM0_CODE_DATA_DIR": str(data_dir()),
"MEM0_PLUGIN_HARNESS": harness["name"],
"MEM0_PLUGIN_ENV_PREFIX": harness["env_prefix"],
"MEM0_PLUGIN_DATA_DIR_NAME": harness["data_dir_name"],
"MEM0_PLUGIN_SOURCE_TAG": harness["source_tag"],
}
)
try:
subprocess.Popen(
[sys.executable, str(worker), str(running_path)],
stdin=subprocess.DEVNULL,
stdout=log_handle, stderr=log_handle,
close_fds=True,
env=child_env,
**detached_process_kwargs(),
)
finally:
log_handle.close()
return True
def recover_pending_handoffs() -> int:
pending_dir = data_dir() / "pending"
pending_dir.mkdir(parents=True, exist_ok=True)
now = time.time()
for running in pending_dir.glob("*.running"):
try:
if now - running.stat().st_mtime > STALE_RUNNING_SECONDS:
running.replace(running.with_suffix(".json"))
except OSError:
continue
recoverable = []
for handoff in pending_dir.glob("*.json"):
try:
age = now - handoff.stat().st_mtime
except OSError:
continue
if age > PENDING_EXPIRY_SECONDS:
handoff.unlink(missing_ok=True)
continue
recoverable.append((age, handoff))
recoverable.sort(key=lambda item: item[0], reverse=True)
launched = 0
for _, handoff in recoverable[:PENDING_LAUNCH_LIMIT]:
launched += int(_launch_handoff(handoff))
return launched
def refresh_pending_handoffs() -> None:
pending_dir = data_dir() / "pending"
if not pending_dir.is_dir():
return
for pattern in ("*.json", "*.running"):
for handoff in pending_dir.glob(pattern):
try:
os.utime(handoff)
except OSError:
continue
def hand_off_flush(
hook_input: dict, reason: str, *, wait_for_inflight: bool = False,
) -> None:
pending_dir = data_dir() / "pending"
pending_dir.mkdir(parents=True, exist_ok=True)
material = (
f"{hook_input.get('cwd', '')}\0{hook_input.get('session_id', '')}\0{reason}"
)
digest = hashlib.sha256(material.encode()).hexdigest()[:24]
handoff_path = pending_dir / f"{digest}-{uuid.uuid4().hex[:8]}.json"
temporary_path = handoff_path.with_suffix(".tmp")
temporary_path.write_text(
json.dumps({
"hook_input": hook_input,
"reason": reason,
"wait_for_inflight": wait_for_inflight,
}),
encoding="utf-8",
)
temporary_path.replace(handoff_path)
_launch_handoff(handoff_path)
def automatic_flush_enabled() -> bool:
return os.environ.get("MEM0_CODE_AUTO_FLUSH", "true").lower() in {
"1", "true", "yes", "on",
}
def schedule_periodic_checkpoint(
store: EvidenceStore, hook_input: dict, repo, session_id: str,
) -> bool:
if (
not automatic_flush_enabled()
or not api_key()
or not store.checkpoint_due(repo.identity, session_id)
):
return False
if store.prepare_flush(repo, session_id, "periodic") is None:
return False
hand_off_flush(hook_input, "periodic")
return True
def _idle_flush_seconds() -> int:
try:
return max(
int(os.environ.get("MEM0_CODE_IDLE_FLUSH_SECONDS", str(DEFAULT_IDLE_FLUSH_SECONDS))),
0,
)
except ValueError:
return DEFAULT_IDLE_FLUSH_SECONDS
def schedule_idle_flush(
store: EvidenceStore, hook_input: dict, repo, session_id: str,
) -> bool:
delay = _idle_flush_seconds()
if delay <= 0 or not automatic_flush_enabled() or not api_key():
return False
if store.has_inflight_flush(repo.identity, session_id):
return False
if not store.has_unflushed_events(repo.identity, session_id):
return False
pending_dir = data_dir() / "pending"
pending_dir.mkdir(parents=True, exist_ok=True)
material = f"idle\0{hook_input.get('cwd', '')}\0{hook_input.get('session_id', '')}"
digest = hashlib.sha256(material.encode()).hexdigest()[:24]
for old in pending_dir.glob(f"idle-{digest}*"):
old.unlink(missing_ok=True)
handoff_path = pending_dir / f"idle-{digest}-{uuid.uuid4().hex[:8]}.json"
temporary_path = handoff_path.with_suffix(".tmp")
temporary_path.write_text(
json.dumps({
"hook_input": hook_input,
"reason": "idle",
"delay_seconds": delay,
}),
encoding="utf-8",
)
temporary_path.replace(handoff_path)
_launch_handoff(handoff_path)
return True
def log_failure(exc: Exception) -> None:
try:
log_path = data_dir() / "plugin-errors.log"
with log_path.open("a", encoding="utf-8") as handle:
handle.write(f"{time.time():.3f} {type(exc).__name__}: {exc}\n")
except OSError:
pass
def run(
*,
record_stop_fn=None,
extra_actions: dict | None = None,
data_dir_env: str = "MEM0_PLUGIN_DATA_DIR",
automatic_flush_reasons: set | None = None,
) -> int:
if record_stop_fn is None:
record_stop_fn = default_record_stop
if automatic_flush_reasons is None:
automatic_flush_reasons = {"session-end"}
base_actions = ["session-start", "user-prompt", "post-tool", "stop", "flush"]
all_actions = base_actions + list((extra_actions or {}).keys())
parser = argparse.ArgumentParser()
parser.add_argument("action", choices=all_actions)
parser.add_argument("--reason", default="manual")
parser.add_argument("--plugin-data-dir", default="")
parser.add_argument("--harness", default="")
args = parser.parse_args()
if args.harness:
configure_harness(args.harness)
telemetry.init(harness=args.harness)
if args.plugin_data_dir:
os.environ[data_dir_env] = args.plugin_data_dir
cache_plugin_api_key()
if args.action == "session-start":
clear_stale_api_key_cache()
hook_input = read_hook_input()
store = EvidenceStore()
try:
if store.is_paused():
if args.action == "session-start":
refresh_pending_handoffs()
telemetry.record("session_start", paused=True)
telemetry.spawn_flush()
return 0
if args.action == "session-start":
if telemetry.is_first_run():
telemetry.record("install")
recovered = recover_pending_handoffs()
record_session_start(store, hook_input)
if recovered:
telemetry.record("handoff_recovered", count=recovered)
telemetry.spawn_flush()
elif args.action == "user-prompt":
output = first_prompt_memory_output(store, hook_input)
if output:
print(json.dumps(output))
elif args.action == "post-tool":
record_tool(store, hook_input)
elif args.action == "stop":
repo, session_id = record_stop_fn(store, hook_input)
if not schedule_periodic_checkpoint(store, hook_input, repo, session_id):
schedule_idle_flush(store, hook_input, repo, session_id)
elif args.action == "flush":
automatic = args.reason in automatic_flush_reasons
if automatic and not automatic_flush_enabled():
return 0
if args.reason == "session-end":
record_stop_fn(store, hook_input)
if os.environ.get("MEM0_CODE_SYNC_FLUSH") == "1":
print(json.dumps(checkpoint_session(store, hook_input, args.reason)))
else:
session_id = str(hook_input.get("session_id") or "unknown-session")
repo = store.repo_for_session(session_id, hook_input.get("cwd"))
already_running = store.has_inflight_flush(repo.identity, session_id)
if already_running and args.reason == "session-end":
hand_off_flush(hook_input, args.reason, wait_for_inflight=True)
elif not already_running and store.prepare_flush(
repo, session_id, args.reason,
) is not None:
hand_off_flush(hook_input, args.reason)
elif extra_actions and args.action in extra_actions:
result = extra_actions[args.action](store, hook_input)
if result:
print(json.dumps(result))
finally:
store.close()
return 0
def entry_point(
*,
record_stop_fn=None,
extra_actions: dict | None = None,
data_dir_env: str = "MEM0_PLUGIN_DATA_DIR",
automatic_flush_reasons: set | None = None,
) -> None:
try:
raise SystemExit(run(
record_stop_fn=record_stop_fn,
extra_actions=extra_actions,
data_dir_env=data_dir_env,
automatic_flush_reasons=automatic_flush_reasons,
))
except Exception as exc:
log_failure(exc)
raise SystemExit(0)
if __name__ == "__main__":
entry_point()
@@ -0,0 +1,247 @@
#!/usr/bin/env python3
"""Expose Mem0's memory search as one local coding-agent tool."""
from __future__ import annotations
import json
import os
import sys
from typing import Any
import telemetry
from memory_core import (
CODING_MEMORY_CATEGORY_NAMES,
PLUGIN_VERSION,
SEARCH_SCOPES,
format_search_result,
resolve_repo,
search_memories,
)
PROTOCOL_VERSION = "2024-11-05"
TOOL_NAME = "search_memories"
TOOL_DESCRIPTION = (
"Search memories from earlier work in this repository. ALWAYS call this "
"tool before answering anything that could depend on prior context: the "
"user's preferences, facts about this codebase, history, people, projects, "
"or earlier decisions. Do not rely on the chat window alone. The "
"repository's memory is shared by everyone who works in it and includes "
"what it took to run, test, or build here, so search before assuming an "
"invocation works. The scope argument changes what is searched: 'repo' "
"(default) is the whole repository's shared memory plus your own "
"preferences, 'dir' narrows the shared part to the directory you are "
"working in, and 'mine' is your preferences alone."
)
TOOL_SCHEMA = {
"type": "object",
"properties": {
"query": {
"type": "string",
"minLength": 1,
"maxLength": 2000,
"description": "A direct question about earlier work in this repository.",
},
"top_k": {
"type": "integer",
"minimum": 1,
"maximum": 20,
"description": "Maximum memories to return. Uses Mem0's configured default when omitted.",
},
"category": {
"type": "string",
"enum": list(CODING_MEMORY_CATEGORY_NAMES),
"description": "Optional memory category. Omit to search every category.",
},
"scope": {
"type": "string",
"enum": list(SEARCH_SCOPES),
"description": (
"Which memories to search. 'repo' (default) is the whole repository's "
"shared memory plus your own preferences, 'dir' narrows the shared "
"part to the current directory, 'mine' is your preferences alone."
),
},
"run_id": {
"type": "string",
"minLength": 1,
"description": (
"Optional coding-agent session ID. With any scope, restricts results to memories "
"saved in that session. Omit to recall memories across sessions."
),
},
},
"required": ["query"],
"additionalProperties": False,
}
class ToolInputError(ValueError):
pass
def _validate_arguments(
arguments: Any,
) -> tuple[str, int | None, str | None, str | None, str | None]:
if not isinstance(arguments, dict):
raise ToolInputError("Search arguments must be an object.")
unknown = set(arguments) - {"query", "top_k", "category", "scope", "run_id"}
if unknown:
raise ToolInputError(f"Unknown search argument: {sorted(unknown)[0]}")
query = arguments.get("query")
if not isinstance(query, str) or not query.strip():
raise ToolInputError("query must be a non-empty string.")
query = query.strip()
if len(query) > 2000:
raise ToolInputError("query must be at most 2,000 characters.")
top_k = arguments.get("top_k")
if top_k is not None and (
isinstance(top_k, bool) or not isinstance(top_k, int) or not 1 <= top_k <= 20
):
raise ToolInputError("top_k must be an integer from 1 to 20.")
category = arguments.get("category")
if category is not None and category not in CODING_MEMORY_CATEGORY_NAMES:
raise ToolInputError("category must be one of Mem0's supported categories.")
scope = arguments.get("scope")
if scope is not None and scope not in SEARCH_SCOPES:
raise ToolInputError(f"scope must be one of {list(SEARCH_SCOPES)}.")
run_id = arguments.get("run_id")
if run_id is not None:
if not isinstance(run_id, str) or not run_id.strip():
raise ToolInputError("run_id must be a non-empty string.")
run_id = run_id.strip()
return query, top_k, category, scope, run_id
def call_search_memories(arguments: Any, cwd: str | None = None) -> str:
query, top_k, category, scope, run_id = _validate_arguments(arguments)
repo = resolve_repo(cwd or os.environ.get("CLAUDE_PROJECT_DIR") or os.getcwd())
result = search_memories(
None,
repo,
None,
query,
top_k=top_k,
category=category,
scope=scope,
run_id=run_id,
operation="mcp-search",
)
return format_search_result(result)
def _workspace_cwd(params: dict[str, Any]) -> str | None:
meta = params.get("_meta")
if not isinstance(meta, dict):
return None
metadata = meta.get("x-codex-turn-metadata")
if not isinstance(metadata, dict):
return None
workspaces = metadata.get("workspaces") or {}
if isinstance(workspaces, dict):
return next((path for path in workspaces if isinstance(path, str) and path), None)
return None
def _tool_response(text: str, *, is_error: bool = False) -> dict[str, Any]:
return {
"content": [{"type": "text", "text": text}],
"isError": is_error,
}
def handle_request(message: Any) -> dict[str, Any] | None:
if not isinstance(message, dict):
return None
request_id = message.get("id")
method = message.get("method")
if method == "notifications/initialized":
return None
if method == "initialize":
requested = (message.get("params") or {}).get("protocolVersion")
return {
"jsonrpc": "2.0",
"id": request_id,
"result": {
"protocolVersion": requested or PROTOCOL_VERSION,
"capabilities": {"tools": {"listChanged": False}},
"serverInfo": {"name": "mem0", "version": PLUGIN_VERSION},
},
}
if method == "ping":
return {"jsonrpc": "2.0", "id": request_id, "result": {}}
if method == "tools/list":
return {
"jsonrpc": "2.0",
"id": request_id,
"result": {
"tools": [
{
"name": TOOL_NAME,
"description": TOOL_DESCRIPTION,
"inputSchema": TOOL_SCHEMA,
"annotations": {
"readOnlyHint": True,
"idempotentHint": True,
"openWorldHint": True,
},
}
]
},
}
if method == "tools/call":
params = message.get("params") or {}
if params.get("name") != TOOL_NAME:
result = _tool_response("Unknown Mem0 tool.", is_error=True)
else:
try:
result = _tool_response(
call_search_memories(params.get("arguments"), _workspace_cwd(params))
)
except ToolInputError as exc:
result = _tool_response(str(exc), is_error=True)
except Exception:
result = _tool_response("Memory search failed.", is_error=True)
return {"jsonrpc": "2.0", "id": request_id, "result": result}
if request_id is None:
return None
return {
"jsonrpc": "2.0",
"id": request_id,
"error": {"code": -32601, "message": "Method not found"},
}
def main() -> int:
for raw_line in sys.stdin:
try:
message = json.loads(raw_line)
response = handle_request(message)
except json.JSONDecodeError:
response = {
"jsonrpc": "2.0",
"id": None,
"error": {"code": -32700, "message": "Parse error"},
}
except Exception:
response = {
"jsonrpc": "2.0",
"id": None,
"error": {"code": -32603, "message": "Internal error"},
}
if response is not None:
sys.stdout.write(json.dumps(response, separators=(",", ":")) + "\n")
sys.stdout.flush()
telemetry.spawn_flush()
return 0
if __name__ == "__main__":
raise SystemExit(main())
@@ -0,0 +1,154 @@
#!/usr/bin/env python3
"""Mem0 diagnostics and user controls."""
from __future__ import annotations
import argparse
import json
import os
import telemetry
from memory_core import (
EvidenceStore,
api_key,
data_dir,
doctor,
forget_remote_repo,
configure_harness,
resolve_repo,
user_id,
)
def _print_status(value: dict) -> None:
last = value.get("last_operation") or {}
print(f"Mem0: {'paused' if value['paused'] else 'active'}")
print(f"Repository: {value['repo_id']}")
print(f"Local data: {value['data_dir']}")
print(f"API key: {'configured' if value['api_key_configured'] else 'missing'}")
print(
"Saved on this computer: "
f"{value['events']} session details, {value['flushes']} memory updates"
)
print(
f"Used in this repository: {value['retrievals']} memories returned, "
f"{value['sidekick_runs']} sidekick runs"
)
if last:
item_label = ""
if last["operation"] in {"flush", "flush-retry"}:
item_label = f", {last['item_count']} memories"
operation = (
"memory update"
if last["operation"] in {"flush", "flush-retry"}
else last["operation"].replace("-", " ")
)
print(
f"Last {operation}: "
f"{'succeeded' if last['success'] else 'failed'} "
f"({last['duration_ms']:.1f} ms{item_label})"
)
sidekick = value.get("last_sidekick") or {}
if sidekick:
state = "finished" if sidekick.get("stopped_at") else "started"
print(
"Last sidekick: "
f"{state}, received {sidekick['context_chars']} characters of memory, "
f"agent {sidekick['agent_id']}"
)
def main() -> int:
parser = argparse.ArgumentParser()
parser.add_argument("--plugin-data-dir", default="")
parser.add_argument("--harness", default="")
subparsers = parser.add_subparsers(dest="command", required=True)
status = subparsers.add_parser("status")
status.add_argument("--json", action="store_true")
doctor_parser = subparsers.add_parser("doctor")
doctor_parser.add_argument("--json", action="store_true")
subparsers.add_parser("pause")
subparsers.add_parser("resume")
forget = subparsers.add_parser("forget")
forget.add_argument("--remote", action="store_true")
forget.add_argument("--yes", action="store_true")
forget.add_argument("--include-project-memory", action="store_true")
args = parser.parse_args()
if args.harness:
source_tag = f"{args.harness.replace('-', '_')}_plugin"
configure_harness(args.harness, source_tag=source_tag)
telemetry.init(harness=args.harness, source_tag=source_tag.upper())
if args.plugin_data_dir:
os.environ["MEM0_CODE_DATA_DIR"] = args.plugin_data_dir
store = EvidenceStore()
try:
repo = resolve_repo(os.getcwd())
telemetry.record("control", repo=repo, action=args.command)
if args.command == "status":
result = {
**store.status(repo.identity),
"repo_id": repo.identity,
"app_id": repo.app_id,
"project_id": repo.project_id,
"directory": repo.directory,
"user_id": user_id(),
"data_dir": str(data_dir()),
"api_key_configured": bool(api_key()),
}
if args.json:
print(json.dumps(result, indent=2, default=str))
else:
_print_status(result)
elif args.command == "doctor":
result = doctor(os.getcwd())
if args.json:
print(json.dumps(result, indent=2, default=str))
else:
for name, check in result["checks"].items():
print(
f"{'PASS' if check['ok'] else 'FAIL'} {name}: {check['detail']}"
)
return 0 if result["ok"] else 1
elif args.command == "pause":
store.set_setting("paused", "true")
print("Mem0 stopped saving and searching memories.")
elif args.command == "resume":
store.set_setting("paused", "false")
print("Mem0 resumed saving and searching memories.")
elif args.command == "forget":
if not args.yes:
print(
"Refusing to delete data without --yes. Add --remote to also "
"delete this user/repository scope from Mem0."
)
return 2
remote_result = (
forget_remote_repo(
repo, include_project_memory=args.include_project_memory
)
if args.remote
else None
)
local_result = store.forget_local_repo(repo.identity)
print(
json.dumps(
{"local": local_result, "remote": remote_result},
indent=2,
default=str,
)
)
if remote_result and remote_result.get("status") == "error":
return 1
finally:
store.close()
telemetry.spawn_flush()
return 0
if __name__ == "__main__":
raise SystemExit(main())
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,397 @@
#!/usr/bin/env python3
"""Anonymous usage telemetry for Mem0 agent plugins.
Hooks run on a 3-6 second budget and fire on every tool call, so recording never
touches the network: `record` appends one JSON line to a local spool and returns.
A detached `python3 telemetry.py` drains the spool in one batched PostHog request,
started once per session and again from the flush worker that is already detached.
Pure stdlib, matching the rest of the plugin. Opt out with MEM0_TELEMETRY=false.
Never sends prompts, memory text, queries, file paths, repository names, or API
keys: only event names, durations, counts, coarse outcomes, and salted hashes.
"""
from __future__ import annotations
import hashlib
import json
import os
import platform
import subprocess
import sys
import time
import urllib.error
import urllib.request
import uuid
from pathlib import Path
from typing import Any
import memory_core
_harness: str = "generic"
_source_tag: str = "MEM0_PLUGIN"
_PRIVATE_KEYS = {
"apikey",
"authorization",
"password",
"query",
"secret",
"prompt",
"token",
"text",
"memory",
"message",
"error",
"path",
"cwd",
"userid",
"agentid",
"runid",
"repoid",
"repositoryid",
"projectid",
"appid",
"filters",
}
def init(harness: str = "generic", source_tag: str = "") -> None:
global _harness, _source_tag
_harness = harness
_source_tag = source_tag or f"MEM0_{harness.upper().replace('-', '_')}_PLUGIN"
POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"
POSTHOG_CAPTURE_URL = "https://us.i.posthog.com/i/v0/e/"
POSTHOG_BATCH_URL = "https://us.i.posthog.com/batch/"
EVENT_PREFIX = "code"
SPOOL_LIMIT_BYTES = 256 * 1024
BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
def is_enabled() -> bool:
"""Whether telemetry is switched on for this process."""
return os.environ.get("MEM0_TELEMETRY", "true").strip().lower() not in {
"false",
"0",
"no",
"off",
}
def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _safe_value(value: Any) -> Any:
if isinstance(value, str):
return memory_core.redact(value)
if isinstance(value, dict):
return {
key: _safe_value(item)
for key, item in value.items()
if "".join(character for character in str(key).lower() if character.isalnum())
not in _PRIVATE_KEYS
}
if isinstance(value, (list, tuple)):
return [_safe_value(item) for item in value]
if value is None or isinstance(value, (bool, int, float)):
return value
return memory_core.redact(value)
def _spool_path() -> Path:
return memory_core.data_dir() / "telemetry.jsonl"
def _identity_path() -> Path:
return memory_core.data_dir() / "telemetry-identity.json"
def _read_identity() -> dict[str, str]:
try:
value = json.loads(_identity_path().read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return {}
return value if isinstance(value, dict) else {}
def _write_identity(identity: dict[str, str]) -> None:
path = _identity_path()
temporary = path.with_suffix(f".{os.getpid()}.tmp")
try:
path.parent.mkdir(parents=True, exist_ok=True)
temporary.write_text(json.dumps(identity), encoding="utf-8")
temporary.replace(path)
except OSError:
try:
temporary.unlink()
except OSError:
pass
def anonymous_id(identity: dict[str, str] | None = None) -> str:
"""Per-machine anonymous identifier, created and persisted on first use."""
identity = _read_identity() if identity is None else identity
existing = identity.get("anonymous_id")
if existing:
return existing
created = f"code-anon-{uuid.uuid4().hex}"
identity["anonymous_id"] = created
_write_identity(identity)
return created
def is_first_run() -> bool:
"""Whether this machine has never recorded a plugin event before."""
return not _identity_path().exists()
def record(
event: str,
*,
repo: Any = None,
session_id: str | None = None,
**properties: Any,
) -> None:
"""Append one event to the local spool. Never blocks and never raises."""
if not is_enabled():
return
try:
spool = _spool_path()
try:
if spool.stat().st_size > SPOOL_LIMIT_BYTES:
return
except OSError:
pass
properties = _safe_value(properties)
properties.update(
harness=_harness,
plugin_version=memory_core.PLUGIN_VERSION,
os=sys.platform,
python_version=platform.python_version(),
)
if repo is not None:
properties["repo_hash"] = _digest(getattr(repo, "identity", ""))
if session_id:
properties["session_hash"] = _digest(session_id)
line = json.dumps(
{
"event": f"{EVENT_PREFIX}.{event}",
"timestamp": memory_core.utc_now(),
"properties": {
key: value for key, value in properties.items() if value is not None
},
},
separators=(",", ":"),
default=str,
)
spool.parent.mkdir(parents=True, exist_ok=True)
with spool.open("a", encoding="utf-8") as handle:
handle.write(line + "\n")
except Exception:
pass
def error_kind(exc: BaseException | str) -> str:
"""Coarse, content-free label for a failure, safe to send."""
text = exc if isinstance(exc, str) else f"{type(exc).__name__}: {exc}"
lowered = text.lower()
if "timed out" in lowered or "timeout" in lowered:
return "timeout"
if "401" in lowered or "403" in lowered or "unauthor" in lowered or "forbidden" in lowered:
return "auth"
if "429" in lowered or "rate limit" in lowered:
return "rate-limited"
if any(code in lowered for code in ("500", "502", "503", "504")):
return "server-error"
if "400" in lowered or "422" in lowered:
return "bad-request"
if isinstance(exc, str):
return "other"
if isinstance(exc, urllib.error.URLError):
return "network"
return type(exc).__name__
def spawn_flush() -> bool:
"""Start the detached sender that drains the spool."""
if not is_enabled():
return False
try:
if not _spool_path().exists() and not any(
memory_core.data_dir().glob("telemetry-*.sending")
):
return False
subprocess.Popen(
[sys.executable, str(Path(__file__).resolve())],
stdin=subprocess.DEVNULL,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
close_fds=True,
**memory_core.detached_process_kwargs(),
)
return True
except Exception:
return False
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
spool = _spool_path()
try:
spool.replace(claim)
return claim
except OSError:
pass
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
try:
orphan.replace(claim)
return claim
except OSError:
continue
return None
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
request = urllib.request.Request(
url, headers={"Authorization": f"Token {key}", "Content-Type": "application/json"}
)
try:
with urllib.request.urlopen(request, timeout=SEND_TIMEOUT) as response:
payload = json.loads(response.read().decode("utf-8"))
except Exception:
return ""
email = payload.get("user_email") if isinstance(payload, dict) else ""
return email if isinstance(email, str) else ""
def _post(payload: dict[str, Any], url: str) -> bool:
request = urllib.request.Request(
url,
data=json.dumps(payload, default=str).encode("utf-8"),
headers={"Content-Type": "application/json"},
)
try:
with urllib.request.urlopen(request, timeout=SEND_TIMEOUT):
return True
except Exception:
return False
def resolve_distinct_id() -> tuple[str, str]:
"""Return the PostHog distinct id and the anonymous id it replaced, if any."""
identity = _read_identity()
email = identity.get("email", "")
if email:
return email, ""
key = memory_core.api_key()
if not key:
return anonymous_id(identity), ""
email = _resolve_email(key)
if not email:
return anonymous_id(identity), ""
previous = identity.get("anonymous_id", "")
identity["email"] = email
_write_identity(identity)
return email, previous
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
if claim is None:
return 0
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
return 0
events = []
for line in lines:
try:
value = json.loads(line)
except json.JSONDecodeError:
continue
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
try:
claim.unlink()
except OSError:
pass
return 0
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
_post(
{
"api_key": POSTHOG_API_KEY,
"event": "$identify",
"distinct_id": distinct_id,
"properties": {
"$anon_distinct_id": aliased_anonymous_id,
"$lib": "posthog-python",
},
},
POSTHOG_CAPTURE_URL,
)
sent = 0
for start in range(0, len(events), BATCH_SIZE):
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
"timestamp": event.get("timestamp"),
"properties": {
"source": _source_tag,
"language": "python",
"$process_person_profile": False,
"$lib": "posthog-python",
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
def main() -> int:
flush()
return 0
if __name__ == "__main__":
try:
raise SystemExit(main())
except Exception:
raise SystemExit(0)