fix(plugins): check the lease before judging a claim exhausted

Review finding from @kartik-mem0 on this PR, and the most serious one: it loses
events, which is what this PR exists to prevent.

_claim_parked judged exhaustion before liveness. Claiming a parked file bumps
its attempt count and refreshes its mtime, so the moment a sender takes the
final attempt the file looks exhausted to every other sender while its owner is
actively draining it. The second sender unlinked it, and everything in that
batch was gone.

The liveness check now runs first, so a batch under a live lease is skipped
whatever its attempt count. The cleanup is deferred, not cancelled: once the
lease lapses, the same exhausted file is reaped on a later run.

Two tests. The first walks a batch to the final attempt and asserts a second
sender neither takes it nor deletes it, and that the events are still in it. The
second asserts an abandoned exhausted batch is still discarded once its lease
lapses, which is the over-correction to guard against. Confirmed the first fails
against the previous ordering.

288 passed, 8 skipped.

Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb
This commit is contained in:
Saket Aryan
2026-09-16 20:35:13 +05:30
parent 26760b00b1
commit d8c99fb405
8 changed files with 117 additions and 21 deletions
@@ -461,6 +461,14 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
@@ -471,9 +479,6 @@ def _claim_parked(directory: Path) -> Path | None:
except OSError:
pass
continue
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it.
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)