diff --git a/integrations/agent-plugin-core/python/hook_runner.py b/integrations/agent-plugin-core/python/hook_runner.py index 2a2adb557..b3d8a3210 100644 --- a/integrations/agent-plugin-core/python/hook_runner.py +++ b/integrations/agent-plugin-core/python/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -305,8 +310,19 @@ def run( return 0 if args.action == "session-start": - if telemetry.is_first_run(): + # Claims the marker atomically and says which event to record, so a + # second session starting alongside this one cannot record it too. + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) + if first_event == "install": telemetry.record("install") + elif first_event == "upgrade": + # First run after a build that never wrote the marker; the + # predecessor version was never recorded anywhere. + telemetry.record("upgrade", from_version="pre-0.3") + else: + previous = telemetry.claim_version_change() + if previous: + telemetry.record("upgrade", from_version=previous) recovered = recover_pending_handoffs() record_session_start(store, hook_input) if recovered: diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index a4375f002..c1fc0084c 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -302,9 +302,176 @@ 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" + + def is_first_run() -> bool: - """Whether this machine has never recorded a plugin event before.""" - return not _identity_path().exists() + """Whether install has never been recorded on this machine. + + Deliberately NOT the identity file. That file is only written by a + successful flush, so an offline or firewalled user recorded code.install on + every single session, forever — and every 0.2.x user recorded one on their + first 0.3.x session because 0.2.x never wrote it at all. + """ + return not _install_state_path().exists() + + +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: + """Claim the one install/upgrade record for this machine, atomically. + + Returns the event to record ("install" or "upgrade"), or None if another + session already claimed it. O_CREAT|O_EXCL so two sessions starting together + cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. + """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + + path = _install_state_path() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return None + except OSError: + return None + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump( + { + "plugin_version": memory_core.PLUGIN_VERSION, + "installed_at": memory_core.utc_now(), + "upgraded": upgrading, + }, + stream, + ) + # Durable before this returns. The O_EXCL open is what makes the + # claim exclusive, so it cannot be replaced by a temp-and-rename + # without losing that, which leaves the content as the thing to make + # safe. A kill between the open and this fsync used to leave a marker + # that exists but parses to nothing: is_first_run reads it as claimed + # and claim_version_change cannot read a version out of it. + stream.flush() + os.fsync(stream.fileno()) + except OSError: + pass + return "upgrade" if upgrading else "install" + + +def _data_dir_has_content() -> bool: + """Whether anything predates this session in the plugin data directory.""" + try: + for entry in memory_core.data_dir().iterdir(): + if entry.name != "install-state.json": + return True + except OSError: + pass + return False + + +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + +def claim_version_change() -> str | None: + """Return the previously recorded version if it differs, updating the marker. + + Only meaningful once the marker exists — the first transition into 0.3.x has + no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting + the marker means the next session sees no change and records nothing. + """ + path = _install_state_path() + try: + state = json.loads(path.read_text(encoding="utf-8")) + except OSError: + return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None + if not isinstance(state, dict): + _repair_install_state(path) + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") + try: + 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 def record( @@ -633,21 +800,72 @@ def _post(payload: dict[str, Any], url: str) -> bool: def resolve_distinct_id() -> tuple[str, str]: - """Return the PostHog distinct id and the anonymous id it replaced, if any.""" + """Return the PostHog distinct id and the anonymous id it replaced, if any. + + The second value becomes a PostHog $identify alias. It is ONLY ever an + anonymous id: aliasing one account email to another merges two real person + profiles and cannot be undone, so a key that now belongs to a different + account re-resolves with no alias. + """ identity = _read_identity() - email = identity.get("email", "") - if email: - return email, "" key = memory_core.api_key() + fingerprint = _digest(key) if key else "" + email = identity.get("email", "") + + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # 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 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) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" - email = _resolve_email(key) - if not email: + + resolved = _resolve_email(key) + if not resolved: + # The key changed and will not resolve (revoked, offline, API down). + # 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), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: diff --git a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py index 14f911bee..7d48dbcc3 100644 --- a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py +++ b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py @@ -177,3 +177,71 @@ def test_source_tag_defaults_agree_between_the_two_modules(): ) left, right = out.split() assert left == right == "KIMI_PLUGIN" + + +def _session_start(core: Path, data_dir: Path) -> list[str]: + """Drive the real hook_runner session-start path and return lifecycle events.""" + recorded = "\n".join( + [ + "import io, json, sys", + f"sys.path.insert(0, {str(core)!r})", + "import telemetry, hook_runner", + "seen = []", + "telemetry.record = lambda event, **kw: seen.append(event) or None", + "telemetry.spawn_flush = lambda: False", + # run() reads sys.argv through argparse; it takes no positional args. + "sys.argv = ['hook_runner', 'session-start']", + "sys.stdin = io.StringIO('{}')", + "hook_runner.run()", + "print(json.dumps([e for e in seen if e in ('install', 'upgrade')]))", + ] + ) + import json as _json + + return _json.loads(_run(core, data_dir, recorded) or "[]") + + +def test_a_fresh_install_reports_install_not_upgrade(): + """The decision must survive the writes hook_runner does before asking. + + claim_install() is reached only after cache_plugin_api_key() has written + `api-key` and EvidenceStore() has created `evidence.sqlite3`. Asking "is the + data dir empty" at that point always saw content, so code.install could + never fire and every new user was counted as an upgrade. + """ + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + assert _session_start(core, data_dir) == ["install"] + + +def test_the_lifecycle_event_fires_exactly_once(): + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + first = _session_start(core, data_dir) + second = _session_start(core, data_dir) + third = _session_start(core, data_dir) + + assert first == ["install"] + assert second == [] + assert third == [] + + +def test_an_existing_data_dir_reports_upgrade(): + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + data_dir.mkdir(parents=True) + # A 0.2.x leftover: the data dir survives the upgrade. + (data_dir / "requirements.txt").write_text("mem0ai\n", encoding="utf-8") + assert _session_start(core, data_dir) == ["upgrade"] diff --git a/integrations/antigravity-plugin/core/hook_runner.py b/integrations/antigravity-plugin/core/hook_runner.py index 2a2adb557..b3d8a3210 100644 --- a/integrations/antigravity-plugin/core/hook_runner.py +++ b/integrations/antigravity-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -305,8 +310,19 @@ def run( return 0 if args.action == "session-start": - if telemetry.is_first_run(): + # Claims the marker atomically and says which event to record, so a + # second session starting alongside this one cannot record it too. + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) + if first_event == "install": telemetry.record("install") + elif first_event == "upgrade": + # First run after a build that never wrote the marker; the + # predecessor version was never recorded anywhere. + telemetry.record("upgrade", from_version="pre-0.3") + else: + previous = telemetry.claim_version_change() + if previous: + telemetry.record("upgrade", from_version=previous) recovered = recover_pending_handoffs() record_session_start(store, hook_input) if recovered: diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index a4375f002..c1fc0084c 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -302,9 +302,176 @@ 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" + + def is_first_run() -> bool: - """Whether this machine has never recorded a plugin event before.""" - return not _identity_path().exists() + """Whether install has never been recorded on this machine. + + Deliberately NOT the identity file. That file is only written by a + successful flush, so an offline or firewalled user recorded code.install on + every single session, forever — and every 0.2.x user recorded one on their + first 0.3.x session because 0.2.x never wrote it at all. + """ + return not _install_state_path().exists() + + +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: + """Claim the one install/upgrade record for this machine, atomically. + + Returns the event to record ("install" or "upgrade"), or None if another + session already claimed it. O_CREAT|O_EXCL so two sessions starting together + cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. + """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + + path = _install_state_path() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return None + except OSError: + return None + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump( + { + "plugin_version": memory_core.PLUGIN_VERSION, + "installed_at": memory_core.utc_now(), + "upgraded": upgrading, + }, + stream, + ) + # Durable before this returns. The O_EXCL open is what makes the + # claim exclusive, so it cannot be replaced by a temp-and-rename + # without losing that, which leaves the content as the thing to make + # safe. A kill between the open and this fsync used to leave a marker + # that exists but parses to nothing: is_first_run reads it as claimed + # and claim_version_change cannot read a version out of it. + stream.flush() + os.fsync(stream.fileno()) + except OSError: + pass + return "upgrade" if upgrading else "install" + + +def _data_dir_has_content() -> bool: + """Whether anything predates this session in the plugin data directory.""" + try: + for entry in memory_core.data_dir().iterdir(): + if entry.name != "install-state.json": + return True + except OSError: + pass + return False + + +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + +def claim_version_change() -> str | None: + """Return the previously recorded version if it differs, updating the marker. + + Only meaningful once the marker exists — the first transition into 0.3.x has + no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting + the marker means the next session sees no change and records nothing. + """ + path = _install_state_path() + try: + state = json.loads(path.read_text(encoding="utf-8")) + except OSError: + return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None + if not isinstance(state, dict): + _repair_install_state(path) + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") + try: + 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 def record( @@ -633,21 +800,72 @@ def _post(payload: dict[str, Any], url: str) -> bool: def resolve_distinct_id() -> tuple[str, str]: - """Return the PostHog distinct id and the anonymous id it replaced, if any.""" + """Return the PostHog distinct id and the anonymous id it replaced, if any. + + The second value becomes a PostHog $identify alias. It is ONLY ever an + anonymous id: aliasing one account email to another merges two real person + profiles and cannot be undone, so a key that now belongs to a different + account re-resolves with no alias. + """ identity = _read_identity() - email = identity.get("email", "") - if email: - return email, "" key = memory_core.api_key() + fingerprint = _digest(key) if key else "" + email = identity.get("email", "") + + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # 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 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) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" - email = _resolve_email(key) - if not email: + + resolved = _resolve_email(key) + if not resolved: + # The key changed and will not resolve (revoked, offline, API down). + # 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), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: diff --git a/integrations/claude-code-plugin/core/hook_runner.py b/integrations/claude-code-plugin/core/hook_runner.py index 2a2adb557..b3d8a3210 100644 --- a/integrations/claude-code-plugin/core/hook_runner.py +++ b/integrations/claude-code-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -305,8 +310,19 @@ def run( return 0 if args.action == "session-start": - if telemetry.is_first_run(): + # Claims the marker atomically and says which event to record, so a + # second session starting alongside this one cannot record it too. + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) + if first_event == "install": telemetry.record("install") + elif first_event == "upgrade": + # First run after a build that never wrote the marker; the + # predecessor version was never recorded anywhere. + telemetry.record("upgrade", from_version="pre-0.3") + else: + previous = telemetry.claim_version_change() + if previous: + telemetry.record("upgrade", from_version=previous) recovered = recover_pending_handoffs() record_session_start(store, hook_input) if recovered: diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index a4375f002..c1fc0084c 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -302,9 +302,176 @@ 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" + + def is_first_run() -> bool: - """Whether this machine has never recorded a plugin event before.""" - return not _identity_path().exists() + """Whether install has never been recorded on this machine. + + Deliberately NOT the identity file. That file is only written by a + successful flush, so an offline or firewalled user recorded code.install on + every single session, forever — and every 0.2.x user recorded one on their + first 0.3.x session because 0.2.x never wrote it at all. + """ + return not _install_state_path().exists() + + +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: + """Claim the one install/upgrade record for this machine, atomically. + + Returns the event to record ("install" or "upgrade"), or None if another + session already claimed it. O_CREAT|O_EXCL so two sessions starting together + cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. + """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + + path = _install_state_path() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return None + except OSError: + return None + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump( + { + "plugin_version": memory_core.PLUGIN_VERSION, + "installed_at": memory_core.utc_now(), + "upgraded": upgrading, + }, + stream, + ) + # Durable before this returns. The O_EXCL open is what makes the + # claim exclusive, so it cannot be replaced by a temp-and-rename + # without losing that, which leaves the content as the thing to make + # safe. A kill between the open and this fsync used to leave a marker + # that exists but parses to nothing: is_first_run reads it as claimed + # and claim_version_change cannot read a version out of it. + stream.flush() + os.fsync(stream.fileno()) + except OSError: + pass + return "upgrade" if upgrading else "install" + + +def _data_dir_has_content() -> bool: + """Whether anything predates this session in the plugin data directory.""" + try: + for entry in memory_core.data_dir().iterdir(): + if entry.name != "install-state.json": + return True + except OSError: + pass + return False + + +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + +def claim_version_change() -> str | None: + """Return the previously recorded version if it differs, updating the marker. + + Only meaningful once the marker exists — the first transition into 0.3.x has + no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting + the marker means the next session sees no change and records nothing. + """ + path = _install_state_path() + try: + state = json.loads(path.read_text(encoding="utf-8")) + except OSError: + return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None + if not isinstance(state, dict): + _repair_install_state(path) + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") + try: + 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 def record( @@ -633,21 +800,72 @@ def _post(payload: dict[str, Any], url: str) -> bool: def resolve_distinct_id() -> tuple[str, str]: - """Return the PostHog distinct id and the anonymous id it replaced, if any.""" + """Return the PostHog distinct id and the anonymous id it replaced, if any. + + The second value becomes a PostHog $identify alias. It is ONLY ever an + anonymous id: aliasing one account email to another merges two real person + profiles and cannot be undone, so a key that now belongs to a different + account re-resolves with no alias. + """ identity = _read_identity() - email = identity.get("email", "") - if email: - return email, "" key = memory_core.api_key() + fingerprint = _digest(key) if key else "" + email = identity.get("email", "") + + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # 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 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) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" - email = _resolve_email(key) - if not email: + + resolved = _resolve_email(key) + if not resolved: + # The key changed and will not resolve (revoked, offline, API down). + # 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), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: diff --git a/integrations/claude-code-plugin/tests/test_telemetry.py b/integrations/claude-code-plugin/tests/test_telemetry.py index de0db10c6..73bf8f121 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -265,12 +265,150 @@ 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_is_first_run_flips_after_the_first_identity_write(isolated_env): +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. + + Keying first-run off it meant an offline user recorded code.install on every + session forever, and every 0.2.x user recorded one on their first 0.3.x run. + """ assert telemetry.is_first_run() telemetry.anonymous_id() + assert telemetry.is_first_run() + + +def test_claiming_install_ends_first_run(isolated_env): + assert telemetry.claim_install() == "install" assert not telemetry.is_first_run() +def test_install_can_only_be_claimed_once(isolated_env): + """Two sessions starting together must not both record an install.""" + assert telemetry.claim_install() == "install" + assert telemetry.claim_install() is None + + +def test_a_populated_data_dir_reads_as_an_upgrade(isolated_env): + """A fresh install has an empty data directory; anything else predates it.""" + data_dir = memory_core.data_dir() + data_dir.mkdir(parents=True, exist_ok=True) + (data_dir / "requirements.txt").write_text("mem0ai\n", encoding="utf-8") + assert telemetry.claim_install() == "upgrade" + + +def test_a_version_change_is_claimed_once(isolated_env): + 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") + + assert telemetry.claim_version_change() == "0.0.1-old" + assert telemetry.claim_version_change() is None + + def test_spawn_flush_does_nothing_without_a_spool(isolated_env): with patch.object(telemetry.subprocess, "Popen") as popen: assert telemetry.spawn_flush() is False diff --git a/integrations/codex-plugin/core/hook_runner.py b/integrations/codex-plugin/core/hook_runner.py index 2a2adb557..b3d8a3210 100644 --- a/integrations/codex-plugin/core/hook_runner.py +++ b/integrations/codex-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -305,8 +310,19 @@ def run( return 0 if args.action == "session-start": - if telemetry.is_first_run(): + # Claims the marker atomically and says which event to record, so a + # second session starting alongside this one cannot record it too. + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) + if first_event == "install": telemetry.record("install") + elif first_event == "upgrade": + # First run after a build that never wrote the marker; the + # predecessor version was never recorded anywhere. + telemetry.record("upgrade", from_version="pre-0.3") + else: + previous = telemetry.claim_version_change() + if previous: + telemetry.record("upgrade", from_version=previous) recovered = recover_pending_handoffs() record_session_start(store, hook_input) if recovered: diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index a4375f002..c1fc0084c 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -302,9 +302,176 @@ 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" + + def is_first_run() -> bool: - """Whether this machine has never recorded a plugin event before.""" - return not _identity_path().exists() + """Whether install has never been recorded on this machine. + + Deliberately NOT the identity file. That file is only written by a + successful flush, so an offline or firewalled user recorded code.install on + every single session, forever — and every 0.2.x user recorded one on their + first 0.3.x session because 0.2.x never wrote it at all. + """ + return not _install_state_path().exists() + + +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: + """Claim the one install/upgrade record for this machine, atomically. + + Returns the event to record ("install" or "upgrade"), or None if another + session already claimed it. O_CREAT|O_EXCL so two sessions starting together + cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. + """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + + path = _install_state_path() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return None + except OSError: + return None + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump( + { + "plugin_version": memory_core.PLUGIN_VERSION, + "installed_at": memory_core.utc_now(), + "upgraded": upgrading, + }, + stream, + ) + # Durable before this returns. The O_EXCL open is what makes the + # claim exclusive, so it cannot be replaced by a temp-and-rename + # without losing that, which leaves the content as the thing to make + # safe. A kill between the open and this fsync used to leave a marker + # that exists but parses to nothing: is_first_run reads it as claimed + # and claim_version_change cannot read a version out of it. + stream.flush() + os.fsync(stream.fileno()) + except OSError: + pass + return "upgrade" if upgrading else "install" + + +def _data_dir_has_content() -> bool: + """Whether anything predates this session in the plugin data directory.""" + try: + for entry in memory_core.data_dir().iterdir(): + if entry.name != "install-state.json": + return True + except OSError: + pass + return False + + +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + +def claim_version_change() -> str | None: + """Return the previously recorded version if it differs, updating the marker. + + Only meaningful once the marker exists — the first transition into 0.3.x has + no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting + the marker means the next session sees no change and records nothing. + """ + path = _install_state_path() + try: + state = json.loads(path.read_text(encoding="utf-8")) + except OSError: + return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None + if not isinstance(state, dict): + _repair_install_state(path) + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") + try: + 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 def record( @@ -633,21 +800,72 @@ def _post(payload: dict[str, Any], url: str) -> bool: def resolve_distinct_id() -> tuple[str, str]: - """Return the PostHog distinct id and the anonymous id it replaced, if any.""" + """Return the PostHog distinct id and the anonymous id it replaced, if any. + + The second value becomes a PostHog $identify alias. It is ONLY ever an + anonymous id: aliasing one account email to another merges two real person + profiles and cannot be undone, so a key that now belongs to a different + account re-resolves with no alias. + """ identity = _read_identity() - email = identity.get("email", "") - if email: - return email, "" key = memory_core.api_key() + fingerprint = _digest(key) if key else "" + email = identity.get("email", "") + + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # 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 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) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" - email = _resolve_email(key) - if not email: + + resolved = _resolve_email(key) + if not resolved: + # The key changed and will not resolve (revoked, offline, API down). + # 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), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: diff --git a/integrations/cursor-plugin/core/hook_runner.py b/integrations/cursor-plugin/core/hook_runner.py index 2a2adb557..b3d8a3210 100644 --- a/integrations/cursor-plugin/core/hook_runner.py +++ b/integrations/cursor-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -305,8 +310,19 @@ def run( return 0 if args.action == "session-start": - if telemetry.is_first_run(): + # Claims the marker atomically and says which event to record, so a + # second session starting alongside this one cannot record it too. + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) + if first_event == "install": telemetry.record("install") + elif first_event == "upgrade": + # First run after a build that never wrote the marker; the + # predecessor version was never recorded anywhere. + telemetry.record("upgrade", from_version="pre-0.3") + else: + previous = telemetry.claim_version_change() + if previous: + telemetry.record("upgrade", from_version=previous) recovered = recover_pending_handoffs() record_session_start(store, hook_input) if recovered: diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index a4375f002..c1fc0084c 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -302,9 +302,176 @@ 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" + + def is_first_run() -> bool: - """Whether this machine has never recorded a plugin event before.""" - return not _identity_path().exists() + """Whether install has never been recorded on this machine. + + Deliberately NOT the identity file. That file is only written by a + successful flush, so an offline or firewalled user recorded code.install on + every single session, forever — and every 0.2.x user recorded one on their + first 0.3.x session because 0.2.x never wrote it at all. + """ + return not _install_state_path().exists() + + +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: + """Claim the one install/upgrade record for this machine, atomically. + + Returns the event to record ("install" or "upgrade"), or None if another + session already claimed it. O_CREAT|O_EXCL so two sessions starting together + cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. + """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + + path = _install_state_path() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return None + except OSError: + return None + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump( + { + "plugin_version": memory_core.PLUGIN_VERSION, + "installed_at": memory_core.utc_now(), + "upgraded": upgrading, + }, + stream, + ) + # Durable before this returns. The O_EXCL open is what makes the + # claim exclusive, so it cannot be replaced by a temp-and-rename + # without losing that, which leaves the content as the thing to make + # safe. A kill between the open and this fsync used to leave a marker + # that exists but parses to nothing: is_first_run reads it as claimed + # and claim_version_change cannot read a version out of it. + stream.flush() + os.fsync(stream.fileno()) + except OSError: + pass + return "upgrade" if upgrading else "install" + + +def _data_dir_has_content() -> bool: + """Whether anything predates this session in the plugin data directory.""" + try: + for entry in memory_core.data_dir().iterdir(): + if entry.name != "install-state.json": + return True + except OSError: + pass + return False + + +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + +def claim_version_change() -> str | None: + """Return the previously recorded version if it differs, updating the marker. + + Only meaningful once the marker exists — the first transition into 0.3.x has + no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting + the marker means the next session sees no change and records nothing. + """ + path = _install_state_path() + try: + state = json.loads(path.read_text(encoding="utf-8")) + except OSError: + return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None + if not isinstance(state, dict): + _repair_install_state(path) + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") + try: + 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 def record( @@ -633,21 +800,72 @@ def _post(payload: dict[str, Any], url: str) -> bool: def resolve_distinct_id() -> tuple[str, str]: - """Return the PostHog distinct id and the anonymous id it replaced, if any.""" + """Return the PostHog distinct id and the anonymous id it replaced, if any. + + The second value becomes a PostHog $identify alias. It is ONLY ever an + anonymous id: aliasing one account email to another merges two real person + profiles and cannot be undone, so a key that now belongs to a different + account re-resolves with no alias. + """ identity = _read_identity() - email = identity.get("email", "") - if email: - return email, "" key = memory_core.api_key() + fingerprint = _digest(key) if key else "" + email = identity.get("email", "") + + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # 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 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) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" - email = _resolve_email(key) - if not email: + + resolved = _resolve_email(key) + if not resolved: + # The key changed and will not resolve (revoked, offline, API down). + # 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), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: diff --git a/integrations/kimi-plugin/core/hook_runner.py b/integrations/kimi-plugin/core/hook_runner.py index 2a2adb557..b3d8a3210 100644 --- a/integrations/kimi-plugin/core/hook_runner.py +++ b/integrations/kimi-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -305,8 +310,19 @@ def run( return 0 if args.action == "session-start": - if telemetry.is_first_run(): + # Claims the marker atomically and says which event to record, so a + # second session starting alongside this one cannot record it too. + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) + if first_event == "install": telemetry.record("install") + elif first_event == "upgrade": + # First run after a build that never wrote the marker; the + # predecessor version was never recorded anywhere. + telemetry.record("upgrade", from_version="pre-0.3") + else: + previous = telemetry.claim_version_change() + if previous: + telemetry.record("upgrade", from_version=previous) recovered = recover_pending_handoffs() record_session_start(store, hook_input) if recovered: diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index a4375f002..c1fc0084c 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -302,9 +302,176 @@ 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" + + def is_first_run() -> bool: - """Whether this machine has never recorded a plugin event before.""" - return not _identity_path().exists() + """Whether install has never been recorded on this machine. + + Deliberately NOT the identity file. That file is only written by a + successful flush, so an offline or firewalled user recorded code.install on + every single session, forever — and every 0.2.x user recorded one on their + first 0.3.x session because 0.2.x never wrote it at all. + """ + return not _install_state_path().exists() + + +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: + """Claim the one install/upgrade record for this machine, atomically. + + Returns the event to record ("install" or "upgrade"), or None if another + session already claimed it. O_CREAT|O_EXCL so two sessions starting together + cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. + """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + + path = _install_state_path() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return None + except OSError: + return None + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump( + { + "plugin_version": memory_core.PLUGIN_VERSION, + "installed_at": memory_core.utc_now(), + "upgraded": upgrading, + }, + stream, + ) + # Durable before this returns. The O_EXCL open is what makes the + # claim exclusive, so it cannot be replaced by a temp-and-rename + # without losing that, which leaves the content as the thing to make + # safe. A kill between the open and this fsync used to leave a marker + # that exists but parses to nothing: is_first_run reads it as claimed + # and claim_version_change cannot read a version out of it. + stream.flush() + os.fsync(stream.fileno()) + except OSError: + pass + return "upgrade" if upgrading else "install" + + +def _data_dir_has_content() -> bool: + """Whether anything predates this session in the plugin data directory.""" + try: + for entry in memory_core.data_dir().iterdir(): + if entry.name != "install-state.json": + return True + except OSError: + pass + return False + + +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + +def claim_version_change() -> str | None: + """Return the previously recorded version if it differs, updating the marker. + + Only meaningful once the marker exists — the first transition into 0.3.x has + no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting + the marker means the next session sees no change and records nothing. + """ + path = _install_state_path() + try: + state = json.loads(path.read_text(encoding="utf-8")) + except OSError: + return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None + if not isinstance(state, dict): + _repair_install_state(path) + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") + try: + 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 def record( @@ -633,21 +800,72 @@ def _post(payload: dict[str, Any], url: str) -> bool: def resolve_distinct_id() -> tuple[str, str]: - """Return the PostHog distinct id and the anonymous id it replaced, if any.""" + """Return the PostHog distinct id and the anonymous id it replaced, if any. + + The second value becomes a PostHog $identify alias. It is ONLY ever an + anonymous id: aliasing one account email to another merges two real person + profiles and cannot be undone, so a key that now belongs to a different + account re-resolves with no alias. + """ identity = _read_identity() - email = identity.get("email", "") - if email: - return email, "" key = memory_core.api_key() + fingerprint = _digest(key) if key else "" + email = identity.get("email", "") + + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # 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 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) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" - email = _resolve_email(key) - if not email: + + resolved = _resolve_email(key) + if not resolved: + # The key changed and will not resolve (revoked, offline, API down). + # 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), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index a4375f002..c1fc0084c 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -302,9 +302,176 @@ 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" + + def is_first_run() -> bool: - """Whether this machine has never recorded a plugin event before.""" - return not _identity_path().exists() + """Whether install has never been recorded on this machine. + + Deliberately NOT the identity file. That file is only written by a + successful flush, so an offline or firewalled user recorded code.install on + every single session, forever — and every 0.2.x user recorded one on their + first 0.3.x session because 0.2.x never wrote it at all. + """ + return not _install_state_path().exists() + + +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: + """Claim the one install/upgrade record for this machine, atomically. + + Returns the event to record ("install" or "upgrade"), or None if another + session already claimed it. O_CREAT|O_EXCL so two sessions starting together + cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. + """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + + path = _install_state_path() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) + try: + path.parent.mkdir(parents=True, exist_ok=True) + handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return None + except OSError: + return None + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump( + { + "plugin_version": memory_core.PLUGIN_VERSION, + "installed_at": memory_core.utc_now(), + "upgraded": upgrading, + }, + stream, + ) + # Durable before this returns. The O_EXCL open is what makes the + # claim exclusive, so it cannot be replaced by a temp-and-rename + # without losing that, which leaves the content as the thing to make + # safe. A kill between the open and this fsync used to leave a marker + # that exists but parses to nothing: is_first_run reads it as claimed + # and claim_version_change cannot read a version out of it. + stream.flush() + os.fsync(stream.fileno()) + except OSError: + pass + return "upgrade" if upgrading else "install" + + +def _data_dir_has_content() -> bool: + """Whether anything predates this session in the plugin data directory.""" + try: + for entry in memory_core.data_dir().iterdir(): + if entry.name != "install-state.json": + return True + except OSError: + pass + return False + + +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + +def claim_version_change() -> str | None: + """Return the previously recorded version if it differs, updating the marker. + + Only meaningful once the marker exists — the first transition into 0.3.x has + no recorded predecessor and reports "pre-0.3" instead. Claiming by rewriting + the marker means the next session sees no change and records nothing. + """ + path = _install_state_path() + try: + state = json.loads(path.read_text(encoding="utf-8")) + except OSError: + return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None + if not isinstance(state, dict): + _repair_install_state(path) + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + temporary = path.with_suffix(f".{os.getpid()}.tmp") + try: + 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 def record( @@ -633,21 +800,72 @@ def _post(payload: dict[str, Any], url: str) -> bool: def resolve_distinct_id() -> tuple[str, str]: - """Return the PostHog distinct id and the anonymous id it replaced, if any.""" + """Return the PostHog distinct id and the anonymous id it replaced, if any. + + The second value becomes a PostHog $identify alias. It is ONLY ever an + anonymous id: aliasing one account email to another merges two real person + profiles and cannot be undone, so a key that now belongs to a different + account re-resolves with no alias. + """ identity = _read_identity() - email = identity.get("email", "") - if email: - return email, "" key = memory_core.api_key() + fingerprint = _digest(key) if key else "" + email = identity.get("email", "") + + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # 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 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) + return _rotate_anonymous_id(identity), "" return anonymous_id(identity), "" - email = _resolve_email(key) - if not email: + + resolved = _resolve_email(key) + if not resolved: + # The key changed and will not resolve (revoked, offline, API down). + # 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), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: