123 lines
4.3 KiB
Python
123 lines
4.3 KiB
Python
"""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
|