From 95d4fc27e23c1e05888915c142051c7c28235a95 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Tue, 15 Sep 2026 00:25:12 +0530 Subject: [PATCH 1/2] fix(plugins): make the telemetry salt stable, its own file, and memoized MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review found three ways the first cut produced worse data than no salt at all. All three came from keeping the salt as a key in the identity dict and doing an unlocked read-modify-write. Hooks are short-lived separate processes firing on every tool call, and people run more than one agent window, so several processes would read {}, each mint its own uuid4, and each hash with it. One repository hashed several ways in the window before a writer won. resolve_distinct_id holds a copy of that same dict across a network call to /v1/ping/ with a 5s timeout, so whichever write landed second erased the other's key: losing the salt changes repo_hash mid-stream, losing the email fires a second $identify and splits the person. _write_identity swallows OSError, and nothing memoized, so on a read-only or full data directory every single event got a brand-new random salt — unbounded cardinality in PostHog, which is strictly worse than the unsalted value it replaced. The salt now lives in its own file claimed with O_CREAT|O_EXCL, so exactly one process wins and the losers read the winner's value, and it is memoized per process. When it cannot be persisted the fallback is derived from the data directory path: stable for the machine rather than random per call. Its own file also means record() no longer creates telemetry-identity.json as a side effect. is_first_run keys off that file, so the first cut would have silently suppressed the install event — a production metric change hidden in a docs PR. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/telemetry.py | 62 ++++++++++++++++--- .../antigravity-plugin/core/telemetry.py | 62 ++++++++++++++++--- .../claude-code-plugin/core/telemetry.py | 62 ++++++++++++++++--- .../tests/test_telemetry.py | 47 ++++++++++++++ integrations/codex-plugin/core/telemetry.py | 62 ++++++++++++++++--- integrations/cursor-plugin/core/telemetry.py | 62 ++++++++++++++++--- integrations/kimi-plugin/core/telemetry.py | 62 ++++++++++++++++--- .../mem0-agent-plugin/core/telemetry.py | 62 ++++++++++++++++--- 8 files changed, 425 insertions(+), 56 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,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: diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,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: diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,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: diff --git a/integrations/claude-code-plugin/tests/test_telemetry.py b/integrations/claude-code-plugin/tests/test_telemetry.py index 0ae280b19..eaa741ac5 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +import os import sys from pathlib import Path from unittest.mock import patch @@ -273,3 +274,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 1b6bd854e..239836604 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,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: diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,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: diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,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: diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,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: From bd17f2b8c92c315c0250336211725c91f9bcd817 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Tue, 15 Sep 2026 00:30:58 +0530 Subject: [PATCH 2/2] fix(plugins): make the retry budget reachable and the rewrite durable MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review found that the first cut traded the duplicate-delivery bug for a worse one, and disproved its own load-bearing safety claim by experiment. Expiry was unreachable. _claim_parked touched the mtime on every re-claim and _release_claim backdated to exactly now minus the stale threshold, so a file's age hovered around 121 seconds and never approached the 7-day expiry. The attempt count in the filename therefore bounded nothing: an undeliverable batch (revoked key, proxy 403, oversized event) lived on disk forever, and because spawn_flush starts a sender whenever a .sending file exists, it spawned a detached Python process on every hook, MCP call and CLI invocation, forever. The old code self-healed here, so this was a regression. Expiry now gates on the attempt budget, which is the thing that actually accumulates; age stays only as a backstop for files that never carried an attempt marker. The attempt parser sniffed for a leading "a", which also matches a hex id like a1234567, so a legacy telemetry--.sending file parsed as attempt 1234567 and was deleted unsent on the first flush after upgrade — precisely the population this PR is meant to protect. Anchored on field position instead. The rewrite was not durable: no fsync before the rename, and _drain unlinked any claim that parsed to zero events. A crash between write and rename left the claim empty, and the next flush deleted it. Now fsynced, and a non-empty file that parses to nothing is quarantined as .corrupt rather than destroyed. read_text raises UnicodeDecodeError on a torn file, which `except OSError` does not catch. flush() runs from a bare `finally:` in flush_worker, so the exception also skipped the handoff cleanup and left it stuck in .running. The per-batch rewrite's return value was discarded, so a failed rewrite let the loop continue as though progress had been recorded — reintroducing the exact duplicate delivery this PR exists to fix. .partial files orphaned by a crash between write and rename matched no glob in the module and were never cleaned up. Also replaces the heartbeat test, which asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS` — two constants, executing none of the code under test. It now drives the real rewrite and watches the mtime move. New tests cover expiry being reachable, legacy filename parsing, torn-claim quarantine, failed-rewrite behaviour and debris sweeping. 64 core tests, 199 host tests. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/telemetry.py | 78 ++++++++++++--- .../tests/test_spool_delivery.py | 98 +++++++++++++++++-- .../antigravity-plugin/core/telemetry.py | 78 ++++++++++++--- .../claude-code-plugin/core/telemetry.py | 78 ++++++++++++--- integrations/codex-plugin/core/telemetry.py | 78 ++++++++++++--- integrations/cursor-plugin/core/telemetry.py | 78 ++++++++++++--- integrations/kimi-plugin/core/telemetry.py | 78 ++++++++++++--- .../mem0-agent-plugin/core/telemetry.py | 78 ++++++++++++--- 8 files changed, 558 insertions(+), 86 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -313,9 +313,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 @@ -349,6 +358,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. @@ -364,7 +388,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: @@ -406,10 +434,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 @@ -498,6 +530,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: @@ -518,7 +551,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: @@ -529,8 +573,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 @@ -580,8 +631,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 55ee7c907..fa26cc9a6 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -313,9 +313,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 @@ -349,6 +358,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. @@ -364,7 +388,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: @@ -406,10 +434,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 @@ -498,6 +530,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: @@ -518,7 +551,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: @@ -529,8 +573,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 @@ -580,8 +631,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 55ee7c907..fa26cc9a6 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -313,9 +313,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 @@ -349,6 +358,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. @@ -364,7 +388,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: @@ -406,10 +434,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 @@ -498,6 +530,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: @@ -518,7 +551,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: @@ -529,8 +573,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 @@ -580,8 +631,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/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -313,9 +313,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 @@ -349,6 +358,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. @@ -364,7 +388,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: @@ -406,10 +434,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 @@ -498,6 +530,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: @@ -518,7 +551,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: @@ -529,8 +573,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 @@ -580,8 +631,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 55ee7c907..fa26cc9a6 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -313,9 +313,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 @@ -349,6 +358,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. @@ -364,7 +388,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: @@ -406,10 +434,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 @@ -498,6 +530,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: @@ -518,7 +551,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: @@ -529,8 +573,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 @@ -580,8 +631,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 55ee7c907..fa26cc9a6 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -313,9 +313,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 @@ -349,6 +358,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. @@ -364,7 +388,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: @@ -406,10 +434,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 @@ -498,6 +530,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: @@ -518,7 +551,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: @@ -529,8 +573,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 @@ -580,8 +631,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 55ee7c907..fa26cc9a6 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -313,9 +313,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 @@ -349,6 +358,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. @@ -364,7 +388,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: @@ -406,10 +434,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 @@ -498,6 +530,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: @@ -518,7 +551,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: @@ -529,8 +573,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 @@ -580,8 +631,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