From 87af5e51396b3c4d05a1120bf2500172740ce11f Mon Sep 17 00:00:00 2001 From: Yiorgis Gozadinos Date: Wed, 27 May 2026 11:46:37 +0300 Subject: [PATCH] Boot-reap stale claims at WorkerPool start --- .../haiku/rag/ingester/workers/pool.py | 7 ++++++ tests/ingester/test_workers.py | 25 +++++++++++++++++++ 2 files changed, 32 insertions(+) diff --git a/haiku_rag_slim/haiku/rag/ingester/workers/pool.py b/haiku_rag_slim/haiku/rag/ingester/workers/pool.py index 70c442b6..14970076 100644 --- a/haiku_rag_slim/haiku/rag/ingester/workers/pool.py +++ b/haiku_rag_slim/haiku/rag/ingester/workers/pool.py @@ -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()) diff --git a/tests/ingester/test_workers.py b/tests/ingester/test_workers.py index e3fe752e..46994394 100644 --- a/tests/ingester/test_workers.py +++ b/tests/ingester/test_workers.py @@ -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