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
This commit is contained in:
@@ -305,8 +305,19 @@ def run(
|
||||
return 0
|
||||
|
||||
if args.action == "session-start":
|
||||
if telemetry.is_first_run():
|
||||
# Claims the marker atomically and says which event to record, so a
|
||||
# second session starting alongside this one cannot record it too.
|
||||
first_event = telemetry.claim_install()
|
||||
if first_event == "install":
|
||||
telemetry.record("install")
|
||||
elif first_event == "upgrade":
|
||||
# First run after a build that never wrote the marker; the
|
||||
# predecessor version was never recorded anywhere.
|
||||
telemetry.record("upgrade", from_version="pre-0.3")
|
||||
else:
|
||||
previous = telemetry.claim_version_change()
|
||||
if previous:
|
||||
telemetry.record("upgrade", from_version=previous)
|
||||
recovered = recover_pending_handoffs()
|
||||
record_session_start(store, hook_input)
|
||||
if recovered:
|
||||
|
||||
@@ -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,10 @@ 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,
|
||||
# query string or X-Mem0-Source header, never from metadata. 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 +2555,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:
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Anonymous usage telemetry for Mem0 agent plugins.
|
||||
"""Usage telemetry for Mem0 agent plugins.
|
||||
|
||||
Events are linked to your Mem0 account email when an API key is configured, and
|
||||
to a random per-machine id otherwise. Not anonymous — the Python SDK and CLI
|
||||
attribute the same way.
|
||||
|
||||
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.
|
||||
@@ -9,7 +13,8 @@ started once per session and again from the flush worker that is already detache
|
||||
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.
|
||||
keys: only event names, durations, counts, coarse outcomes, and repo/session
|
||||
identifiers hashed with a random per-install salt.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -29,8 +34,23 @@ from typing import Any
|
||||
|
||||
import memory_core
|
||||
|
||||
_harness: str = "generic"
|
||||
_source_tag: str = "MEM0_PLUGIN"
|
||||
# Seeded from the per-host module the build generates into core/. Two processes
|
||||
# in this pipeline never call init() — mcp_server.py, and the detached
|
||||
# `python3 telemetry.py` sender that spawn_flush() starts — so a module default
|
||||
# was what every one of their events got labelled with.
|
||||
try: # pragma: no cover - absent only in the un-built shared source tree
|
||||
from _harness_id import HARNESS_ID as _DEFAULT_HARNESS
|
||||
from _harness_id import PLATFORM_APPLICATION as _PLATFORM_APPLICATION
|
||||
from _harness_id import PLATFORM_SOURCE as _PLATFORM_SOURCE
|
||||
from _harness_id import SOURCE_TAG as _DEFAULT_SOURCE_TAG
|
||||
except ImportError:
|
||||
_DEFAULT_HARNESS = "generic"
|
||||
_DEFAULT_SOURCE_TAG = "MEM0_PLUGIN"
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
"apikey",
|
||||
"authorization",
|
||||
@@ -56,10 +76,19 @@ _PRIVATE_KEYS = {
|
||||
}
|
||||
|
||||
|
||||
def init(harness: str = "generic", source_tag: str = "") -> None:
|
||||
def init(harness: str = "", source_tag: str = "") -> None:
|
||||
"""Override the generated identity. Optional — core/_harness_id.py is the default.
|
||||
|
||||
The fallback shape matches memory_core.configure_harness's (``<HOST>_PLUGIN``).
|
||||
It used to be ``MEM0_<HOST>_PLUGIN`` here and ``<host>_plugin`` there, which
|
||||
meant one plugin could emit three different source values depending on which
|
||||
process happened to send the batch.
|
||||
"""
|
||||
global _harness, _source_tag
|
||||
_harness = harness
|
||||
_source_tag = source_tag or f"MEM0_{harness.upper().replace('-', '_')}_PLUGIN"
|
||||
_harness = harness or _DEFAULT_HARNESS
|
||||
_source_tag = source_tag or (
|
||||
f"{_harness.upper().replace('-', '_')}_PLUGIN" if harness else _DEFAULT_SOURCE_TAG
|
||||
)
|
||||
|
||||
POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"
|
||||
POSTHOG_CAPTURE_URL = "https://us.i.posthog.com/i/v0/e/"
|
||||
@@ -70,6 +99,11 @@ BATCH_SIZE = 100
|
||||
SEND_TIMEOUT = 5
|
||||
CLAIM_STALE_SECONDS = 120
|
||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
||||
# A batch is only discarded once it has genuinely been retried this many times.
|
||||
MAX_CLAIM_ATTEMPTS = 3
|
||||
# Parked claims drained per run, after the live spool. Bounded so a long backlog
|
||||
# cannot turn one flush into an unbounded send loop.
|
||||
MAX_PARKED_PER_RUN = 3
|
||||
|
||||
|
||||
def is_enabled() -> bool:
|
||||
@@ -83,9 +117,36 @@ def is_enabled() -> bool:
|
||||
|
||||
|
||||
def _digest(value: str, length: int = 16) -> str:
|
||||
"""Unsalted digest. Only for values that are already secrets (API keys)."""
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
"""Salted digest for values drawn from a guessable space.
|
||||
|
||||
repo.identity is a git remote URL, or ``local:<absolute path>`` when there is
|
||||
no remote — which normally contains the account username. Sixteen unsalted
|
||||
hex characters over that input space is enumerable, so this is not a
|
||||
privacy control without the salt. Salting per install keeps every
|
||||
within-account join the analytics actually use and gives up only
|
||||
cross-machine joins on the same repository, which nothing computes.
|
||||
"""
|
||||
if not value:
|
||||
return ""
|
||||
return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _safe_value(value: Any) -> Any:
|
||||
if isinstance(value, str):
|
||||
return memory_core.redact(value)
|
||||
@@ -145,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str:
|
||||
return created
|
||||
|
||||
|
||||
def _install_state_path() -> Path:
|
||||
return memory_core.data_dir() / "install-state.json"
|
||||
|
||||
|
||||
def is_first_run() -> bool:
|
||||
"""Whether this machine has never recorded a plugin event before."""
|
||||
return not _identity_path().exists()
|
||||
"""Whether install has never been recorded on this machine.
|
||||
|
||||
Deliberately NOT the identity file. That file is only written by a
|
||||
successful flush, so an offline or firewalled user recorded code.install on
|
||||
every single session, forever — and every 0.2.x user recorded one on their
|
||||
first 0.3.x session because 0.2.x never wrote it at all.
|
||||
"""
|
||||
return not _install_state_path().exists()
|
||||
|
||||
|
||||
def claim_install() -> str | None:
|
||||
"""Claim the one install/upgrade record for this machine, atomically.
|
||||
|
||||
Returns the event to record ("install" or "upgrade"), or None if another
|
||||
session already claimed it. O_CREAT|O_EXCL so two sessions starting together
|
||||
cannot both win.
|
||||
"""
|
||||
path = _install_state_path()
|
||||
# A fresh install has an empty data directory. Anything already there —
|
||||
# a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read
|
||||
# before the marker is created, since creating it would itself be content.
|
||||
upgrading = _data_dir_has_content()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
except FileExistsError:
|
||||
return None
|
||||
except OSError:
|
||||
return None
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
json.dump(
|
||||
{
|
||||
"plugin_version": memory_core.PLUGIN_VERSION,
|
||||
"installed_at": memory_core.utc_now(),
|
||||
"upgraded": upgrading,
|
||||
},
|
||||
stream,
|
||||
)
|
||||
except OSError:
|
||||
pass
|
||||
return "upgrade" if upgrading else "install"
|
||||
|
||||
|
||||
def _data_dir_has_content() -> bool:
|
||||
"""Whether anything predates this session in the plugin data directory."""
|
||||
try:
|
||||
for entry in memory_core.data_dir().iterdir():
|
||||
if entry.name != "install-state.json":
|
||||
return True
|
||||
except OSError:
|
||||
pass
|
||||
return False
|
||||
|
||||
|
||||
def claim_version_change() -> str | None:
|
||||
"""Return the previously recorded version if it differs, updating the marker.
|
||||
|
||||
Only meaningful once the marker exists — the first transition into 0.3.x has
|
||||
no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting
|
||||
the marker means the next session sees no change and records nothing.
|
||||
"""
|
||||
path = _install_state_path()
|
||||
try:
|
||||
state = json.loads(path.read_text(encoding="utf-8"))
|
||||
except (OSError, json.JSONDecodeError):
|
||||
return None
|
||||
if not isinstance(state, dict):
|
||||
return None
|
||||
previous = str(state.get("plugin_version") or "")
|
||||
if not previous or previous == memory_core.PLUGIN_VERSION:
|
||||
return None
|
||||
state["plugin_version"] = memory_core.PLUGIN_VERSION
|
||||
state["upgraded_at"] = memory_core.utc_now()
|
||||
try:
|
||||
temporary = path.with_suffix(f".{os.getpid()}.tmp")
|
||||
temporary.write_text(json.dumps(state), encoding="utf-8")
|
||||
temporary.replace(path)
|
||||
except OSError:
|
||||
return None
|
||||
return previous
|
||||
|
||||
|
||||
def record(
|
||||
@@ -168,19 +312,25 @@ def record(
|
||||
except OSError:
|
||||
pass
|
||||
properties = _safe_value(properties)
|
||||
# Stamped in the RECORDING process, beside harness. `source` used to be
|
||||
# read in the sending process from a module global, so whichever process
|
||||
# drained the spool named every event in it. flush() spreads per-event
|
||||
# properties last, so this now wins over any sender's default.
|
||||
properties.update(
|
||||
harness=_harness,
|
||||
source=_source_tag,
|
||||
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", ""))
|
||||
properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", ""))
|
||||
if session_id:
|
||||
properties["session_hash"] = _digest(session_id)
|
||||
properties["session_hash"] = _scoped_digest(session_id)
|
||||
line = json.dumps(
|
||||
{
|
||||
"event": f"{EVENT_PREFIX}.{event}",
|
||||
"uuid": str(uuid.uuid4()),
|
||||
"timestamp": memory_core.utc_now(),
|
||||
"properties": {
|
||||
key: value for key, value in properties.items() if value is not None
|
||||
@@ -239,38 +389,139 @@ def spawn_flush() -> bool:
|
||||
return False
|
||||
|
||||
|
||||
def _claim_name(attempt: int = 0) -> str:
|
||||
"""Claim filename. The attempt count rides in the name so the 7-day expiry
|
||||
only ever discards a batch that was actually retried and failed."""
|
||||
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
|
||||
|
||||
def _touch(path: Path) -> None:
|
||||
"""Refresh mtime so a claim's age measures time since it was claimed.
|
||||
|
||||
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
|
||||
after a quiet minute inherited the spool's last-write time and looked
|
||||
abandoned the instant it was made. A second sender would then take it over
|
||||
while the first was still posting, and both would deliver the batch.
|
||||
"""
|
||||
try:
|
||||
os.utime(path, None)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
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"
|
||||
claim = directory / _claim_name()
|
||||
spool = _spool_path()
|
||||
try:
|
||||
spool.replace(claim)
|
||||
_touch(claim)
|
||||
return claim
|
||||
except OSError:
|
||||
pass
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
Kept separate from the live spool so flush() can drain both in one run.
|
||||
Previously parked batches were only reachable when no spool existed at all,
|
||||
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,
|
||||
even though its own presence is what started the sender.
|
||||
"""
|
||||
now = time.time()
|
||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
||||
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||
try:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS:
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
continue
|
||||
if age < CLAIM_STALE_SECONDS:
|
||||
# Someone else holds a live lease on it.
|
||||
continue
|
||||
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||
try:
|
||||
orphan.replace(claim)
|
||||
_touch(claim)
|
||||
return claim
|
||||
except OSError:
|
||||
continue
|
||||
return None
|
||||
|
||||
|
||||
def _safe_mtime(path: Path) -> float:
|
||||
try:
|
||||
return path.stat().st_mtime
|
||||
except OSError:
|
||||
return 0.0
|
||||
|
||||
|
||||
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
"""Persist the unsent remainder, atomically, and refresh the lease.
|
||||
|
||||
Called after every successful batch. Two jobs: a retry resumes where the
|
||||
send stopped instead of re-posting from the top, and the rewrite doubles as
|
||||
the lease heartbeat, so a slow sender does not have its claim stolen
|
||||
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
|
||||
"""
|
||||
if not remaining:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
except OSError:
|
||||
try:
|
||||
temporary.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return False
|
||||
|
||||
|
||||
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
|
||||
"""Persist the remainder and drop the lease, because this sender has given up.
|
||||
|
||||
Distinct from the per-batch heartbeat: heartbeating on the way out would
|
||||
make an abandoned batch look actively owned for a further
|
||||
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
|
||||
threshold lets the next flush pick it up immediately, while the attempt
|
||||
count in the filename still bounds how many times that can happen.
|
||||
"""
|
||||
if not _rewrite_claim(claim, remaining):
|
||||
return
|
||||
try:
|
||||
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||
os.utime(claim, (released, released))
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
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/"
|
||||
@@ -300,34 +551,76 @@ def _post(payload: dict[str, Any], url: str) -> bool:
|
||||
|
||||
|
||||
def resolve_distinct_id() -> tuple[str, str]:
|
||||
"""Return the PostHog distinct id and the anonymous id it replaced, if any."""
|
||||
"""Return the PostHog distinct id and the anonymous id it replaced, if any.
|
||||
|
||||
The second value becomes a PostHog $identify alias. It is ONLY ever an
|
||||
anonymous id: aliasing one account email to another merges two real person
|
||||
profiles and cannot be undone, so a key that now belongs to a different
|
||||
account re-resolves with no alias.
|
||||
"""
|
||||
identity = _read_identity()
|
||||
email = identity.get("email", "")
|
||||
if email:
|
||||
return email, ""
|
||||
key = memory_core.api_key()
|
||||
fingerprint = _digest(key) if key else ""
|
||||
email = identity.get("email", "")
|
||||
|
||||
if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint:
|
||||
return email, ""
|
||||
|
||||
if not key:
|
||||
# No key to verify the account with; do not keep attributing to it.
|
||||
if email:
|
||||
identity.pop("email", None)
|
||||
identity.pop("key_fingerprint", None)
|
||||
_write_identity(identity)
|
||||
return anonymous_id(identity), ""
|
||||
email = _resolve_email(key)
|
||||
if not email:
|
||||
return anonymous_id(identity), ""
|
||||
previous = identity.get("anonymous_id", "")
|
||||
identity["email"] = email
|
||||
|
||||
resolved = _resolve_email(key)
|
||||
if not resolved:
|
||||
return (email, "") if email else (anonymous_id(identity), "")
|
||||
|
||||
# Alias only when going anonymous -> email for the first time.
|
||||
previous = "" if email else identity.get("anonymous_id", "")
|
||||
identity["email"] = resolved
|
||||
identity["key_fingerprint"] = fingerprint
|
||||
_write_identity(identity)
|
||||
return email, previous
|
||||
return resolved, previous
|
||||
|
||||
|
||||
def flush() -> int:
|
||||
"""Drain claimed spools to PostHog and return the number of events sent."""
|
||||
"""Drain the live spool, then any parked claims, and return events sent."""
|
||||
if not is_enabled():
|
||||
return 0
|
||||
claim = _claim_spool()
|
||||
sent, delivered = _drain(_claim_spool())
|
||||
if not delivered:
|
||||
# The network is failing. Retrying other batches now would only burn
|
||||
# their attempt budget against the same broken connection.
|
||||
return sent
|
||||
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
break
|
||||
count, delivered = _drain(parked)
|
||||
sent += count
|
||||
if not delivered:
|
||||
break
|
||||
return sent
|
||||
|
||||
|
||||
def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
"""Post one claimed batch file, recording progress after every batch.
|
||||
|
||||
Returns (events sent, whether everything was delivered).
|
||||
"""
|
||||
if claim is None:
|
||||
return 0
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
return 0
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
try:
|
||||
@@ -341,7 +634,7 @@ def flush() -> int:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0
|
||||
return 0, True
|
||||
|
||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||
if aliased_anonymous_id:
|
||||
@@ -360,12 +653,17 @@ def flush() -> int:
|
||||
|
||||
sent = 0
|
||||
for start in range(0, len(events), BATCH_SIZE):
|
||||
chunk = events[start : start + BATCH_SIZE]
|
||||
batch = [
|
||||
{
|
||||
"event": event["event"],
|
||||
"distinct_id": distinct_id,
|
||||
# Carried through from record() so a resend can be collapsed.
|
||||
"uuid": event.get("uuid"),
|
||||
"timestamp": event.get("timestamp"),
|
||||
"properties": {
|
||||
# Fallback only: events recorded by a build before source
|
||||
# moved into record() have none of their own.
|
||||
"source": _source_tag,
|
||||
"language": "python",
|
||||
"$process_person_profile": False,
|
||||
@@ -373,16 +671,19 @@ def flush() -> int:
|
||||
**(event.get("properties") or {}),
|
||||
},
|
||||
}
|
||||
for event in events[start : start + BATCH_SIZE]
|
||||
for event in chunk
|
||||
]
|
||||
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
|
||||
# Keep only what has not been delivered, and release the lease.
|
||||
# Previously the whole file was kept and the retry re-posted every
|
||||
# batch, including the ones that had already arrived.
|
||||
_release_claim(claim, events[start:])
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
return sent, True
|
||||
|
||||
|
||||
def main() -> int:
|
||||
|
||||
Reference in New Issue
Block a user