Initial parser admin implementation
This commit is contained in:
@@ -0,0 +1,98 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from .constants import JOB_STATUS_IN_PROGRESS, JOB_STATUS_PENDING, JOB_STATUS_RETRY
|
||||
|
||||
|
||||
async def is_worker_enabled(pool, name: str) -> bool:
|
||||
value = await pool.fetchval("SELECT enabled FROM worker_controls WHERE name=$1", name)
|
||||
return bool(value)
|
||||
|
||||
|
||||
async def claim_job(pool, job_type: str, worker_id: str) -> dict | None:
|
||||
row = await pool.fetchrow(
|
||||
"""
|
||||
WITH cte AS (
|
||||
SELECT id
|
||||
FROM jobs
|
||||
WHERE type=$1
|
||||
AND status IN ($2, $3)
|
||||
AND next_run_at <= NOW()
|
||||
ORDER BY next_run_at ASC, id ASC
|
||||
FOR UPDATE SKIP LOCKED
|
||||
LIMIT 1
|
||||
)
|
||||
UPDATE jobs j
|
||||
SET status=$4,
|
||||
locked_by=$5,
|
||||
locked_at=NOW(),
|
||||
updated_at=NOW()
|
||||
FROM cte
|
||||
WHERE j.id=cte.id
|
||||
RETURNING j.*
|
||||
""",
|
||||
job_type,
|
||||
JOB_STATUS_PENDING,
|
||||
JOB_STATUS_RETRY,
|
||||
JOB_STATUS_IN_PROGRESS,
|
||||
worker_id,
|
||||
)
|
||||
return dict(row) if row else None
|
||||
|
||||
|
||||
async def ack_done(pool, job_id: int) -> None:
|
||||
await pool.execute(
|
||||
"""
|
||||
UPDATE jobs
|
||||
SET status='done',
|
||||
locked_by=NULL,
|
||||
locked_at=NULL,
|
||||
last_error=NULL,
|
||||
updated_at=NOW()
|
||||
WHERE id=$1
|
||||
""",
|
||||
job_id,
|
||||
)
|
||||
|
||||
|
||||
async def ack_retry(pool, job: dict, error: str, delay_sec: int = 60) -> None:
|
||||
attempt = int(job.get("attempts") or 0) + 1
|
||||
max_attempts = int(job.get("max_attempts") or 5)
|
||||
status = "dead" if attempt >= max_attempts else "retry"
|
||||
await pool.execute(
|
||||
"""
|
||||
UPDATE jobs
|
||||
SET status=$2,
|
||||
attempts=$3,
|
||||
next_run_at=CASE WHEN $2='retry' THEN NOW() + ($4 * INTERVAL '1 second') ELSE next_run_at END,
|
||||
locked_by=NULL,
|
||||
locked_at=NULL,
|
||||
last_error=$5,
|
||||
updated_at=NOW()
|
||||
WHERE id=$1
|
||||
""",
|
||||
int(job["id"]),
|
||||
status,
|
||||
attempt,
|
||||
int(delay_sec),
|
||||
error[:1000],
|
||||
)
|
||||
|
||||
|
||||
async def recover_stale_jobs(pool, job_type: str, stale_minutes: int = 20) -> int:
|
||||
result = await pool.execute(
|
||||
"""
|
||||
UPDATE jobs
|
||||
SET status='retry',
|
||||
locked_by=NULL,
|
||||
locked_at=NULL,
|
||||
next_run_at=NOW(),
|
||||
last_error=COALESCE(last_error, 'recovered stale lock'),
|
||||
updated_at=NOW()
|
||||
WHERE type=$1
|
||||
AND status='in_progress'
|
||||
AND locked_at < NOW() - ($2 * INTERVAL '1 minute')
|
||||
""",
|
||||
job_type,
|
||||
stale_minutes,
|
||||
)
|
||||
return int(str(result).split()[-1])
|
||||
Reference in New Issue
Block a user