Files
tabdeal-job-queue/backend/jobs/tests/test_services.py

186 lines
6.8 KiB
Python

from concurrent.futures import ThreadPoolExecutor
from datetime import timedelta
import pytest
from django.conf import settings
from django.db import close_old_connections, connection
from django.utils import timezone
from jobs.models import Job, JobEvent
from jobs.services import (
claim_next_job,
cleanup_expired_jobs,
complete_job,
create_job,
fail_job,
renew_lease,
retry_failed_job,
)
from jobs.worker import JobWorkerRunner
@pytest.mark.django_db
def test_successful_claim_sets_lock_and_event():
job = create_job(job_type="demo.success").job
claimed = claim_next_job(worker_id="worker-a", lease_seconds=30)
assert claimed.id == job.id
assert claimed.status == Job.Status.RUNNING
assert claimed.locked_by == "worker-a"
assert claimed.locked_until is not None
assert claimed.attempts == 1
assert JobEvent.objects.filter(job=job, type=JobEvent.Type.CLAIMED, worker_id="worker-a").count() == 1
@pytest.mark.django_db(transaction=True)
def test_concurrent_claim_prevention_with_postgresql():
if connection.vendor != "postgresql":
pytest.skip("SKIP LOCKED concurrency is a PostgreSQL behavior.")
create_job(job_type="demo.success")
def claim(index):
close_old_connections()
try:
job = claim_next_job(worker_id=f"worker-{index}", lease_seconds=30)
return str(job.id) if job else None
finally:
close_old_connections()
with ThreadPoolExecutor(max_workers=10) as executor:
claimed_ids = list(executor.map(claim, range(10)))
assert len([job_id for job_id in claimed_ids if job_id]) == 1
assert JobEvent.objects.filter(type=JobEvent.Type.CLAIMED).count() == 1
@pytest.mark.django_db
def test_priority_and_available_at_ordering_are_deterministic():
future = timezone.now() + timedelta(hours=1)
low = create_job(job_type="demo.success", priority=10).job
high = create_job(job_type="demo.success", priority=90).job
create_job(job_type="demo.success", priority=100, available_at=future)
claimed = claim_next_job(worker_id="worker-a", lease_seconds=30)
assert claimed.id == high.id
next_claim = claim_next_job(worker_id="worker-b", lease_seconds=30)
assert next_claim.id == low.id
@pytest.mark.django_db
def test_idempotent_create_returns_existing_job():
first = create_job(job_type="demo.success", idempotency_key="invoice-123")
second = create_job(job_type="demo.success", idempotency_key="invoice-123")
assert first.created is True
assert second.created is False
assert second.job.id == first.job.id
assert Job.objects.count() == 1
assert JobEvent.objects.filter(type=JobEvent.Type.CREATED).count() == 1
@pytest.mark.django_db
def test_failure_schedules_exponential_retry():
job = create_job(job_type="demo.fail", max_attempts=3).job
claimed = claim_next_job(worker_id="worker-a", lease_seconds=30)
before = timezone.now()
updated = fail_job(claimed.id, worker_id="worker-a", attempt=1, error="boom")
assert updated.status == Job.Status.QUEUED
assert updated.attempts == 1
assert updated.available_at >= before + timedelta(seconds=settings.JOB_WORKER_BACKOFF_BASE_SECONDS)
assert updated.locked_by is None
assert JobEvent.objects.filter(job=job, type=JobEvent.Type.RETRY_SCHEDULED).exists()
@pytest.mark.django_db
def test_exhausted_attempts_become_failed():
job = create_job(job_type="demo.fail", max_attempts=1).job
claimed = claim_next_job(worker_id="worker-a", lease_seconds=30)
updated = fail_job(claimed.id, worker_id="worker-a", attempt=1, error="boom")
assert updated.status == Job.Status.FAILED
assert updated.finished_at is not None
assert updated.locked_by is None
assert JobEvent.objects.filter(job=job, type=JobEvent.Type.FAILED).exists()
@pytest.mark.django_db
def test_expired_lease_cleanup_requeues_or_fails():
requeue = create_job(job_type="demo.timeout", max_attempts=2).job
fail = create_job(job_type="demo.timeout", max_attempts=1).job
requeue_claim = claim_next_job(worker_id="worker-a", lease_seconds=30)
fail_claim = claim_next_job(worker_id="worker-b", lease_seconds=30)
Job.objects.filter(id__in=[requeue_claim.id, fail_claim.id]).update(
locked_until=timezone.now() - timedelta(seconds=1)
)
assert cleanup_expired_jobs(batch_size=10) == 2
requeue.refresh_from_db()
fail.refresh_from_db()
assert requeue.status == Job.Status.QUEUED
assert fail.status == Job.Status.FAILED
assert JobEvent.objects.filter(job=requeue, type=JobEvent.Type.TIMEOUT_REQUEUED).exists()
assert JobEvent.objects.filter(job=fail, type=JobEvent.Type.FAILED).exists()
@pytest.mark.django_db
def test_stale_worker_cannot_complete_old_attempt():
job = create_job(job_type="demo.timeout", max_attempts=3).job
first = claim_next_job(worker_id="worker-a", lease_seconds=30)
Job.objects.filter(id=first.id).update(locked_until=timezone.now() - timedelta(seconds=1))
cleanup_expired_jobs(batch_size=10)
Job.objects.filter(id=job.id).update(available_at=timezone.now() - timedelta(seconds=1))
second = claim_next_job(worker_id="worker-b", lease_seconds=30)
stale_result = complete_job(job.id, worker_id="worker-a", attempt=1, result={"ok": True})
assert stale_result is None
job.refresh_from_db()
assert job.status == Job.Status.RUNNING
assert job.locked_by == "worker-b"
assert job.attempts == second.attempts == 2
@pytest.mark.django_db
def test_lease_renewal_only_current_owner_and_attempt():
job = create_job(job_type="demo.slow").job
claimed = claim_next_job(worker_id="worker-a", lease_seconds=5)
old_locked_until = claimed.locked_until
assert renew_lease(job.id, worker_id="worker-b", attempt=1, lease_seconds=60) is None
assert renew_lease(job.id, worker_id="worker-a", attempt=2, lease_seconds=60) is None
renewed = renew_lease(job.id, worker_id="worker-a", attempt=1, lease_seconds=60)
assert renewed is not None
assert renewed.locked_until > old_locked_until
assert JobEvent.objects.filter(job=job, type=JobEvent.Type.LEASE_RENEWED).count() == 1
@pytest.mark.django_db
def test_manual_retry_resets_failed_job():
job = create_job(job_type="demo.fail", max_attempts=1).job
claimed = claim_next_job(worker_id="worker-a", lease_seconds=30)
fail_job(claimed.id, worker_id="worker-a", attempt=1, error="boom")
retried = retry_failed_job(job.id)
assert retried.status == Job.Status.QUEUED
assert retried.attempts == 0
assert retried.finished_at is None
assert JobEvent.objects.filter(job=job, type=JobEvent.Type.MANUAL_RETRY).exists()
def test_runner_request_stop_prevents_more_claiming():
runner = JobWorkerRunner(thread_count=1)
assert runner.stop_event.is_set() is False
runner.request_stop()
assert runner.stop_event.is_set() is True