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():