"""Two workers against the real database, at the same instant."""
import os, pathlib, sys, threading
sys.path.insert(0, r"C:\taxpilot-ai-work")
for line in pathlib.Path(r"C:\taxpilot-ai-staging\shared\.env").read_text(encoding="utf-8").splitlines():
    line = line.strip()
    if line and not line.startswith("#") and "=" in line:
        k, v = line.split("=", 1)
        os.environ.setdefault(k.strip(), v.strip().strip('"'))

from datetime import UTC, datetime
from app.runtime.container import Container
from app.database.migrator import Migrator
from app.intake import PendingRegistration

c = Container()
print("   applied:", Migrator(c.connect).apply() or "already up to date")

box = c.intake_outbox
for item in box.all():
    box.done(item)

for i in range(20):
    box.add(PendingRegistration(
        source="whatsapp", provider_message_id=f"RACE-{i:03d}",
        original_filename=f"race-{i}.pdf", received_at=datetime.now(UTC),
    ))
print("   queued:", len(box))

taken, lock = [], threading.Lock()

def worker(name):
    own = Container().intake_outbox          # its own connection, as a real worker has
    while True:
        got = own.claim(name)
        if got is None:
            return
        with lock:
            taken.append(got.provider_message_id)

threads = [threading.Thread(target=worker, args=(f"w{i}",)) for i in range(4)]
for t in threads: t.start()
for t in threads: t.join()

print(f"   4 workers claimed {len(taken)} items; distinct = {len(set(taken))}")
print("   duplicates:", len(taken) - len(set(taken)))

for item in Container().intake_outbox.all():
    Container().intake_outbox.done(item)
