From 0fefd044baf6f87de7ebb12a2ef8eb2c80a8f046 Mon Sep 17 00:00:00 2001 From: avi Date: Mon, 14 Sep 2026 10:37:09 -0500 Subject: [PATCH] Engine: inline queue reclaims running-job corpses at startup (single-process engine has no live workers) --- backend/shonar/services/inline_queue.py | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/backend/shonar/services/inline_queue.py b/backend/shonar/services/inline_queue.py index 39f5c48..501711c 100644 --- a/backend/shonar/services/inline_queue.py +++ b/backend/shonar/services/inline_queue.py @@ -43,6 +43,22 @@ async def start() -> None: """Start the consumer and re-run anything the DB says is pending.""" global _consumer, _retention _get_queue() + # Inline engine = single process: any row still marked `running` at + # startup is a corpse from the previous process (the worker died with + # it). sweep_stale's 2h live-worker grace — correct for multi-worker + # arq deployments — would starve these jobs, so reclaim them first. + from sqlalchemy import update + + from shonar.db.models import JobStatus, ProcessingJob + from shonar.db.session import session_factory + + async with session_factory()() as s: + await s.execute( + update(ProcessingJob) + .where(ProcessingJob.status == JobStatus.running) + .values(status=JobStatus.queued, started_at=None) + ) + await s.commit() if _consumer is None or _consumer.done(): _consumer = asyncio.create_task(_consume(), name="shonar-inline-queue") if _retention is None or _retention.done():