Files

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"