Files
Lyra/worker/lyra_worker/claim.py
T
Jonathan bff590a00f feat(worker): claim_next skips per-item paused jobs
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-15 12:01:38 +02:00

58 lines
2.0 KiB
Python

import psycopg
# States a job passes through while a worker is actively processing it. On startup this
# fresh worker owns none of them, so any job left here was orphaned by a crash/restart.
_IN_FLIGHT_STATES = ("matching", "matched", "downloading", "tagging")
def reclaim_stuck_jobs(conn: psycopg.Connection) -> int:
"""Reset jobs left mid-pipeline by a prior worker crash/restart back to 'requested' so
they get re-claimed. Safe to call once at startup: nothing is processing them, staging is
cleared beforehand, and the pipeline redoes intake→import idempotently (atomic import).
Returns how many were reclaimed. Terminal jobs (imported/needs_attention) are untouched."""
with conn.cursor() as cur:
cur.execute(
"""
UPDATE "Job"
SET state = 'requested', "currentStage" = 'intake', "claimedAt" = NULL,
error = NULL, "downloadProgress" = 0, "updatedAt" = now()
WHERE state IN ('matching', 'matched', 'downloading', 'tagging')
"""
)
n = cur.rowcount
conn.commit()
return n
def claim_next(conn: psycopg.Connection) -> str | None:
"""Atomically claim the oldest queued job. Returns its id or None."""
with conn.cursor() as cur:
cur.execute(
"""
SELECT id FROM "Job"
WHERE state = 'requested' AND NOT paused
ORDER BY "createdAt" ASC
FOR UPDATE SKIP LOCKED
LIMIT 1
"""
)
row = cur.fetchone()
if row is None:
conn.rollback()
return None
job_id = row[0]
cur.execute(
"""
UPDATE "Job"
SET state = 'matching',
"currentStage" = 'match',
"claimedAt" = now(),
attempts = attempts + 1,
"updatedAt" = now()
WHERE id = %s
""",
(job_id,),
)
conn.commit()
return job_id