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
This commit is contained in:
Saket Aryan
2026-09-16 22:14:42 +05:30
parent 0377b9a85e
commit 40287f6f04
8 changed files with 244 additions and 29 deletions
@@ -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:
@@ -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-<pid>-<hex>.sending, and hex can start with 'a'."""
assert telemetry._claim_attempt(Path("telemetry-999-deadbeef.sending")) == 0
@@ -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:
@@ -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:
+22 -4
View File
@@ -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:
+22 -4
View File
@@ -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:
+22 -4
View File
@@ -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:
@@ -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: