diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 8cce9cbc7..274174a9e 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -145,36 +145,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 @@ -187,10 +200,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,6 +272,25 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: return created +def _rotate_anonymous_id(identity: dict[str, str]) -> str: + """Mint a fresh anonymous id because the account context is gone. + + The previous id may already have been merged into a person profile by an + $identify, and that merge is permanent. Reusing it after a logout or a key + change attributes everything that follows to the account that just went + away, which is the same misattribution the key fingerprint exists to stop, + only arriving through the anonymous path instead. + + `aliased` is cleared with it: the new id has never been merged, so it is + eligible to be aliased into whatever account comes next. + """ + created = f"code-anon-{uuid.uuid4().hex}" + identity["anonymous_id"] = created + identity.pop("aliased", None) + _write_identity(identity) + return created + + def _install_state_path() -> Path: return memory_core.data_dir() / "install-state.json" @@ -380,11 +419,19 @@ def claim_version_change() -> str | None: state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") try: - temporary = path.with_suffix(f".{os.getpid()}.tmp") temporary.write_text(json.dumps(state), encoding="utf-8") temporary.replace(path) except OSError: + # Release the claim. The marker still records the old version, so + # without this the sentinel makes claim_version_change return early on + # every later run and this version's upgrade is never recorded again. + for leftover in (sentinel, temporary): + try: + leftover.unlink() + except OSError: + pass return None return previous @@ -418,10 +465,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}", @@ -566,6 +620,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 @@ -576,9 +638,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) @@ -695,26 +754,43 @@ def resolve_distinct_id() -> tuple[str, str]: 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. + # Rows written before fingerprints existed. Verify rather than + # adopt: a key changed before the upgrade would otherwise bind the + # new key to the previous account's email, permanently, and the + # fingerprint would then agree with itself forever after. + verified = _resolve_email(key) + if not verified: + # Offline, firewalled, or the API is down. Keep the previous + # behaviour and retry on the next flush rather than dropping a + # real account attribution. Safe because the same network that + # failed /v1/ping/ is about to fail the PostHog POST, so nothing + # is delivered under the unverified identity in the meantime. + return email, "" + identity["email"] = verified identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, "" + return verified, "" if not key: # No key to verify the account with; do not keep attributing to it. if email: identity.pop("email", None) identity.pop("key_fingerprint", None) - _write_identity(identity) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" resolved = _resolve_email(key) if not resolved: # The key changed and will not resolve (revoked, offline, API down). - # Do not keep attributing to the previous account. + # Reaching here with an email means the recorded fingerprint disagreed, + # so the key really did change. Drop the account and rotate: the stored + # anonymous id may already be merged into that account's person, and + # reusing it would keep the events on the profile we are trying to + # leave. + if email: + identity.pop("email", None) + identity.pop("key_fingerprint", None) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" # Alias only when going anonymous -> email for the first time. Once an anon 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 8cce9cbc7..274174a9e 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -145,36 +145,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 @@ -187,10 +200,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,6 +272,25 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: return created +def _rotate_anonymous_id(identity: dict[str, str]) -> str: + """Mint a fresh anonymous id because the account context is gone. + + The previous id may already have been merged into a person profile by an + $identify, and that merge is permanent. Reusing it after a logout or a key + change attributes everything that follows to the account that just went + away, which is the same misattribution the key fingerprint exists to stop, + only arriving through the anonymous path instead. + + `aliased` is cleared with it: the new id has never been merged, so it is + eligible to be aliased into whatever account comes next. + """ + created = f"code-anon-{uuid.uuid4().hex}" + identity["anonymous_id"] = created + identity.pop("aliased", None) + _write_identity(identity) + return created + + def _install_state_path() -> Path: return memory_core.data_dir() / "install-state.json" @@ -380,11 +419,19 @@ def claim_version_change() -> str | None: state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") try: - temporary = path.with_suffix(f".{os.getpid()}.tmp") temporary.write_text(json.dumps(state), encoding="utf-8") temporary.replace(path) except OSError: + # Release the claim. The marker still records the old version, so + # without this the sentinel makes claim_version_change return early on + # every later run and this version's upgrade is never recorded again. + for leftover in (sentinel, temporary): + try: + leftover.unlink() + except OSError: + pass return None return previous @@ -418,10 +465,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}", @@ -566,6 +620,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 @@ -576,9 +638,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) @@ -695,26 +754,43 @@ def resolve_distinct_id() -> tuple[str, str]: 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. + # Rows written before fingerprints existed. Verify rather than + # adopt: a key changed before the upgrade would otherwise bind the + # new key to the previous account's email, permanently, and the + # fingerprint would then agree with itself forever after. + verified = _resolve_email(key) + if not verified: + # Offline, firewalled, or the API is down. Keep the previous + # behaviour and retry on the next flush rather than dropping a + # real account attribution. Safe because the same network that + # failed /v1/ping/ is about to fail the PostHog POST, so nothing + # is delivered under the unverified identity in the meantime. + return email, "" + identity["email"] = verified identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, "" + return verified, "" if not key: # No key to verify the account with; do not keep attributing to it. if email: identity.pop("email", None) identity.pop("key_fingerprint", None) - _write_identity(identity) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" resolved = _resolve_email(key) if not resolved: # The key changed and will not resolve (revoked, offline, API down). - # Do not keep attributing to the previous account. + # Reaching here with an email means the recorded fingerprint disagreed, + # so the key really did change. Drop the account and rotate: the stored + # anonymous id may already be merged into that account's person, and + # reusing it would keep the events on the profile we are trying to + # leave. + if email: + identity.pop("email", None) + identity.pop("key_fingerprint", None) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" # Alias only when going anonymous -> email for the first time. Once an anon diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 8cce9cbc7..274174a9e 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -145,36 +145,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 @@ -187,10 +200,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,6 +272,25 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: return created +def _rotate_anonymous_id(identity: dict[str, str]) -> str: + """Mint a fresh anonymous id because the account context is gone. + + The previous id may already have been merged into a person profile by an + $identify, and that merge is permanent. Reusing it after a logout or a key + change attributes everything that follows to the account that just went + away, which is the same misattribution the key fingerprint exists to stop, + only arriving through the anonymous path instead. + + `aliased` is cleared with it: the new id has never been merged, so it is + eligible to be aliased into whatever account comes next. + """ + created = f"code-anon-{uuid.uuid4().hex}" + identity["anonymous_id"] = created + identity.pop("aliased", None) + _write_identity(identity) + return created + + def _install_state_path() -> Path: return memory_core.data_dir() / "install-state.json" @@ -380,11 +419,19 @@ def claim_version_change() -> str | None: state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") try: - temporary = path.with_suffix(f".{os.getpid()}.tmp") temporary.write_text(json.dumps(state), encoding="utf-8") temporary.replace(path) except OSError: + # Release the claim. The marker still records the old version, so + # without this the sentinel makes claim_version_change return early on + # every later run and this version's upgrade is never recorded again. + for leftover in (sentinel, temporary): + try: + leftover.unlink() + except OSError: + pass return None return previous @@ -418,10 +465,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}", @@ -566,6 +620,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 @@ -576,9 +638,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) @@ -695,26 +754,43 @@ def resolve_distinct_id() -> tuple[str, str]: 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. + # Rows written before fingerprints existed. Verify rather than + # adopt: a key changed before the upgrade would otherwise bind the + # new key to the previous account's email, permanently, and the + # fingerprint would then agree with itself forever after. + verified = _resolve_email(key) + if not verified: + # Offline, firewalled, or the API is down. Keep the previous + # behaviour and retry on the next flush rather than dropping a + # real account attribution. Safe because the same network that + # failed /v1/ping/ is about to fail the PostHog POST, so nothing + # is delivered under the unverified identity in the meantime. + return email, "" + identity["email"] = verified identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, "" + return verified, "" if not key: # No key to verify the account with; do not keep attributing to it. if email: identity.pop("email", None) identity.pop("key_fingerprint", None) - _write_identity(identity) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" resolved = _resolve_email(key) if not resolved: # The key changed and will not resolve (revoked, offline, API down). - # Do not keep attributing to the previous account. + # Reaching here with an email means the recorded fingerprint disagreed, + # so the key really did change. Drop the account and rotate: the stored + # anonymous id may already be merged into that account's person, and + # reusing it would keep the events on the profile we are trying to + # leave. + if email: + identity.pop("email", None) + identity.pop("key_fingerprint", None) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" # Alias only when going anonymous -> email for the first time. Once an anon diff --git a/integrations/claude-code-plugin/tests/test_telemetry.py b/integrations/claude-code-plugin/tests/test_telemetry.py index efd2e11c0..8d83e4a7b 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -265,6 +265,109 @@ def test_an_unresolvable_key_falls_back_to_the_anonymous_id(isolated_env, monkey assert telemetry.resolve_distinct_id()[0].startswith("code-anon-") +def test_logging_out_does_not_leave_events_on_the_previous_account(isolated_env, monkeypatch): + """Review finding: clearing the email kept an id already merged into a person. + + The anonymous id is offered to PostHog as $anon_distinct_id on first sign-in, + and that merge is permanent. Keeping it after the key goes away means every + later anonymous event lands on the account that just left. + """ + # Run anonymously first, which is the only way an id exists to be merged. + merged = telemetry.anonymous_id() + + monkeypatch.setenv("MEM0_API_KEY", "key-for-account-a") + with patch.object(telemetry, "_resolve_email", lambda key: "a@example.com"): + identified, alias = telemetry.resolve_distinct_id() + assert identified == "a@example.com" + assert alias == merged, "the anonymous id was merged into this account" + + monkeypatch.delenv("MEM0_API_KEY", raising=False) + after_logout, logout_alias = telemetry.resolve_distinct_id() + + assert after_logout.startswith("code-anon-") + assert after_logout != merged, "reused an id already merged into the previous account" + assert logout_alias == "" + assert "aliased" not in telemetry._read_identity(), "rotated id must be aliasable again" + + +def test_a_changed_key_that_will_not_resolve_rotates_the_anonymous_id(isolated_env, monkeypatch): + """Same leak by the other route: fingerprint disagrees and the lookup fails.""" + merged = telemetry.anonymous_id() + monkeypatch.setenv("MEM0_API_KEY", "key-for-account-a") + with patch.object(telemetry, "_resolve_email", lambda key: "a@example.com"): + telemetry.resolve_distinct_id() + + monkeypatch.setenv("MEM0_API_KEY", "key-for-account-b") + with patch.object(telemetry, "_resolve_email", lambda key: ""): + after, alias = telemetry.resolve_distinct_id() + + assert after.startswith("code-anon-") + assert after != merged + assert alias == "" + assert "email" not in telemetry._read_identity() + + +def test_a_legacy_cached_email_is_verified_before_the_key_is_bound(isolated_env, monkeypatch): + """Review finding: a key changed before upgrading bound the wrong account. + + Rows written before fingerprints existed carry an email and no fingerprint. + Adopting the current key without checking pinned that key to the previous + account's email, and every run after that agreed with itself. + """ + telemetry._write_identity({"email": "old@example.com", "anonymous_id": "code-anon-seed"}) + monkeypatch.setenv("MEM0_API_KEY", "key-for-account-b") + + with patch.object(telemetry, "_resolve_email", lambda key: "new@example.com"): + resolved, alias = telemetry.resolve_distinct_id() + + assert resolved == "new@example.com" + assert alias == "", "email to email must never alias; it merges two real people" + stored = telemetry._read_identity() + assert stored["email"] == "new@example.com" + assert stored["key_fingerprint"] == telemetry._digest("key-for-account-b") + + +def test_a_legacy_row_keeps_working_when_the_account_cannot_be_checked(isolated_env, monkeypatch): + """Firewalled users must not lose attribution, and must not bind unverified. + + The same network that fails /v1/ping/ fails the PostHog POST, so nothing is + delivered under the unverified identity while this holds. + """ + telemetry._write_identity({"email": "old@example.com"}) + monkeypatch.setenv("MEM0_API_KEY", "key-for-account-b") + + with patch.object(telemetry, "_resolve_email", lambda key: ""): + resolved, _ = telemetry.resolve_distinct_id() + + assert resolved == "old@example.com" + assert "key_fingerprint" not in telemetry._read_identity(), "bound an unverified key" + + +def test_a_failed_upgrade_claim_can_be_retried(isolated_env, monkeypatch): + """Review finding: a failed rewrite left the sentinel and suppressed forever. + + claim_version_change returns early on FileExistsError, and the marker still + holds the old version, so the upgrade for that version was never recorded + again on that machine. + """ + telemetry.claim_install() + state_path = memory_core.data_dir() / "install-state.json" + state = json.loads(state_path.read_text()) + state["plugin_version"] = "0.0.1-old" + state_path.write_text(json.dumps(state), encoding="utf-8") + + real_replace = Path.replace + + def failing_replace(self, target): + raise OSError("disk full") + + monkeypatch.setattr(Path, "replace", failing_replace) + assert telemetry.claim_version_change() is None + + monkeypatch.setattr(Path, "replace", real_replace) + assert telemetry.claim_version_change() == "0.0.1-old", "sentinel suppressed the retry" + + def test_first_run_is_not_flipped_by_writing_the_identity_file(isolated_env): """The identity file is written by a successful flush, not by recording. @@ -350,14 +453,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 8cce9cbc7..274174a9e 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -145,36 +145,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 @@ -187,10 +200,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,6 +272,25 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: return created +def _rotate_anonymous_id(identity: dict[str, str]) -> str: + """Mint a fresh anonymous id because the account context is gone. + + The previous id may already have been merged into a person profile by an + $identify, and that merge is permanent. Reusing it after a logout or a key + change attributes everything that follows to the account that just went + away, which is the same misattribution the key fingerprint exists to stop, + only arriving through the anonymous path instead. + + `aliased` is cleared with it: the new id has never been merged, so it is + eligible to be aliased into whatever account comes next. + """ + created = f"code-anon-{uuid.uuid4().hex}" + identity["anonymous_id"] = created + identity.pop("aliased", None) + _write_identity(identity) + return created + + def _install_state_path() -> Path: return memory_core.data_dir() / "install-state.json" @@ -380,11 +419,19 @@ def claim_version_change() -> str | None: state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") try: - temporary = path.with_suffix(f".{os.getpid()}.tmp") temporary.write_text(json.dumps(state), encoding="utf-8") temporary.replace(path) except OSError: + # Release the claim. The marker still records the old version, so + # without this the sentinel makes claim_version_change return early on + # every later run and this version's upgrade is never recorded again. + for leftover in (sentinel, temporary): + try: + leftover.unlink() + except OSError: + pass return None return previous @@ -418,10 +465,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}", @@ -566,6 +620,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 @@ -576,9 +638,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) @@ -695,26 +754,43 @@ def resolve_distinct_id() -> tuple[str, str]: 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. + # Rows written before fingerprints existed. Verify rather than + # adopt: a key changed before the upgrade would otherwise bind the + # new key to the previous account's email, permanently, and the + # fingerprint would then agree with itself forever after. + verified = _resolve_email(key) + if not verified: + # Offline, firewalled, or the API is down. Keep the previous + # behaviour and retry on the next flush rather than dropping a + # real account attribution. Safe because the same network that + # failed /v1/ping/ is about to fail the PostHog POST, so nothing + # is delivered under the unverified identity in the meantime. + return email, "" + identity["email"] = verified identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, "" + return verified, "" if not key: # No key to verify the account with; do not keep attributing to it. if email: identity.pop("email", None) identity.pop("key_fingerprint", None) - _write_identity(identity) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" resolved = _resolve_email(key) if not resolved: # The key changed and will not resolve (revoked, offline, API down). - # Do not keep attributing to the previous account. + # Reaching here with an email means the recorded fingerprint disagreed, + # so the key really did change. Drop the account and rotate: the stored + # anonymous id may already be merged into that account's person, and + # reusing it would keep the events on the profile we are trying to + # leave. + if email: + identity.pop("email", None) + identity.pop("key_fingerprint", None) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" # Alias only when going anonymous -> email for the first time. Once an anon diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 8cce9cbc7..274174a9e 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -145,36 +145,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 @@ -187,10 +200,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,6 +272,25 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: return created +def _rotate_anonymous_id(identity: dict[str, str]) -> str: + """Mint a fresh anonymous id because the account context is gone. + + The previous id may already have been merged into a person profile by an + $identify, and that merge is permanent. Reusing it after a logout or a key + change attributes everything that follows to the account that just went + away, which is the same misattribution the key fingerprint exists to stop, + only arriving through the anonymous path instead. + + `aliased` is cleared with it: the new id has never been merged, so it is + eligible to be aliased into whatever account comes next. + """ + created = f"code-anon-{uuid.uuid4().hex}" + identity["anonymous_id"] = created + identity.pop("aliased", None) + _write_identity(identity) + return created + + def _install_state_path() -> Path: return memory_core.data_dir() / "install-state.json" @@ -380,11 +419,19 @@ def claim_version_change() -> str | None: state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") try: - temporary = path.with_suffix(f".{os.getpid()}.tmp") temporary.write_text(json.dumps(state), encoding="utf-8") temporary.replace(path) except OSError: + # Release the claim. The marker still records the old version, so + # without this the sentinel makes claim_version_change return early on + # every later run and this version's upgrade is never recorded again. + for leftover in (sentinel, temporary): + try: + leftover.unlink() + except OSError: + pass return None return previous @@ -418,10 +465,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}", @@ -566,6 +620,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 @@ -576,9 +638,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) @@ -695,26 +754,43 @@ def resolve_distinct_id() -> tuple[str, str]: 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. + # Rows written before fingerprints existed. Verify rather than + # adopt: a key changed before the upgrade would otherwise bind the + # new key to the previous account's email, permanently, and the + # fingerprint would then agree with itself forever after. + verified = _resolve_email(key) + if not verified: + # Offline, firewalled, or the API is down. Keep the previous + # behaviour and retry on the next flush rather than dropping a + # real account attribution. Safe because the same network that + # failed /v1/ping/ is about to fail the PostHog POST, so nothing + # is delivered under the unverified identity in the meantime. + return email, "" + identity["email"] = verified identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, "" + return verified, "" if not key: # No key to verify the account with; do not keep attributing to it. if email: identity.pop("email", None) identity.pop("key_fingerprint", None) - _write_identity(identity) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" resolved = _resolve_email(key) if not resolved: # The key changed and will not resolve (revoked, offline, API down). - # Do not keep attributing to the previous account. + # Reaching here with an email means the recorded fingerprint disagreed, + # so the key really did change. Drop the account and rotate: the stored + # anonymous id may already be merged into that account's person, and + # reusing it would keep the events on the profile we are trying to + # leave. + if email: + identity.pop("email", None) + identity.pop("key_fingerprint", None) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" # Alias only when going anonymous -> email for the first time. Once an anon diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 8cce9cbc7..274174a9e 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -145,36 +145,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 @@ -187,10 +200,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,6 +272,25 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: return created +def _rotate_anonymous_id(identity: dict[str, str]) -> str: + """Mint a fresh anonymous id because the account context is gone. + + The previous id may already have been merged into a person profile by an + $identify, and that merge is permanent. Reusing it after a logout or a key + change attributes everything that follows to the account that just went + away, which is the same misattribution the key fingerprint exists to stop, + only arriving through the anonymous path instead. + + `aliased` is cleared with it: the new id has never been merged, so it is + eligible to be aliased into whatever account comes next. + """ + created = f"code-anon-{uuid.uuid4().hex}" + identity["anonymous_id"] = created + identity.pop("aliased", None) + _write_identity(identity) + return created + + def _install_state_path() -> Path: return memory_core.data_dir() / "install-state.json" @@ -380,11 +419,19 @@ def claim_version_change() -> str | None: state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") try: - temporary = path.with_suffix(f".{os.getpid()}.tmp") temporary.write_text(json.dumps(state), encoding="utf-8") temporary.replace(path) except OSError: + # Release the claim. The marker still records the old version, so + # without this the sentinel makes claim_version_change return early on + # every later run and this version's upgrade is never recorded again. + for leftover in (sentinel, temporary): + try: + leftover.unlink() + except OSError: + pass return None return previous @@ -418,10 +465,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}", @@ -566,6 +620,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 @@ -576,9 +638,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) @@ -695,26 +754,43 @@ def resolve_distinct_id() -> tuple[str, str]: 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. + # Rows written before fingerprints existed. Verify rather than + # adopt: a key changed before the upgrade would otherwise bind the + # new key to the previous account's email, permanently, and the + # fingerprint would then agree with itself forever after. + verified = _resolve_email(key) + if not verified: + # Offline, firewalled, or the API is down. Keep the previous + # behaviour and retry on the next flush rather than dropping a + # real account attribution. Safe because the same network that + # failed /v1/ping/ is about to fail the PostHog POST, so nothing + # is delivered under the unverified identity in the meantime. + return email, "" + identity["email"] = verified identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, "" + return verified, "" if not key: # No key to verify the account with; do not keep attributing to it. if email: identity.pop("email", None) identity.pop("key_fingerprint", None) - _write_identity(identity) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" resolved = _resolve_email(key) if not resolved: # The key changed and will not resolve (revoked, offline, API down). - # Do not keep attributing to the previous account. + # Reaching here with an email means the recorded fingerprint disagreed, + # so the key really did change. Drop the account and rotate: the stored + # anonymous id may already be merged into that account's person, and + # reusing it would keep the events on the profile we are trying to + # leave. + if email: + identity.pop("email", None) + identity.pop("key_fingerprint", None) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" # Alias only when going anonymous -> email for the first time. Once an anon diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 8cce9cbc7..274174a9e 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -145,36 +145,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 @@ -187,10 +200,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,6 +272,25 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: return created +def _rotate_anonymous_id(identity: dict[str, str]) -> str: + """Mint a fresh anonymous id because the account context is gone. + + The previous id may already have been merged into a person profile by an + $identify, and that merge is permanent. Reusing it after a logout or a key + change attributes everything that follows to the account that just went + away, which is the same misattribution the key fingerprint exists to stop, + only arriving through the anonymous path instead. + + `aliased` is cleared with it: the new id has never been merged, so it is + eligible to be aliased into whatever account comes next. + """ + created = f"code-anon-{uuid.uuid4().hex}" + identity["anonymous_id"] = created + identity.pop("aliased", None) + _write_identity(identity) + return created + + def _install_state_path() -> Path: return memory_core.data_dir() / "install-state.json" @@ -380,11 +419,19 @@ def claim_version_change() -> str | None: state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") try: - temporary = path.with_suffix(f".{os.getpid()}.tmp") temporary.write_text(json.dumps(state), encoding="utf-8") temporary.replace(path) except OSError: + # Release the claim. The marker still records the old version, so + # without this the sentinel makes claim_version_change return early on + # every later run and this version's upgrade is never recorded again. + for leftover in (sentinel, temporary): + try: + leftover.unlink() + except OSError: + pass return None return previous @@ -418,10 +465,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}", @@ -566,6 +620,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 @@ -576,9 +638,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) @@ -695,26 +754,43 @@ def resolve_distinct_id() -> tuple[str, str]: 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. + # Rows written before fingerprints existed. Verify rather than + # adopt: a key changed before the upgrade would otherwise bind the + # new key to the previous account's email, permanently, and the + # fingerprint would then agree with itself forever after. + verified = _resolve_email(key) + if not verified: + # Offline, firewalled, or the API is down. Keep the previous + # behaviour and retry on the next flush rather than dropping a + # real account attribution. Safe because the same network that + # failed /v1/ping/ is about to fail the PostHog POST, so nothing + # is delivered under the unverified identity in the meantime. + return email, "" + identity["email"] = verified identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, "" + return verified, "" if not key: # No key to verify the account with; do not keep attributing to it. if email: identity.pop("email", None) identity.pop("key_fingerprint", None) - _write_identity(identity) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" resolved = _resolve_email(key) if not resolved: # The key changed and will not resolve (revoked, offline, API down). - # Do not keep attributing to the previous account. + # Reaching here with an email means the recorded fingerprint disagreed, + # so the key really did change. Drop the account and rotate: the stored + # anonymous id may already be merged into that account's person, and + # reusing it would keep the events on the profile we are trying to + # leave. + if email: + identity.pop("email", None) + identity.pop("key_fingerprint", None) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" # Alias only when going anonymous -> email for the first time. Once an anon