Merge branch 'pr4/install-marker-and-identity' into pr5/surface-headers
# Conflicts: # integrations/agent-plugin-core/tests/test_uninitialised_identity.py
This commit is contained in:
@@ -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":
|
||||
|
||||
@@ -49,6 +49,7 @@ except ImportError:
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_salt_cache: str = ""
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _salt_path() -> Path:
|
||||
return memory_core.data_dir() / "telemetry-salt"
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
"""Random per-install salt, created once and memoized for the process.
|
||||
|
||||
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
|
||||
the identity file. Three reasons, all of which produced wrong data when this
|
||||
lived in the identity dict:
|
||||
|
||||
- Hooks are short-lived separate processes firing on every tool call, and
|
||||
people run more than one agent window. A read-modify-write would let each
|
||||
process mint its own salt, so one repository would hash several ways in the
|
||||
window before a writer won.
|
||||
- resolve_distinct_id holds a copy of the identity dict across a network call
|
||||
to /v1/ping/, so whichever write landed second erased the other's key —
|
||||
losing either the salt (repo_hash changes mid-stream) or the email (a
|
||||
second $identify, splitting the person).
|
||||
- Touching the identity file from record() would create it, and is_first_run
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
@@ -221,18 +267,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 +329,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 +352,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:
|
||||
@@ -396,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
|
||||
|
||||
Anchored on field position, not on a leading "a": the legacy shape is
|
||||
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
|
||||
otherwise parse as attempt 1234567 and be discarded unsent on the first
|
||||
flush after an upgrade.
|
||||
"""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
parts = stem.split("-")
|
||||
if len(parts) != 4:
|
||||
return 0
|
||||
tail = parts[3]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
@@ -432,6 +536,21 @@ def _claim_spool() -> Path | None:
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _sweep_debris(directory: Path) -> None:
|
||||
"""Remove temp files orphaned by a crash between write and rename.
|
||||
|
||||
Neither glob in this module matches *.partial, so nothing else would ever
|
||||
clean them up.
|
||||
"""
|
||||
now = time.time()
|
||||
for debris in directory.glob("telemetry-*.partial"):
|
||||
try:
|
||||
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
|
||||
debris.unlink()
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
@@ -447,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
# Attempts, not age. Every re-claim touches the mtime and every release
|
||||
# backdates it by a fixed amount, so age is pinned near the stale
|
||||
# threshold and never reaches the expiry. Age stays only as a backstop
|
||||
# for files that never carried an attempt marker.
|
||||
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
@@ -489,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
|
||||
# fsync before the rename: without it the rename can land while the
|
||||
# bytes have not, and the claim comes back empty or truncated after a
|
||||
# crash. _drain then reads zero events and unlinks it.
|
||||
with open(temporary, "w", encoding="utf-8") as handle:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
@@ -563,8 +690,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 +713,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)
|
||||
@@ -599,6 +742,7 @@ def flush() -> int:
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
_sweep_debris(directory)
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
@@ -619,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
except (OSError, ValueError):
|
||||
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
|
||||
# rather than retry: flush() runs from a bare `finally:` in
|
||||
# flush_worker, so raising here also skips the handoff cleanup, and an
|
||||
# undecodable file would otherwise be re-read on every flush forever.
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt"))
|
||||
except OSError:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
@@ -630,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
if isinstance(value, dict) and value.get("event"):
|
||||
events.append(value)
|
||||
if not events:
|
||||
# Only delete when the file really is empty. A non-empty file that
|
||||
# parses to nothing is a torn write, and its contents are the unsent
|
||||
# remainder — deleting it is the data loss this PR exists to prevent.
|
||||
try:
|
||||
claim.unlink()
|
||||
empty = claim.stat().st_size == 0
|
||||
except OSError:
|
||||
empty = True
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
@@ -681,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
# a crash repeats at most one batch instead of the entire file. If the
|
||||
# rewrite fails the claim still holds delivered events, so stop rather
|
||||
# than carry on as though progress were recorded — continuing is how the
|
||||
# duplicate delivery this PR fixes would come back.
|
||||
if not _rewrite_claim(claim, events[start + len(chunk) :]):
|
||||
_release_claim(claim, events[start + len(chunk) :])
|
||||
return sent, False
|
||||
return sent, True
|
||||
|
||||
|
||||
|
||||
@@ -112,16 +112,19 @@ def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry):
|
||||
assert names == {"code.parked", "code.fresh"}
|
||||
|
||||
|
||||
def test_an_untried_batch_is_not_expired_by_age_alone(telemetry):
|
||||
"""Expiry should discard what failed, not what never got a turn."""
|
||||
def test_a_batch_is_retried_until_the_budget_is_spent_not_discarded(telemetry):
|
||||
"""Expiry discards what failed repeatedly, not what merely sat for a while.
|
||||
|
||||
The budget is the attempt count, because age cannot be one: every re-claim
|
||||
touches the mtime and every release backdates it, so age never accumulates.
|
||||
"""
|
||||
telemetry.record("parked")
|
||||
telemetry._post = lambda payload, url: False
|
||||
telemetry.flush()
|
||||
|
||||
parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
|
||||
assert len(parked) == 1
|
||||
ancient = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60)
|
||||
os.utime(parked[0], (ancient, ancient))
|
||||
assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS
|
||||
|
||||
sent: list[dict] = []
|
||||
telemetry._post = lambda payload, url: sent.append(payload) or True
|
||||
@@ -154,11 +157,88 @@ def test_progress_is_recorded_after_every_batch(telemetry):
|
||||
assert json.loads(remaining[0])["properties"]["index"] == 200
|
||||
|
||||
|
||||
def test_the_heartbeat_stays_well_inside_the_lease(telemetry):
|
||||
def test_the_heartbeat_actually_refreshes_the_lease(telemetry):
|
||||
"""The claim rewrite doubles as the lease heartbeat.
|
||||
|
||||
_post makes a single attempt with SEND_TIMEOUT and no retry, so a heartbeat
|
||||
lands at least that often. If a retry loop is ever added to _post, this is
|
||||
the assertion that catches a sender losing its claim mid-flight.
|
||||
Previously asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS`, which compares
|
||||
two constants and executes none of the code under test. Drive the real
|
||||
rewrite and watch the mtime move instead.
|
||||
"""
|
||||
assert telemetry.SEND_TIMEOUT * 4 < telemetry.CLAIM_STALE_SECONDS
|
||||
for index in range(150):
|
||||
telemetry.record("search", index=index)
|
||||
claim = telemetry._claim_spool()
|
||||
assert claim is not None
|
||||
|
||||
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
|
||||
os.utime(claim, (stale, stale))
|
||||
assert time.time() - claim.stat().st_mtime > telemetry.CLAIM_STALE_SECONDS
|
||||
|
||||
telemetry._rewrite_claim(claim, [{"event": "code.x", "properties": {}}])
|
||||
assert time.time() - claim.stat().st_mtime < telemetry.CLAIM_STALE_SECONDS
|
||||
|
||||
|
||||
def test_an_undeliverable_batch_is_eventually_given_up_on(telemetry):
|
||||
"""Expiry has to be reachable from a state the state machine can produce.
|
||||
|
||||
It was not: every re-claim touched the mtime and every release backdated it
|
||||
by a fixed amount, so age hovered near the stale threshold and the 7-day
|
||||
expiry never fired. An undeliverable batch lived on disk forever, and
|
||||
spawn_flush saw it and started a sender on every hook.
|
||||
"""
|
||||
telemetry.record("doomed")
|
||||
telemetry._post = lambda payload, url: False
|
||||
|
||||
for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3):
|
||||
telemetry.flush()
|
||||
|
||||
leftover = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
|
||||
assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}"
|
||||
|
||||
|
||||
def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry):
|
||||
"""The old shape is telemetry-<pid>-<hex>.sending, and hex can start with 'a'."""
|
||||
assert telemetry._claim_attempt(Path("telemetry-999-deadbeef.sending")) == 0
|
||||
assert telemetry._claim_attempt(Path("telemetry-999-a1234567.sending")) == 0
|
||||
assert telemetry._claim_attempt(Path("telemetry-999-deadbeef-a2.sending")) == 2
|
||||
|
||||
|
||||
def test_a_torn_claim_is_quarantined_not_deleted(telemetry):
|
||||
"""A non-empty file that parses to nothing is the remainder, not garbage."""
|
||||
telemetry.record("search")
|
||||
claim = telemetry._claim_spool()
|
||||
claim.write_bytes(b"\xff\xfe not utf-8 at all")
|
||||
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
|
||||
os.utime(claim, (stale, stale))
|
||||
|
||||
sent = telemetry.flush()
|
||||
|
||||
assert sent == 0
|
||||
assert not claim.exists()
|
||||
quarantined = list(telemetry.memory_core.data_dir().glob("*.corrupt"))
|
||||
assert len(quarantined) == 1, "torn claim was destroyed instead of kept"
|
||||
|
||||
|
||||
def test_a_failed_rewrite_stops_instead_of_redelivering(telemetry):
|
||||
"""Ignoring the rewrite result reintroduced the duplicates this PR fixes."""
|
||||
for index in range(250):
|
||||
telemetry.record("search", index=index)
|
||||
|
||||
telemetry._rewrite_claim = lambda claim, remaining: False
|
||||
delivered = []
|
||||
telemetry._post = lambda payload, url: delivered.extend(payload.get("batch", [])) or True
|
||||
|
||||
telemetry.flush()
|
||||
assert len(delivered) == 100, f"kept going after a failed rewrite: {len(delivered)}"
|
||||
|
||||
|
||||
def test_partial_files_are_swept(telemetry):
|
||||
"""Nothing else globs *.partial, so a crash mid-rename orphans one forever."""
|
||||
data_dir = telemetry.memory_core.data_dir()
|
||||
data_dir.mkdir(parents=True, exist_ok=True)
|
||||
debris = data_dir / "telemetry-1-abc-a0.1.partial"
|
||||
debris.write_text("x", encoding="utf-8")
|
||||
old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
|
||||
os.utime(debris, (old, old))
|
||||
|
||||
telemetry.flush()
|
||||
assert not debris.exists()
|
||||
|
||||
@@ -181,6 +181,37 @@ def test_source_tag_defaults_agree_between_the_two_modules():
|
||||
|
||||
def test_the_plugin_declares_its_surface_in_the_body_and_the_headers():
|
||||
"""Body and headers both, because only the body works on every backend."""
|
||||
|
||||
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")
|
||||
@@ -204,3 +235,35 @@ def test_the_plugin_declares_its_surface_in_the_body_and_the_headers():
|
||||
# The transport headers the three call sites relied on must survive.
|
||||
assert headers["auth"] == "Token k"
|
||||
assert headers["ctype"] == "application/json"
|
||||
|
||||
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"]
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -49,6 +49,7 @@ except ImportError:
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_salt_cache: str = ""
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _salt_path() -> Path:
|
||||
return memory_core.data_dir() / "telemetry-salt"
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
"""Random per-install salt, created once and memoized for the process.
|
||||
|
||||
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
|
||||
the identity file. Three reasons, all of which produced wrong data when this
|
||||
lived in the identity dict:
|
||||
|
||||
- Hooks are short-lived separate processes firing on every tool call, and
|
||||
people run more than one agent window. A read-modify-write would let each
|
||||
process mint its own salt, so one repository would hash several ways in the
|
||||
window before a writer won.
|
||||
- resolve_distinct_id holds a copy of the identity dict across a network call
|
||||
to /v1/ping/, so whichever write landed second erased the other's key —
|
||||
losing either the salt (repo_hash changes mid-stream) or the email (a
|
||||
second $identify, splitting the person).
|
||||
- Touching the identity file from record() would create it, and is_first_run
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
@@ -221,18 +267,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 +329,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 +352,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:
|
||||
@@ -396,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
|
||||
|
||||
Anchored on field position, not on a leading "a": the legacy shape is
|
||||
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
|
||||
otherwise parse as attempt 1234567 and be discarded unsent on the first
|
||||
flush after an upgrade.
|
||||
"""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
parts = stem.split("-")
|
||||
if len(parts) != 4:
|
||||
return 0
|
||||
tail = parts[3]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
@@ -432,6 +536,21 @@ def _claim_spool() -> Path | None:
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _sweep_debris(directory: Path) -> None:
|
||||
"""Remove temp files orphaned by a crash between write and rename.
|
||||
|
||||
Neither glob in this module matches *.partial, so nothing else would ever
|
||||
clean them up.
|
||||
"""
|
||||
now = time.time()
|
||||
for debris in directory.glob("telemetry-*.partial"):
|
||||
try:
|
||||
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
|
||||
debris.unlink()
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
@@ -447,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
# Attempts, not age. Every re-claim touches the mtime and every release
|
||||
# backdates it by a fixed amount, so age is pinned near the stale
|
||||
# threshold and never reaches the expiry. Age stays only as a backstop
|
||||
# for files that never carried an attempt marker.
|
||||
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
@@ -489,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
|
||||
# fsync before the rename: without it the rename can land while the
|
||||
# bytes have not, and the claim comes back empty or truncated after a
|
||||
# crash. _drain then reads zero events and unlinks it.
|
||||
with open(temporary, "w", encoding="utf-8") as handle:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
@@ -563,8 +690,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 +713,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)
|
||||
@@ -599,6 +742,7 @@ def flush() -> int:
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
_sweep_debris(directory)
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
@@ -619,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
except (OSError, ValueError):
|
||||
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
|
||||
# rather than retry: flush() runs from a bare `finally:` in
|
||||
# flush_worker, so raising here also skips the handoff cleanup, and an
|
||||
# undecodable file would otherwise be re-read on every flush forever.
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt"))
|
||||
except OSError:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
@@ -630,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
if isinstance(value, dict) and value.get("event"):
|
||||
events.append(value)
|
||||
if not events:
|
||||
# Only delete when the file really is empty. A non-empty file that
|
||||
# parses to nothing is a torn write, and its contents are the unsent
|
||||
# remainder — deleting it is the data loss this PR exists to prevent.
|
||||
try:
|
||||
claim.unlink()
|
||||
empty = claim.stat().st_size == 0
|
||||
except OSError:
|
||||
empty = True
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
@@ -681,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
# a crash repeats at most one batch instead of the entire file. If the
|
||||
# rewrite fails the claim still holds delivered events, so stop rather
|
||||
# than carry on as though progress were recorded — continuing is how the
|
||||
# duplicate delivery this PR fixes would come back.
|
||||
if not _rewrite_claim(claim, events[start + len(chunk) :]):
|
||||
_release_claim(claim, events[start + len(chunk) :])
|
||||
return sent, False
|
||||
return sent, True
|
||||
|
||||
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -49,6 +49,7 @@ except ImportError:
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_salt_cache: str = ""
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _salt_path() -> Path:
|
||||
return memory_core.data_dir() / "telemetry-salt"
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
"""Random per-install salt, created once and memoized for the process.
|
||||
|
||||
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
|
||||
the identity file. Three reasons, all of which produced wrong data when this
|
||||
lived in the identity dict:
|
||||
|
||||
- Hooks are short-lived separate processes firing on every tool call, and
|
||||
people run more than one agent window. A read-modify-write would let each
|
||||
process mint its own salt, so one repository would hash several ways in the
|
||||
window before a writer won.
|
||||
- resolve_distinct_id holds a copy of the identity dict across a network call
|
||||
to /v1/ping/, so whichever write landed second erased the other's key —
|
||||
losing either the salt (repo_hash changes mid-stream) or the email (a
|
||||
second $identify, splitting the person).
|
||||
- Touching the identity file from record() would create it, and is_first_run
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
@@ -221,18 +267,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 +329,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 +352,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:
|
||||
@@ -396,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
|
||||
|
||||
Anchored on field position, not on a leading "a": the legacy shape is
|
||||
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
|
||||
otherwise parse as attempt 1234567 and be discarded unsent on the first
|
||||
flush after an upgrade.
|
||||
"""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
parts = stem.split("-")
|
||||
if len(parts) != 4:
|
||||
return 0
|
||||
tail = parts[3]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
@@ -432,6 +536,21 @@ def _claim_spool() -> Path | None:
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _sweep_debris(directory: Path) -> None:
|
||||
"""Remove temp files orphaned by a crash between write and rename.
|
||||
|
||||
Neither glob in this module matches *.partial, so nothing else would ever
|
||||
clean them up.
|
||||
"""
|
||||
now = time.time()
|
||||
for debris in directory.glob("telemetry-*.partial"):
|
||||
try:
|
||||
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
|
||||
debris.unlink()
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
@@ -447,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
# Attempts, not age. Every re-claim touches the mtime and every release
|
||||
# backdates it by a fixed amount, so age is pinned near the stale
|
||||
# threshold and never reaches the expiry. Age stays only as a backstop
|
||||
# for files that never carried an attempt marker.
|
||||
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
@@ -489,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
|
||||
# fsync before the rename: without it the rename can land while the
|
||||
# bytes have not, and the claim comes back empty or truncated after a
|
||||
# crash. _drain then reads zero events and unlinks it.
|
||||
with open(temporary, "w", encoding="utf-8") as handle:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
@@ -563,8 +690,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 +713,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)
|
||||
@@ -599,6 +742,7 @@ def flush() -> int:
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
_sweep_debris(directory)
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
@@ -619,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
except (OSError, ValueError):
|
||||
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
|
||||
# rather than retry: flush() runs from a bare `finally:` in
|
||||
# flush_worker, so raising here also skips the handoff cleanup, and an
|
||||
# undecodable file would otherwise be re-read on every flush forever.
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt"))
|
||||
except OSError:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
@@ -630,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
if isinstance(value, dict) and value.get("event"):
|
||||
events.append(value)
|
||||
if not events:
|
||||
# Only delete when the file really is empty. A non-empty file that
|
||||
# parses to nothing is a torn write, and its contents are the unsent
|
||||
# remainder — deleting it is the data loss this PR exists to prevent.
|
||||
try:
|
||||
claim.unlink()
|
||||
empty = claim.stat().st_size == 0
|
||||
except OSError:
|
||||
empty = True
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
@@ -681,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
# a crash repeats at most one batch instead of the entire file. If the
|
||||
# rewrite fails the claim still holds delivered events, so stop rather
|
||||
# than carry on as though progress were recorded — continuing is how the
|
||||
# duplicate delivery this PR fixes would come back.
|
||||
if not _rewrite_claim(claim, events[start + len(chunk) :]):
|
||||
_release_claim(claim, events[start + len(chunk) :])
|
||||
return sent, False
|
||||
return sent, True
|
||||
|
||||
|
||||
|
||||
@@ -315,3 +315,49 @@ def test_spawn_flush_does_nothing_without_a_spool(isolated_env):
|
||||
with patch.object(telemetry.subprocess, "Popen") as popen:
|
||||
assert telemetry.spawn_flush() is True
|
||||
popen.assert_called_once()
|
||||
|
||||
|
||||
def test_salt_is_stable_across_processes(isolated_env):
|
||||
"""Hooks are separate short-lived processes; one repo must hash one way.
|
||||
|
||||
An unlocked read-modify-write let each process mint its own salt, so a
|
||||
repository hashed several ways in the window before one writer won.
|
||||
"""
|
||||
import subprocess as sp
|
||||
|
||||
core = str(Path(__file__).resolve().parents[1] / "core")
|
||||
script = (
|
||||
f"import sys; sys.path.insert(0, {core!r})\n"
|
||||
"import telemetry\n"
|
||||
"print(telemetry._install_salt())"
|
||||
)
|
||||
env = {**os.environ, "MEM0_CODE_DATA_DIR": str(memory_core.data_dir())}
|
||||
salts = {
|
||||
sp.run([sys.executable, "-c", script], capture_output=True, text=True, env=env).stdout.strip()
|
||||
for _ in range(4)
|
||||
}
|
||||
assert len(salts) == 1, f"one repo hashed {len(salts)} ways: {salts}"
|
||||
|
||||
|
||||
def test_salt_does_not_touch_the_identity_file(isolated_env):
|
||||
"""The identity file is is_first_run's marker and the sender's email store.
|
||||
|
||||
Writing the salt into it would create it from record(), suppressing the
|
||||
install event, and would race resolve_distinct_id, which holds a stale copy
|
||||
of that dict across a network call.
|
||||
"""
|
||||
telemetry._install_salt()
|
||||
assert not telemetry._identity_path().exists()
|
||||
|
||||
|
||||
def test_salt_is_stable_when_it_cannot_be_persisted(isolated_env, monkeypatch):
|
||||
"""A read-only data dir must degrade to a weaker salt, not to random-per-call.
|
||||
|
||||
Random per call is unbounded cardinality in PostHog, which is worse than no
|
||||
salt at all.
|
||||
"""
|
||||
telemetry._salt_cache = ""
|
||||
monkeypatch.setattr(telemetry.os, "open", lambda *a, **k: (_ for _ in ()).throw(OSError("read-only")))
|
||||
first = telemetry._install_salt()
|
||||
telemetry._salt_cache = ""
|
||||
assert telemetry._install_salt() == first
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -49,6 +49,7 @@ except ImportError:
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_salt_cache: str = ""
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _salt_path() -> Path:
|
||||
return memory_core.data_dir() / "telemetry-salt"
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
"""Random per-install salt, created once and memoized for the process.
|
||||
|
||||
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
|
||||
the identity file. Three reasons, all of which produced wrong data when this
|
||||
lived in the identity dict:
|
||||
|
||||
- Hooks are short-lived separate processes firing on every tool call, and
|
||||
people run more than one agent window. A read-modify-write would let each
|
||||
process mint its own salt, so one repository would hash several ways in the
|
||||
window before a writer won.
|
||||
- resolve_distinct_id holds a copy of the identity dict across a network call
|
||||
to /v1/ping/, so whichever write landed second erased the other's key —
|
||||
losing either the salt (repo_hash changes mid-stream) or the email (a
|
||||
second $identify, splitting the person).
|
||||
- Touching the identity file from record() would create it, and is_first_run
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
@@ -221,18 +267,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 +329,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 +352,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:
|
||||
@@ -396,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
|
||||
|
||||
Anchored on field position, not on a leading "a": the legacy shape is
|
||||
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
|
||||
otherwise parse as attempt 1234567 and be discarded unsent on the first
|
||||
flush after an upgrade.
|
||||
"""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
parts = stem.split("-")
|
||||
if len(parts) != 4:
|
||||
return 0
|
||||
tail = parts[3]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
@@ -432,6 +536,21 @@ def _claim_spool() -> Path | None:
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _sweep_debris(directory: Path) -> None:
|
||||
"""Remove temp files orphaned by a crash between write and rename.
|
||||
|
||||
Neither glob in this module matches *.partial, so nothing else would ever
|
||||
clean them up.
|
||||
"""
|
||||
now = time.time()
|
||||
for debris in directory.glob("telemetry-*.partial"):
|
||||
try:
|
||||
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
|
||||
debris.unlink()
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
@@ -447,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
# Attempts, not age. Every re-claim touches the mtime and every release
|
||||
# backdates it by a fixed amount, so age is pinned near the stale
|
||||
# threshold and never reaches the expiry. Age stays only as a backstop
|
||||
# for files that never carried an attempt marker.
|
||||
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
@@ -489,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
|
||||
# fsync before the rename: without it the rename can land while the
|
||||
# bytes have not, and the claim comes back empty or truncated after a
|
||||
# crash. _drain then reads zero events and unlinks it.
|
||||
with open(temporary, "w", encoding="utf-8") as handle:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
@@ -563,8 +690,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 +713,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)
|
||||
@@ -599,6 +742,7 @@ def flush() -> int:
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
_sweep_debris(directory)
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
@@ -619,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
except (OSError, ValueError):
|
||||
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
|
||||
# rather than retry: flush() runs from a bare `finally:` in
|
||||
# flush_worker, so raising here also skips the handoff cleanup, and an
|
||||
# undecodable file would otherwise be re-read on every flush forever.
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt"))
|
||||
except OSError:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
@@ -630,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
if isinstance(value, dict) and value.get("event"):
|
||||
events.append(value)
|
||||
if not events:
|
||||
# Only delete when the file really is empty. A non-empty file that
|
||||
# parses to nothing is a torn write, and its contents are the unsent
|
||||
# remainder — deleting it is the data loss this PR exists to prevent.
|
||||
try:
|
||||
claim.unlink()
|
||||
empty = claim.stat().st_size == 0
|
||||
except OSError:
|
||||
empty = True
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
@@ -681,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
# a crash repeats at most one batch instead of the entire file. If the
|
||||
# rewrite fails the claim still holds delivered events, so stop rather
|
||||
# than carry on as though progress were recorded — continuing is how the
|
||||
# duplicate delivery this PR fixes would come back.
|
||||
if not _rewrite_claim(claim, events[start + len(chunk) :]):
|
||||
_release_claim(claim, events[start + len(chunk) :])
|
||||
return sent, False
|
||||
return sent, True
|
||||
|
||||
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -49,6 +49,7 @@ except ImportError:
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_salt_cache: str = ""
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _salt_path() -> Path:
|
||||
return memory_core.data_dir() / "telemetry-salt"
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
"""Random per-install salt, created once and memoized for the process.
|
||||
|
||||
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
|
||||
the identity file. Three reasons, all of which produced wrong data when this
|
||||
lived in the identity dict:
|
||||
|
||||
- Hooks are short-lived separate processes firing on every tool call, and
|
||||
people run more than one agent window. A read-modify-write would let each
|
||||
process mint its own salt, so one repository would hash several ways in the
|
||||
window before a writer won.
|
||||
- resolve_distinct_id holds a copy of the identity dict across a network call
|
||||
to /v1/ping/, so whichever write landed second erased the other's key —
|
||||
losing either the salt (repo_hash changes mid-stream) or the email (a
|
||||
second $identify, splitting the person).
|
||||
- Touching the identity file from record() would create it, and is_first_run
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
@@ -221,18 +267,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 +329,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 +352,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:
|
||||
@@ -396,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
|
||||
|
||||
Anchored on field position, not on a leading "a": the legacy shape is
|
||||
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
|
||||
otherwise parse as attempt 1234567 and be discarded unsent on the first
|
||||
flush after an upgrade.
|
||||
"""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
parts = stem.split("-")
|
||||
if len(parts) != 4:
|
||||
return 0
|
||||
tail = parts[3]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
@@ -432,6 +536,21 @@ def _claim_spool() -> Path | None:
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _sweep_debris(directory: Path) -> None:
|
||||
"""Remove temp files orphaned by a crash between write and rename.
|
||||
|
||||
Neither glob in this module matches *.partial, so nothing else would ever
|
||||
clean them up.
|
||||
"""
|
||||
now = time.time()
|
||||
for debris in directory.glob("telemetry-*.partial"):
|
||||
try:
|
||||
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
|
||||
debris.unlink()
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
@@ -447,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
# Attempts, not age. Every re-claim touches the mtime and every release
|
||||
# backdates it by a fixed amount, so age is pinned near the stale
|
||||
# threshold and never reaches the expiry. Age stays only as a backstop
|
||||
# for files that never carried an attempt marker.
|
||||
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
@@ -489,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
|
||||
# fsync before the rename: without it the rename can land while the
|
||||
# bytes have not, and the claim comes back empty or truncated after a
|
||||
# crash. _drain then reads zero events and unlinks it.
|
||||
with open(temporary, "w", encoding="utf-8") as handle:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
@@ -563,8 +690,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 +713,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)
|
||||
@@ -599,6 +742,7 @@ def flush() -> int:
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
_sweep_debris(directory)
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
@@ -619,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
except (OSError, ValueError):
|
||||
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
|
||||
# rather than retry: flush() runs from a bare `finally:` in
|
||||
# flush_worker, so raising here also skips the handoff cleanup, and an
|
||||
# undecodable file would otherwise be re-read on every flush forever.
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt"))
|
||||
except OSError:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
@@ -630,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
if isinstance(value, dict) and value.get("event"):
|
||||
events.append(value)
|
||||
if not events:
|
||||
# Only delete when the file really is empty. A non-empty file that
|
||||
# parses to nothing is a torn write, and its contents are the unsent
|
||||
# remainder — deleting it is the data loss this PR exists to prevent.
|
||||
try:
|
||||
claim.unlink()
|
||||
empty = claim.stat().st_size == 0
|
||||
except OSError:
|
||||
empty = True
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
@@ -681,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
# a crash repeats at most one batch instead of the entire file. If the
|
||||
# rewrite fails the claim still holds delivered events, so stop rather
|
||||
# than carry on as though progress were recorded — continuing is how the
|
||||
# duplicate delivery this PR fixes would come back.
|
||||
if not _rewrite_claim(claim, events[start + len(chunk) :]):
|
||||
_release_claim(claim, events[start + len(chunk) :])
|
||||
return sent, False
|
||||
return sent, True
|
||||
|
||||
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -49,6 +49,7 @@ except ImportError:
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_salt_cache: str = ""
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _salt_path() -> Path:
|
||||
return memory_core.data_dir() / "telemetry-salt"
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
"""Random per-install salt, created once and memoized for the process.
|
||||
|
||||
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
|
||||
the identity file. Three reasons, all of which produced wrong data when this
|
||||
lived in the identity dict:
|
||||
|
||||
- Hooks are short-lived separate processes firing on every tool call, and
|
||||
people run more than one agent window. A read-modify-write would let each
|
||||
process mint its own salt, so one repository would hash several ways in the
|
||||
window before a writer won.
|
||||
- resolve_distinct_id holds a copy of the identity dict across a network call
|
||||
to /v1/ping/, so whichever write landed second erased the other's key —
|
||||
losing either the salt (repo_hash changes mid-stream) or the email (a
|
||||
second $identify, splitting the person).
|
||||
- Touching the identity file from record() would create it, and is_first_run
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
@@ -221,18 +267,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 +329,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 +352,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:
|
||||
@@ -396,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
|
||||
|
||||
Anchored on field position, not on a leading "a": the legacy shape is
|
||||
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
|
||||
otherwise parse as attempt 1234567 and be discarded unsent on the first
|
||||
flush after an upgrade.
|
||||
"""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
parts = stem.split("-")
|
||||
if len(parts) != 4:
|
||||
return 0
|
||||
tail = parts[3]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
@@ -432,6 +536,21 @@ def _claim_spool() -> Path | None:
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _sweep_debris(directory: Path) -> None:
|
||||
"""Remove temp files orphaned by a crash between write and rename.
|
||||
|
||||
Neither glob in this module matches *.partial, so nothing else would ever
|
||||
clean them up.
|
||||
"""
|
||||
now = time.time()
|
||||
for debris in directory.glob("telemetry-*.partial"):
|
||||
try:
|
||||
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
|
||||
debris.unlink()
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
@@ -447,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
# Attempts, not age. Every re-claim touches the mtime and every release
|
||||
# backdates it by a fixed amount, so age is pinned near the stale
|
||||
# threshold and never reaches the expiry. Age stays only as a backstop
|
||||
# for files that never carried an attempt marker.
|
||||
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
@@ -489,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
|
||||
# fsync before the rename: without it the rename can land while the
|
||||
# bytes have not, and the claim comes back empty or truncated after a
|
||||
# crash. _drain then reads zero events and unlinks it.
|
||||
with open(temporary, "w", encoding="utf-8") as handle:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
@@ -563,8 +690,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 +713,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)
|
||||
@@ -599,6 +742,7 @@ def flush() -> int:
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
_sweep_debris(directory)
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
@@ -619,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
except (OSError, ValueError):
|
||||
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
|
||||
# rather than retry: flush() runs from a bare `finally:` in
|
||||
# flush_worker, so raising here also skips the handoff cleanup, and an
|
||||
# undecodable file would otherwise be re-read on every flush forever.
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt"))
|
||||
except OSError:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
@@ -630,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
if isinstance(value, dict) and value.get("event"):
|
||||
events.append(value)
|
||||
if not events:
|
||||
# Only delete when the file really is empty. A non-empty file that
|
||||
# parses to nothing is a torn write, and its contents are the unsent
|
||||
# remainder — deleting it is the data loss this PR exists to prevent.
|
||||
try:
|
||||
claim.unlink()
|
||||
empty = claim.stat().st_size == 0
|
||||
except OSError:
|
||||
empty = True
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
@@ -681,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
# a crash repeats at most one batch instead of the entire file. If the
|
||||
# rewrite fails the claim still holds delivered events, so stop rather
|
||||
# than carry on as though progress were recorded — continuing is how the
|
||||
# duplicate delivery this PR fixes would come back.
|
||||
if not _rewrite_claim(claim, events[start + len(chunk) :]):
|
||||
_release_claim(claim, events[start + len(chunk) :])
|
||||
return sent, False
|
||||
return sent, True
|
||||
|
||||
|
||||
|
||||
@@ -49,6 +49,7 @@ except ImportError:
|
||||
_PLATFORM_SOURCE = "MEM0_PLUGIN"
|
||||
_PLATFORM_APPLICATION = ""
|
||||
|
||||
_salt_cache: str = ""
|
||||
_harness: str = _DEFAULT_HARNESS
|
||||
_source_tag: str = _DEFAULT_SOURCE_TAG
|
||||
_PRIVATE_KEYS = {
|
||||
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
|
||||
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _salt_path() -> Path:
|
||||
return memory_core.data_dir() / "telemetry-salt"
|
||||
|
||||
|
||||
def _install_salt() -> str:
|
||||
"""Random per-install salt, created on first use and kept in the identity file."""
|
||||
identity = _read_identity()
|
||||
salt = identity.get("salt")
|
||||
if not salt:
|
||||
salt = uuid.uuid4().hex
|
||||
identity["salt"] = salt
|
||||
_write_identity(identity)
|
||||
return salt
|
||||
"""Random per-install salt, created once and memoized for the process.
|
||||
|
||||
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
|
||||
the identity file. Three reasons, all of which produced wrong data when this
|
||||
lived in the identity dict:
|
||||
|
||||
- Hooks are short-lived separate processes firing on every tool call, and
|
||||
people run more than one agent window. A read-modify-write would let each
|
||||
process mint its own salt, so one repository would hash several ways in the
|
||||
window before a writer won.
|
||||
- resolve_distinct_id holds a copy of the identity dict across a network call
|
||||
to /v1/ping/, so whichever write landed second erased the other's key —
|
||||
losing either the salt (repo_hash changes mid-stream) or the email (a
|
||||
second $identify, splitting the person).
|
||||
- Touching the identity file from record() would create it, and is_first_run
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
@@ -221,18 +267,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 +329,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 +352,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:
|
||||
@@ -396,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
|
||||
|
||||
|
||||
def _claim_attempt(claim: Path) -> int:
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
|
||||
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
|
||||
|
||||
Anchored on field position, not on a leading "a": the legacy shape is
|
||||
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
|
||||
otherwise parse as attempt 1234567 and be discarded unsent on the first
|
||||
flush after an upgrade.
|
||||
"""
|
||||
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||
tail = stem.rsplit("-", 1)[-1]
|
||||
parts = stem.split("-")
|
||||
if len(parts) != 4:
|
||||
return 0
|
||||
tail = parts[3]
|
||||
if tail.startswith("a") and tail[1:].isdigit():
|
||||
return int(tail[1:])
|
||||
return 0
|
||||
@@ -432,6 +536,21 @@ def _claim_spool() -> Path | None:
|
||||
return _claim_parked(directory)
|
||||
|
||||
|
||||
def _sweep_debris(directory: Path) -> None:
|
||||
"""Remove temp files orphaned by a crash between write and rename.
|
||||
|
||||
Neither glob in this module matches *.partial, so nothing else would ever
|
||||
clean them up.
|
||||
"""
|
||||
now = time.time()
|
||||
for debris in directory.glob("telemetry-*.partial"):
|
||||
try:
|
||||
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
|
||||
debris.unlink()
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
|
||||
def _claim_parked(directory: Path) -> Path | None:
|
||||
"""Take the oldest abandoned claim, if any lease has actually expired.
|
||||
|
||||
@@ -447,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||
# Attempts, not age. Every re-claim touches the mtime and every release
|
||||
# backdates it by a fixed amount, so age is pinned near the stale
|
||||
# threshold and never reaches the expiry. Age stays only as a backstop
|
||||
# for files that never carried an attempt marker.
|
||||
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
|
||||
try:
|
||||
orphan.unlink()
|
||||
except OSError:
|
||||
@@ -489,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
|
||||
return True
|
||||
temporary = claim.with_suffix(f".{os.getpid()}.partial")
|
||||
try:
|
||||
temporary.write_text(
|
||||
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||
encoding="utf-8",
|
||||
)
|
||||
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
|
||||
# fsync before the rename: without it the rename can land while the
|
||||
# bytes have not, and the claim comes back empty or truncated after a
|
||||
# crash. _drain then reads zero events and unlinks it.
|
||||
with open(temporary, "w", encoding="utf-8") as handle:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
temporary.replace(claim)
|
||||
_touch(claim)
|
||||
return True
|
||||
@@ -563,8 +690,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 +713,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)
|
||||
@@ -599,6 +742,7 @@ def flush() -> int:
|
||||
# Parked batches used to starve behind the live spool indefinitely. Bounded
|
||||
# per run so a long backlog cannot turn one flush into an unbounded loop.
|
||||
directory = memory_core.data_dir()
|
||||
_sweep_debris(directory)
|
||||
for _ in range(MAX_PARKED_PER_RUN):
|
||||
parked = _claim_parked(directory)
|
||||
if parked is None:
|
||||
@@ -619,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return 0, True
|
||||
try:
|
||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||
except OSError:
|
||||
except (OSError, ValueError):
|
||||
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
|
||||
# rather than retry: flush() runs from a bare `finally:` in
|
||||
# flush_worker, so raising here also skips the handoff cleanup, and an
|
||||
# undecodable file would otherwise be re-read on every flush forever.
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt"))
|
||||
except OSError:
|
||||
try:
|
||||
claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
events = []
|
||||
for line in lines:
|
||||
@@ -630,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
if isinstance(value, dict) and value.get("event"):
|
||||
events.append(value)
|
||||
if not events:
|
||||
# Only delete when the file really is empty. A non-empty file that
|
||||
# parses to nothing is a torn write, and its contents are the unsent
|
||||
# remainder — deleting it is the data loss this PR exists to prevent.
|
||||
try:
|
||||
claim.unlink()
|
||||
empty = claim.stat().st_size == 0
|
||||
except OSError:
|
||||
empty = True
|
||||
try:
|
||||
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return 0, True
|
||||
@@ -681,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
|
||||
return sent, False
|
||||
sent += len(chunk)
|
||||
# Record progress and refresh the lease after each successful batch, so
|
||||
# a crash repeats at most one batch instead of the entire file.
|
||||
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||
# a crash repeats at most one batch instead of the entire file. If the
|
||||
# rewrite fails the claim still holds delivered events, so stop rather
|
||||
# than carry on as though progress were recorded — continuing is how the
|
||||
# duplicate delivery this PR fixes would come back.
|
||||
if not _rewrite_claim(claim, events[start + len(chunk) :]):
|
||||
_release_claim(claim, events[start + len(chunk) :])
|
||||
return sent, False
|
||||
return sent, True
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user