From 7710a4e180de38813779fae0322914ca8b01a46e Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Wed, 16 Sep 2026 20:33:54 +0530 Subject: [PATCH 1/2] fix(plugins): publish the salt atomically, and omit the hash when there is none Review finding from @kartik-mem0 on this PR. O_CREAT|O_EXCL then write leaves a window where the salt file exists and is empty. Hooks are short-lived processes firing on every tool call and people run several agent windows, so a concurrent reader lands in that window, reads nothing, and falls back to a digest of the salt file's own path, memoized for its whole run. That path is guessable, so the race silently replaced the privacy control with something an attacker can compute, and hashed the same repository two ways depending on timing. The value is now written to a private temp file, fsynced, and published with os.link, which is atomic and fails if another process already published one. Link rather than replace, so losing the race adopts their salt instead of clobbering it. The temp file is removed either way. The derived fallback is gone rather than fixed. _scoped_digest returns "" when there is no salt and record() omits the property, because an unsalted digest over a git remote or a home-directory path is close to plaintext, and shipping one under a name that says hash is worse than sending nothing. Three tests: the racing reader never sees the name half-written, a second writer adopts the first's salt and leaves no temp file, and an unwritable data directory drops the property instead of emitting a weak one. The old test asserted the fallback behaviour and is replaced. 265 passed, 8 skipped. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/telemetry.py | 61 +++++++++++++------ .../antigravity-plugin/core/telemetry.py | 61 +++++++++++++------ .../claude-code-plugin/core/telemetry.py | 61 +++++++++++++------ .../tests/test_telemetry.py | 57 +++++++++++++++-- integrations/codex-plugin/core/telemetry.py | 61 +++++++++++++------ integrations/cursor-plugin/core/telemetry.py | 61 +++++++++++++------ integrations/kimi-plugin/core/telemetry.py | 61 +++++++++++++------ .../mem0-agent-plugin/core/telemetry.py | 61 +++++++++++++------ 8 files changed, 359 insertions(+), 125 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 239836604..2cb1056e0 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -116,36 +116,49 @@ def _install_salt() -> str: 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. + Published atomically, and there is deliberately no derived fallback. Creating + the file with O_CREAT|O_EXCL and then writing into it leaves a window where + the file exists and is empty, and a concurrent hook that reads it in that + window gets nothing. Falling back to a digest of the path would hand that + process a salt an attacker can compute, memoized for its whole run, which is + the privacy control this function exists to provide silently turning itself + off under load. The salt is written to a private temp file first and linked + into place, so the name either does not exist or already has the full value. + + Returns "" when it genuinely cannot persist. Callers omit the hash entirely + rather than emit an unsalted one. """ global _salt_cache if _salt_cache: return _salt_cache path = _salt_path() + temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp") try: path.parent.mkdir(parents=True, exist_ok=True) - handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + stream.flush() + os.fsync(stream.fileno()) try: - with os.fdopen(handle, "w", encoding="utf-8") as stream: - stream.write(uuid.uuid4().hex) + # Atomic claim: fails if another process already published one. + # os.link rather than replace, which would clobber theirs. + os.link(temporary, path) + except FileExistsError: + pass + except OSError: + pass + finally: + try: + temporary.unlink() 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 @@ -158,10 +171,17 @@ def _scoped_digest(value: str, length: int = 16) -> str: privacy control without the salt. Salting per install keeps every within-account join the analytics actually use and gives up only cross-machine joins on the same repository, which nothing computes. + + Returns "" when there is no salt, so record() omits the property. An + unsalted digest over this input space is close to plaintext, and emitting one + under a name that implies it is hashed is worse than sending nothing. """ if not value: return "" - return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length] + salt = _install_salt() + if not salt: + return "" + return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length] def _safe_value(value: Any) -> Any: @@ -252,10 +272,17 @@ def record( os=sys.platform, python_version=platform.python_version(), ) + # Assigned only when the digest is real. _scoped_digest returns "" when + # the salt could not be persisted, and an empty property is worse than an + # absent one: it survives the None filter below and reads as a value. if repo is not None: - properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", "")) + repo_hash = _scoped_digest(getattr(repo, "identity", "")) + if repo_hash: + properties["repo_hash"] = repo_hash if session_id: - properties["session_hash"] = _scoped_digest(session_id) + session_hash = _scoped_digest(session_id) + if session_hash: + properties["session_hash"] = session_hash line = json.dumps( { "event": f"{EVENT_PREFIX}.{event}", diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 239836604..2cb1056e0 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -116,36 +116,49 @@ def _install_salt() -> str: 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. + Published atomically, and there is deliberately no derived fallback. Creating + the file with O_CREAT|O_EXCL and then writing into it leaves a window where + the file exists and is empty, and a concurrent hook that reads it in that + window gets nothing. Falling back to a digest of the path would hand that + process a salt an attacker can compute, memoized for its whole run, which is + the privacy control this function exists to provide silently turning itself + off under load. The salt is written to a private temp file first and linked + into place, so the name either does not exist or already has the full value. + + Returns "" when it genuinely cannot persist. Callers omit the hash entirely + rather than emit an unsalted one. """ global _salt_cache if _salt_cache: return _salt_cache path = _salt_path() + temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp") try: path.parent.mkdir(parents=True, exist_ok=True) - handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + stream.flush() + os.fsync(stream.fileno()) try: - with os.fdopen(handle, "w", encoding="utf-8") as stream: - stream.write(uuid.uuid4().hex) + # Atomic claim: fails if another process already published one. + # os.link rather than replace, which would clobber theirs. + os.link(temporary, path) + except FileExistsError: + pass + except OSError: + pass + finally: + try: + temporary.unlink() 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 @@ -158,10 +171,17 @@ def _scoped_digest(value: str, length: int = 16) -> str: privacy control without the salt. Salting per install keeps every within-account join the analytics actually use and gives up only cross-machine joins on the same repository, which nothing computes. + + Returns "" when there is no salt, so record() omits the property. An + unsalted digest over this input space is close to plaintext, and emitting one + under a name that implies it is hashed is worse than sending nothing. """ if not value: return "" - return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length] + salt = _install_salt() + if not salt: + return "" + return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length] def _safe_value(value: Any) -> Any: @@ -252,10 +272,17 @@ def record( os=sys.platform, python_version=platform.python_version(), ) + # Assigned only when the digest is real. _scoped_digest returns "" when + # the salt could not be persisted, and an empty property is worse than an + # absent one: it survives the None filter below and reads as a value. if repo is not None: - properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", "")) + repo_hash = _scoped_digest(getattr(repo, "identity", "")) + if repo_hash: + properties["repo_hash"] = repo_hash if session_id: - properties["session_hash"] = _scoped_digest(session_id) + session_hash = _scoped_digest(session_id) + if session_hash: + properties["session_hash"] = session_hash line = json.dumps( { "event": f"{EVENT_PREFIX}.{event}", diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 239836604..2cb1056e0 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -116,36 +116,49 @@ def _install_salt() -> str: 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. + Published atomically, and there is deliberately no derived fallback. Creating + the file with O_CREAT|O_EXCL and then writing into it leaves a window where + the file exists and is empty, and a concurrent hook that reads it in that + window gets nothing. Falling back to a digest of the path would hand that + process a salt an attacker can compute, memoized for its whole run, which is + the privacy control this function exists to provide silently turning itself + off under load. The salt is written to a private temp file first and linked + into place, so the name either does not exist or already has the full value. + + Returns "" when it genuinely cannot persist. Callers omit the hash entirely + rather than emit an unsalted one. """ global _salt_cache if _salt_cache: return _salt_cache path = _salt_path() + temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp") try: path.parent.mkdir(parents=True, exist_ok=True) - handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + stream.flush() + os.fsync(stream.fileno()) try: - with os.fdopen(handle, "w", encoding="utf-8") as stream: - stream.write(uuid.uuid4().hex) + # Atomic claim: fails if another process already published one. + # os.link rather than replace, which would clobber theirs. + os.link(temporary, path) + except FileExistsError: + pass + except OSError: + pass + finally: + try: + temporary.unlink() 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 @@ -158,10 +171,17 @@ def _scoped_digest(value: str, length: int = 16) -> str: privacy control without the salt. Salting per install keeps every within-account join the analytics actually use and gives up only cross-machine joins on the same repository, which nothing computes. + + Returns "" when there is no salt, so record() omits the property. An + unsalted digest over this input space is close to plaintext, and emitting one + under a name that implies it is hashed is worse than sending nothing. """ if not value: return "" - return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length] + salt = _install_salt() + if not salt: + return "" + return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length] def _safe_value(value: Any) -> Any: @@ -252,10 +272,17 @@ def record( os=sys.platform, python_version=platform.python_version(), ) + # Assigned only when the digest is real. _scoped_digest returns "" when + # the salt could not be persisted, and an empty property is worse than an + # absent one: it survives the None filter below and reads as a value. if repo is not None: - properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", "")) + repo_hash = _scoped_digest(getattr(repo, "identity", "")) + if repo_hash: + properties["repo_hash"] = repo_hash if session_id: - properties["session_hash"] = _scoped_digest(session_id) + session_hash = _scoped_digest(session_id) + if session_hash: + properties["session_hash"] = session_hash line = json.dumps( { "event": f"{EVENT_PREFIX}.{event}", diff --git a/integrations/claude-code-plugin/tests/test_telemetry.py b/integrations/claude-code-plugin/tests/test_telemetry.py index eaa741ac5..b816b629e 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -309,14 +309,59 @@ def test_salt_does_not_touch_the_identity_file(isolated_env): 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. +def test_no_salt_means_no_hash_rather_than_an_unsalted_one(isolated_env, monkeypatch): + """A read-only data dir drops the property; it must not emit a weak digest. - Random per call is unbounded cardinality in PostHog, which is worse than no - salt at all. + The previous fallback was a digest of the salt file's own path, which an + attacker can compute, memoized for the whole process. A property named + repo_hash carrying an effectively unsalted digest is worse than no property: + it reads as protected and is not. """ telemetry._salt_cache = "" monkeypatch.setattr(telemetry.os, "open", lambda *a, **k: (_ for _ in ()).throw(OSError("read-only"))) - first = telemetry._install_salt() + + assert telemetry._install_salt() == "" + assert telemetry._scoped_digest("git@github.com:acme/secret.git") == "" + + +def test_a_half_written_salt_is_never_visible_to_another_process(isolated_env, monkeypatch): + """The window this closes: file created, value not yet written. + + O_CREAT|O_EXCL then write leaves the name present and empty in between. A + hook reading it there used to get "", fall back to the path digest and cache + that for its whole run, so the same repo hashed two ways depending on timing. + Publishing by link means the name either does not exist or is complete. + """ telemetry._salt_cache = "" - assert telemetry._install_salt() == first + salt_path = telemetry._salt_path() + observed = [] + + real_link = telemetry.os.link + + def observing_link(source, target): + # Stand where the racing reader stands: after the temp file is written, + # before the real name exists. + observed.append(salt_path.exists()) + return real_link(source, target) + + monkeypatch.setattr(telemetry.os, "link", observing_link) + salt = telemetry._install_salt() + + assert observed == [False], "the salt name existed before it held a value" + assert len(salt) == 32 + assert salt_path.read_text(encoding="utf-8").strip() == salt + + +def test_a_concurrent_writer_does_not_clobber_the_published_salt(isolated_env): + """Second process to finish must adopt the first one's salt, not replace it. + + os.link rather than os.replace is what makes losing the race harmless. + """ + telemetry._salt_cache = "" + first = telemetry._install_salt() + + telemetry._salt_cache = "" + second = telemetry._install_salt() + + assert second == first + assert not list(telemetry._salt_path().parent.glob("telemetry-salt.*.tmp")), "temp file left behind" diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 239836604..2cb1056e0 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -116,36 +116,49 @@ def _install_salt() -> str: 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. + Published atomically, and there is deliberately no derived fallback. Creating + the file with O_CREAT|O_EXCL and then writing into it leaves a window where + the file exists and is empty, and a concurrent hook that reads it in that + window gets nothing. Falling back to a digest of the path would hand that + process a salt an attacker can compute, memoized for its whole run, which is + the privacy control this function exists to provide silently turning itself + off under load. The salt is written to a private temp file first and linked + into place, so the name either does not exist or already has the full value. + + Returns "" when it genuinely cannot persist. Callers omit the hash entirely + rather than emit an unsalted one. """ global _salt_cache if _salt_cache: return _salt_cache path = _salt_path() + temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp") try: path.parent.mkdir(parents=True, exist_ok=True) - handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + stream.flush() + os.fsync(stream.fileno()) try: - with os.fdopen(handle, "w", encoding="utf-8") as stream: - stream.write(uuid.uuid4().hex) + # Atomic claim: fails if another process already published one. + # os.link rather than replace, which would clobber theirs. + os.link(temporary, path) + except FileExistsError: + pass + except OSError: + pass + finally: + try: + temporary.unlink() 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 @@ -158,10 +171,17 @@ def _scoped_digest(value: str, length: int = 16) -> str: privacy control without the salt. Salting per install keeps every within-account join the analytics actually use and gives up only cross-machine joins on the same repository, which nothing computes. + + Returns "" when there is no salt, so record() omits the property. An + unsalted digest over this input space is close to plaintext, and emitting one + under a name that implies it is hashed is worse than sending nothing. """ if not value: return "" - return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length] + salt = _install_salt() + if not salt: + return "" + return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length] def _safe_value(value: Any) -> Any: @@ -252,10 +272,17 @@ def record( os=sys.platform, python_version=platform.python_version(), ) + # Assigned only when the digest is real. _scoped_digest returns "" when + # the salt could not be persisted, and an empty property is worse than an + # absent one: it survives the None filter below and reads as a value. if repo is not None: - properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", "")) + repo_hash = _scoped_digest(getattr(repo, "identity", "")) + if repo_hash: + properties["repo_hash"] = repo_hash if session_id: - properties["session_hash"] = _scoped_digest(session_id) + session_hash = _scoped_digest(session_id) + if session_hash: + properties["session_hash"] = session_hash line = json.dumps( { "event": f"{EVENT_PREFIX}.{event}", diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 239836604..2cb1056e0 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -116,36 +116,49 @@ def _install_salt() -> str: 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. + Published atomically, and there is deliberately no derived fallback. Creating + the file with O_CREAT|O_EXCL and then writing into it leaves a window where + the file exists and is empty, and a concurrent hook that reads it in that + window gets nothing. Falling back to a digest of the path would hand that + process a salt an attacker can compute, memoized for its whole run, which is + the privacy control this function exists to provide silently turning itself + off under load. The salt is written to a private temp file first and linked + into place, so the name either does not exist or already has the full value. + + Returns "" when it genuinely cannot persist. Callers omit the hash entirely + rather than emit an unsalted one. """ global _salt_cache if _salt_cache: return _salt_cache path = _salt_path() + temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp") try: path.parent.mkdir(parents=True, exist_ok=True) - handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + stream.flush() + os.fsync(stream.fileno()) try: - with os.fdopen(handle, "w", encoding="utf-8") as stream: - stream.write(uuid.uuid4().hex) + # Atomic claim: fails if another process already published one. + # os.link rather than replace, which would clobber theirs. + os.link(temporary, path) + except FileExistsError: + pass + except OSError: + pass + finally: + try: + temporary.unlink() 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 @@ -158,10 +171,17 @@ def _scoped_digest(value: str, length: int = 16) -> str: privacy control without the salt. Salting per install keeps every within-account join the analytics actually use and gives up only cross-machine joins on the same repository, which nothing computes. + + Returns "" when there is no salt, so record() omits the property. An + unsalted digest over this input space is close to plaintext, and emitting one + under a name that implies it is hashed is worse than sending nothing. """ if not value: return "" - return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length] + salt = _install_salt() + if not salt: + return "" + return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length] def _safe_value(value: Any) -> Any: @@ -252,10 +272,17 @@ def record( os=sys.platform, python_version=platform.python_version(), ) + # Assigned only when the digest is real. _scoped_digest returns "" when + # the salt could not be persisted, and an empty property is worse than an + # absent one: it survives the None filter below and reads as a value. if repo is not None: - properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", "")) + repo_hash = _scoped_digest(getattr(repo, "identity", "")) + if repo_hash: + properties["repo_hash"] = repo_hash if session_id: - properties["session_hash"] = _scoped_digest(session_id) + session_hash = _scoped_digest(session_id) + if session_hash: + properties["session_hash"] = session_hash line = json.dumps( { "event": f"{EVENT_PREFIX}.{event}", diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 239836604..2cb1056e0 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -116,36 +116,49 @@ def _install_salt() -> str: 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. + Published atomically, and there is deliberately no derived fallback. Creating + the file with O_CREAT|O_EXCL and then writing into it leaves a window where + the file exists and is empty, and a concurrent hook that reads it in that + window gets nothing. Falling back to a digest of the path would hand that + process a salt an attacker can compute, memoized for its whole run, which is + the privacy control this function exists to provide silently turning itself + off under load. The salt is written to a private temp file first and linked + into place, so the name either does not exist or already has the full value. + + Returns "" when it genuinely cannot persist. Callers omit the hash entirely + rather than emit an unsalted one. """ global _salt_cache if _salt_cache: return _salt_cache path = _salt_path() + temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp") try: path.parent.mkdir(parents=True, exist_ok=True) - handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + stream.flush() + os.fsync(stream.fileno()) try: - with os.fdopen(handle, "w", encoding="utf-8") as stream: - stream.write(uuid.uuid4().hex) + # Atomic claim: fails if another process already published one. + # os.link rather than replace, which would clobber theirs. + os.link(temporary, path) + except FileExistsError: + pass + except OSError: + pass + finally: + try: + temporary.unlink() 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 @@ -158,10 +171,17 @@ def _scoped_digest(value: str, length: int = 16) -> str: privacy control without the salt. Salting per install keeps every within-account join the analytics actually use and gives up only cross-machine joins on the same repository, which nothing computes. + + Returns "" when there is no salt, so record() omits the property. An + unsalted digest over this input space is close to plaintext, and emitting one + under a name that implies it is hashed is worse than sending nothing. """ if not value: return "" - return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length] + salt = _install_salt() + if not salt: + return "" + return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length] def _safe_value(value: Any) -> Any: @@ -252,10 +272,17 @@ def record( os=sys.platform, python_version=platform.python_version(), ) + # Assigned only when the digest is real. _scoped_digest returns "" when + # the salt could not be persisted, and an empty property is worse than an + # absent one: it survives the None filter below and reads as a value. if repo is not None: - properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", "")) + repo_hash = _scoped_digest(getattr(repo, "identity", "")) + if repo_hash: + properties["repo_hash"] = repo_hash if session_id: - properties["session_hash"] = _scoped_digest(session_id) + session_hash = _scoped_digest(session_id) + if session_hash: + properties["session_hash"] = session_hash line = json.dumps( { "event": f"{EVENT_PREFIX}.{event}", diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 239836604..2cb1056e0 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -116,36 +116,49 @@ def _install_salt() -> str: 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. + Published atomically, and there is deliberately no derived fallback. Creating + the file with O_CREAT|O_EXCL and then writing into it leaves a window where + the file exists and is empty, and a concurrent hook that reads it in that + window gets nothing. Falling back to a digest of the path would hand that + process a salt an attacker can compute, memoized for its whole run, which is + the privacy control this function exists to provide silently turning itself + off under load. The salt is written to a private temp file first and linked + into place, so the name either does not exist or already has the full value. + + Returns "" when it genuinely cannot persist. Callers omit the hash entirely + rather than emit an unsalted one. """ global _salt_cache if _salt_cache: return _salt_cache path = _salt_path() + temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp") try: path.parent.mkdir(parents=True, exist_ok=True) - handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + with os.fdopen(handle, "w", encoding="utf-8") as stream: + stream.write(uuid.uuid4().hex) + stream.flush() + os.fsync(stream.fileno()) try: - with os.fdopen(handle, "w", encoding="utf-8") as stream: - stream.write(uuid.uuid4().hex) + # Atomic claim: fails if another process already published one. + # os.link rather than replace, which would clobber theirs. + os.link(temporary, path) + except FileExistsError: + pass + except OSError: + pass + finally: + try: + temporary.unlink() 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 @@ -158,10 +171,17 @@ def _scoped_digest(value: str, length: int = 16) -> str: privacy control without the salt. Salting per install keeps every within-account join the analytics actually use and gives up only cross-machine joins on the same repository, which nothing computes. + + Returns "" when there is no salt, so record() omits the property. An + unsalted digest over this input space is close to plaintext, and emitting one + under a name that implies it is hashed is worse than sending nothing. """ if not value: return "" - return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length] + salt = _install_salt() + if not salt: + return "" + return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length] def _safe_value(value: Any) -> Any: @@ -252,10 +272,17 @@ def record( os=sys.platform, python_version=platform.python_version(), ) + # Assigned only when the digest is real. _scoped_digest returns "" when + # the salt could not be persisted, and an empty property is worse than an + # absent one: it survives the None filter below and reads as a value. if repo is not None: - properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", "")) + repo_hash = _scoped_digest(getattr(repo, "identity", "")) + if repo_hash: + properties["repo_hash"] = repo_hash if session_id: - properties["session_hash"] = _scoped_digest(session_id) + session_hash = _scoped_digest(session_id) + if session_hash: + properties["session_hash"] = session_hash line = json.dumps( { "event": f"{EVENT_PREFIX}.{event}", From d8c99fb405d59adc28225ce7c603d0b13fe29533 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Wed, 16 Sep 2026 20:35:13 +0530 Subject: [PATCH 2/2] fix(plugins): check the lease before judging a claim exhausted Review finding from @kartik-mem0 on this PR, and the most serious one: it loses events, which is what this PR exists to prevent. _claim_parked judged exhaustion before liveness. Claiming a parked file bumps its attempt count and refreshes its mtime, so the moment a sender takes the final attempt the file looks exhausted to every other sender while its owner is actively draining it. The second sender unlinked it, and everything in that batch was gone. The liveness check now runs first, so a batch under a live lease is skipped whatever its attempt count. The cleanup is deferred, not cancelled: once the lease lapses, the same exhausted file is reaped on a later run. Two tests. The first walks a batch to the final attempt and asserts a second sender neither takes it nor deletes it, and that the events are still in it. The second asserts an abandoned exhausted batch is still discarded once its lease lapses, which is the over-correction to guard against. Confirmed the first fails against the previous ordering. 288 passed, 8 skipped. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/telemetry.py | 11 +++- .../tests/test_spool_delivery.py | 61 +++++++++++++++++++ .../antigravity-plugin/core/telemetry.py | 11 +++- .../claude-code-plugin/core/telemetry.py | 11 +++- integrations/codex-plugin/core/telemetry.py | 11 +++- integrations/cursor-plugin/core/telemetry.py | 11 +++- integrations/kimi-plugin/core/telemetry.py | 11 +++- .../mem0-agent-plugin/core/telemetry.py | 11 +++- 8 files changed, 117 insertions(+), 21 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue # 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 @@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None: except OSError: pass continue - if age < CLAIM_STALE_SECONDS: - # Someone else holds a live lease on it. - continue claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) diff --git a/integrations/agent-plugin-core/tests/test_spool_delivery.py b/integrations/agent-plugin-core/tests/test_spool_delivery.py index a51027904..a7646cc31 100644 --- a/integrations/agent-plugin-core/tests/test_spool_delivery.py +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -104,6 +104,67 @@ def test_a_fresh_claim_is_not_immediately_stealable(telemetry): assert telemetry._claim_parked(first.parent) is None +def test_a_live_final_attempt_is_not_deleted_by_another_sender(telemetry): + """Review finding: exhaustion was judged before liveness, so owners lost batches. + + Claiming a parked file bumps its attempt count and refreshes its mtime. Once + the count reaches the budget, the owner draining it looked exhausted to every + other sender, which unlinked the file out from under it. Everything in that + batch was gone, which is precisely the loss this PR exists to stop. + """ + telemetry.record("search", reason="owned-by-the-first-sender") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + + claim = telemetry._claim_spool() + assert claim is not None + + # Walk it to the final attempt, ageing it each round so it can be re-claimed. + # _claim_spool hands back a0 and _release_claim keeps the name, so it takes + # one full round per attempt to reach the budget. + for _ in range(telemetry.MAX_CLAIM_ATTEMPTS): + # Carry the marker through each rewrite so the final assertion proves the + # events survived, not merely that some file with the right name did. + telemetry._release_claim(claim, [{"event": "code.search", "uuid": "owned-by-the-first-sender"}]) + parked = sorted(claim.parent.glob("telemetry-*.sending")) + assert parked, "the batch was dropped while still inside its budget" + os.utime(parked[0], (stale, stale)) + claim = telemetry._claim_parked(claim.parent) + assert claim is not None + + assert telemetry._claim_attempt(claim) >= telemetry.MAX_CLAIM_ATTEMPTS + assert claim.exists() + + # The owner is draining it right now: fresh mtime, live lease. + second_sender = telemetry._claim_parked(claim.parent) + + assert second_sender is None, "a second sender took a batch under a live lease" + assert claim.exists(), "a second sender deleted a batch its owner was draining" + assert "owned-by-the-first-sender" in claim.read_text(encoding="utf-8") + + +def test_an_exhausted_batch_is_still_discarded_once_its_lease_lapses(telemetry): + """The liveness check must defer the cleanup, not cancel it. + + Guards the obvious over-correction: skipping live claims is only safe if an + abandoned one at the same attempt count is still reaped on a later run. + """ + telemetry.record("search") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + + claim = telemetry._claim_spool() + assert claim is not None + exhausted = claim.parent / telemetry._claim_name(telemetry.MAX_CLAIM_ATTEMPTS) + claim.replace(exhausted) + os.utime(exhausted, (stale, stale)) + + assert telemetry._claim_parked(exhausted.parent) is None + assert not exhausted.exists(), "an abandoned exhausted batch was left behind forever" + + def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry): """Defect 6: parked claims were only reachable when no spool existed. diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue # 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 @@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None: except OSError: pass continue - if age < CLAIM_STALE_SECONDS: - # Someone else holds a live lease on it. - continue claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue # 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 @@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None: except OSError: pass continue - if age < CLAIM_STALE_SECONDS: - # Someone else holds a live lease on it. - continue claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue # 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 @@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None: except OSError: pass continue - if age < CLAIM_STALE_SECONDS: - # Someone else holds a live lease on it. - continue claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue # 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 @@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None: except OSError: pass continue - if age < CLAIM_STALE_SECONDS: - # Someone else holds a live lease on it. - continue claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue # 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 @@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None: except OSError: pass continue - if age < CLAIM_STALE_SECONDS: - # Someone else holds a live lease on it. - continue claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue # 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 @@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None: except OSError: pass continue - if age < CLAIM_STALE_SECONDS: - # Someone else holds a live lease on it. - continue claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim)