diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/agent-plugin-core/tests/test_spool_delivery.py b/integrations/agent-plugin-core/tests/test_spool_delivery.py index 754307454..08266cc98 100644 --- a/integrations/agent-plugin-core/tests/test_spool_delivery.py +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -112,16 +112,19 @@ def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry): assert names == {"code.parked", "code.fresh"} -def test_an_untried_batch_is_not_expired_by_age_alone(telemetry): - """Expiry should discard what failed, not what never got a turn.""" +def test_a_batch_is_retried_until_the_budget_is_spent_not_discarded(telemetry): + """Expiry discards what failed repeatedly, not what merely sat for a while. + + The budget is the attempt count, because age cannot be one: every re-claim + touches the mtime and every release backdates it, so age never accumulates. + """ telemetry.record("parked") telemetry._post = lambda payload, url: False telemetry.flush() parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) assert len(parked) == 1 - ancient = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60) - os.utime(parked[0], (ancient, ancient)) + assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS sent: list[dict] = [] telemetry._post = lambda payload, url: sent.append(payload) or True @@ -154,11 +157,88 @@ def test_progress_is_recorded_after_every_batch(telemetry): assert json.loads(remaining[0])["properties"]["index"] == 200 -def test_the_heartbeat_stays_well_inside_the_lease(telemetry): +def test_the_heartbeat_actually_refreshes_the_lease(telemetry): """The claim rewrite doubles as the lease heartbeat. - _post makes a single attempt with SEND_TIMEOUT and no retry, so a heartbeat - lands at least that often. If a retry loop is ever added to _post, this is - the assertion that catches a sender losing its claim mid-flight. + Previously asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS`, which compares + two constants and executes none of the code under test. Drive the real + rewrite and watch the mtime move instead. """ - assert telemetry.SEND_TIMEOUT * 4 < telemetry.CLAIM_STALE_SECONDS + for index in range(150): + telemetry.record("search", index=index) + claim = telemetry._claim_spool() + assert claim is not None + + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(claim, (stale, stale)) + assert time.time() - claim.stat().st_mtime > telemetry.CLAIM_STALE_SECONDS + + telemetry._rewrite_claim(claim, [{"event": "code.x", "properties": {}}]) + assert time.time() - claim.stat().st_mtime < telemetry.CLAIM_STALE_SECONDS + + +def test_an_undeliverable_batch_is_eventually_given_up_on(telemetry): + """Expiry has to be reachable from a state the state machine can produce. + + It was not: every re-claim touched the mtime and every release backdated it + by a fixed amount, so age hovered near the stale threshold and the 7-day + expiry never fired. An undeliverable batch lived on disk forever, and + spawn_flush saw it and started a sender on every hook. + """ + telemetry.record("doomed") + telemetry._post = lambda payload, url: False + + for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3): + telemetry.flush() + + leftover = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) + assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}" + + +def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry): + """The old shape is telemetry--.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() diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 55ee7c907..fa26cc9a6 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -313,9 +313,18 @@ def _claim_name(attempt: int = 0) -> str: def _claim_attempt(claim: Path) -> int: - """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.""" + """Attempts recorded in a claim filename; 0 for the pre-attempt-count shape. + + Anchored on field position, not on a leading "a": the legacy shape is + ``telemetry--.sending`` and a hex id such as ``a1234567`` would + otherwise parse as attempt 1234567 and be discarded unsent on the first + flush after an upgrade. + """ stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name - tail = stem.rsplit("-", 1)[-1] + parts = stem.split("-") + if len(parts) != 4: + return 0 + tail = parts[3] if tail.startswith("a") and tail[1:].isdigit(): return int(tail[1:]) return 0 @@ -349,6 +358,21 @@ def _claim_spool() -> Path | None: return _claim_parked(directory) +def _sweep_debris(directory: Path) -> None: + """Remove temp files orphaned by a crash between write and rename. + + Neither glob in this module matches *.partial, so nothing else would ever + clean them up. + """ + now = time.time() + for debris in directory.glob("telemetry-*.partial"): + try: + if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS: + debris.unlink() + except OSError: + continue + + def _claim_parked(directory: Path) -> Path | None: """Take the oldest abandoned claim, if any lease has actually expired. @@ -364,7 +388,11 @@ def _claim_parked(directory: Path) -> Path | None: age = now - orphan.stat().st_mtime except OSError: continue - if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS: + # Attempts, not age. Every re-claim touches the mtime and every release + # backdates it by a fixed amount, so age is pinned near the stale + # threshold and never reaches the expiry. Age stays only as a backstop + # for files that never carried an attempt marker. + if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS: try: orphan.unlink() except OSError: @@ -406,10 +434,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool: return True temporary = claim.with_suffix(f".{os.getpid()}.partial") try: - temporary.write_text( - "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining), - encoding="utf-8", - ) + payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining) + # fsync before the rename: without it the rename can land while the + # bytes have not, and the claim comes back empty or truncated after a + # crash. _drain then reads zero events and unlinks it. + with open(temporary, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) temporary.replace(claim) _touch(claim) return True @@ -498,6 +530,7 @@ def flush() -> int: # Parked batches used to starve behind the live spool indefinitely. Bounded # per run so a long backlog cannot turn one flush into an unbounded loop. directory = memory_core.data_dir() + _sweep_debris(directory) for _ in range(MAX_PARKED_PER_RUN): parked = _claim_parked(directory) if parked is None: @@ -518,7 +551,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except OSError: + except (OSError, ValueError): + # ValueError covers UnicodeDecodeError from a torn write. Quarantine + # rather than retry: flush() runs from a bare `finally:` in + # flush_worker, so raising here also skips the handoff cleanup, and an + # undecodable file would otherwise be re-read on every flush forever. + try: + claim.replace(claim.with_suffix(".corrupt")) + except OSError: + try: + claim.unlink() + except OSError: + pass return 0, True events = [] for line in lines: @@ -529,8 +573,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]: if isinstance(value, dict) and value.get("event"): events.append(value) if not events: + # Only delete when the file really is empty. A non-empty file that + # parses to nothing is a torn write, and its contents are the unsent + # remainder — deleting it is the data loss this PR exists to prevent. try: - claim.unlink() + empty = claim.stat().st_size == 0 + except OSError: + empty = True + try: + claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink() except OSError: pass return 0, True @@ -580,8 +631,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return sent, False sent += len(chunk) # Record progress and refresh the lease after each successful batch, so - # a crash repeats at most one batch instead of the entire file. - _rewrite_claim(claim, events[start + len(chunk) :]) + # a crash repeats at most one batch instead of the entire file. If the + # rewrite fails the claim still holds delivered events, so stop rather + # than carry on as though progress were recorded — continuing is how the + # duplicate delivery this PR fixes would come back. + if not _rewrite_claim(claim, events[start + len(chunk) :]): + _release_claim(claim, events[start + len(chunk) :]) + return sent, False return sent, True