Boot-reap stale claims at WorkerPool start

This commit is contained in:
Yiorgis Gozadinos 2026-05-27 11:46:37 +03:00
parent 1b36452629
commit 87af5e5139
No known key found for this signature in database
2 changed files with 32 additions and 0 deletions

View file

@ -83,6 +83,13 @@ class WorkerPool:
if self._workers:
raise RuntimeError("WorkerPool already started")
self._stop.clear()
# Any rows in `claimed` at start time are owned by workers from a
# previous process that didn't get to release them (SIGKILL, OOM,
# host reboot). Reset them so fresh workers can claim immediately
# instead of waiting on the reaper's claim_timeout_s.
reset = await self._jobs.reap_stale(claim_timeout_seconds=0)
if reset:
logger.info("Boot-reaped %d stale claim(s) from previous process", reset)
for i in range(self._worker_count):
self._workers.append(asyncio.create_task(self._worker_loop(f"worker-{i}")))
self._reaper = asyncio.create_task(self._reaper_loop())

View file

@ -566,6 +566,31 @@ async def test_breaker_ignores_permanent_errors(client, jobs, sync):
# --- reaper ---
@pytest.mark.asyncio
async def test_boot_reap_resets_pre_existing_claims(client, jobs, sync):
"""A SIGKILL'd previous process leaves rows in `claimed` state. The new
WorkerPool.start() must reset them immediately so fresh workers can
claim them, instead of waiting on the periodic reaper's claim_timeout_s
window (default 1800s)."""
await jobs.enqueue("src", "u", JobOp.UPSERT)
pre_claimed = await jobs.claim_next("ghost-worker")
assert pre_claimed is not None
assert pre_claimed.status is JobStatus.CLAIMED
pool = _pool(client, jobs, sync, worker_count=0)
await pool.start()
try:
refreshed = await jobs.get_job(pre_claimed.id)
assert refreshed is not None
assert refreshed.status is JobStatus.QUEUED
assert refreshed.claimed_by is None
# attempts was incremented by claim_next; reap decrements it so the
# next claim doesn't see a consumed retry it never actually used.
assert refreshed.attempts == 0
finally:
await pool.stop()
@pytest.mark.asyncio
async def test_reaper_resets_stale_claims(client, jobs, sync, conn):
from datetime import UTC, datetime, timedelta