diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 157db297d..55ee7c907 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -99,6 +99,11 @@ 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 def is_enabled() -> bool: @@ -301,38 +306,139 @@ 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.""" + stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name + tail = stem.rsplit("-", 1)[-1] + 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 _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")): + 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_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: try: orphan.unlink() except OSError: pass continue if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. 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: + temporary.write_text( + "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), + encoding="utf-8", + ) + 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: + released = time.time() - CLAIM_STALE_SECONDS - 1 + 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/" @@ -380,16 +486,40 @@ 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() + 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 OSError: - return 0 + return 0, True events = [] for line in lines: try: @@ -403,7 +533,7 @@ def flush() -> int: claim.unlink() except OSError: pass - return 0 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -422,10 +552,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 @@ -437,16 +570,19 @@ 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. + _rewrite_claim(claim, events[start + len(chunk) :]) + 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..754307454 --- /dev/null +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -0,0 +1,164 @@ +"""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): + monkeypatch.setenv("MEM0_CODE_DATA_DIR", str(tmp_path / "data")) + monkeypatch.syspath_prepend(str(HOST_CORE)) + for name in ("telemetry", "memory_core", "_harness_id"): + sys.modules.pop(name, None) + module = importlib.import_module("telemetry") + monkeypatch.setattr(module, "resolve_distinct_id", lambda: ("tester@example.com", "")) + yield module + for name in ("telemetry", "memory_core", "_harness_id"): + sys.modules.pop(name, None) + + +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_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_an_untried_batch_is_not_expired_by_age_alone(telemetry): + """Expiry should discard what failed, not what never got a turn.""" + telemetry.record("parked") + telemetry._post = lambda payload, url: False + telemetry.flush() + + parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) + assert len(parked) == 1 + ancient = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60) + os.utime(parked[0], (ancient, ancient)) + + 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_stays_well_inside_the_lease(telemetry): + """The claim rewrite doubles as the lease heartbeat. + + _post makes a single attempt with SEND_TIMEOUT and no retry, so a heartbeat + lands at least that often. If a retry loop is ever added to _post, this is + the assertion that catches a sender losing its claim mid-flight. + """ + assert telemetry.SEND_TIMEOUT * 4 < telemetry.CLAIM_STALE_SECONDS diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 157db297d..55ee7c907 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -99,6 +99,11 @@ 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 def is_enabled() -> bool: @@ -301,38 +306,139 @@ 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.""" + stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name + tail = stem.rsplit("-", 1)[-1] + 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 _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")): + 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_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: try: orphan.unlink() except OSError: pass continue if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. 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: + temporary.write_text( + "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), + encoding="utf-8", + ) + 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: + released = time.time() - CLAIM_STALE_SECONDS - 1 + 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/" @@ -380,16 +486,40 @@ 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() + 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 OSError: - return 0 + return 0, True events = [] for line in lines: try: @@ -403,7 +533,7 @@ def flush() -> int: claim.unlink() except OSError: pass - return 0 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -422,10 +552,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 @@ -437,16 +570,19 @@ 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. + _rewrite_claim(claim, events[start + len(chunk) :]) + 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 157db297d..55ee7c907 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -99,6 +99,11 @@ 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 def is_enabled() -> bool: @@ -301,38 +306,139 @@ 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.""" + stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name + tail = stem.rsplit("-", 1)[-1] + 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 _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")): + 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_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: try: orphan.unlink() except OSError: pass continue if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. 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: + temporary.write_text( + "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), + encoding="utf-8", + ) + 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: + released = time.time() - CLAIM_STALE_SECONDS - 1 + 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/" @@ -380,16 +486,40 @@ 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() + 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 OSError: - return 0 + return 0, True events = [] for line in lines: try: @@ -403,7 +533,7 @@ def flush() -> int: claim.unlink() except OSError: pass - return 0 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -422,10 +552,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 @@ -437,16 +570,19 @@ 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. + _rewrite_claim(claim, events[start + len(chunk) :]) + 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 0ae280b19..580bc1100 100644 --- a/integrations/claude-code-plugin/tests/test_telemetry.py +++ b/integrations/claude-code-plugin/tests/test_telemetry.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +import os import sys from pathlib import Path from unittest.mock import patch @@ -198,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 @@ -210,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 157db297d..55ee7c907 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -99,6 +99,11 @@ 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 def is_enabled() -> bool: @@ -301,38 +306,139 @@ 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.""" + stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name + tail = stem.rsplit("-", 1)[-1] + 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 _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")): + 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_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: try: orphan.unlink() except OSError: pass continue if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. 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: + temporary.write_text( + "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), + encoding="utf-8", + ) + 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: + released = time.time() - CLAIM_STALE_SECONDS - 1 + 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/" @@ -380,16 +486,40 @@ 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() + 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 OSError: - return 0 + return 0, True events = [] for line in lines: try: @@ -403,7 +533,7 @@ def flush() -> int: claim.unlink() except OSError: pass - return 0 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -422,10 +552,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 @@ -437,16 +570,19 @@ 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. + _rewrite_claim(claim, events[start + len(chunk) :]) + return sent, True def main() -> int: diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 157db297d..55ee7c907 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -99,6 +99,11 @@ 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 def is_enabled() -> bool: @@ -301,38 +306,139 @@ 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.""" + stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name + tail = stem.rsplit("-", 1)[-1] + 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 _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")): + 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_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: try: orphan.unlink() except OSError: pass continue if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. 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: + temporary.write_text( + "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), + encoding="utf-8", + ) + 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: + released = time.time() - CLAIM_STALE_SECONDS - 1 + 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/" @@ -380,16 +486,40 @@ 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() + 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 OSError: - return 0 + return 0, True events = [] for line in lines: try: @@ -403,7 +533,7 @@ def flush() -> int: claim.unlink() except OSError: pass - return 0 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -422,10 +552,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 @@ -437,16 +570,19 @@ 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. + _rewrite_claim(claim, events[start + len(chunk) :]) + return sent, True def main() -> int: diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 157db297d..55ee7c907 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -99,6 +99,11 @@ 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 def is_enabled() -> bool: @@ -301,38 +306,139 @@ 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.""" + stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name + tail = stem.rsplit("-", 1)[-1] + 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 _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")): + 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_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: try: orphan.unlink() except OSError: pass continue if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. 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: + temporary.write_text( + "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), + encoding="utf-8", + ) + 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: + released = time.time() - CLAIM_STALE_SECONDS - 1 + 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/" @@ -380,16 +486,40 @@ 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() + 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 OSError: - return 0 + return 0, True events = [] for line in lines: try: @@ -403,7 +533,7 @@ def flush() -> int: claim.unlink() except OSError: pass - return 0 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -422,10 +552,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 @@ -437,16 +570,19 @@ 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. + _rewrite_claim(claim, events[start + len(chunk) :]) + 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 157db297d..55ee7c907 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -99,6 +99,11 @@ 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 def is_enabled() -> bool: @@ -301,38 +306,139 @@ 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.""" + stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name + tail = stem.rsplit("-", 1)[-1] + 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 _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")): + 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_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: try: orphan.unlink() except OSError: pass continue if age < CLAIM_STALE_SECONDS: + # Someone else holds a live lease on it. 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: + temporary.write_text( + "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), + encoding="utf-8", + ) + 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: + released = time.time() - CLAIM_STALE_SECONDS - 1 + 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/" @@ -380,16 +486,40 @@ 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() + 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 OSError: - return 0 + return 0, True events = [] for line in lines: try: @@ -403,7 +533,7 @@ def flush() -> int: claim.unlink() except OSError: pass - return 0 + return 0, True distinct_id, aliased_anonymous_id = resolve_distinct_id() if aliased_anonymous_id: @@ -422,10 +552,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 @@ -437,16 +570,19 @@ 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. + _rewrite_claim(claim, events[start + len(chunk) :]) + return sent, True def main() -> int: