diff --git a/integrations/agent-plugin-core/python/hook_runner.py b/integrations/agent-plugin-core/python/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/agent-plugin-core/python/hook_runner.py +++ b/integrations/agent-plugin-core/python/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index d7933bfbc..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: @@ -221,18 +267,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +329,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +352,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -396,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 @@ -432,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. @@ -447,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: @@ -489,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 @@ -563,8 +690,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +713,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) @@ -599,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: @@ -619,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: @@ -630,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 @@ -681,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/agent-plugin-core/tests/test_uninitialised_identity.py b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py index 1c742b518..3a733d0bd 100644 --- a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py +++ b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py @@ -181,6 +181,37 @@ def test_source_tag_defaults_agree_between_the_two_modules(): def test_the_plugin_declares_its_surface_in_the_body_and_the_headers(): """Body and headers both, because only the body works on every backend.""" + +def _session_start(core: Path, data_dir: Path) -> list[str]: + """Drive the real hook_runner session-start path and return lifecycle events.""" + recorded = "\n".join( + [ + "import io, json, sys", + f"sys.path.insert(0, {str(core)!r})", + "import telemetry, hook_runner", + "seen = []", + "telemetry.record = lambda event, **kw: seen.append(event) or None", + "telemetry.spawn_flush = lambda: False", + # run() reads sys.argv through argparse; it takes no positional args. + "sys.argv = ['hook_runner', 'session-start']", + "sys.stdin = io.StringIO('{}')", + "hook_runner.run()", + "print(json.dumps([e for e in seen if e in ('install', 'upgrade')]))", + ] + ) + import json as _json + + return _json.loads(_run(core, data_dir, recorded) or "[]") + + +def test_a_fresh_install_reports_install_not_upgrade(): + """The decision must survive the writes hook_runner does before asking. + + claim_install() is reached only after cache_plugin_api_key() has written + `api-key` and EvidenceStore() has created `evidence.sqlite3`. Asking "is the + data dir empty" at that point always saw content, so code.install could + never fire and every new user was counted as an upgrade. + """ core = _core_dir("claude-code-plugin") if not core.exists(): pytest.skip("claude-code-plugin is not built in this tree") @@ -204,3 +235,35 @@ def test_the_plugin_declares_its_surface_in_the_body_and_the_headers(): # The transport headers the three call sites relied on must survive. assert headers["auth"] == "Token k" assert headers["ctype"] == "application/json" + + data_dir = Path(tmp) / "data" + assert _session_start(core, data_dir) == ["install"] + + +def test_the_lifecycle_event_fires_exactly_once(): + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + first = _session_start(core, data_dir) + second = _session_start(core, data_dir) + third = _session_start(core, data_dir) + + assert first == ["install"] + assert second == [] + assert third == [] + + +def test_an_existing_data_dir_reports_upgrade(): + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + data_dir.mkdir(parents=True) + # A 0.2.x leftover: the data dir survives the upgrade. + (data_dir / "requirements.txt").write_text("mem0ai\n", encoding="utf-8") + assert _session_start(core, data_dir) == ["upgrade"] diff --git a/integrations/antigravity-plugin/core/hook_runner.py b/integrations/antigravity-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/antigravity-plugin/core/hook_runner.py +++ b/integrations/antigravity-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index d7933bfbc..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: @@ -221,18 +267,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +329,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +352,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -396,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 @@ -432,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. @@ -447,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: @@ -489,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 @@ -563,8 +690,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +713,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) @@ -599,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: @@ -619,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: @@ -630,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 @@ -681,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/hook_runner.py b/integrations/claude-code-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/claude-code-plugin/core/hook_runner.py +++ b/integrations/claude-code-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index d7933bfbc..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: @@ -221,18 +267,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +329,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +352,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -396,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 @@ -432,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. @@ -447,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: @@ -489,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 @@ -563,8 +690,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +713,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) @@ -599,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: @@ -619,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: @@ -630,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 @@ -681,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/hook_runner.py b/integrations/codex-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/codex-plugin/core/hook_runner.py +++ b/integrations/codex-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index d7933bfbc..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: @@ -221,18 +267,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +329,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +352,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -396,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 @@ -432,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. @@ -447,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: @@ -489,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 @@ -563,8 +690,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +713,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) @@ -599,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: @@ -619,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: @@ -630,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 @@ -681,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/hook_runner.py b/integrations/cursor-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/cursor-plugin/core/hook_runner.py +++ b/integrations/cursor-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index d7933bfbc..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: @@ -221,18 +267,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +329,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +352,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -396,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 @@ -432,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. @@ -447,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: @@ -489,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 @@ -563,8 +690,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +713,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) @@ -599,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: @@ -619,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: @@ -630,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 @@ -681,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/hook_runner.py b/integrations/kimi-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/kimi-plugin/core/hook_runner.py +++ b/integrations/kimi-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index d7933bfbc..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: @@ -221,18 +267,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +329,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +352,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -396,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 @@ -432,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. @@ -447,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: @@ -489,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 @@ -563,8 +690,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +713,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) @@ -599,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: @@ -619,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: @@ -630,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 @@ -681,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 d7933bfbc..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: @@ -221,18 +267,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +329,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +352,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -396,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 @@ -432,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. @@ -447,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: @@ -489,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 @@ -563,8 +690,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +713,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) @@ -599,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: @@ -619,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: @@ -630,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 @@ -681,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