65 lines
1.9 KiB
Python
65 lines
1.9 KiB
Python
"""Concurrent one-shot session start tests."""
|
|
from __future__ import annotations
|
|
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
from pathlib import Path
|
|
from threading import Barrier
|
|
|
|
from app.services.sessions import SessionStore
|
|
|
|
|
|
def test_synchronized_starts_share_one_active_session(tmp_path: Path):
|
|
store = SessionStore(tmp_path)
|
|
barrier = Barrier(8)
|
|
|
|
def start_one():
|
|
barrier.wait(timeout=5)
|
|
return store.start(
|
|
org_id="org-1",
|
|
user_id="user-1",
|
|
group_id="group-a",
|
|
persona_id="persona-01",
|
|
persona_name="Customer",
|
|
)
|
|
|
|
with ThreadPoolExecutor(max_workers=8) as pool:
|
|
results = list(pool.map(lambda _: start_one(), range(8)))
|
|
|
|
sessions = [session for session, _ in results]
|
|
assert len({session["id"] for session in sessions}) == 1
|
|
assert sum(1 for _, resumed in results if resumed is False) == 1
|
|
assert sum(1 for _, resumed in results if resumed is True) == 7
|
|
|
|
|
|
def test_finished_scope_remains_blocked_after_concurrent_start(tmp_path: Path):
|
|
store = SessionStore(tmp_path)
|
|
session, _ = store.start(
|
|
org_id="org-1",
|
|
user_id="user-1",
|
|
group_id="group-a",
|
|
persona_id="persona-01",
|
|
persona_name="Customer",
|
|
)
|
|
store.update(session["id"], status="finished", outcome="won")
|
|
|
|
barrier = Barrier(8)
|
|
|
|
def try_start():
|
|
barrier.wait(timeout=5)
|
|
try:
|
|
store.start(
|
|
org_id="org-1",
|
|
user_id="user-1",
|
|
group_id="group-a",
|
|
persona_id="persona-01",
|
|
persona_name="Customer",
|
|
)
|
|
except ValueError as exc:
|
|
return str(exc)
|
|
return "unexpected-success"
|
|
|
|
with ThreadPoolExecutor(max_workers=8) as pool:
|
|
results = list(pool.map(lambda _: try_start(), range(8)))
|
|
|
|
assert results == ["you have already trained on this persona (one-shot)"] * 8
|