diff --git a/integrations/agent-plugin-core/python/hook_runner.py b/integrations/agent-plugin-core/python/hook_runner.py index 2a2adb557..559702102 100644 --- a/integrations/agent-plugin-core/python/hook_runner.py +++ b/integrations/agent-plugin-core/python/hook_runner.py @@ -305,8 +305,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() + 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 55ee7c907..d7933bfbc 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -206,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: 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 claim_install() -> 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. + """ + path = _install_state_path() + # A fresh install has an empty data directory. Anything already there — + # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read + # before the marker is created, since creating it would itself be content. + upgrading = _data_dir_has_content() + 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, + ) + 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 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, json.JSONDecodeError): + return None + if not isinstance(state, dict): + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text(json.dumps(state), encoding="utf-8") + temporary.replace(path) + except OSError: + return None + return previous def record( @@ -468,21 +551,39 @@ 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 identity.get("key_fingerprint", "") == fingerprint and fingerprint: + return email, "" + 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 anonymous_id(identity), "" - email = _resolve_email(key) - if not email: - return anonymous_id(identity), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + resolved = _resolve_email(key) + if not resolved: + return (email, "") if email else (anonymous_id(identity), "") + + # Alias only when going anonymous -> email for the first time. + previous = "" if email else identity.get("anonymous_id", "") + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: diff --git a/integrations/antigravity-plugin/core/hook_runner.py b/integrations/antigravity-plugin/core/hook_runner.py index 2a2adb557..559702102 100644 --- a/integrations/antigravity-plugin/core/hook_runner.py +++ b/integrations/antigravity-plugin/core/hook_runner.py @@ -305,8 +305,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() + 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 55ee7c907..d7933bfbc 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -206,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: 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 claim_install() -> 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. + """ + path = _install_state_path() + # A fresh install has an empty data directory. Anything already there — + # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read + # before the marker is created, since creating it would itself be content. + upgrading = _data_dir_has_content() + 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, + ) + 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 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, json.JSONDecodeError): + return None + if not isinstance(state, dict): + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text(json.dumps(state), encoding="utf-8") + temporary.replace(path) + except OSError: + return None + return previous def record( @@ -468,21 +551,39 @@ 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 identity.get("key_fingerprint", "") == fingerprint and fingerprint: + return email, "" + 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 anonymous_id(identity), "" - email = _resolve_email(key) - if not email: - return anonymous_id(identity), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + resolved = _resolve_email(key) + if not resolved: + return (email, "") if email else (anonymous_id(identity), "") + + # Alias only when going anonymous -> email for the first time. + previous = "" if email else identity.get("anonymous_id", "") + 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..559702102 100644 --- a/integrations/claude-code-plugin/core/hook_runner.py +++ b/integrations/claude-code-plugin/core/hook_runner.py @@ -305,8 +305,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() + 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 55ee7c907..d7933bfbc 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -206,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: 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 claim_install() -> 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. + """ + path = _install_state_path() + # A fresh install has an empty data directory. Anything already there — + # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read + # before the marker is created, since creating it would itself be content. + upgrading = _data_dir_has_content() + 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, + ) + 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 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, json.JSONDecodeError): + return None + if not isinstance(state, dict): + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text(json.dumps(state), encoding="utf-8") + temporary.replace(path) + except OSError: + return None + return previous def record( @@ -468,21 +551,39 @@ 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 identity.get("key_fingerprint", "") == fingerprint and fingerprint: + return email, "" + 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 anonymous_id(identity), "" - email = _resolve_email(key) - if not email: - return anonymous_id(identity), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + resolved = _resolve_email(key) + if not resolved: + return (email, "") if email else (anonymous_id(identity), "") + + # Alias only when going anonymous -> email for the first time. + previous = "" if email else identity.get("anonymous_id", "") + 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 580bc1100..4476a9dd3 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -265,12 +265,47 @@ 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_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..559702102 100644 --- a/integrations/codex-plugin/core/hook_runner.py +++ b/integrations/codex-plugin/core/hook_runner.py @@ -305,8 +305,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() + 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 55ee7c907..d7933bfbc 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -206,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: 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 claim_install() -> 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. + """ + path = _install_state_path() + # A fresh install has an empty data directory. Anything already there — + # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read + # before the marker is created, since creating it would itself be content. + upgrading = _data_dir_has_content() + 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, + ) + 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 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, json.JSONDecodeError): + return None + if not isinstance(state, dict): + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text(json.dumps(state), encoding="utf-8") + temporary.replace(path) + except OSError: + return None + return previous def record( @@ -468,21 +551,39 @@ 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 identity.get("key_fingerprint", "") == fingerprint and fingerprint: + return email, "" + 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 anonymous_id(identity), "" - email = _resolve_email(key) - if not email: - return anonymous_id(identity), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + resolved = _resolve_email(key) + if not resolved: + return (email, "") if email else (anonymous_id(identity), "") + + # Alias only when going anonymous -> email for the first time. + previous = "" if email else identity.get("anonymous_id", "") + 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..559702102 100644 --- a/integrations/cursor-plugin/core/hook_runner.py +++ b/integrations/cursor-plugin/core/hook_runner.py @@ -305,8 +305,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() + 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 55ee7c907..d7933bfbc 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -206,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: 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 claim_install() -> 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. + """ + path = _install_state_path() + # A fresh install has an empty data directory. Anything already there — + # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read + # before the marker is created, since creating it would itself be content. + upgrading = _data_dir_has_content() + 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, + ) + 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 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, json.JSONDecodeError): + return None + if not isinstance(state, dict): + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text(json.dumps(state), encoding="utf-8") + temporary.replace(path) + except OSError: + return None + return previous def record( @@ -468,21 +551,39 @@ 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 identity.get("key_fingerprint", "") == fingerprint and fingerprint: + return email, "" + 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 anonymous_id(identity), "" - email = _resolve_email(key) - if not email: - return anonymous_id(identity), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + resolved = _resolve_email(key) + if not resolved: + return (email, "") if email else (anonymous_id(identity), "") + + # Alias only when going anonymous -> email for the first time. + previous = "" if email else identity.get("anonymous_id", "") + 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..559702102 100644 --- a/integrations/kimi-plugin/core/hook_runner.py +++ b/integrations/kimi-plugin/core/hook_runner.py @@ -305,8 +305,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() + 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 55ee7c907..d7933bfbc 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -206,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: 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 claim_install() -> 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. + """ + path = _install_state_path() + # A fresh install has an empty data directory. Anything already there — + # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read + # before the marker is created, since creating it would itself be content. + upgrading = _data_dir_has_content() + 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, + ) + 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 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, json.JSONDecodeError): + return None + if not isinstance(state, dict): + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text(json.dumps(state), encoding="utf-8") + temporary.replace(path) + except OSError: + return None + return previous def record( @@ -468,21 +551,39 @@ 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 identity.get("key_fingerprint", "") == fingerprint and fingerprint: + return email, "" + 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 anonymous_id(identity), "" - email = _resolve_email(key) - if not email: - return anonymous_id(identity), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + resolved = _resolve_email(key) + if not resolved: + return (email, "") if email else (anonymous_id(identity), "") + + # Alias only when going anonymous -> email for the first time. + previous = "" if email else identity.get("anonymous_id", "") + 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 55ee7c907..d7933bfbc 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -206,9 +206,92 @@ def anonymous_id(identity: dict[str, str] | None = None) -> str: 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 claim_install() -> 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. + """ + path = _install_state_path() + # A fresh install has an empty data directory. Anything already there — + # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read + # before the marker is created, since creating it would itself be content. + upgrading = _data_dir_has_content() + 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, + ) + 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 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, json.JSONDecodeError): + return None + if not isinstance(state, dict): + return None + previous = str(state.get("plugin_version") or "") + if not previous or previous == memory_core.PLUGIN_VERSION: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION + state["upgraded_at"] = memory_core.utc_now() + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text(json.dumps(state), encoding="utf-8") + temporary.replace(path) + except OSError: + return None + return previous def record( @@ -468,21 +551,39 @@ 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 identity.get("key_fingerprint", "") == fingerprint and fingerprint: + return email, "" + 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 anonymous_id(identity), "" - email = _resolve_email(key) - if not email: - return anonymous_id(identity), "" - previous = identity.get("anonymous_id", "") - identity["email"] = email + + resolved = _resolve_email(key) + if not resolved: + return (email, "") if email else (anonymous_id(identity), "") + + # Alias only when going anonymous -> email for the first time. + previous = "" if email else identity.get("anonymous_id", "") + identity["email"] = resolved + identity["key_fingerprint"] = fingerprint _write_identity(identity) - return email, previous + return resolved, previous def flush() -> int: