diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index a24c97264..8cce9cbc7 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -49,6 +49,7 @@ except ImportError: _PLATFORM_SOURCE = "MEM0_PLUGIN" _PLATFORM_APPLICATION = "" +_salt_cache: str = "" _harness: str = _DEFAULT_HARNESS _source_tag: str = _DEFAULT_SOURCE_TAG _PRIVATE_KEYS = { @@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + 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 + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: @@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -481,6 +536,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: 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", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -664,6 +742,7 @@ def flush() -> int: # 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() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: 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) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/agent-plugin-core/tests/test_spool_delivery.py b/integrations/agent-plugin-core/tests/test_spool_delivery.py index 754307454..08266cc98 100644 --- a/integrations/agent-plugin-core/tests/test_spool_delivery.py +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -112,16 +112,19 @@ def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry): assert names == {"code.parked", "code.fresh"} -def test_an_untried_batch_is_not_expired_by_age_alone(telemetry): - """Expiry should discard what failed, not what never got a turn.""" +def test_a_batch_is_retried_until_the_budget_is_spent_not_discarded(telemetry): + """Expiry discards what failed repeatedly, not what merely sat for a while. + + The budget is the attempt count, because age cannot be one: every re-claim + touches the mtime and every release backdates it, so age never accumulates. + """ telemetry.record("parked") telemetry._post = lambda payload, url: False telemetry.flush() parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) assert len(parked) == 1 - ancient = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60) - os.utime(parked[0], (ancient, ancient)) + assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS sent: list[dict] = [] telemetry._post = lambda payload, url: sent.append(payload) or True @@ -154,11 +157,88 @@ def test_progress_is_recorded_after_every_batch(telemetry): assert json.loads(remaining[0])["properties"]["index"] == 200 -def test_the_heartbeat_stays_well_inside_the_lease(telemetry): +def test_the_heartbeat_actually_refreshes_the_lease(telemetry): """The claim rewrite doubles as the lease heartbeat. - _post makes a single attempt with SEND_TIMEOUT and no retry, so a heartbeat - lands at least that often. If a retry loop is ever added to _post, this is - the assertion that catches a sender losing its claim mid-flight. + Previously asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS`, which compares + two constants and executes none of the code under test. Drive the real + rewrite and watch the mtime move instead. """ - assert telemetry.SEND_TIMEOUT * 4 < telemetry.CLAIM_STALE_SECONDS + for index in range(150): + telemetry.record("search", index=index) + claim = telemetry._claim_spool() + assert claim is not None + + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(claim, (stale, stale)) + assert time.time() - claim.stat().st_mtime > telemetry.CLAIM_STALE_SECONDS + + telemetry._rewrite_claim(claim, [{"event": "code.x", "properties": {}}]) + assert time.time() - claim.stat().st_mtime < telemetry.CLAIM_STALE_SECONDS + + +def test_an_undeliverable_batch_is_eventually_given_up_on(telemetry): + """Expiry has to be reachable from a state the state machine can produce. + + It was not: every re-claim touched the mtime and every release backdated it + by a fixed amount, so age hovered near the stale threshold and the 7-day + expiry never fired. An undeliverable batch lived on disk forever, and + spawn_flush saw it and started a sender on every hook. + """ + telemetry.record("doomed") + telemetry._post = lambda payload, url: False + + for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3): + telemetry.flush() + + leftover = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) + assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}" + + +def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry): + """The old shape is telemetry--.sending, and hex can start with 'a'.""" + assert telemetry._claim_attempt(Path("telemetry-999-deadbeef.sending")) == 0 + assert telemetry._claim_attempt(Path("telemetry-999-a1234567.sending")) == 0 + assert telemetry._claim_attempt(Path("telemetry-999-deadbeef-a2.sending")) == 2 + + +def test_a_torn_claim_is_quarantined_not_deleted(telemetry): + """A non-empty file that parses to nothing is the remainder, not garbage.""" + telemetry.record("search") + claim = telemetry._claim_spool() + claim.write_bytes(b"\xff\xfe not utf-8 at all") + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(claim, (stale, stale)) + + sent = telemetry.flush() + + assert sent == 0 + assert not claim.exists() + quarantined = list(telemetry.memory_core.data_dir().glob("*.corrupt")) + assert len(quarantined) == 1, "torn claim was destroyed instead of kept" + + +def test_a_failed_rewrite_stops_instead_of_redelivering(telemetry): + """Ignoring the rewrite result reintroduced the duplicates this PR fixes.""" + for index in range(250): + telemetry.record("search", index=index) + + telemetry._rewrite_claim = lambda claim, remaining: False + delivered = [] + telemetry._post = lambda payload, url: delivered.extend(payload.get("batch", [])) or True + + telemetry.flush() + assert len(delivered) == 100, f"kept going after a failed rewrite: {len(delivered)}" + + +def test_partial_files_are_swept(telemetry): + """Nothing else globs *.partial, so a crash mid-rename orphans one forever.""" + data_dir = telemetry.memory_core.data_dir() + data_dir.mkdir(parents=True, exist_ok=True) + debris = data_dir / "telemetry-1-abc-a0.1.partial" + debris.write_text("x", encoding="utf-8") + old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(debris, (old, old)) + + telemetry.flush() + assert not debris.exists() diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index a24c97264..8cce9cbc7 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -49,6 +49,7 @@ except ImportError: _PLATFORM_SOURCE = "MEM0_PLUGIN" _PLATFORM_APPLICATION = "" +_salt_cache: str = "" _harness: str = _DEFAULT_HARNESS _source_tag: str = _DEFAULT_SOURCE_TAG _PRIVATE_KEYS = { @@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + 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 + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: @@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -481,6 +536,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: 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", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -664,6 +742,7 @@ def flush() -> int: # 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() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: 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) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index a24c97264..8cce9cbc7 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -49,6 +49,7 @@ except ImportError: _PLATFORM_SOURCE = "MEM0_PLUGIN" _PLATFORM_APPLICATION = "" +_salt_cache: str = "" _harness: str = _DEFAULT_HARNESS _source_tag: str = _DEFAULT_SOURCE_TAG _PRIVATE_KEYS = { @@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + 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 + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: @@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -481,6 +536,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: 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", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -664,6 +742,7 @@ def flush() -> int: # 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() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: 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) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/claude-code-plugin/tests/test_telemetry.py b/integrations/claude-code-plugin/tests/test_telemetry.py index 4476a9dd3..efd2e11c0 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -315,3 +315,49 @@ def test_spawn_flush_does_nothing_without_a_spool(isolated_env): with patch.object(telemetry.subprocess, "Popen") as popen: assert telemetry.spawn_flush() is True popen.assert_called_once() + + +def test_salt_is_stable_across_processes(isolated_env): + """Hooks are separate short-lived processes; one repo must hash one way. + + An unlocked read-modify-write let each process mint its own salt, so a + repository hashed several ways in the window before one writer won. + """ + import subprocess as sp + + core = str(Path(__file__).resolve().parents[1] / "core") + script = ( + f"import sys; sys.path.insert(0, {core!r})\n" + "import telemetry\n" + "print(telemetry._install_salt())" + ) + env = {**os.environ, "MEM0_CODE_DATA_DIR": str(memory_core.data_dir())} + salts = { + sp.run([sys.executable, "-c", script], capture_output=True, text=True, env=env).stdout.strip() + for _ in range(4) + } + assert len(salts) == 1, f"one repo hashed {len(salts)} ways: {salts}" + + +def test_salt_does_not_touch_the_identity_file(isolated_env): + """The identity file is is_first_run's marker and the sender's email store. + + Writing the salt into it would create it from record(), suppressing the + install event, and would race resolve_distinct_id, which holds a stale copy + of that dict across a network call. + """ + telemetry._install_salt() + assert not telemetry._identity_path().exists() + + +def test_salt_is_stable_when_it_cannot_be_persisted(isolated_env, monkeypatch): + """A read-only data dir must degrade to a weaker salt, not to random-per-call. + + Random per call is unbounded cardinality in PostHog, which is worse than no + salt at all. + """ + telemetry._salt_cache = "" + monkeypatch.setattr(telemetry.os, "open", lambda *a, **k: (_ for _ in ()).throw(OSError("read-only"))) + first = telemetry._install_salt() + telemetry._salt_cache = "" + assert telemetry._install_salt() == first diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index a24c97264..8cce9cbc7 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -49,6 +49,7 @@ except ImportError: _PLATFORM_SOURCE = "MEM0_PLUGIN" _PLATFORM_APPLICATION = "" +_salt_cache: str = "" _harness: str = _DEFAULT_HARNESS _source_tag: str = _DEFAULT_SOURCE_TAG _PRIVATE_KEYS = { @@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + 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 + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: @@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -481,6 +536,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: 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", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -664,6 +742,7 @@ def flush() -> int: # 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() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: 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) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index a24c97264..8cce9cbc7 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -49,6 +49,7 @@ except ImportError: _PLATFORM_SOURCE = "MEM0_PLUGIN" _PLATFORM_APPLICATION = "" +_salt_cache: str = "" _harness: str = _DEFAULT_HARNESS _source_tag: str = _DEFAULT_SOURCE_TAG _PRIVATE_KEYS = { @@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + 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 + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: @@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -481,6 +536,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: 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", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -664,6 +742,7 @@ def flush() -> int: # 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() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: 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) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index a24c97264..8cce9cbc7 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -49,6 +49,7 @@ except ImportError: _PLATFORM_SOURCE = "MEM0_PLUGIN" _PLATFORM_APPLICATION = "" +_salt_cache: str = "" _harness: str = _DEFAULT_HARNESS _source_tag: str = _DEFAULT_SOURCE_TAG _PRIVATE_KEYS = { @@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + 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 + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: @@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -481,6 +536,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: 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", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -664,6 +742,7 @@ def flush() -> int: # 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() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: 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) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index a24c97264..8cce9cbc7 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -49,6 +49,7 @@ except ImportError: _PLATFORM_SOURCE = "MEM0_PLUGIN" _PLATFORM_APPLICATION = "" +_salt_cache: str = "" _harness: str = _DEFAULT_HARNESS _source_tag: str = _DEFAULT_SOURCE_TAG _PRIVATE_KEYS = { @@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + 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 + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: @@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -481,6 +536,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: 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", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -664,6 +742,7 @@ def flush() -> int: # 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() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: 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) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True