diff --git a/integrations/agent-plugin-core/python/hook_runner.py b/integrations/agent-plugin-core/python/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/agent-plugin-core/python/hook_runner.py +++ b/integrations/agent-plugin-core/python/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -221,18 +221,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +283,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +306,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -563,8 +612,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +635,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) diff --git a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py index 14f911bee..7d48dbcc3 100644 --- a/integrations/agent-plugin-core/tests/test_uninitialised_identity.py +++ b/integrations/agent-plugin-core/tests/test_uninitialised_identity.py @@ -177,3 +177,71 @@ def test_source_tag_defaults_agree_between_the_two_modules(): ) left, right = out.split() assert left == right == "KIMI_PLUGIN" + + +def _session_start(core: Path, data_dir: Path) -> list[str]: + """Drive the real hook_runner session-start path and return lifecycle events.""" + recorded = "\n".join( + [ + "import io, json, sys", + f"sys.path.insert(0, {str(core)!r})", + "import telemetry, hook_runner", + "seen = []", + "telemetry.record = lambda event, **kw: seen.append(event) or None", + "telemetry.spawn_flush = lambda: False", + # run() reads sys.argv through argparse; it takes no positional args. + "sys.argv = ['hook_runner', 'session-start']", + "sys.stdin = io.StringIO('{}')", + "hook_runner.run()", + "print(json.dumps([e for e in seen if e in ('install', 'upgrade')]))", + ] + ) + import json as _json + + return _json.loads(_run(core, data_dir, recorded) or "[]") + + +def test_a_fresh_install_reports_install_not_upgrade(): + """The decision must survive the writes hook_runner does before asking. + + claim_install() is reached only after cache_plugin_api_key() has written + `api-key` and EvidenceStore() has created `evidence.sqlite3`. Asking "is the + data dir empty" at that point always saw content, so code.install could + never fire and every new user was counted as an upgrade. + """ + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + assert _session_start(core, data_dir) == ["install"] + + +def test_the_lifecycle_event_fires_exactly_once(): + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + first = _session_start(core, data_dir) + second = _session_start(core, data_dir) + third = _session_start(core, data_dir) + + assert first == ["install"] + assert second == [] + assert third == [] + + +def test_an_existing_data_dir_reports_upgrade(): + core = _core_dir("claude-code-plugin") + if not core.exists(): + pytest.skip("claude-code-plugin is not built in this tree") + + with tempfile.TemporaryDirectory() as tmp: + data_dir = Path(tmp) / "data" + data_dir.mkdir(parents=True) + # A 0.2.x leftover: the data dir survives the upgrade. + (data_dir / "requirements.txt").write_text("mem0ai\n", encoding="utf-8") + assert _session_start(core, data_dir) == ["upgrade"] diff --git a/integrations/antigravity-plugin/core/hook_runner.py b/integrations/antigravity-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/antigravity-plugin/core/hook_runner.py +++ b/integrations/antigravity-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -221,18 +221,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +283,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +306,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -563,8 +612,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +635,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) diff --git a/integrations/claude-code-plugin/core/hook_runner.py b/integrations/claude-code-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/claude-code-plugin/core/hook_runner.py +++ b/integrations/claude-code-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -221,18 +221,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +283,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +306,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -563,8 +612,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +635,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) diff --git a/integrations/codex-plugin/core/hook_runner.py b/integrations/codex-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/codex-plugin/core/hook_runner.py +++ b/integrations/codex-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -221,18 +221,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +283,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +306,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -563,8 +612,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +635,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) diff --git a/integrations/cursor-plugin/core/hook_runner.py b/integrations/cursor-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/cursor-plugin/core/hook_runner.py +++ b/integrations/cursor-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -221,18 +221,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +283,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +306,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -563,8 +612,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +635,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) diff --git a/integrations/kimi-plugin/core/hook_runner.py b/integrations/kimi-plugin/core/hook_runner.py index 559702102..b3d8a3210 100644 --- a/integrations/kimi-plugin/core/hook_runner.py +++ b/integrations/kimi-plugin/core/hook_runner.py @@ -290,6 +290,11 @@ def run( if args.plugin_data_dir: os.environ[data_dir_env] = args.plugin_data_dir + # Snapshot BEFORE anything writes to the data dir: cache_plugin_api_key + # writes `api-key` and EvidenceStore creates `evidence.sqlite3`, so asking + # after them always saw content and every fresh install reported an upgrade. + data_dir_was_empty = telemetry.data_dir_was_empty() + cache_plugin_api_key() if args.action == "session-start": clear_stale_api_key_cache() @@ -307,7 +312,7 @@ def run( if args.action == "session-start": # Claims the marker atomically and says which event to record, so a # second session starting alongside this one cannot record it too. - first_event = telemetry.claim_install() + first_event = telemetry.claim_install(was_empty=data_dir_was_empty) if first_event == "install": telemetry.record("install") elif first_event == "upgrade": diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -221,18 +221,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +283,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +306,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -563,8 +612,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +635,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity) diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index d7933bfbc..a24c97264 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -221,18 +221,35 @@ def is_first_run() -> bool: return not _install_state_path().exists() -def claim_install() -> str | None: +def data_dir_was_empty() -> bool: + """Whether the data directory is untouched. Call BEFORE anything writes to it. + + hook_runner reaches claim_install() only after cache_plugin_api_key() has + written `api-key` and EvidenceStore() has created `evidence.sqlite3`, so + asking at claim time always saw content and every fresh install reported an + upgrade. The caller snapshots this at the top of the run instead. + """ + return not _data_dir_has_content() + + +def claim_install(was_empty: bool | None = None) -> str | None: """Claim the one install/upgrade record for this machine, atomically. Returns the event to record ("install" or "upgrade"), or None if another session already claimed it. O_CREAT|O_EXCL so two sessions starting together cannot both win. + + `was_empty` must come from data_dir_was_empty() called before this process + wrote anything. Omitting it falls back to checking now, which is only + correct for a caller that has touched nothing. """ + if not is_enabled(): + # Never consume the one-shot claim while the user is opted out, or they + # would silently lose their install event if they later opt in. + return None + path = _install_state_path() - # A fresh install has an empty data directory. Anything already there — - # a 0.2.x venv, an evidence db, a spool — means this is an upgrade. Read - # before the marker is created, since creating it would itself be content. - upgrading = _data_dir_has_content() + upgrading = not (data_dir_was_empty() if was_empty is None else was_empty) try: path.parent.mkdir(parents=True, exist_ok=True) handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) @@ -266,6 +283,19 @@ def _data_dir_has_content() -> bool: return False +def _repair_install_state(path: Path) -> None: + """Rewrite an unparseable marker so version tracking can resume.""" + try: + temporary = path.with_suffix(f".{os.getpid()}.tmp") + temporary.write_text( + json.dumps({"plugin_version": memory_core.PLUGIN_VERSION, "repaired_at": memory_core.utc_now()}), + encoding="utf-8", + ) + temporary.replace(path) + except OSError: + pass + + def claim_version_change() -> str | None: """Return the previously recorded version if it differs, updating the marker. @@ -276,13 +306,32 @@ def claim_version_change() -> str | None: path = _install_state_path() try: state = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): + except OSError: return None + except json.JSONDecodeError: + # A crash between O_EXCL and the write leaves an empty marker. Left + # alone it disables every future upgrade event on this machine, because + # claim_install sees the file and this function cannot parse it. + state = None if not isinstance(state, dict): + _repair_install_state(path) return None previous = str(state.get("plugin_version") or "") if not previous or previous == memory_core.PLUGIN_VERSION: return None + # Claim the transition with an exclusive sentinel before rewriting the + # marker. A plain read-modify-write let every concurrently starting session + # observe the old version and each record its own upgrade — and the first + # session after a version bump is exactly when several agent windows restart + # together. + sentinel = path.with_name(f"upgraded-{memory_core.PLUGIN_VERSION}") + try: + os.close(os.open(sentinel, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + return None + except OSError: + return None + state["plugin_version"] = memory_core.PLUGIN_VERSION state["upgraded_at"] = memory_core.utc_now() try: @@ -563,8 +612,18 @@ def resolve_distinct_id() -> tuple[str, str]: fingerprint = _digest(key) if key else "" email = identity.get("email", "") - if email and identity.get("key_fingerprint", "") == fingerprint and fingerprint: - return email, "" + if email and fingerprint: + recorded = identity.get("key_fingerprint", "") + if recorded == fingerprint: + return email, "" + if not recorded: + # Rows written before fingerprints existed. Adopt the current key + # rather than re-resolving: otherwise every existing user pays an + # uncached /v1/ping/ on every flush, forever, and a firewalled one + # pays the full timeout each time. + identity["key_fingerprint"] = fingerprint + _write_identity(identity) + return email, "" if not key: # No key to verify the account with; do not keep attributing to it. @@ -576,10 +635,16 @@ def resolve_distinct_id() -> tuple[str, str]: resolved = _resolve_email(key) if not resolved: - return (email, "") if email else (anonymous_id(identity), "") + # The key changed and will not resolve (revoked, offline, API down). + # Do not keep attributing to the previous account. + return anonymous_id(identity), "" - # Alias only when going anonymous -> email for the first time. - previous = "" if email else identity.get("anonymous_id", "") + # Alias only when going anonymous -> email for the first time. Once an anon + # id has been merged into an account it must never be offered again: an + # alias naming an already-identified id is what could link two real people. + previous = "" if (email or identity.get("aliased")) else identity.get("anonymous_id", "") + if previous: + identity["aliased"] = True identity["email"] = resolved identity["key_fingerprint"] = fingerprint _write_identity(identity)