From 40287f6f04d969eb7ac540485e8b5cf630bcc7a6 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Wed, 16 Sep 2026 22:14:42 +0530 Subject: [PATCH] fix(plugins): a batch that could not be read is not a delivered batch Two review findings from @karthik-indla on this PR. _drain returned (0, True) on any read failure, so a claim nothing was posted from counted as fully delivered. flush() then carried on to the next claim as though this one had arrived, and the single signal that says the run went badly never fired. The two cases are now separated: undecodable content is still quarantined and reported delivered, because there is nothing left to send and the rest of the run should continue, while an OSError leaves the file exactly where it is and reports undelivered. Quarantining there would discard events over a transient filesystem error, and nothing ever re-globs .corrupt. Retries had no time backoff. _release_claim backdated straight to immediately-reclaimable, so two senders meeting one momentary failure could walk a batch from attempt 0 to the limit within seconds and discard it, when a retry a minute later would have delivered. Releases now carry a cooldown that grows with the attempts already spent, clamped so the mtime never lands in the future and reads as a live lease. Four tests: an unreadable batch is neither delivered nor quarantined, undecodable content still is quarantined so one torn file cannot block every later claim, and attempts cannot be burned without waiting. The expiry test now ages the file between flushes, which is the wall time a real retry waits. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .../agent-plugin-core/python/telemetry.py | 26 +++++- .../tests/test_spool_delivery.py | 91 ++++++++++++++++++- .../antigravity-plugin/core/telemetry.py | 26 +++++- .../claude-code-plugin/core/telemetry.py | 26 +++++- integrations/codex-plugin/core/telemetry.py | 26 +++++- integrations/cursor-plugin/core/telemetry.py | 26 +++++- integrations/kimi-plugin/core/telemetry.py | 26 +++++- .../mem0-agent-plugin/core/telemetry.py | 26 +++++- 8 files changed, 244 insertions(+), 29 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index a8f018f8e..a5107d430 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -105,6 +105,11 @@ MAX_CLAIM_ATTEMPTS = 3 # Parked claims drained per run, after the live spool. Bounded so a long backlog # cannot turn one flush into an unbounded send loop. MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -556,7 +561,11 @@ def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: if not _rewrite_claim(claim, remaining): return try: - released = time.time() - CLAIM_STALE_SECONDS - 1 + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown os.utime(claim, (released, released)) except OSError: pass @@ -642,11 +651,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except (OSError, ValueError): - # ValueError covers UnicodeDecodeError from a torn write. Quarantine - # rather than retry: flush() runs from a bare `finally:` in + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # quarantine rather than retry. flush() runs from a bare `finally:` in # flush_worker, so raising here also skips the handoff cleanup, and an # undecodable file would otherwise be re-read on every flush forever. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. try: claim.replace(claim.with_suffix(".corrupt")) except OSError: @@ -655,6 +666,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: except OSError: pass return 0, True + except OSError: + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: diff --git a/integrations/agent-plugin-core/tests/test_spool_delivery.py b/integrations/agent-plugin-core/tests/test_spool_delivery.py index a7646cc31..ac0bf6b22 100644 --- a/integrations/agent-plugin-core/tests/test_spool_delivery.py +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -266,13 +266,102 @@ def test_an_undeliverable_batch_is_eventually_given_up_on(telemetry): telemetry.record("doomed") telemetry._post = lambda payload, url: False + directory = telemetry.memory_core.data_dir() for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3): telemetry.flush() + # Attempts now carry a cooldown, so a released claim is not instantly + # reclaimable. Age it to stand in for the wall time a real retry waits; + # without this the loop spins inside one cooldown and proves nothing. + for parked in directory.glob("telemetry-*.sending"): + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(parked, (stale, stale)) - leftover = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending")) + leftover = list(directory.glob("telemetry-*.sending")) assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}" +def test_a_batch_that_cannot_be_read_is_not_counted_as_delivered(telemetry): + """Review finding: a read failure reported 'everything delivered'. + + Nothing was posted, so calling it delivered lets flush() carry on to other + claims as though this batch had arrived, and hides the failure from the one + signal that says the run went badly. It also must not quarantine: a briefly + unreadable file is retryable, and moving it to .corrupt discards the events + over a transient filesystem error, because nothing ever re-globs .corrupt. + """ + telemetry.record("search") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + claim = telemetry._claim_spool() + assert claim is not None + + original = Path.read_text + + def unreadable(self, *args, **kwargs): + if self == claim: + raise OSError(5, "I/O error") + return original(self, *args, **kwargs) + + Path.read_text = unreadable + try: + sent, delivered = telemetry._drain(claim) + finally: + Path.read_text = original + + assert sent == 0 + assert delivered is False, "an unread batch was reported as delivered" + assert claim.exists(), "a transient read error discarded the batch" + assert not list(claim.parent.glob("*.corrupt")), "quarantined over a transient error" + + +def test_undecodable_content_is_still_quarantined_and_the_run_continues(telemetry): + """The other half: genuinely unrecoverable content must not block the run. + + Guards the over-correction. If every read problem returned undelivered, one + torn file would stop every later claim on every flush, forever. + """ + telemetry.record("search") + spool = telemetry._spool_path() + stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60) + os.utime(spool, (stale, stale)) + claim = telemetry._claim_spool() + assert claim is not None + claim.write_bytes(b"\xff\xfe torn \x00 write") + + sent, delivered = telemetry._drain(claim) + + assert (sent, delivered) == (0, True) + assert not claim.exists() + assert list(claim.parent.glob("*.corrupt")), "unrecoverable content was not quarantined" + + +def test_retries_are_spread_over_real_time_not_burned_at_once(telemetry): + """Review finding: releasing straight to reclaimable spent the budget instantly. + + Two senders hitting one momentary failure could walk a batch from attempt 0 + to the limit within seconds and discard it, when a retry a minute later would + have delivered. Each release now has to age past a cooldown that grows with + the attempts already spent. + """ + telemetry.record("doomed") + telemetry._post = lambda payload, url: False + + directory = telemetry.memory_core.data_dir() + telemetry.flush() + + parked = list(directory.glob("telemetry-*.sending")) + assert parked, "the batch was discarded on its first failure" + assert telemetry._claim_attempt(parked[0]) == 0 + + # Second sender, immediately: the cooldown has not elapsed, so it must not + # be able to spend another attempt. + telemetry.flush() + still = list(directory.glob("telemetry-*.sending")) + assert len(still) == 1 + assert telemetry._claim_attempt(still[0]) <= 1, "burned attempts without waiting" + + def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry): """The old shape is telemetry--.sending, and hex can start with 'a'.""" assert telemetry._claim_attempt(Path("telemetry-999-deadbeef.sending")) == 0 diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index a8f018f8e..a5107d430 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -105,6 +105,11 @@ MAX_CLAIM_ATTEMPTS = 3 # Parked claims drained per run, after the live spool. Bounded so a long backlog # cannot turn one flush into an unbounded send loop. MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -556,7 +561,11 @@ def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: if not _rewrite_claim(claim, remaining): return try: - released = time.time() - CLAIM_STALE_SECONDS - 1 + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown os.utime(claim, (released, released)) except OSError: pass @@ -642,11 +651,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except (OSError, ValueError): - # ValueError covers UnicodeDecodeError from a torn write. Quarantine - # rather than retry: flush() runs from a bare `finally:` in + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # quarantine rather than retry. flush() runs from a bare `finally:` in # flush_worker, so raising here also skips the handoff cleanup, and an # undecodable file would otherwise be re-read on every flush forever. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. try: claim.replace(claim.with_suffix(".corrupt")) except OSError: @@ -655,6 +666,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: except OSError: pass return 0, True + except OSError: + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index a8f018f8e..a5107d430 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -105,6 +105,11 @@ MAX_CLAIM_ATTEMPTS = 3 # Parked claims drained per run, after the live spool. Bounded so a long backlog # cannot turn one flush into an unbounded send loop. MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -556,7 +561,11 @@ def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: if not _rewrite_claim(claim, remaining): return try: - released = time.time() - CLAIM_STALE_SECONDS - 1 + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown os.utime(claim, (released, released)) except OSError: pass @@ -642,11 +651,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except (OSError, ValueError): - # ValueError covers UnicodeDecodeError from a torn write. Quarantine - # rather than retry: flush() runs from a bare `finally:` in + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # quarantine rather than retry. flush() runs from a bare `finally:` in # flush_worker, so raising here also skips the handoff cleanup, and an # undecodable file would otherwise be re-read on every flush forever. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. try: claim.replace(claim.with_suffix(".corrupt")) except OSError: @@ -655,6 +666,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: except OSError: pass return 0, True + except OSError: + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index a8f018f8e..a5107d430 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -105,6 +105,11 @@ MAX_CLAIM_ATTEMPTS = 3 # Parked claims drained per run, after the live spool. Bounded so a long backlog # cannot turn one flush into an unbounded send loop. MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -556,7 +561,11 @@ def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: if not _rewrite_claim(claim, remaining): return try: - released = time.time() - CLAIM_STALE_SECONDS - 1 + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown os.utime(claim, (released, released)) except OSError: pass @@ -642,11 +651,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except (OSError, ValueError): - # ValueError covers UnicodeDecodeError from a torn write. Quarantine - # rather than retry: flush() runs from a bare `finally:` in + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # quarantine rather than retry. flush() runs from a bare `finally:` in # flush_worker, so raising here also skips the handoff cleanup, and an # undecodable file would otherwise be re-read on every flush forever. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. try: claim.replace(claim.with_suffix(".corrupt")) except OSError: @@ -655,6 +666,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: except OSError: pass return 0, True + except OSError: + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index a8f018f8e..a5107d430 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -105,6 +105,11 @@ MAX_CLAIM_ATTEMPTS = 3 # Parked claims drained per run, after the live spool. Bounded so a long backlog # cannot turn one flush into an unbounded send loop. MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -556,7 +561,11 @@ def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: if not _rewrite_claim(claim, remaining): return try: - released = time.time() - CLAIM_STALE_SECONDS - 1 + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown os.utime(claim, (released, released)) except OSError: pass @@ -642,11 +651,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except (OSError, ValueError): - # ValueError covers UnicodeDecodeError from a torn write. Quarantine - # rather than retry: flush() runs from a bare `finally:` in + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # quarantine rather than retry. flush() runs from a bare `finally:` in # flush_worker, so raising here also skips the handoff cleanup, and an # undecodable file would otherwise be re-read on every flush forever. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. try: claim.replace(claim.with_suffix(".corrupt")) except OSError: @@ -655,6 +666,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: except OSError: pass return 0, True + except OSError: + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index a8f018f8e..a5107d430 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -105,6 +105,11 @@ MAX_CLAIM_ATTEMPTS = 3 # Parked claims drained per run, after the live spool. Bounded so a long backlog # cannot turn one flush into an unbounded send loop. MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -556,7 +561,11 @@ def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: if not _rewrite_claim(claim, remaining): return try: - released = time.time() - CLAIM_STALE_SECONDS - 1 + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown os.utime(claim, (released, released)) except OSError: pass @@ -642,11 +651,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except (OSError, ValueError): - # ValueError covers UnicodeDecodeError from a torn write. Quarantine - # rather than retry: flush() runs from a bare `finally:` in + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # quarantine rather than retry. flush() runs from a bare `finally:` in # flush_worker, so raising here also skips the handoff cleanup, and an # undecodable file would otherwise be re-read on every flush forever. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. try: claim.replace(claim.with_suffix(".corrupt")) except OSError: @@ -655,6 +666,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: except OSError: pass return 0, True + except OSError: + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index a8f018f8e..a5107d430 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -105,6 +105,11 @@ MAX_CLAIM_ATTEMPTS = 3 # Parked claims drained per run, after the live spool. Bounded so a long backlog # cannot turn one flush into an unbounded send loop. MAX_PARKED_PER_RUN = 3 +# Added to the wait before a released claim becomes reclaimable, per attempt +# already spent. Releasing straight to "reclaimable now" let two senders burn the +# whole budget within seconds of one another on a single momentary failure, and +# discard a batch a retry a minute later would have delivered. +RETRY_COOLDOWN_SECONDS = 60 def is_enabled() -> bool: @@ -556,7 +561,11 @@ def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None: if not _rewrite_claim(claim, remaining): return try: - released = time.time() - CLAIM_STALE_SECONDS - 1 + # Backdate past the stale threshold so the next flush can pick it up, + # minus a cooldown that grows with the attempts already spent. Clamped so + # the mtime never lands in the future, which would read as a live lease. + cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS) + released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown os.utime(claim, (released, released)) except OSError: pass @@ -642,11 +651,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: return 0, True try: lines = claim.read_text(encoding="utf-8").splitlines() - except (OSError, ValueError): - # ValueError covers UnicodeDecodeError from a torn write. Quarantine - # rather than retry: flush() runs from a bare `finally:` in + except ValueError: + # UnicodeDecodeError from a torn write: the content is unrecoverable, so + # quarantine rather than retry. flush() runs from a bare `finally:` in # flush_worker, so raising here also skips the handoff cleanup, and an # undecodable file would otherwise be re-read on every flush forever. + # Reported as delivered because there is nothing left to deliver and the + # rest of the run should continue. try: claim.replace(claim.with_suffix(".corrupt")) except OSError: @@ -655,6 +666,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]: except OSError: pass return 0, True + except OSError: + # Could not read it, which is not the same as having nothing to send. + # The file is left exactly where it is: a vanished or briefly unreadable + # claim is retryable, and quarantining it here would discard events over + # a transient filesystem error. Reported as undelivered so the run stops + # instead of counting a batch nothing was posted from as delivered. + return 0, False events = [] for line in lines: try: