diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 6d7aeecae..a4375f002 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -100,6 +100,16 @@ BATCH_SIZE = 100 SEND_TIMEOUT = 5 CLAIM_STALE_SECONDS = 120 CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60 +# A batch is only discarded once it has genuinely been retried this many times. +MAX_CLAIM_ATTEMPTS = 3 +# Parked claims drained per run, after the live spool. Bounded so a long backlog +# cannot turn one flush into an unbounded send loop. +MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -399,38 +409,201 @@ def spawn_flush() -> bool: return False +def _claim_name(attempt: int = 0) -> str: + """Claim filename. The attempt count rides in the name so the 7-day expiry + only ever discards a batch that was actually retried and failed.""" + return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending" + + +def _claim_attempt(claim: Path) -> int: + """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--.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 + 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 + + +def _touch(path: Path) -> None: + """Refresh mtime so a claim's age measures time since it was claimed. + + ``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created + after a quiet minute inherited the spool's last-write time and looked + abandoned the instant it was made. A second sender would then take it over + while the first was still posting, and both would deliver the batch. + """ + try: + os.utime(path, None) + except OSError: + pass + + def _claim_spool() -> Path | None: """Rename the spool aside so exactly one sender owns each batch.""" directory = memory_core.data_dir() - claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending" + claim = directory / _claim_name() spool = _spool_path() try: spool.replace(claim) + _touch(claim) return claim except OSError: pass + return _claim_parked(directory) + + +def _sweep_debris(directory: Path) -> None: + """Remove files nothing else will ever pick up again. + + *.partial is a temp file orphaned by a crash between write and rename. + *.corrupt is a batch quarantined for undecodable content. No glob in this + module matches either, so without this they accumulate on disk for the life + of the install. + + Quarantined batches are kept far longer than debris: they are the only + evidence left of events that could not be delivered, and someone diagnosing + a report of missing telemetry has to be able to find one. + """ now = time.time() - for orphan in sorted(directory.glob("telemetry-*.sending")): + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + for quarantined in directory.glob("telemetry-*.corrupt"): + try: + if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS: + quarantined.unlink() + except OSError: + continue + # The same reasoning covers *.tmp. _write_identity and _install_salt both + # create one and unlink it in a finally, which a SIGKILL skips, and no glob + # in this module matches the leftovers either. + for temporary in directory.glob("telemetry-*.tmp"): + try: + if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS: + temporary.unlink() + except OSError: + continue + + +def _claim_parked(directory: Path) -> Path | None: + """Take the oldest abandoned claim, if any lease has actually expired. + + Kept separate from the live spool so flush() can drain both in one run. + Previously parked batches were only reachable when no spool existed at all, + and because sessions keep recording there usually was one — so a batch + parked by a failed send waited until the 7-day expiry deleted it unsent, + even though its own presence is what started the sender. + """ + now = time.time() + for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime): try: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS: + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # 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: pass continue - if age < CLAIM_STALE_SECONDS: - continue + claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) + _touch(claim) return claim except OSError: continue return None +def _safe_mtime(path: Path) -> float: + try: + return path.stat().st_mtime + except OSError: + return 0.0 + + +def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: + """Persist the unsent remainder, atomically, and refresh the lease. + + Called after every successful batch. Two jobs: a retry resumes where the + send stopped instead of re-posting from the top, and the rewrite doubles as + the lease heartbeat, so a slow sender does not have its claim stolen + mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS. + """ + if not remaining: + try: + claim.unlink() + except OSError: + pass + return True + temporary = claim.with_suffix(f".{os.getpid()}.partial") + try: + 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 + except OSError: + try: + temporary.unlink() + except OSError: + pass + return False + + +def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: + """Persist the remainder and drop the lease, because this sender has given up. + + Distinct from the per-batch heartbeat: heartbeating on the way out would + make an abandoned batch look actively owned for a further + CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the + threshold lets the next flush pick it up immediately, while the attempt + count in the filename still bounds how many times that can happen. + """ + if not _rewrite_claim(claim, remaining): + return + try: + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown + os.utime(claim, (released, released)) + except OSError: + pass + + def _resolve_email(key: str) -> str: """Trade the API key for the account email so events join other Mem0 surfaces.""" url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/" @@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]: def flush() -> int: - """Drain claimed spools to PostHog and return the number of events sent.""" + """Drain the live spool, then any parked claims, and return events sent.""" if not is_enabled(): return 0 - claim = _claim_spool() + sent, delivered = _drain(_claim_spool()) + if not delivered: + # The network is failing. Retrying other batches now would only burn + # their attempt budget against the same broken connection. + return sent + + # 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: + break + count, delivered = _drain(parked) + sent += count + if not delivered: + break + return sent + + +def _drain(claim: Path | None) -> tuple[int, bool]: + """Post one claimed batch file, recording progress after every batch. + + Returns (events sent, whether everything was delivered). + """ if claim is None: - return 0 + return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # 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. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass + return 0, True except OSError: - return 0 + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: @@ -497,11 +715,18 @@ def flush() -> int: 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 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -520,10 +745,13 @@ def flush() -> int: sent = 0 for start in range(0, len(events), BATCH_SIZE): + chunk = events[start : start + BATCH_SIZE] batch = [ { "event": event["event"], "distinct_id": distinct_id, + # Carried through from record() so a resend can be collapsed. + "uuid": event.get("uuid"), "timestamp": event.get("timestamp"), "properties": { # Fallback only: events recorded by a build before source @@ -535,16 +763,24 @@ def flush() -> int: **(event.get("properties") or {}), }, } - for event in events[start : start + BATCH_SIZE] + for event in chunk ] if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL): - return sent - sent += len(batch) - try: - claim.unlink() - except OSError: - pass - return sent + # Keep only what has not been delivered, and release the lease. + # Previously the whole file was kept and the retry re-posted every + # batch, including the ones that had already arrived. + _release_claim(claim, events[start:]) + 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. 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 def main() -> int: diff --git a/integrations/agent-plugin-core/tests/test_spool_delivery.py b/integrations/agent-plugin-core/tests/test_spool_delivery.py new file mode 100644 index 000000000..accd631a6 --- /dev/null +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -0,0 +1,446 @@ +"""Delivery semantics of the telemetry spool: no duplicates, no starvation. + +These run against a built host's core in-process (not a subprocess) because they +need to inject failures into ``_post``. The identity tests next door cover the +uninitialised-process case that needs a real interpreter. +""" + +from __future__ import annotations + +import importlib +import json +import os +import sys +import time +from pathlib import Path + +import pytest + +CORE_ROOT = Path(__file__).resolve().parents[1] +REPOSITORY_ROOT = CORE_ROOT.parents[1] +HOST_CORE = REPOSITORY_ROOT / "integrations" / "claude-code-plugin" / "core" + +pytestmark = pytest.mark.skipif(not HOST_CORE.exists(), reason="claude-code-plugin is not built") + + +@pytest.fixture() +def telemetry(tmp_path, monkeypatch): + # CI runs this directory and claude-code-plugin/tests in ONE pytest process, + # and that suite's conftest sets MEM0_TELEMETRY=false at import, process-wide. + # Without this the whole file silently no-ops: record() returns early and + # every assertion sees an empty spool. Do not rely on ambient env. + monkeypatch.setenv("MEM0_TELEMETRY", "true") + monkeypatch.setenv("MEM0_CODE_DATA_DIR", str(tmp_path / "data")) + monkeypatch.syspath_prepend(str(HOST_CORE)) + + # Save and RESTORE rather than delete. claude-code-plugin/tests/conftest.py + # imports memory_core once at collection and calls configure_harness() on it; + # dropping the module left a later re-import with default harness config, so + # tests in that suite failed depending on collection order. + names = ("telemetry", "memory_core", "_harness_id") + saved = {name: sys.modules.get(name) for name in names} + for name in names: + sys.modules.pop(name, None) + + module = importlib.import_module("telemetry") + monkeypatch.setattr(module, "resolve_distinct_id", lambda: ("tester@example.com", "")) + try: + yield module + finally: + for name in names: + sys.modules.pop(name, None) + if saved[name] is not None: + sys.modules[name] = saved[name] + + +def _delivered(payloads): + return [event for payload in payloads if "batch" in payload for event in payload["batch"]] + + +def test_a_partial_failure_does_not_redeliver_what_already_arrived(telemetry): + """Defect 2a: flush kept the whole claim on failure and retried from the top. + + 150 events across two batches, the second failing, previously delivered 250. + """ + for index in range(150): + telemetry.record("search", index=index) + + sent: list[dict] = [] + calls = {"n": 0} + + def flaky(payload, url): + calls["n"] += 1 + if calls["n"] == 2: # second batch fails + return False + sent.append(payload) + return True + + telemetry._post = flaky + telemetry.flush() + + telemetry._post = lambda payload, url: sent.append(payload) or True + telemetry.flush() + + events = _delivered(sent) + assert len(events) == 150 + assert len({event["uuid"] for event in events}) == 150 + + +def test_a_fresh_claim_is_not_immediately_stealable(telemetry): + """Defect 2b: rename preserves mtime, so a claim inherited the spool's age. + + With the last write older than the stale threshold, a claim made now looked + abandoned the instant it existed and a second sender took it over. + """ + telemetry.record("search") + spool = telemetry._spool_path() + old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (old, old)) + + first = telemetry._claim_spool() + assert first is not None + + # A second sender starting right now must find nothing to take. + assert telemetry._claim_parked(first.parent) is None + + +def test_a_live_final_attempt_is_not_deleted_by_another_sender(telemetry): + """Review finding: exhaustion was judged before liveness, so owners lost batches. + + Claiming a parked file bumps its attempt count and refreshes its mtime. Once + the count reaches the budget, the owner draining it looked exhausted to every + other sender, which unlinked the file out from under it. Everything in that + batch was gone, which is precisely the loss this PR exists to stop. + """ + telemetry.record("search", reason="owned-by-the-first-sender") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + + claim = telemetry._claim_spool() + assert claim is not None + + # Walk it to the final attempt, ageing it each round so it can be re-claimed. + # _claim_spool hands back a0 and _release_claim keeps the name, so it takes + # one full round per attempt to reach the budget. + for _ in range(telemetry.MAX_CLAIM_ATTEMPTS): + # Carry the marker through each rewrite so the final assertion proves the + # events survived, not merely that some file with the right name did. + telemetry._release_claim(claim, [{"event": "code.search", "uuid": "owned-by-the-first-sender"}]) + parked = sorted(claim.parent.glob("telemetry-*.sending")) + assert parked, "the batch was dropped while still inside its budget" + os.utime(parked[0], (stale, stale)) + claim = telemetry._claim_parked(claim.parent) + assert claim is not None + + assert telemetry._claim_attempt(claim) >= telemetry.MAX_CLAIM_ATTEMPTS + assert claim.exists() + + # The owner is draining it right now: fresh mtime, live lease. + second_sender = telemetry._claim_parked(claim.parent) + + assert second_sender is None, "a second sender took a batch under a live lease" + assert claim.exists(), "a second sender deleted a batch its owner was draining" + assert "owned-by-the-first-sender" in claim.read_text(encoding="utf-8") + + +def test_an_exhausted_batch_is_still_discarded_once_its_lease_lapses(telemetry): + """The liveness check must defer the cleanup, not cancel it. + + Guards the obvious over-correction: skipping live claims is only safe if an + abandoned one at the same attempt count is still reaped on a later run. + """ + telemetry.record("search") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + + claim = telemetry._claim_spool() + assert claim is not None + exhausted = claim.parent / telemetry._claim_name(telemetry.MAX_CLAIM_ATTEMPTS) + claim.replace(exhausted) + os.utime(exhausted, (stale, stale)) + + assert telemetry._claim_parked(exhausted.parent) is None + assert not exhausted.exists(), "an abandoned exhausted batch was left behind forever" + + +def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry): + """Defect 6: parked claims were only reachable when no spool existed. + + Because sessions keep recording there usually was one, so a batch parked by + a failed send waited until the 7-day expiry deleted it unsent — even though + its own presence is what starts the sender. + """ + 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 + old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(parked[0], (old, old)) + + telemetry.record("fresh") + sent: list[dict] = [] + telemetry._post = lambda payload, url: sent.append(payload) or True + telemetry.flush() + + names = {event["event"] for event in _delivered(sent)} + assert names == {"code.parked", "code.fresh"} + + +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 + assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS + + sent: list[dict] = [] + telemetry._post = lambda payload, url: sent.append(payload) or True + telemetry.flush() + + assert [event["event"] for event in _delivered(sent)] == ["code.parked"] + + +def test_progress_is_recorded_after_every_batch(telemetry): + """A crash repeats at most one batch, not the whole file.""" + for index in range(250): + telemetry.record("search", index=index) + + calls = {"n": 0} + + def die_after_two(payload, url): + calls["n"] += 1 + if calls["n"] > 2: + return False + return True + + telemetry._post = die_after_two + telemetry.flush() + + parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) + assert len(parked) == 1 + remaining = parked[0].read_text(encoding="utf-8").strip().splitlines() + # Two batches of 100 landed; only the last 50 should still be pending. + assert len(remaining) == 50 + assert json.loads(remaining[0])["properties"]["index"] == 200 + + +def test_the_heartbeat_actually_refreshes_the_lease(telemetry): + """The claim rewrite doubles as the lease heartbeat. + + 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. + """ + 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 + + directory = telemetry.memory_core.data_dir() + for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3): + telemetry.flush() + # Attempts now carry a cooldown, so a released claim is not instantly + # reclaimable. Age it to stand in for the wall time a real retry waits; + # without this the loop spins inside one cooldown and proves nothing. + for parked in directory.glob("telemetry-*.sending"): + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(parked, (stale, stale)) + + leftover = list(directory.glob("telemetry-*.sending")) + assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}" + + +def test_a_batch_that_cannot_be_read_is_not_counted_as_delivered(telemetry): + """Review finding: a read failure reported 'everything delivered'. + + Nothing was posted, so calling it delivered lets flush() carry on to other + claims as though this batch had arrived, and hides the failure from the one + signal that says the run went badly. It also must not quarantine: a briefly + unreadable file is retryable, and moving it to .corrupt discards the events + over a transient filesystem error, because nothing ever re-globs .corrupt. + """ + telemetry.record("search") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + claim = telemetry._claim_spool() + assert claim is not None + + original = Path.read_text + + def unreadable(self, *args, **kwargs): + if self == claim: + raise OSError(5, "I/O error") + return original(self, *args, **kwargs) + + Path.read_text = unreadable + try: + sent, delivered = telemetry._drain(claim) + finally: + Path.read_text = original + + assert sent == 0 + assert delivered is False, "an unread batch was reported as delivered" + assert claim.exists(), "a transient read error discarded the batch" + assert not list(claim.parent.glob("*.corrupt")), "quarantined over a transient error" + + +def test_undecodable_content_is_still_quarantined_and_the_run_continues(telemetry): + """The other half: genuinely unrecoverable content must not block the run. + + Guards the over-correction. If every read problem returned undelivered, one + torn file would stop every later claim on every flush, forever. + """ + telemetry.record("search") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + claim = telemetry._claim_spool() + assert claim is not None + claim.write_bytes(b"\xff\xfe torn \x00 write") + + sent, delivered = telemetry._drain(claim) + + assert (sent, delivered) == (0, True) + assert not claim.exists() + assert list(claim.parent.glob("*.corrupt")), "unrecoverable content was not quarantined" + + +def test_retries_are_spread_over_real_time_not_burned_at_once(telemetry): + """Review finding: releasing straight to reclaimable spent the budget instantly. + + Two senders hitting one momentary failure could walk a batch from attempt 0 + to the limit within seconds and discard it, when a retry a minute later would + have delivered. Each release now has to age past a cooldown that grows with + the attempts already spent. + """ + telemetry.record("doomed") + telemetry._post = lambda payload, url: False + + directory = telemetry.memory_core.data_dir() + telemetry.flush() + + parked = list(directory.glob("telemetry-*.sending")) + assert parked, "the batch was discarded on its first failure" + assert telemetry._claim_attempt(parked[0]) == 0 + + # Second sender, immediately: the cooldown has not elapsed, so it must not + # be able to spend another attempt. + telemetry.flush() + still = list(directory.glob("telemetry-*.sending")) + assert len(still) == 1 + assert telemetry._claim_attempt(still[0]) <= 1, "burned attempts without waiting" + + +def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry): + """The old shape is telemetry--.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() + + +def test_quarantined_batches_are_eventually_collected(telemetry): + """Nothing re-globs .corrupt, so without a sweep they live on disk forever. + + Kept much longer than .partial debris on purpose: a quarantined batch is the + only remaining evidence of events that could not be delivered. + """ + directory = telemetry.memory_core.data_dir() + directory.mkdir(parents=True, exist_ok=True) + fresh = directory / "telemetry-1-aaaaaaaa-a0.corrupt" + old = directory / "telemetry-2-bbbbbbbb-a0.corrupt" + for path in (fresh, old): + path.write_text("torn", encoding="utf-8") + expired = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60) + os.utime(old, (expired, expired)) + + telemetry._sweep_debris(directory) + + assert fresh.exists(), "a recent quarantine was discarded before anyone could look at it" + assert not old.exists(), "an expired quarantine was left on disk forever" + + +def test_temp_files_orphaned_by_a_kill_are_collected(telemetry): + """_write_identity and _install_salt unlink in a finally, which SIGKILL skips.""" + directory = telemetry.memory_core.data_dir() + directory.mkdir(parents=True, exist_ok=True) + orphan = directory / "telemetry-salt.999.tmp" + orphan.write_text("abandoned", encoding="utf-8") + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(orphan, (stale, stale)) + + telemetry._sweep_debris(directory) + + assert not orphan.exists(), "a killed process left a temp file on disk forever" diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 6d7aeecae..a4375f002 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -100,6 +100,16 @@ BATCH_SIZE = 100 SEND_TIMEOUT = 5 CLAIM_STALE_SECONDS = 120 CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60 +# A batch is only discarded once it has genuinely been retried this many times. +MAX_CLAIM_ATTEMPTS = 3 +# Parked claims drained per run, after the live spool. Bounded so a long backlog +# cannot turn one flush into an unbounded send loop. +MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -399,38 +409,201 @@ def spawn_flush() -> bool: return False +def _claim_name(attempt: int = 0) -> str: + """Claim filename. The attempt count rides in the name so the 7-day expiry + only ever discards a batch that was actually retried and failed.""" + return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending" + + +def _claim_attempt(claim: Path) -> int: + """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--.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 + 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 + + +def _touch(path: Path) -> None: + """Refresh mtime so a claim's age measures time since it was claimed. + + ``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created + after a quiet minute inherited the spool's last-write time and looked + abandoned the instant it was made. A second sender would then take it over + while the first was still posting, and both would deliver the batch. + """ + try: + os.utime(path, None) + except OSError: + pass + + def _claim_spool() -> Path | None: """Rename the spool aside so exactly one sender owns each batch.""" directory = memory_core.data_dir() - claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending" + claim = directory / _claim_name() spool = _spool_path() try: spool.replace(claim) + _touch(claim) return claim except OSError: pass + return _claim_parked(directory) + + +def _sweep_debris(directory: Path) -> None: + """Remove files nothing else will ever pick up again. + + *.partial is a temp file orphaned by a crash between write and rename. + *.corrupt is a batch quarantined for undecodable content. No glob in this + module matches either, so without this they accumulate on disk for the life + of the install. + + Quarantined batches are kept far longer than debris: they are the only + evidence left of events that could not be delivered, and someone diagnosing + a report of missing telemetry has to be able to find one. + """ now = time.time() - for orphan in sorted(directory.glob("telemetry-*.sending")): + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + for quarantined in directory.glob("telemetry-*.corrupt"): + try: + if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS: + quarantined.unlink() + except OSError: + continue + # The same reasoning covers *.tmp. _write_identity and _install_salt both + # create one and unlink it in a finally, which a SIGKILL skips, and no glob + # in this module matches the leftovers either. + for temporary in directory.glob("telemetry-*.tmp"): + try: + if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS: + temporary.unlink() + except OSError: + continue + + +def _claim_parked(directory: Path) -> Path | None: + """Take the oldest abandoned claim, if any lease has actually expired. + + Kept separate from the live spool so flush() can drain both in one run. + Previously parked batches were only reachable when no spool existed at all, + and because sessions keep recording there usually was one — so a batch + parked by a failed send waited until the 7-day expiry deleted it unsent, + even though its own presence is what started the sender. + """ + now = time.time() + for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime): try: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS: + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # 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: pass continue - if age < CLAIM_STALE_SECONDS: - continue + claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) + _touch(claim) return claim except OSError: continue return None +def _safe_mtime(path: Path) -> float: + try: + return path.stat().st_mtime + except OSError: + return 0.0 + + +def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: + """Persist the unsent remainder, atomically, and refresh the lease. + + Called after every successful batch. Two jobs: a retry resumes where the + send stopped instead of re-posting from the top, and the rewrite doubles as + the lease heartbeat, so a slow sender does not have its claim stolen + mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS. + """ + if not remaining: + try: + claim.unlink() + except OSError: + pass + return True + temporary = claim.with_suffix(f".{os.getpid()}.partial") + try: + 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 + except OSError: + try: + temporary.unlink() + except OSError: + pass + return False + + +def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: + """Persist the remainder and drop the lease, because this sender has given up. + + Distinct from the per-batch heartbeat: heartbeating on the way out would + make an abandoned batch look actively owned for a further + CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the + threshold lets the next flush pick it up immediately, while the attempt + count in the filename still bounds how many times that can happen. + """ + if not _rewrite_claim(claim, remaining): + return + try: + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown + os.utime(claim, (released, released)) + except OSError: + pass + + def _resolve_email(key: str) -> str: """Trade the API key for the account email so events join other Mem0 surfaces.""" url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/" @@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]: def flush() -> int: - """Drain claimed spools to PostHog and return the number of events sent.""" + """Drain the live spool, then any parked claims, and return events sent.""" if not is_enabled(): return 0 - claim = _claim_spool() + sent, delivered = _drain(_claim_spool()) + if not delivered: + # The network is failing. Retrying other batches now would only burn + # their attempt budget against the same broken connection. + return sent + + # 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: + break + count, delivered = _drain(parked) + sent += count + if not delivered: + break + return sent + + +def _drain(claim: Path | None) -> tuple[int, bool]: + """Post one claimed batch file, recording progress after every batch. + + Returns (events sent, whether everything was delivered). + """ if claim is None: - return 0 + return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # 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. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass + return 0, True except OSError: - return 0 + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: @@ -497,11 +715,18 @@ def flush() -> int: 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 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -520,10 +745,13 @@ def flush() -> int: sent = 0 for start in range(0, len(events), BATCH_SIZE): + chunk = events[start : start + BATCH_SIZE] batch = [ { "event": event["event"], "distinct_id": distinct_id, + # Carried through from record() so a resend can be collapsed. + "uuid": event.get("uuid"), "timestamp": event.get("timestamp"), "properties": { # Fallback only: events recorded by a build before source @@ -535,16 +763,24 @@ def flush() -> int: **(event.get("properties") or {}), }, } - for event in events[start : start + BATCH_SIZE] + for event in chunk ] if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL): - return sent - sent += len(batch) - try: - claim.unlink() - except OSError: - pass - return sent + # Keep only what has not been delivered, and release the lease. + # Previously the whole file was kept and the retry re-posted every + # batch, including the ones that had already arrived. + _release_claim(claim, events[start:]) + 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. 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 def main() -> int: diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 6d7aeecae..a4375f002 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -100,6 +100,16 @@ BATCH_SIZE = 100 SEND_TIMEOUT = 5 CLAIM_STALE_SECONDS = 120 CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60 +# A batch is only discarded once it has genuinely been retried this many times. +MAX_CLAIM_ATTEMPTS = 3 +# Parked claims drained per run, after the live spool. Bounded so a long backlog +# cannot turn one flush into an unbounded send loop. +MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -399,38 +409,201 @@ def spawn_flush() -> bool: return False +def _claim_name(attempt: int = 0) -> str: + """Claim filename. The attempt count rides in the name so the 7-day expiry + only ever discards a batch that was actually retried and failed.""" + return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending" + + +def _claim_attempt(claim: Path) -> int: + """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--.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 + 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 + + +def _touch(path: Path) -> None: + """Refresh mtime so a claim's age measures time since it was claimed. + + ``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created + after a quiet minute inherited the spool's last-write time and looked + abandoned the instant it was made. A second sender would then take it over + while the first was still posting, and both would deliver the batch. + """ + try: + os.utime(path, None) + except OSError: + pass + + def _claim_spool() -> Path | None: """Rename the spool aside so exactly one sender owns each batch.""" directory = memory_core.data_dir() - claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending" + claim = directory / _claim_name() spool = _spool_path() try: spool.replace(claim) + _touch(claim) return claim except OSError: pass + return _claim_parked(directory) + + +def _sweep_debris(directory: Path) -> None: + """Remove files nothing else will ever pick up again. + + *.partial is a temp file orphaned by a crash between write and rename. + *.corrupt is a batch quarantined for undecodable content. No glob in this + module matches either, so without this they accumulate on disk for the life + of the install. + + Quarantined batches are kept far longer than debris: they are the only + evidence left of events that could not be delivered, and someone diagnosing + a report of missing telemetry has to be able to find one. + """ now = time.time() - for orphan in sorted(directory.glob("telemetry-*.sending")): + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + for quarantined in directory.glob("telemetry-*.corrupt"): + try: + if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS: + quarantined.unlink() + except OSError: + continue + # The same reasoning covers *.tmp. _write_identity and _install_salt both + # create one and unlink it in a finally, which a SIGKILL skips, and no glob + # in this module matches the leftovers either. + for temporary in directory.glob("telemetry-*.tmp"): + try: + if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS: + temporary.unlink() + except OSError: + continue + + +def _claim_parked(directory: Path) -> Path | None: + """Take the oldest abandoned claim, if any lease has actually expired. + + Kept separate from the live spool so flush() can drain both in one run. + Previously parked batches were only reachable when no spool existed at all, + and because sessions keep recording there usually was one — so a batch + parked by a failed send waited until the 7-day expiry deleted it unsent, + even though its own presence is what started the sender. + """ + now = time.time() + for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime): try: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS: + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # 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: pass continue - if age < CLAIM_STALE_SECONDS: - continue + claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) + _touch(claim) return claim except OSError: continue return None +def _safe_mtime(path: Path) -> float: + try: + return path.stat().st_mtime + except OSError: + return 0.0 + + +def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: + """Persist the unsent remainder, atomically, and refresh the lease. + + Called after every successful batch. Two jobs: a retry resumes where the + send stopped instead of re-posting from the top, and the rewrite doubles as + the lease heartbeat, so a slow sender does not have its claim stolen + mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS. + """ + if not remaining: + try: + claim.unlink() + except OSError: + pass + return True + temporary = claim.with_suffix(f".{os.getpid()}.partial") + try: + 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 + except OSError: + try: + temporary.unlink() + except OSError: + pass + return False + + +def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: + """Persist the remainder and drop the lease, because this sender has given up. + + Distinct from the per-batch heartbeat: heartbeating on the way out would + make an abandoned batch look actively owned for a further + CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the + threshold lets the next flush pick it up immediately, while the attempt + count in the filename still bounds how many times that can happen. + """ + if not _rewrite_claim(claim, remaining): + return + try: + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown + os.utime(claim, (released, released)) + except OSError: + pass + + def _resolve_email(key: str) -> str: """Trade the API key for the account email so events join other Mem0 surfaces.""" url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/" @@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]: def flush() -> int: - """Drain claimed spools to PostHog and return the number of events sent.""" + """Drain the live spool, then any parked claims, and return events sent.""" if not is_enabled(): return 0 - claim = _claim_spool() + sent, delivered = _drain(_claim_spool()) + if not delivered: + # The network is failing. Retrying other batches now would only burn + # their attempt budget against the same broken connection. + return sent + + # 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: + break + count, delivered = _drain(parked) + sent += count + if not delivered: + break + return sent + + +def _drain(claim: Path | None) -> tuple[int, bool]: + """Post one claimed batch file, recording progress after every batch. + + Returns (events sent, whether everything was delivered). + """ if claim is None: - return 0 + return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # 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. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass + return 0, True except OSError: - return 0 + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: @@ -497,11 +715,18 @@ def flush() -> int: 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 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -520,10 +745,13 @@ def flush() -> int: sent = 0 for start in range(0, len(events), BATCH_SIZE): + chunk = events[start : start + BATCH_SIZE] batch = [ { "event": event["event"], "distinct_id": distinct_id, + # Carried through from record() so a resend can be collapsed. + "uuid": event.get("uuid"), "timestamp": event.get("timestamp"), "properties": { # Fallback only: events recorded by a build before source @@ -535,16 +763,24 @@ def flush() -> int: **(event.get("properties") or {}), }, } - for event in events[start : start + BATCH_SIZE] + for event in chunk ] if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL): - return sent - sent += len(batch) - try: - claim.unlink() - except OSError: - pass - return sent + # Keep only what has not been delivered, and release the lease. + # Previously the whole file was kept and the retry re-posted every + # batch, including the ones that had already arrived. + _release_claim(claim, events[start:]) + 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. 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 def main() -> int: diff --git a/integrations/claude-code-plugin/tests/test_telemetry.py b/integrations/claude-code-plugin/tests/test_telemetry.py index 6029fa98a..de0db10c6 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -199,9 +199,11 @@ def test_a_stale_claim_is_reclaimed(isolated_env, monkeypatch): telemetry.record("search") orphan = telemetry._claim_spool() assert orphan is not None - monkeypatch.setattr( - telemetry.time, "time", lambda: orphan.stat().st_mtime + telemetry.CLAIM_STALE_SECONDS + 1 - ) + # Frozen rather than re-stat'd per call: flush() drains the live spool and + # then looks for parked claims in the same run, so by the second look this + # file no longer exists. + stale_now = orphan.stat().st_mtime + telemetry.CLAIM_STALE_SECONDS + 1 + monkeypatch.setattr(telemetry.time, "time", lambda: stale_now) with patch.object(telemetry, "_post", lambda payload, url: True): assert telemetry.flush() == 1 @@ -211,9 +213,13 @@ def test_an_expired_claim_is_dropped(isolated_env, monkeypatch): telemetry.record("search") orphan = telemetry._claim_spool() assert orphan is not None - monkeypatch.setattr( - telemetry.time, "time", lambda: orphan.stat().st_mtime + telemetry.CLAIM_EXPIRY_SECONDS + 1 - ) + expired_now = orphan.stat().st_mtime + telemetry.CLAIM_EXPIRY_SECONDS + 1 + monkeypatch.setattr(telemetry.time, "time", lambda: expired_now) + # Expiry now only discards a batch that was genuinely retried and failed, + # so age alone is not enough — age it past the attempt budget too. + retried = orphan.parent / orphan.name.replace("-a0.", f"-a{telemetry.MAX_CLAIM_ATTEMPTS}.") + orphan.replace(retried) + os.utime(retried, (expired_now, expired_now - telemetry.CLAIM_EXPIRY_SECONDS - 1)) assert telemetry._claim_spool() is None assert not list(memory_core.data_dir().glob("telemetry-*.sending")) diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 6d7aeecae..a4375f002 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -100,6 +100,16 @@ BATCH_SIZE = 100 SEND_TIMEOUT = 5 CLAIM_STALE_SECONDS = 120 CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60 +# A batch is only discarded once it has genuinely been retried this many times. +MAX_CLAIM_ATTEMPTS = 3 +# Parked claims drained per run, after the live spool. Bounded so a long backlog +# cannot turn one flush into an unbounded send loop. +MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -399,38 +409,201 @@ def spawn_flush() -> bool: return False +def _claim_name(attempt: int = 0) -> str: + """Claim filename. The attempt count rides in the name so the 7-day expiry + only ever discards a batch that was actually retried and failed.""" + return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending" + + +def _claim_attempt(claim: Path) -> int: + """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--.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 + 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 + + +def _touch(path: Path) -> None: + """Refresh mtime so a claim's age measures time since it was claimed. + + ``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created + after a quiet minute inherited the spool's last-write time and looked + abandoned the instant it was made. A second sender would then take it over + while the first was still posting, and both would deliver the batch. + """ + try: + os.utime(path, None) + except OSError: + pass + + def _claim_spool() -> Path | None: """Rename the spool aside so exactly one sender owns each batch.""" directory = memory_core.data_dir() - claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending" + claim = directory / _claim_name() spool = _spool_path() try: spool.replace(claim) + _touch(claim) return claim except OSError: pass + return _claim_parked(directory) + + +def _sweep_debris(directory: Path) -> None: + """Remove files nothing else will ever pick up again. + + *.partial is a temp file orphaned by a crash between write and rename. + *.corrupt is a batch quarantined for undecodable content. No glob in this + module matches either, so without this they accumulate on disk for the life + of the install. + + Quarantined batches are kept far longer than debris: they are the only + evidence left of events that could not be delivered, and someone diagnosing + a report of missing telemetry has to be able to find one. + """ now = time.time() - for orphan in sorted(directory.glob("telemetry-*.sending")): + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + for quarantined in directory.glob("telemetry-*.corrupt"): + try: + if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS: + quarantined.unlink() + except OSError: + continue + # The same reasoning covers *.tmp. _write_identity and _install_salt both + # create one and unlink it in a finally, which a SIGKILL skips, and no glob + # in this module matches the leftovers either. + for temporary in directory.glob("telemetry-*.tmp"): + try: + if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS: + temporary.unlink() + except OSError: + continue + + +def _claim_parked(directory: Path) -> Path | None: + """Take the oldest abandoned claim, if any lease has actually expired. + + Kept separate from the live spool so flush() can drain both in one run. + Previously parked batches were only reachable when no spool existed at all, + and because sessions keep recording there usually was one — so a batch + parked by a failed send waited until the 7-day expiry deleted it unsent, + even though its own presence is what started the sender. + """ + now = time.time() + for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime): try: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS: + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # 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: pass continue - if age < CLAIM_STALE_SECONDS: - continue + claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) + _touch(claim) return claim except OSError: continue return None +def _safe_mtime(path: Path) -> float: + try: + return path.stat().st_mtime + except OSError: + return 0.0 + + +def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: + """Persist the unsent remainder, atomically, and refresh the lease. + + Called after every successful batch. Two jobs: a retry resumes where the + send stopped instead of re-posting from the top, and the rewrite doubles as + the lease heartbeat, so a slow sender does not have its claim stolen + mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS. + """ + if not remaining: + try: + claim.unlink() + except OSError: + pass + return True + temporary = claim.with_suffix(f".{os.getpid()}.partial") + try: + 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 + except OSError: + try: + temporary.unlink() + except OSError: + pass + return False + + +def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: + """Persist the remainder and drop the lease, because this sender has given up. + + Distinct from the per-batch heartbeat: heartbeating on the way out would + make an abandoned batch look actively owned for a further + CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the + threshold lets the next flush pick it up immediately, while the attempt + count in the filename still bounds how many times that can happen. + """ + if not _rewrite_claim(claim, remaining): + return + try: + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown + os.utime(claim, (released, released)) + except OSError: + pass + + def _resolve_email(key: str) -> str: """Trade the API key for the account email so events join other Mem0 surfaces.""" url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/" @@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]: def flush() -> int: - """Drain claimed spools to PostHog and return the number of events sent.""" + """Drain the live spool, then any parked claims, and return events sent.""" if not is_enabled(): return 0 - claim = _claim_spool() + sent, delivered = _drain(_claim_spool()) + if not delivered: + # The network is failing. Retrying other batches now would only burn + # their attempt budget against the same broken connection. + return sent + + # 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: + break + count, delivered = _drain(parked) + sent += count + if not delivered: + break + return sent + + +def _drain(claim: Path | None) -> tuple[int, bool]: + """Post one claimed batch file, recording progress after every batch. + + Returns (events sent, whether everything was delivered). + """ if claim is None: - return 0 + return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # 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. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass + return 0, True except OSError: - return 0 + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: @@ -497,11 +715,18 @@ def flush() -> int: 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 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -520,10 +745,13 @@ def flush() -> int: sent = 0 for start in range(0, len(events), BATCH_SIZE): + chunk = events[start : start + BATCH_SIZE] batch = [ { "event": event["event"], "distinct_id": distinct_id, + # Carried through from record() so a resend can be collapsed. + "uuid": event.get("uuid"), "timestamp": event.get("timestamp"), "properties": { # Fallback only: events recorded by a build before source @@ -535,16 +763,24 @@ def flush() -> int: **(event.get("properties") or {}), }, } - for event in events[start : start + BATCH_SIZE] + for event in chunk ] if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL): - return sent - sent += len(batch) - try: - claim.unlink() - except OSError: - pass - return sent + # Keep only what has not been delivered, and release the lease. + # Previously the whole file was kept and the retry re-posted every + # batch, including the ones that had already arrived. + _release_claim(claim, events[start:]) + 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. 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 def main() -> int: diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 6d7aeecae..a4375f002 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -100,6 +100,16 @@ BATCH_SIZE = 100 SEND_TIMEOUT = 5 CLAIM_STALE_SECONDS = 120 CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60 +# A batch is only discarded once it has genuinely been retried this many times. +MAX_CLAIM_ATTEMPTS = 3 +# Parked claims drained per run, after the live spool. Bounded so a long backlog +# cannot turn one flush into an unbounded send loop. +MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -399,38 +409,201 @@ def spawn_flush() -> bool: return False +def _claim_name(attempt: int = 0) -> str: + """Claim filename. The attempt count rides in the name so the 7-day expiry + only ever discards a batch that was actually retried and failed.""" + return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending" + + +def _claim_attempt(claim: Path) -> int: + """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--.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 + 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 + + +def _touch(path: Path) -> None: + """Refresh mtime so a claim's age measures time since it was claimed. + + ``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created + after a quiet minute inherited the spool's last-write time and looked + abandoned the instant it was made. A second sender would then take it over + while the first was still posting, and both would deliver the batch. + """ + try: + os.utime(path, None) + except OSError: + pass + + def _claim_spool() -> Path | None: """Rename the spool aside so exactly one sender owns each batch.""" directory = memory_core.data_dir() - claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending" + claim = directory / _claim_name() spool = _spool_path() try: spool.replace(claim) + _touch(claim) return claim except OSError: pass + return _claim_parked(directory) + + +def _sweep_debris(directory: Path) -> None: + """Remove files nothing else will ever pick up again. + + *.partial is a temp file orphaned by a crash between write and rename. + *.corrupt is a batch quarantined for undecodable content. No glob in this + module matches either, so without this they accumulate on disk for the life + of the install. + + Quarantined batches are kept far longer than debris: they are the only + evidence left of events that could not be delivered, and someone diagnosing + a report of missing telemetry has to be able to find one. + """ now = time.time() - for orphan in sorted(directory.glob("telemetry-*.sending")): + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + for quarantined in directory.glob("telemetry-*.corrupt"): + try: + if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS: + quarantined.unlink() + except OSError: + continue + # The same reasoning covers *.tmp. _write_identity and _install_salt both + # create one and unlink it in a finally, which a SIGKILL skips, and no glob + # in this module matches the leftovers either. + for temporary in directory.glob("telemetry-*.tmp"): + try: + if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS: + temporary.unlink() + except OSError: + continue + + +def _claim_parked(directory: Path) -> Path | None: + """Take the oldest abandoned claim, if any lease has actually expired. + + Kept separate from the live spool so flush() can drain both in one run. + Previously parked batches were only reachable when no spool existed at all, + and because sessions keep recording there usually was one — so a batch + parked by a failed send waited until the 7-day expiry deleted it unsent, + even though its own presence is what started the sender. + """ + now = time.time() + for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime): try: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS: + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # 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: pass continue - if age < CLAIM_STALE_SECONDS: - continue + claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) + _touch(claim) return claim except OSError: continue return None +def _safe_mtime(path: Path) -> float: + try: + return path.stat().st_mtime + except OSError: + return 0.0 + + +def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: + """Persist the unsent remainder, atomically, and refresh the lease. + + Called after every successful batch. Two jobs: a retry resumes where the + send stopped instead of re-posting from the top, and the rewrite doubles as + the lease heartbeat, so a slow sender does not have its claim stolen + mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS. + """ + if not remaining: + try: + claim.unlink() + except OSError: + pass + return True + temporary = claim.with_suffix(f".{os.getpid()}.partial") + try: + 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 + except OSError: + try: + temporary.unlink() + except OSError: + pass + return False + + +def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: + """Persist the remainder and drop the lease, because this sender has given up. + + Distinct from the per-batch heartbeat: heartbeating on the way out would + make an abandoned batch look actively owned for a further + CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the + threshold lets the next flush pick it up immediately, while the attempt + count in the filename still bounds how many times that can happen. + """ + if not _rewrite_claim(claim, remaining): + return + try: + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown + os.utime(claim, (released, released)) + except OSError: + pass + + def _resolve_email(key: str) -> str: """Trade the API key for the account email so events join other Mem0 surfaces.""" url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/" @@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]: def flush() -> int: - """Drain claimed spools to PostHog and return the number of events sent.""" + """Drain the live spool, then any parked claims, and return events sent.""" if not is_enabled(): return 0 - claim = _claim_spool() + sent, delivered = _drain(_claim_spool()) + if not delivered: + # The network is failing. Retrying other batches now would only burn + # their attempt budget against the same broken connection. + return sent + + # 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: + break + count, delivered = _drain(parked) + sent += count + if not delivered: + break + return sent + + +def _drain(claim: Path | None) -> tuple[int, bool]: + """Post one claimed batch file, recording progress after every batch. + + Returns (events sent, whether everything was delivered). + """ if claim is None: - return 0 + return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # 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. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass + return 0, True except OSError: - return 0 + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: @@ -497,11 +715,18 @@ def flush() -> int: 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 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -520,10 +745,13 @@ def flush() -> int: sent = 0 for start in range(0, len(events), BATCH_SIZE): + chunk = events[start : start + BATCH_SIZE] batch = [ { "event": event["event"], "distinct_id": distinct_id, + # Carried through from record() so a resend can be collapsed. + "uuid": event.get("uuid"), "timestamp": event.get("timestamp"), "properties": { # Fallback only: events recorded by a build before source @@ -535,16 +763,24 @@ def flush() -> int: **(event.get("properties") or {}), }, } - for event in events[start : start + BATCH_SIZE] + for event in chunk ] if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL): - return sent - sent += len(batch) - try: - claim.unlink() - except OSError: - pass - return sent + # Keep only what has not been delivered, and release the lease. + # Previously the whole file was kept and the retry re-posted every + # batch, including the ones that had already arrived. + _release_claim(claim, events[start:]) + 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. 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 def main() -> int: diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 6d7aeecae..a4375f002 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -100,6 +100,16 @@ BATCH_SIZE = 100 SEND_TIMEOUT = 5 CLAIM_STALE_SECONDS = 120 CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60 +# A batch is only discarded once it has genuinely been retried this many times. +MAX_CLAIM_ATTEMPTS = 3 +# Parked claims drained per run, after the live spool. Bounded so a long backlog +# cannot turn one flush into an unbounded send loop. +MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -399,38 +409,201 @@ def spawn_flush() -> bool: return False +def _claim_name(attempt: int = 0) -> str: + """Claim filename. The attempt count rides in the name so the 7-day expiry + only ever discards a batch that was actually retried and failed.""" + return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending" + + +def _claim_attempt(claim: Path) -> int: + """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--.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 + 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 + + +def _touch(path: Path) -> None: + """Refresh mtime so a claim's age measures time since it was claimed. + + ``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created + after a quiet minute inherited the spool's last-write time and looked + abandoned the instant it was made. A second sender would then take it over + while the first was still posting, and both would deliver the batch. + """ + try: + os.utime(path, None) + except OSError: + pass + + def _claim_spool() -> Path | None: """Rename the spool aside so exactly one sender owns each batch.""" directory = memory_core.data_dir() - claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending" + claim = directory / _claim_name() spool = _spool_path() try: spool.replace(claim) + _touch(claim) return claim except OSError: pass + return _claim_parked(directory) + + +def _sweep_debris(directory: Path) -> None: + """Remove files nothing else will ever pick up again. + + *.partial is a temp file orphaned by a crash between write and rename. + *.corrupt is a batch quarantined for undecodable content. No glob in this + module matches either, so without this they accumulate on disk for the life + of the install. + + Quarantined batches are kept far longer than debris: they are the only + evidence left of events that could not be delivered, and someone diagnosing + a report of missing telemetry has to be able to find one. + """ now = time.time() - for orphan in sorted(directory.glob("telemetry-*.sending")): + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + for quarantined in directory.glob("telemetry-*.corrupt"): + try: + if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS: + quarantined.unlink() + except OSError: + continue + # The same reasoning covers *.tmp. _write_identity and _install_salt both + # create one and unlink it in a finally, which a SIGKILL skips, and no glob + # in this module matches the leftovers either. + for temporary in directory.glob("telemetry-*.tmp"): + try: + if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS: + temporary.unlink() + except OSError: + continue + + +def _claim_parked(directory: Path) -> Path | None: + """Take the oldest abandoned claim, if any lease has actually expired. + + Kept separate from the live spool so flush() can drain both in one run. + Previously parked batches were only reachable when no spool existed at all, + and because sessions keep recording there usually was one — so a batch + parked by a failed send waited until the 7-day expiry deleted it unsent, + even though its own presence is what started the sender. + """ + now = time.time() + for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime): try: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS: + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # 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: pass continue - if age < CLAIM_STALE_SECONDS: - continue + claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) + _touch(claim) return claim except OSError: continue return None +def _safe_mtime(path: Path) -> float: + try: + return path.stat().st_mtime + except OSError: + return 0.0 + + +def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: + """Persist the unsent remainder, atomically, and refresh the lease. + + Called after every successful batch. Two jobs: a retry resumes where the + send stopped instead of re-posting from the top, and the rewrite doubles as + the lease heartbeat, so a slow sender does not have its claim stolen + mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS. + """ + if not remaining: + try: + claim.unlink() + except OSError: + pass + return True + temporary = claim.with_suffix(f".{os.getpid()}.partial") + try: + 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 + except OSError: + try: + temporary.unlink() + except OSError: + pass + return False + + +def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: + """Persist the remainder and drop the lease, because this sender has given up. + + Distinct from the per-batch heartbeat: heartbeating on the way out would + make an abandoned batch look actively owned for a further + CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the + threshold lets the next flush pick it up immediately, while the attempt + count in the filename still bounds how many times that can happen. + """ + if not _rewrite_claim(claim, remaining): + return + try: + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown + os.utime(claim, (released, released)) + except OSError: + pass + + def _resolve_email(key: str) -> str: """Trade the API key for the account email so events join other Mem0 surfaces.""" url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/" @@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]: def flush() -> int: - """Drain claimed spools to PostHog and return the number of events sent.""" + """Drain the live spool, then any parked claims, and return events sent.""" if not is_enabled(): return 0 - claim = _claim_spool() + sent, delivered = _drain(_claim_spool()) + if not delivered: + # The network is failing. Retrying other batches now would only burn + # their attempt budget against the same broken connection. + return sent + + # 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: + break + count, delivered = _drain(parked) + sent += count + if not delivered: + break + return sent + + +def _drain(claim: Path | None) -> tuple[int, bool]: + """Post one claimed batch file, recording progress after every batch. + + Returns (events sent, whether everything was delivered). + """ if claim is None: - return 0 + return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # 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. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass + return 0, True except OSError: - return 0 + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: @@ -497,11 +715,18 @@ def flush() -> int: 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 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -520,10 +745,13 @@ def flush() -> int: sent = 0 for start in range(0, len(events), BATCH_SIZE): + chunk = events[start : start + BATCH_SIZE] batch = [ { "event": event["event"], "distinct_id": distinct_id, + # Carried through from record() so a resend can be collapsed. + "uuid": event.get("uuid"), "timestamp": event.get("timestamp"), "properties": { # Fallback only: events recorded by a build before source @@ -535,16 +763,24 @@ def flush() -> int: **(event.get("properties") or {}), }, } - for event in events[start : start + BATCH_SIZE] + for event in chunk ] if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL): - return sent - sent += len(batch) - try: - claim.unlink() - except OSError: - pass - return sent + # Keep only what has not been delivered, and release the lease. + # Previously the whole file was kept and the retry re-posted every + # batch, including the ones that had already arrived. + _release_claim(claim, events[start:]) + 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. 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 def main() -> int: diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 6d7aeecae..a4375f002 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -100,6 +100,16 @@ BATCH_SIZE = 100 SEND_TIMEOUT = 5 CLAIM_STALE_SECONDS = 120 CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60 +# A batch is only discarded once it has genuinely been retried this many times. +MAX_CLAIM_ATTEMPTS = 3 +# Parked claims drained per run, after the live spool. Bounded so a long backlog +# cannot turn one flush into an unbounded send loop. +MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -399,38 +409,201 @@ def spawn_flush() -> bool: return False +def _claim_name(attempt: int = 0) -> str: + """Claim filename. The attempt count rides in the name so the 7-day expiry + only ever discards a batch that was actually retried and failed.""" + return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending" + + +def _claim_attempt(claim: Path) -> int: + """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--.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 + 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 + + +def _touch(path: Path) -> None: + """Refresh mtime so a claim's age measures time since it was claimed. + + ``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created + after a quiet minute inherited the spool's last-write time and looked + abandoned the instant it was made. A second sender would then take it over + while the first was still posting, and both would deliver the batch. + """ + try: + os.utime(path, None) + except OSError: + pass + + def _claim_spool() -> Path | None: """Rename the spool aside so exactly one sender owns each batch.""" directory = memory_core.data_dir() - claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending" + claim = directory / _claim_name() spool = _spool_path() try: spool.replace(claim) + _touch(claim) return claim except OSError: pass + return _claim_parked(directory) + + +def _sweep_debris(directory: Path) -> None: + """Remove files nothing else will ever pick up again. + + *.partial is a temp file orphaned by a crash between write and rename. + *.corrupt is a batch quarantined for undecodable content. No glob in this + module matches either, so without this they accumulate on disk for the life + of the install. + + Quarantined batches are kept far longer than debris: they are the only + evidence left of events that could not be delivered, and someone diagnosing + a report of missing telemetry has to be able to find one. + """ now = time.time() - for orphan in sorted(directory.glob("telemetry-*.sending")): + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + for quarantined in directory.glob("telemetry-*.corrupt"): + try: + if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS: + quarantined.unlink() + except OSError: + continue + # The same reasoning covers *.tmp. _write_identity and _install_salt both + # create one and unlink it in a finally, which a SIGKILL skips, and no glob + # in this module matches the leftovers either. + for temporary in directory.glob("telemetry-*.tmp"): + try: + if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS: + temporary.unlink() + except OSError: + continue + + +def _claim_parked(directory: Path) -> Path | None: + """Take the oldest abandoned claim, if any lease has actually expired. + + Kept separate from the live spool so flush() can drain both in one run. + Previously parked batches were only reachable when no spool existed at all, + and because sessions keep recording there usually was one — so a batch + parked by a failed send waited until the 7-day expiry deleted it unsent, + even though its own presence is what started the sender. + """ + now = time.time() + for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime): try: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS: + if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. This check has to come + # first. Claiming a file bumps its attempt count and refreshes its + # mtime, so a sender that has just taken the final attempt looks + # exhausted to everyone else while it is actively draining. Judging + # exhaustion before liveness let a second sender unlink a batch out + # from under its owner, losing every event in it. + continue + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # 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: pass continue - if age < CLAIM_STALE_SECONDS: - continue + claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1) try: orphan.replace(claim) + _touch(claim) return claim except OSError: continue return None +def _safe_mtime(path: Path) -> float: + try: + return path.stat().st_mtime + except OSError: + return 0.0 + + +def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: + """Persist the unsent remainder, atomically, and refresh the lease. + + Called after every successful batch. Two jobs: a retry resumes where the + send stopped instead of re-posting from the top, and the rewrite doubles as + the lease heartbeat, so a slow sender does not have its claim stolen + mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS. + """ + if not remaining: + try: + claim.unlink() + except OSError: + pass + return True + temporary = claim.with_suffix(f".{os.getpid()}.partial") + try: + 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 + except OSError: + try: + temporary.unlink() + except OSError: + pass + return False + + +def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: + """Persist the remainder and drop the lease, because this sender has given up. + + Distinct from the per-batch heartbeat: heartbeating on the way out would + make an abandoned batch look actively owned for a further + CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the + threshold lets the next flush pick it up immediately, while the attempt + count in the filename still bounds how many times that can happen. + """ + if not _rewrite_claim(claim, remaining): + return + try: + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown + os.utime(claim, (released, released)) + except OSError: + pass + + def _resolve_email(key: str) -> str: """Trade the API key for the account email so events join other Mem0 surfaces.""" url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/" @@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]: def flush() -> int: - """Drain claimed spools to PostHog and return the number of events sent.""" + """Drain the live spool, then any parked claims, and return events sent.""" if not is_enabled(): return 0 - claim = _claim_spool() + sent, delivered = _drain(_claim_spool()) + if not delivered: + # The network is failing. Retrying other batches now would only burn + # their attempt budget against the same broken connection. + return sent + + # 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: + break + count, delivered = _drain(parked) + sent += count + if not delivered: + break + return sent + + +def _drain(claim: Path | None) -> tuple[int, bool]: + """Post one claimed batch file, recording progress after every batch. + + Returns (events sent, whether everything was delivered). + """ if claim is None: - return 0 + return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # 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. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass + return 0, True except OSError: - return 0 + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: @@ -497,11 +715,18 @@ def flush() -> int: 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 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -520,10 +745,13 @@ def flush() -> int: sent = 0 for start in range(0, len(events), BATCH_SIZE): + chunk = events[start : start + BATCH_SIZE] batch = [ { "event": event["event"], "distinct_id": distinct_id, + # Carried through from record() so a resend can be collapsed. + "uuid": event.get("uuid"), "timestamp": event.get("timestamp"), "properties": { # Fallback only: events recorded by a build before source @@ -535,16 +763,24 @@ def flush() -> int: **(event.get("properties") or {}), }, } - for event in events[start : start + BATCH_SIZE] + for event in chunk ] if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL): - return sent - sent += len(batch) - try: - claim.unlink() - except OSError: - pass - return sent + # Keep only what has not been delivered, and release the lease. + # Previously the whole file was kept and the retry re-posted every + # batch, including the ones that had already arrived. + _release_claim(claim, events[start:]) + 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. 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 def main() -> int: