From d8c99fb405d59adc28225ce7c603d0b13fe29533 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Wed, 16 Sep 2026 20:35:13 +0530 Subject: [PATCH] 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 --- .../agent-plugin-core/python/telemetry.py | 11 +++- .../tests/test_spool_delivery.py | 61 +++++++++++++++++++ .../antigravity-plugin/core/telemetry.py | 11 +++- .../claude-code-plugin/core/telemetry.py | 11 +++- integrations/codex-plugin/core/telemetry.py | 11 +++- integrations/cursor-plugin/core/telemetry.py | 11 +++- integrations/kimi-plugin/core/telemetry.py | 11 +++- .../mem0-agent-plugin/core/telemetry.py | 11 +++- 8 files changed, 117 insertions(+), 21 deletions(-) diff --git a/integrations/agent-plugin-core/python/telemetry.py b/integrations/agent-plugin-core/python/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/agent-plugin-core/python/telemetry.py +++ b/integrations/agent-plugin-core/python/telemetry.py @@ -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) diff --git a/integrations/agent-plugin-core/tests/test_spool_delivery.py b/integrations/agent-plugin-core/tests/test_spool_delivery.py index a51027904..a7646cc31 100644 --- a/integrations/agent-plugin-core/tests/test_spool_delivery.py +++ b/integrations/agent-plugin-core/tests/test_spool_delivery.py @@ -104,6 +104,67 @@ def test_a_fresh_claim_is_not_immediately_stealable(telemetry): assert telemetry._claim_parked(first.parent) is None +def test_a_live_final_attempt_is_not_deleted_by_another_sender(telemetry): + """Review finding: exhaustion was judged before liveness, so owners lost batches. + + Claiming a parked file bumps its attempt count and refreshes its mtime. Once + the count reaches the budget, the owner draining it looked exhausted to every + other sender, which unlinked the file out from under it. Everything in that + batch was gone, which is precisely the loss this PR exists to stop. + """ + telemetry.record("search", reason="owned-by-the-first-sender") + 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 + + # Walk it to the final attempt, ageing it each round so it can be re-claimed. + # _claim_spool hands back a0 and _release_claim keeps the name, so it takes + # one full round per attempt to reach the budget. + for _ in range(telemetry.MAX_CLAIM_ATTEMPTS): + # Carry the marker through each rewrite so the final assertion proves the + # events survived, not merely that some file with the right name did. + telemetry._release_claim(claim, [{"event": "code.search", "uuid": "owned-by-the-first-sender"}]) + parked = sorted(claim.parent.glob("telemetry-*.sending")) + assert parked, "the batch was dropped while still inside its budget" + os.utime(parked[0], (stale, stale)) + claim = telemetry._claim_parked(claim.parent) + assert claim is not None + + assert telemetry._claim_attempt(claim) >= telemetry.MAX_CLAIM_ATTEMPTS + assert claim.exists() + + # The owner is draining it right now: fresh mtime, live lease. + second_sender = telemetry._claim_parked(claim.parent) + + assert second_sender is None, "a second sender took a batch under a live lease" + assert claim.exists(), "a second sender deleted a batch its owner was draining" + assert "owned-by-the-first-sender" in claim.read_text(encoding="utf-8") + + +def test_an_exhausted_batch_is_still_discarded_once_its_lease_lapses(telemetry): + """The liveness check must defer the cleanup, not cancel it. + + Guards the obvious over-correction: skipping live claims is only safe if an + abandoned one at the same attempt count is still reaped on a later run. + """ + 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 + exhausted = claim.parent / telemetry._claim_name(telemetry.MAX_CLAIM_ATTEMPTS) + claim.replace(exhausted) + os.utime(exhausted, (stale, stale)) + + assert telemetry._claim_parked(exhausted.parent) is None + assert not exhausted.exists(), "an abandoned exhausted batch was left behind forever" + + def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry): """Defect 6: parked claims were only reachable when no spool existed. diff --git a/integrations/antigravity-plugin/core/telemetry.py b/integrations/antigravity-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/antigravity-plugin/core/telemetry.py +++ b/integrations/antigravity-plugin/core/telemetry.py @@ -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) diff --git a/integrations/claude-code-plugin/core/telemetry.py b/integrations/claude-code-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/claude-code-plugin/core/telemetry.py +++ b/integrations/claude-code-plugin/core/telemetry.py @@ -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) diff --git a/integrations/codex-plugin/core/telemetry.py b/integrations/codex-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/codex-plugin/core/telemetry.py +++ b/integrations/codex-plugin/core/telemetry.py @@ -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) diff --git a/integrations/cursor-plugin/core/telemetry.py b/integrations/cursor-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/cursor-plugin/core/telemetry.py +++ b/integrations/cursor-plugin/core/telemetry.py @@ -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) diff --git a/integrations/kimi-plugin/core/telemetry.py b/integrations/kimi-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/kimi-plugin/core/telemetry.py +++ b/integrations/kimi-plugin/core/telemetry.py @@ -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) diff --git a/integrations/mem0-agent-plugin/core/telemetry.py b/integrations/mem0-agent-plugin/core/telemetry.py index 792a7f2eb..48d6660aa 100644 --- a/integrations/mem0-agent-plugin/core/telemetry.py +++ b/integrations/mem0-agent-plugin/core/telemetry.py @@ -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)