66 lines
2.5 KiB
Python
66 lines
2.5 KiB
Python
"""High-contention concurrency test verifying atomic claims with 8+ parallel workers."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import concurrent.futures
|
|
import threading
|
|
from pathlib import Path
|
|
from typing import List
|
|
|
|
from src.runtime.storage.sqlite_store import SQLiteStore
|
|
|
|
|
|
def test_concurrent_claims_with_8_workers(tmp_path: Path):
|
|
"""Executes 8 parallel threads attempting to claim the same article fingerprint simultaneously.
|
|
|
|
Asserts:
|
|
1. Exactly 1 worker successfully obtains the claim and transitions to completed_text.
|
|
2. 7 workers receive active_claim or reuse existing result without duplicate writes.
|
|
3. Zero SQLite deadlock / database locked errors occur.
|
|
"""
|
|
db_path = tmp_path / "test_concurrency.db"
|
|
store = SQLiteStore(db_path)
|
|
fingerprint = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"
|
|
|
|
barrier = threading.Barrier(8)
|
|
results: List[str] = []
|
|
lock = threading.Lock()
|
|
|
|
def worker_action(worker_id: int):
|
|
# Synchronize all 8 workers at the starting line
|
|
barrier.wait()
|
|
worker_store = SQLiteStore(db_path)
|
|
try:
|
|
is_new, record = worker_store.claim_or_get_execution(
|
|
fingerprint=fingerprint,
|
|
source_url="https://example.com/test",
|
|
selected_extractor="trafilatura",
|
|
config_version="1.0.0",
|
|
)
|
|
with lock:
|
|
results.append(f"worker_{worker_id}:{'claimed' if is_new else 'reused'}")
|
|
if is_new:
|
|
# Worker simulates processing and records completion
|
|
worker_store.record_transition(
|
|
fingerprint,
|
|
"completed_text",
|
|
reason="Completed by winner worker",
|
|
extra_fields={"final_status": "completed_text"},
|
|
)
|
|
except Exception as e:
|
|
with lock:
|
|
results.append(f"worker_{worker_id}:error:{e}")
|
|
|
|
with concurrent.futures.ThreadPoolExecutor(max_workers=8) as executor:
|
|
futures = [executor.submit(worker_action, i) for i in range(8)]
|
|
concurrent.futures.wait(futures)
|
|
|
|
# Exactly 1 claimed
|
|
claimed_count = sum(1 for r in results if ":claimed" in r)
|
|
assert claimed_count == 1, f"Expected exactly 1 claim, got: {results}"
|
|
|
|
# Final state in DB must be completed_text
|
|
final_record = store.get_execution(fingerprint)
|
|
assert final_record is not None
|
|
assert final_record["current_status"] == "completed_text"
|