"""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"