Engine: inline queue reclaims running-job corpses at startup (single-process engine has no live workers)
This commit is contained in:
parent
7f447188b5
commit
0fefd044ba
1 changed files with 16 additions and 0 deletions
|
|
@ -43,6 +43,22 @@ async def start() -> None:
|
||||||
"""Start the consumer and re-run anything the DB says is pending."""
|
"""Start the consumer and re-run anything the DB says is pending."""
|
||||||
global _consumer, _retention
|
global _consumer, _retention
|
||||||
_get_queue()
|
_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():
|
if _consumer is None or _consumer.done():
|
||||||
_consumer = asyncio.create_task(_consume(), name="shonar-inline-queue")
|
_consumer = asyncio.create_task(_consume(), name="shonar-inline-queue")
|
||||||
if _retention is None or _retention.done():
|
if _retention is None or _retention.done():
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue