186 lines
6.8 KiB
Python
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
|