"""Concurrent job and outbox lease behavior backed by PostgreSQL.""" from concurrent.futures import ThreadPoolExecutor from datetime import datetime, timedelta, timezone import threading import uuid def _claimConcurrently(claim, firstWorker, secondWorker): barrier = threading.Barrier(2) def run(workerID): barrier.wait(timeout=10) return workerID, claim(workerID) with ThreadPoolExecutor(max_workers=2) as executor: futures = [ executor.submit(run, firstWorker), executor.submit(run, secondWorker), ] return dict(future.result(timeout=20) for future in futures) def test_job_claims_are_disjoint_and_support_retry_and_cancel( migratedDatabase ): from core import jobs, users users.registerUser("job-owner", "job-owner-password") userUUID = users.getUserUUID("job-owner") due = datetime.now(timezone.utc) - timedelta(minutes=1) created = [ jobs.create_job( "integration.work", {"sequence": index}, due, user_uuid=userUUID, idempotency_key=f"job-{uuid.uuid4()}", ) for index in range(10) ] claims = _claimConcurrently( lambda worker: jobs.claim_due_jobs(worker, limit=5, lease_seconds=60), "job-worker-one", "job-worker-two", ) firstIDs = {str(item["id"]) for item in claims["job-worker-one"]} secondIDs = {str(item["id"]) for item in claims["job-worker-two"]} assert firstIDs.isdisjoint(secondIDs) assert firstIDs | secondIDs == {str(item["id"]) for item in created} assert jobs.claim_due_jobs("job-worker-three", limit=10) == [] retriedID = str(claims["job-worker-one"][0]["id"]) assert jobs.fail_job(retriedID, "wrong-worker", "must not update") is None retried = jobs.fail_job( retriedID, "job-worker-one", "temporary job failure", retry_seconds=1 ) assert retried["status"] == "pending" assert retried["attempts"] == 1 assert retried["last_error"] == "temporary job failure" assert retried["run_at"] > datetime.now(timezone.utc) cancelledID = str(claims["job-worker-two"][0]["id"]) assert jobs.cancel_job(cancelledID, user_uuid=uuid.uuid4()) is None cancelled = jobs.cancel_job(cancelledID, user_uuid=userUUID) assert cancelled["status"] == "cancelled" assert cancelled["leased_by"] is None def test_outbox_claims_are_disjoint_and_support_retry_and_cancel( migratedDatabase ): from core import outbox, users users.registerUser("outbox-owner", "outbox-owner-password") userUUID = users.getUserUUID("outbox-owner") due = datetime.now(timezone.utc) - timedelta(minutes=1) created = [ outbox.enqueue_message( userUUID, "discord_dm", {"content": f"message {index}"}, idempotency_key=f"outbox-{uuid.uuid4()}", available_at=due, ) for index in range(10) ] claims = _claimConcurrently( lambda worker: outbox.claim_messages( worker, channel="discord_dm", limit=5, lease_seconds=60 ), "outbox-worker-one", "outbox-worker-two", ) firstIDs = {str(item["id"]) for item in claims["outbox-worker-one"]} secondIDs = {str(item["id"]) for item in claims["outbox-worker-two"]} assert firstIDs.isdisjoint(secondIDs) assert firstIDs | secondIDs == {str(item["id"]) for item in created} assert outbox.claim_messages("outbox-worker-three", limit=10) == [] retriedID = str(claims["outbox-worker-one"][0]["id"]) assert outbox.retry_message( retriedID, "wrong-worker", "must not update" ) is None retried = outbox.retry_message( retriedID, "outbox-worker-one", "temporary delivery failure", retry_seconds=1, ) assert retried["status"] == "pending" assert retried["attempts"] == 1 assert retried["last_error"] == "temporary delivery failure" assert retried["available_at"] > datetime.now(timezone.utc) cancelledID = str(claims["outbox-worker-two"][0]["id"]) assert outbox.cancel_message(cancelledID, user_uuid=uuid.uuid4()) is None cancelled = outbox.cancel_message(cancelledID, user_uuid=userUUID) assert cancelled["status"] == "cancelled" assert cancelled["leased_by"] is None