Files
tabdeal-job-queue/backend/jobs/services.py

394 lines
13 KiB
Python

from __future__ import annotations
import logging
from dataclasses import dataclass
from datetime import timedelta
from django.conf import settings
from django.db import IntegrityError, connection, transaction
from django.db.models import Count, F
from django.utils import timezone
from jobs.models import Job, JobEvent
logger = logging.getLogger(__name__)
class JobOwnershipLost(Exception):
pass
@dataclass(frozen=True)
class JobCreateResult:
job: Job
created: bool
def default_max_attempts() -> int:
return max(1, int(settings.JOB_WORKER_MAX_ATTEMPTS_DEFAULT))
def backoff_delay(attempt: int) -> timedelta:
exponent = max(0, attempt - 1)
seconds = min(
int(settings.JOB_WORKER_BACKOFF_BASE_SECONDS) * (2**exponent),
int(settings.JOB_WORKER_BACKOFF_MAX_SECONDS),
)
return timedelta(seconds=seconds)
def emit_event(
job: Job,
event_type: str,
*,
attempt: int | None = None,
worker_id: str | None = None,
message: str = "",
data: dict | None = None,
) -> JobEvent:
return JobEvent.objects.create(
job=job,
type=event_type,
attempt=job.attempts if attempt is None else attempt,
worker_id=worker_id,
message=message,
data=data or {},
)
@transaction.atomic
def create_job(
*,
job_type: str,
payload: dict | None = None,
priority: int = 50,
available_at=None,
max_attempts: int | None = None,
idempotency_key: str | None = None,
) -> JobCreateResult:
normalized_key = idempotency_key.strip() if idempotency_key else None
if normalized_key:
existing = Job.objects.filter(idempotency_key=normalized_key).first()
if existing is not None:
logger.info("Idempotent create returned existing job %s key=%s.", existing.id, normalized_key)
return JobCreateResult(existing, created=False)
try:
job = Job.objects.create(
type=job_type,
payload=payload or {},
priority=priority,
available_at=available_at or timezone.now(),
max_attempts=max_attempts or default_max_attempts(),
idempotency_key=normalized_key,
)
except IntegrityError:
if not normalized_key:
raise
existing = Job.objects.get(idempotency_key=normalized_key)
logger.info("Idempotent create raced and returned existing job %s key=%s.", existing.id, normalized_key)
return JobCreateResult(existing, created=False)
emit_event(
job,
JobEvent.Type.CREATED,
message="Job created",
data={
"type": job.type,
"priority": job.priority,
"available_at": job.available_at.isoformat(),
"idempotency_key": normalized_key,
},
)
logger.info("Created job %s type=%s priority=%s available_at=%s.", job.id, job.type, job.priority, job.available_at)
return JobCreateResult(job, created=True)
@transaction.atomic
def claim_next_job(*, worker_id: str, lease_seconds: int | None = None) -> Job | None:
now = timezone.now()
lease_seconds = lease_seconds or settings.JOB_WORKER_LEASE_SECONDS
queryset = Job.objects.filter(
status=Job.Status.QUEUED,
available_at__lte=now,
attempts__lt=F("max_attempts"),
).order_by("-priority", "available_at", "created_at", "id")
if connection.features.has_select_for_update_skip_locked:
queryset = queryset.select_for_update(skip_locked=True)
else:
queryset = queryset.select_for_update()
job = queryset.first()
if job is None:
logger.debug("No claimable queued job found for worker=%s at %s.", worker_id, now.isoformat())
return None
job.status = Job.Status.RUNNING
job.attempts += 1
job.locked_by = worker_id
job.locked_until = now + timedelta(seconds=lease_seconds)
job.last_error = ""
job.result = None
job.save(update_fields=["status", "attempts", "locked_by", "locked_until", "last_error", "result", "updated_at"])
emit_event(
job,
JobEvent.Type.CLAIMED,
worker_id=worker_id,
message="Job claimed",
data={"locked_until": job.locked_until.isoformat()},
)
logger.info(
"Claimed job %s type=%s attempt=%s worker=%s locked_until=%s.",
job.id,
job.type,
job.attempts,
worker_id,
job.locked_until,
)
return job
@transaction.atomic
def renew_lease(job_id, *, worker_id: str, attempt: int, lease_seconds: int | None = None) -> Job | None:
now = timezone.now()
lease_seconds = lease_seconds or settings.JOB_WORKER_LEASE_SECONDS
locked_until = now + timedelta(seconds=lease_seconds)
updated = Job.objects.filter(
id=job_id,
status=Job.Status.RUNNING,
locked_by=worker_id,
attempts=attempt,
).update(locked_until=locked_until, updated_at=now)
if not updated:
logger.debug("Lease renewal rejected job=%s worker=%s attempt=%s.", job_id, worker_id, attempt)
return None
job = Job.objects.get(id=job_id)
emit_event(
job,
JobEvent.Type.LEASE_RENEWED,
attempt=attempt,
worker_id=worker_id,
message="Lease renewed",
data={"locked_until": locked_until.isoformat()},
)
logger.debug("Renewed lease job=%s worker=%s attempt=%s locked_until=%s.", job_id, worker_id, attempt, locked_until)
return job
@transaction.atomic
def emit_progress(job_id, *, worker_id: str, attempt: int, percent: int, message: str = "") -> bool:
job = (
Job.objects.filter(
id=job_id,
status=Job.Status.RUNNING,
locked_by=worker_id,
attempts=attempt,
)
.select_for_update()
.first()
)
if job is None:
logger.debug("Progress rejected job=%s worker=%s attempt=%s percent=%s.", job_id, worker_id, attempt, percent)
return False
emit_event(
job,
JobEvent.Type.PROGRESS,
attempt=attempt,
worker_id=worker_id,
message=message or f"{percent}% complete",
data={"percent": percent},
)
logger.debug("Recorded progress job=%s worker=%s attempt=%s percent=%s.", job_id, worker_id, attempt, percent)
return True
@transaction.atomic
def complete_job(job_id, *, worker_id: str, attempt: int, result: dict | None = None) -> Job | None:
now = timezone.now()
updated = Job.objects.filter(
id=job_id,
status=Job.Status.RUNNING,
locked_by=worker_id,
attempts=attempt,
).update(
status=Job.Status.SUCCEEDED,
result=result or {},
locked_by=None,
locked_until=None,
finished_at=now,
updated_at=now,
)
if not updated:
logger.debug("Completion rejected job=%s worker=%s attempt=%s.", job_id, worker_id, attempt)
return None
job = Job.objects.get(id=job_id)
emit_event(
job,
JobEvent.Type.SUCCEEDED,
attempt=attempt,
worker_id=worker_id,
message="Job succeeded",
data={"result": job.result},
)
logger.info("Marked job %s succeeded worker=%s attempt=%s.", job.id, worker_id, attempt)
return job
@transaction.atomic
def fail_job(job_id, *, worker_id: str, attempt: int, error: str) -> Job | None:
now = timezone.now()
job = (
Job.objects.select_for_update()
.filter(id=job_id, status=Job.Status.RUNNING, locked_by=worker_id, attempts=attempt)
.first()
)
if job is None:
logger.debug("Failure rejected job=%s worker=%s attempt=%s error=%s.", job_id, worker_id, attempt, error)
return None
if job.attempts < job.max_attempts:
delay = backoff_delay(job.attempts)
job.status = Job.Status.QUEUED
job.available_at = now + delay
job.locked_by = None
job.locked_until = None
job.last_error = error
job.save(update_fields=["status", "available_at", "locked_by", "locked_until", "last_error", "updated_at"])
emit_event(
job,
JobEvent.Type.RETRY_SCHEDULED,
attempt=attempt,
worker_id=worker_id,
message=error,
data={"delay_seconds": int(delay.total_seconds()), "available_at": job.available_at.isoformat()},
)
logger.info(
"Scheduled retry job=%s attempt=%s/%s worker=%s delay_seconds=%s error=%s.",
job.id,
job.attempts,
job.max_attempts,
worker_id,
int(delay.total_seconds()),
error,
)
return job
job.status = Job.Status.FAILED
job.locked_by = None
job.locked_until = None
job.last_error = error
job.finished_at = now
job.save(update_fields=["status", "locked_by", "locked_until", "last_error", "finished_at", "updated_at"])
emit_event(job, JobEvent.Type.FAILED, attempt=attempt, worker_id=worker_id, message=error, data={"error": error})
logger.info("Marked job %s failed after attempt=%s worker=%s error=%s.", job.id, attempt, worker_id, error)
return job
@transaction.atomic
def cleanup_expired_jobs(*, batch_size: int | None = None) -> int:
now = timezone.now()
batch_size = batch_size or settings.JOB_WORKER_CLEANUP_BATCH_SIZE
queryset = Job.objects.filter(status=Job.Status.RUNNING, locked_until__lt=now).order_by("locked_until")
if connection.features.has_select_for_update_skip_locked:
queryset = queryset.select_for_update(skip_locked=True)
else:
queryset = queryset.select_for_update()
expired_jobs = list(queryset[:batch_size])
logger.debug("Found %s expired running job(s) for cleanup.", len(expired_jobs))
for job in expired_jobs:
worker_id = job.locked_by
attempt = job.attempts
if job.attempts < job.max_attempts:
delay = backoff_delay(job.attempts)
job.status = Job.Status.QUEUED
job.available_at = now + delay
job.locked_by = None
job.locked_until = None
job.last_error = "Worker lease expired"
job.save(update_fields=["status", "available_at", "locked_by", "locked_until", "last_error", "updated_at"])
emit_event(
job,
JobEvent.Type.TIMEOUT_REQUEUED,
attempt=attempt,
worker_id=worker_id,
message="Worker lease expired; job requeued",
data={"delay_seconds": int(delay.total_seconds()), "available_at": job.available_at.isoformat()},
)
logger.info(
"Requeued expired job=%s attempt=%s/%s previous_worker=%s delay_seconds=%s.",
job.id,
attempt,
job.max_attempts,
worker_id,
int(delay.total_seconds()),
)
continue
job.status = Job.Status.FAILED
job.locked_by = None
job.locked_until = None
job.last_error = "Worker lease expired"
job.finished_at = now
job.save(update_fields=["status", "locked_by", "locked_until", "last_error", "finished_at", "updated_at"])
emit_event(
job,
JobEvent.Type.FAILED,
attempt=attempt,
worker_id=worker_id,
message="Worker lease expired; attempts exhausted",
data={"error": "Worker lease expired"},
)
logger.info("Failed expired job=%s attempt=%s previous_worker=%s.", job.id, attempt, worker_id)
return len(expired_jobs)
@transaction.atomic
def retry_failed_job(job_id) -> Job:
job = Job.objects.select_for_update().get(id=job_id)
if job.status != Job.Status.FAILED:
raise ValueError("Only failed jobs can be retried.")
job.status = Job.Status.QUEUED
job.available_at = timezone.now()
job.attempts = 0
job.locked_by = None
job.locked_until = None
job.last_error = ""
job.result = None
job.finished_at = None
job.save(
update_fields=[
"status",
"available_at",
"attempts",
"locked_by",
"locked_until",
"last_error",
"result",
"finished_at",
"updated_at",
]
)
emit_event(job, JobEvent.Type.MANUAL_RETRY, message="Failed job manually retried")
logger.info("Manually retried failed job %s.", job.id)
return job
def job_stats() -> dict:
status_counts = dict(Job.objects.order_by().values_list("status").annotate(count=Count("id")))
now = timezone.now()
oldest_queued = Job.objects.filter(status=Job.Status.QUEUED).order_by("created_at").first()
return {
"total": Job.objects.count(),
"by_status": {status: status_counts.get(status, 0) for status in Job.Status.values},
"overdue_running": Job.objects.filter(status=Job.Status.RUNNING, locked_until__lt=now).count(),
"retries_pending": Job.objects.filter(status=Job.Status.QUEUED, attempts__gt=0).count(),
"oldest_queued_at": oldest_queued.created_at.isoformat() if oldest_queued else None,
}