From 95d4fc27e23c1e05888915c142051c7c28235a95 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Tue, 15 Sep 2026 00:25:12 +0530 Subject: [PATCH 1/3] fix(plugins): make the telemetry salt stable, its own file, and memoized MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review found three ways the first cut produced worse data than no salt at all. All three came from keeping the salt as a key in the identity dict and doing an unlocked read-modify-write. Hooks are short-lived separate processes firing on every tool call, and people run more than one agent window, so several processes would read {}, each mint its own uuid4, and each hash with it. One repository hashed several ways in the window before a writer won. resolve_distinct_id holds a copy of that same dict across a network call to /v1/ping/ with a 5s timeout, so whichever write landed second erased the other's key: losing the salt changes repo_hash mid-stream, losing the email fires a second $identify and splits the person. _write_identity swallows OSError, and nothing memoized, so on a read-only or full data directory every single event got a brand-new random salt — unbounded cardinality in PostHog, which is strictly worse than the unsalted value it replaced. The salt now lives in its own file claimed with O_CREAT|O_EXCL, so exactly one process wins and the losers read the winner's value, and it is memoized per process. When it cannot be persisted the fallback is derived from the data directory path: stable for the machine rather than random per call. Its own file also means record() no longer creates telemetry-identity.json as a side effect. is_first_run keys off that file, so the first cut would have silently suppressed the install event — a production metric change hidden in a docs PR. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/telemetry.py | 62 ++++++++++++++++--- .../antigravity-plugin/core/telemetry.py | 62 ++++++++++++++++--- .../claude-code-plugin/core/telemetry.py | 62 ++++++++++++++++--- .../tests/test_telemetry.py | 47 ++++++++++++++ integrations/codex-plugin/core/telemetry.py | 62 ++++++++++++++++--- integrations/cursor-plugin/core/telemetry.py | 62 ++++++++++++++++--- integrations/kimi-plugin/core/telemetry.py | 62 ++++++++++++++++--- .../mem0-agent-plugin/core/telemetry.py | 62 ++++++++++++++++--- 8 files changed, 425 insertions(+), 56 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + def _install_salt() -> str: - """Random per-install salt, created on first use and kept in the identity file.""" - identity = _read_identity() - salt = identity.get("salt") - if not salt: - salt = uuid.uuid4().hex - identity["salt"] = salt - _write_identity(identity) - return salt + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + def _install_salt() -> str: - """Random per-install salt, created on first use and kept in the identity file.""" - identity = _read_identity() - salt = identity.get("salt") - if not salt: - salt = uuid.uuid4().hex - identity["salt"] = salt - _write_identity(identity) - return salt + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + def _install_salt() -> str: - """Random per-install salt, created on first use and kept in the identity file.""" - identity = _read_identity() - salt = identity.get("salt") - if not salt: - salt = uuid.uuid4().hex - identity["salt"] = salt - _write_identity(identity) - return salt + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: diff --git a/integrations/claude-code-plugin/tests/test_telemetry.py b/integrations/claude-code-plugin/tests/test_telemetry.py index 0ae280b19..eaa741ac5 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +import os import sys from pathlib import Path from unittest.mock import patch @@ -273,3 +274,49 @@ def test_spawn_flush_does_nothing_without_a_spool(isolated_env): with patch.object(telemetry.subprocess, "Popen") as popen: assert telemetry.spawn_flush() is True popen.assert_called_once() + + +def test_salt_is_stable_across_processes(isolated_env): + """Hooks are separate short-lived processes; one repo must hash one way. + + An unlocked read-modify-write let each process mint its own salt, so a + repository hashed several ways in the window before one writer won. + """ + import subprocess as sp + + core = str(Path(__file__).resolve().parents[1] / "core") + script = ( + f"import sys; sys.path.insert(0, {core!r})\n" + "import telemetry\n" + "print(telemetry._install_salt())" + ) + env = {**os.environ, "MEM0_CODE_DATA_DIR": str(memory_core.data_dir())} + salts = { + sp.run([sys.executable, "-c", script], capture_output=True, text=True, env=env).stdout.strip() + for _ in range(4) + } + assert len(salts) == 1, f"one repo hashed {len(salts)} ways: {salts}" + + +def test_salt_does_not_touch_the_identity_file(isolated_env): + """The identity file is is_first_run's marker and the sender's email store. + + Writing the salt into it would create it from record(), suppressing the + install event, and would race resolve_distinct_id, which holds a stale copy + of that dict across a network call. + """ + telemetry._install_salt() + assert not telemetry._identity_path().exists() + + +def test_salt_is_stable_when_it_cannot_be_persisted(isolated_env, monkeypatch): + """A read-only data dir must degrade to a weaker salt, not to random-per-call. + + Random per call is unbounded cardinality in PostHog, which is worse than no + salt at all. + """ + telemetry._salt_cache = "" + monkeypatch.setattr(telemetry.os, "open", lambda *a, **k: (_ for _ in ()).throw(OSError("read-only"))) + first = telemetry._install_salt() + telemetry._salt_cache = "" + assert telemetry._install_salt() == first diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + def _install_salt() -> str: - """Random per-install salt, created on first use and kept in the identity file.""" - identity = _read_identity() - salt = identity.get("salt") - if not salt: - salt = uuid.uuid4().hex - identity["salt"] = salt - _write_identity(identity) - return salt + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + def _install_salt() -> str: - """Random per-install salt, created on first use and kept in the identity file.""" - identity = _read_identity() - salt = identity.get("salt") - if not salt: - salt = uuid.uuid4().hex - identity["salt"] = salt - _write_identity(identity) - return salt + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + def _install_salt() -> str: - """Random per-install salt, created on first use and kept in the identity file.""" - identity = _read_identity() - salt = identity.get("salt") - if not salt: - salt = uuid.uuid4().hex - identity["salt"] = salt - _write_identity(identity) - return salt + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 1b6bd854e..239836604 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -34,6 +34,7 @@ from typing import Any import memory_core +_salt_cache: str = "" _harness: str = "generic" _source_tag: str = "MEM0_PLUGIN" _PRIVATE_KEYS = { @@ -92,15 +93,60 @@ def _digest(value: str, length: int = 16) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length] +def _salt_path() -> Path: + return memory_core.data_dir() / "telemetry-salt" + + def _install_salt() -> str: - """Random per-install salt, created on first use and kept in the identity file.""" - identity = _read_identity() - salt = identity.get("salt") - if not salt: - salt = uuid.uuid4().hex - identity["salt"] = salt - _write_identity(identity) - return salt + """Random per-install salt, created once and memoized for the process. + + Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in + the identity file. Three reasons, all of which produced wrong data when this + lived in the identity dict: + + - Hooks are short-lived separate processes firing on every tool call, and + people run more than one agent window. A read-modify-write would let each + process mint its own salt, so one repository would hash several ways in the + window before a writer won. + - resolve_distinct_id holds a copy of the identity dict across a network call + to /v1/ping/, so whichever write landed second erased the other's key — + losing either the salt (repo_hash changes mid-stream) or the email (a + second $identify, splitting the person). + - Touching the identity file from record() would create it, and is_first_run + keys off that file, so recording an event would silently suppress the + install event. + + On a read-only or full data directory the fallback is derived from the data + directory path: stable for the machine rather than random per call, so the + failure mode is a weaker salt and not unbounded cardinality in PostHog. + """ + global _salt_cache + if _salt_cache: + return _salt_cache + + path = _salt_path() + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + except OSError: + pass + except FileExistsError: + pass + except OSError: + # Cannot persist. Stable-per-machine beats random-per-call. + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache + + try: + _salt_cache = path.read_text(encoding="utf-8").strip() + except OSError: + _salt_cache = "" + if not _salt_cache: + _salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest() + return _salt_cache def _scoped_digest(value: str, length: int = 16) -> str: From bd17f2b8c92c315c0250336211725c91f9bcd817 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Tue, 15 Sep 2026 00:30:58 +0530 Subject: [PATCH 2/3] fix(plugins): make the retry budget reachable and the rewrite durable MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review found that the first cut traded the duplicate-delivery bug for a worse one, and disproved its own load-bearing safety claim by experiment. Expiry was unreachable. _claim_parked touched the mtime on every re-claim and _release_claim backdated to exactly now minus the stale threshold, so a file's age hovered around 121 seconds and never approached the 7-day expiry. The attempt count in the filename therefore bounded nothing: an undeliverable batch (revoked key, proxy 403, oversized event) lived on disk forever, and because spawn_flush starts a sender whenever a .sending file exists, it spawned a detached Python process on every hook, MCP call and CLI invocation, forever. The old code self-healed here, so this was a regression. Expiry now gates on the attempt budget, which is the thing that actually accumulates; age stays only as a backstop for files that never carried an attempt marker. The attempt parser sniffed for a leading "a", which also matches a hex id like a1234567, so a legacy telemetry--.sending file parsed as attempt 1234567 and was deleted unsent on the first flush after upgrade — precisely the population this PR is meant to protect. Anchored on field position instead. The rewrite was not durable: no fsync before the rename, and _drain unlinked any claim that parsed to zero events. A crash between write and rename left the claim empty, and the next flush deleted it. Now fsynced, and a non-empty file that parses to nothing is quarantined as .corrupt rather than destroyed. read_text raises UnicodeDecodeError on a torn file, which `except OSError` does not catch. flush() runs from a bare `finally:` in flush_worker, so the exception also skipped the handoff cleanup and left it stuck in .running. The per-batch rewrite's return value was discarded, so a failed rewrite let the loop continue as though progress had been recorded — reintroducing the exact duplicate delivery this PR exists to fix. .partial files orphaned by a crash between write and rename matched no glob in the module and were never cleaned up. Also replaces the heartbeat test, which asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS` — two constants, executing none of the code under test. It now drives the real rewrite and watches the mtime move. New tests cover expiry being reachable, legacy filename parsing, torn-claim quarantine, failed-rewrite behaviour and debris sweeping. 64 core tests, 199 host tests. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/telemetry.py | 78 ++++++++++++--- .../tests/test_spool_delivery.py | 98 +++++++++++++++++-- .../antigravity-plugin/core/telemetry.py | 78 ++++++++++++--- .../claude-code-plugin/core/telemetry.py | 78 ++++++++++++--- integrations/codex-plugin/core/telemetry.py | 78 ++++++++++++--- integrations/cursor-plugin/core/telemetry.py | 78 ++++++++++++--- integrations/kimi-plugin/core/telemetry.py | 78 ++++++++++++--- .../mem0-agent-plugin/core/telemetry.py | 78 ++++++++++++--- 8 files changed, 558 insertions(+), 86 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/agent-plugin-core/tests/test_spool_delivery.py b/integrations/agent-plugin-core/tests/test_spool_delivery.py index 754307454..08266cc98 100644 --- a/integrations/agent-plugin-core/tests/test_spool_delivery.py +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -112,16 +112,19 @@ def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry): assert names == {"code.parked", "code.fresh"} -def test_an_untried_batch_is_not_expired_by_age_alone(telemetry): - """Expiry should discard what failed, not what never got a turn.""" +def test_a_batch_is_retried_until_the_budget_is_spent_not_discarded(telemetry): + """Expiry discards what failed repeatedly, not what merely sat for a while. + + The budget is the attempt count, because age cannot be one: every re-claim + touches the mtime and every release backdates it, so age never accumulates. + """ telemetry.record("parked") telemetry._post = lambda payload, url: False telemetry.flush() parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) assert len(parked) == 1 - ancient = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60) - os.utime(parked[0], (ancient, ancient)) + assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS sent: list[dict] = [] telemetry._post = lambda payload, url: sent.append(payload) or True @@ -154,11 +157,88 @@ def test_progress_is_recorded_after_every_batch(telemetry): assert json.loads(remaining[0])["properties"]["index"] == 200 -def test_the_heartbeat_stays_well_inside_the_lease(telemetry): +def test_the_heartbeat_actually_refreshes_the_lease(telemetry): """The claim rewrite doubles as the lease heartbeat. - _post makes a single attempt with SEND_TIMEOUT and no retry, so a heartbeat - lands at least that often. If a retry loop is ever added to _post, this is - the assertion that catches a sender losing its claim mid-flight. + Previously asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS`, which compares + two constants and executes none of the code under test. Drive the real + rewrite and watch the mtime move instead. """ - assert telemetry.SEND_TIMEOUT * 4 < telemetry.CLAIM_STALE_SECONDS + for index in range(150): + telemetry.record("search", index=index) + claim = telemetry._claim_spool() + assert claim is not None + + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(claim, (stale, stale)) + assert time.time() - claim.stat().st_mtime > telemetry.CLAIM_STALE_SECONDS + + telemetry._rewrite_claim(claim, [{"event": "code.x", "properties": {}}]) + assert time.time() - claim.stat().st_mtime < telemetry.CLAIM_STALE_SECONDS + + +def test_an_undeliverable_batch_is_eventually_given_up_on(telemetry): + """Expiry has to be reachable from a state the state machine can produce. + + It was not: every re-claim touched the mtime and every release backdated it + by a fixed amount, so age hovered near the stale threshold and the 7-day + expiry never fired. An undeliverable batch lived on disk forever, and + spawn_flush saw it and started a sender on every hook. + """ + telemetry.record("doomed") + telemetry._post = lambda payload, url: False + + for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3): + telemetry.flush() + + leftover = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) + assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}" + + +def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry): + """The old shape is telemetry--.sending, and hex can start with 'a'.""" + assert telemetry._claim_attempt(Path("telemetry-999-deadbeef.sending")) == 0 + assert telemetry._claim_attempt(Path("telemetry-999-a1234567.sending")) == 0 + assert telemetry._claim_attempt(Path("telemetry-999-deadbeef-a2.sending")) == 2 + + +def test_a_torn_claim_is_quarantined_not_deleted(telemetry): + """A non-empty file that parses to nothing is the remainder, not garbage.""" + telemetry.record("search") + claim = telemetry._claim_spool() + claim.write_bytes(b"\xff\xfe not utf-8 at all") + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(claim, (stale, stale)) + + sent = telemetry.flush() + + assert sent == 0 + assert not claim.exists() + quarantined = list(telemetry.memory_core.data_dir().glob("*.corrupt")) + assert len(quarantined) == 1, "torn claim was destroyed instead of kept" + + +def test_a_failed_rewrite_stops_instead_of_redelivering(telemetry): + """Ignoring the rewrite result reintroduced the duplicates this PR fixes.""" + for index in range(250): + telemetry.record("search", index=index) + + telemetry._rewrite_claim = lambda claim, remaining: False + delivered = [] + telemetry._post = lambda payload, url: delivered.extend(payload.get("batch", [])) or True + + telemetry.flush() + assert len(delivered) == 100, f"kept going after a failed rewrite: {len(delivered)}" + + +def test_partial_files_are_swept(telemetry): + """Nothing else globs *.partial, so a crash mid-rename orphans one forever.""" + data_dir = telemetry.memory_core.data_dir() + data_dir.mkdir(parents=True, exist_ok=True) + debris = data_dir / "telemetry-1-abc-a0.1.partial" + debris.write_text("x", encoding="utf-8") + old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(debris, (old, old)) + + telemetry.flush() + assert not debris.exists() diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True From 2c885fdcd7c246f59011ba517e35729bf6606a04 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Tue, 15 Sep 2026 00:31:52 +0530 Subject: [PATCH 3/3] fix(plugins): make code.install reachable, and stop pinging on every flush MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review found the headline fix inverted: code.install could never fire, so every fresh install reported an upgrade and the two cohorts became indistinguishable — strictly worse than the bug being fixed. hook_runner reaches claim_install() only after cache_plugin_api_key() has written `api-key` and EvidenceStore() has created `evidence.sqlite3` and its WAL files. Asking "is the data directory empty" at that point always saw content. The caller now snapshots emptiness at the top of the run, before anything writes, and passes it in. Also caught by review, all in the same file: - claim_version_change was an unsynchronized read-modify-write, so several concurrently starting sessions each observed the old version and each recorded an upgrade. The first session after a version bump is exactly when a user's open agent windows all restart together. The transition is now claimed with an exclusive per-version sentinel. - A crash between O_EXCL and the write left an empty marker, which disabled every future upgrade event on that machine: claim_install saw the file and claim_version_change could not parse it. An unparseable marker is now repaired. - claim_install consumed the one-shot claim even under MEM0_TELEMETRY=false, so a user who opted out for their first sessions would never report install after opting in. - Existing users have an email but no key fingerprint, so the fast path always missed and every flush paid an uncached /v1/ping/ — a 5s timeout each time for the offline users this stack keeps citing. Legacy rows now adopt the current key's fingerprint instead of re-resolving. - A key that will not resolve (revoked, offline) kept attributing to the previous account's email, which is the bug this was meant to fix. It now falls back to the anonymous id. - The anonymous id was never rotated, so once it had been merged into one account it was still offered as the alias for the next one. An alias naming an already-identified id is what could link two real people; it is now offered once. The gap that let this ship was that no test drove hook_runner's session-start path — the decision was only ever tested by calling claim_install() directly on a directory nothing had touched. Adds subprocess tests that run the real entrypoint: fresh install, exactly-once, and an existing data dir. 62 core tests, 203 host tests. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/hook_runner.py | 7 +- .../agent-plugin-core/python/telemetry.py | 87 ++++++++++++++++--- .../tests/test_uninitialised_identity.py | 68 +++++++++++++++ .../antigravity-plugin/core/hook_runner.py | 7 +- .../antigravity-plugin/core/telemetry.py | 87 ++++++++++++++++--- .../claude-code-plugin/core/hook_runner.py | 7 +- .../claude-code-plugin/core/telemetry.py | 87 ++++++++++++++++--- integrations/codex-plugin/core/hook_runner.py | 7 +- integrations/codex-plugin/core/telemetry.py | 87 ++++++++++++++++--- .../cursor-plugin/core/hook_runner.py | 7 +- integrations/cursor-plugin/core/telemetry.py | 87 ++++++++++++++++--- integrations/kimi-plugin/core/hook_runner.py | 7 +- integrations/kimi-plugin/core/telemetry.py | 87 ++++++++++++++++--- .../mem0-agent-plugin/core/telemetry.py | 87 ++++++++++++++++--- 14 files changed, 636 insertions(+), 83 deletions(-) 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..a24c97264 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -221,18 +221,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 +283,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 +306,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: @@ -563,8 +612,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 +635,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) diff --git a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py index 14f911bee..7d48dbcc3 100644 --- a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py +++ b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py @@ -177,3 +177,71 @@ def test_source_tag_defaults_agree_between_the_two_modules(): ) left, right = out.split() assert left == right == "KIMI_PLUGIN" + + +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") + + with tempfile.TemporaryDirectory() as tmp: + 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..a24c97264 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -221,18 +221,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 +283,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 +306,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: @@ -563,8 +612,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 +635,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) 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..a24c97264 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -221,18 +221,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 +283,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 +306,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: @@ -563,8 +612,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 +635,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) 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..a24c97264 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -221,18 +221,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 +283,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 +306,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: @@ -563,8 +612,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 +635,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) 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..a24c97264 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -221,18 +221,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 +283,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 +306,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: @@ -563,8 +612,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 +635,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) 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..a24c97264 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -221,18 +221,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 +283,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 +306,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: @@ -563,8 +612,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 +635,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) diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -221,18 +221,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 +283,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 +306,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: @@ -563,8 +612,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 +635,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)