mirror of
https://github.com/turnstonelabs/turnstone.git
synced 2026-08-12 23:12:23 -06:00
b8dd5041c4
The force branch cleared worker ownership and emitted idle from the route thread, while the abandon latch and the queue demote ran only in the stuck worker's own exception handler — a thread force-cancel abandons precisely because it is not making progress. Subscribers on the IDLE fan-out therefore saw an operator-forced idle with the latch unset: the idle observer's operator-Stop gate did not suppress advice, and wake-eligible entries survived un-demoted, so a nudge wake could resume a workstream seconds after the operator forced it to stop. The route now runs the session's abandon machinery first; the abandoned thread re-running it at its eventual death is idempotent.
2549 lines
105 KiB
Python
2549 lines
105 KiB
Python
"""HTTP-boundary authorization tests for turnstone-server.
|
|
|
|
Covers the ownership gates added in PR #2 (sec-1 through sec-9 +
|
|
sec-11) and the kind-validation branches that PR #1 tightened but
|
|
never had Starlette-level regression coverage. Each test crosses
|
|
the middleware → handler boundary via ``TestClient`` so the JWT
|
|
decoding, scope extraction, and audit-context wiring are all exercised.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import queue
|
|
import threading
|
|
from typing import Any
|
|
from unittest.mock import MagicMock
|
|
|
|
import pytest
|
|
from starlette.testclient import TestClient
|
|
|
|
from tests._helpers import wait_until
|
|
from turnstone.core.workstream import INTERJECTION_CAP_CHARS
|
|
|
|
_TEST_JWT_SECRET = "test-jwt-secret-minimum-32-chars!"
|
|
|
|
|
|
# Default permission set for test JWTs. Mirrors what builtin-operator
|
|
# carries: enough perms to exercise create/close/approve gates without
|
|
# turning every existing test into a re-authorization round. Tests
|
|
# negating these gates pass ``permissions=frozenset()`` explicitly.
|
|
_DEFAULT_TEST_PERMS = frozenset(
|
|
{"workstreams.create", "workstreams.close", "tools.approve", "conversation.modify"}
|
|
)
|
|
|
|
|
|
def _make_jwt(
|
|
user_id: str,
|
|
*,
|
|
scopes: frozenset[str] | None = None,
|
|
permissions: frozenset[str] | None = None,
|
|
) -> str:
|
|
from turnstone.core.auth import JWT_AUD_SERVER, create_jwt
|
|
|
|
return create_jwt(
|
|
user_id=user_id,
|
|
scopes=scopes or frozenset({"read", "write", "approve"}),
|
|
source="test",
|
|
secret=_TEST_JWT_SECRET,
|
|
audience=JWT_AUD_SERVER,
|
|
permissions=_DEFAULT_TEST_PERMS if permissions is None else permissions,
|
|
)
|
|
|
|
|
|
def _auth(
|
|
user: str,
|
|
*,
|
|
scopes: frozenset[str] | None = None,
|
|
permissions: frozenset[str] | None = None,
|
|
) -> dict[str, str]:
|
|
return {"Authorization": f"Bearer {_make_jwt(user, scopes=scopes, permissions=permissions)}"}
|
|
|
|
|
|
class TestAssignableScopes:
|
|
"""``service`` scope is a cross-tenant bypass and must never be
|
|
GRANTED via a user-facing token mint (admin API or CLI) — otherwise an
|
|
``admin.users`` holder could self-mint it and see every private
|
|
project's workstreams. Both mint paths route through
|
|
:func:`reject_unassignable_scopes`."""
|
|
|
|
def test_service_scope_rejected(self) -> None:
|
|
from turnstone.core.auth import reject_unassignable_scopes
|
|
|
|
assert reject_unassignable_scopes("service") is not None
|
|
assert reject_unassignable_scopes("read,service") is not None
|
|
assert reject_unassignable_scopes("read,write,approve,service") is not None
|
|
|
|
def test_service_not_in_assignable_set(self) -> None:
|
|
from turnstone.core.auth import ASSIGNABLE_SCOPES, VALID_SCOPES
|
|
|
|
assert "service" in VALID_SCOPES # still a valid runtime scope
|
|
assert "service" not in ASSIGNABLE_SCOPES # but not user-assignable
|
|
|
|
def test_ordinary_scopes_accepted(self) -> None:
|
|
from turnstone.core.auth import reject_unassignable_scopes
|
|
|
|
assert reject_unassignable_scopes("read") is None
|
|
assert reject_unassignable_scopes("read,write,approve") is None
|
|
|
|
def test_empty_and_unknown_rejected(self) -> None:
|
|
from turnstone.core.auth import reject_unassignable_scopes
|
|
|
|
assert reject_unassignable_scopes("") is not None
|
|
assert reject_unassignable_scopes("bogus") is not None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# FakeUI / FakeSession doubles — match the shape the create handler expects
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class _FakeUI:
|
|
def __init__(self, ws_id: str = "", user_id: str = "", **_kw: Any) -> None:
|
|
self.ws_id = ws_id
|
|
self._user_id = user_id
|
|
self.auto_approve = False
|
|
self.auto_approve_tools: set[str] = set()
|
|
self._enqueued: list[dict[str, Any]] = []
|
|
self.states: list[str] = []
|
|
self.infos: list[str] = []
|
|
self.errors: list[str] = []
|
|
self._listeners: list[queue.Queue[dict[str, Any]]] = []
|
|
self._listeners_lock = threading.Lock()
|
|
self._pending_approval: dict[str, Any] | None = None
|
|
self._approval_event = threading.Event()
|
|
self._fg_event = threading.Event()
|
|
self._ws_lock = threading.Lock()
|
|
# Dashboard handler reads these fields under _ws_lock to build
|
|
# per-ws summary rows; keep them zero/empty for the fake so the
|
|
# handler doesn't need to special-case.
|
|
self._ws_prompt_tokens = 0
|
|
self._ws_completion_tokens = 0
|
|
self._ws_tool_calls: dict[str, int] = {}
|
|
self._ws_context_ratio = 0.0
|
|
self._ws_current_activity = ""
|
|
self._ws_activity_state = ""
|
|
self._ws_messages = 0
|
|
self._ws_turn_tool_calls = 0
|
|
self._llm_verdicts: dict[str, dict[str, Any]] = {}
|
|
|
|
def serialize_pending_approval_details(self) -> list[dict[str, Any]]:
|
|
# Mirrors SessionUIBase.serialize_pending_approval_details —
|
|
# the fake is monkeypatched in for ``WebUI`` and the dashboard
|
|
# handler reads this method during projection. Real subclasses
|
|
# iterate their approval-cycle registry (one entry per live
|
|
# cycle); the fake models a single slot, so the list carries
|
|
# zero or one entries.
|
|
pending = self._pending_approval
|
|
if pending is None:
|
|
return []
|
|
items = pending.get("items") or []
|
|
if not items:
|
|
return []
|
|
call_ids = [item.get("call_id", "") for item in items]
|
|
# Match the real impl's pattern (session_ui_base.py): snapshot
|
|
# references under the lock, copy after release. Writers only
|
|
# assign — never mutate — so the reference snapshot is stable
|
|
# outside the lock window.
|
|
with self._ws_lock:
|
|
verdict_refs = {
|
|
cid: self._llm_verdicts[cid]
|
|
for cid in call_ids
|
|
if cid and cid in self._llm_verdicts
|
|
}
|
|
verdicts = {cid: dict(v) for cid, v in verdict_refs.items()}
|
|
serialized: list[dict[str, Any]] = []
|
|
for item in items:
|
|
cid = item.get("call_id", "")
|
|
serialized.append(
|
|
{
|
|
"call_id": cid,
|
|
"header": item.get("header", ""),
|
|
"preview": item.get("preview", ""),
|
|
"func_name": item.get("func_name", ""),
|
|
"approval_label": item.get("approval_label", ""),
|
|
"needs_approval": item.get("needs_approval", False),
|
|
"error": item.get("error"),
|
|
"heuristic_verdict": item.get("verdict"),
|
|
"judge_verdict": verdicts.get(cid),
|
|
}
|
|
)
|
|
# Primary call_id must mirror the real serializer: first
|
|
# *non-empty* in list order, not just first. Aligning the
|
|
# fake here keeps test-vs-prod behavioural drift from
|
|
# masking a real-shape regression.
|
|
primary = next((cid for cid in call_ids if cid), "")
|
|
return [
|
|
{
|
|
"cycle_id": pending.get("cycle_id", ""),
|
|
"call_id": primary,
|
|
"judge_pending": bool(pending.get("judge_pending", False)),
|
|
"items": serialized,
|
|
}
|
|
]
|
|
|
|
def serialize_recent_auto_approvals(self) -> list[dict[str, Any]]:
|
|
# Empty buffer for tests that don't exercise the auto-approve
|
|
# visibility path. /dashboard handler reads this method
|
|
# unconditionally now (paired with serialize_pending_approval_detail);
|
|
# returning [] keeps the row payload compatible without
|
|
# modeling the full ring buffer in the fake.
|
|
return []
|
|
|
|
def _register_listener(self) -> queue.Queue[dict[str, Any]]:
|
|
q: queue.Queue[dict[str, Any]] = queue.Queue()
|
|
with self._listeners_lock:
|
|
self._listeners.append(q)
|
|
return q
|
|
|
|
def _enqueue(self, ev: dict[str, Any]) -> None:
|
|
self._enqueued.append(ev)
|
|
|
|
def on_stream_end(self) -> None:
|
|
pass
|
|
|
|
def on_state_change(self, state: str) -> None:
|
|
self.states.append(state)
|
|
|
|
def on_info(self, msg: str) -> None:
|
|
self.infos.append(msg)
|
|
|
|
def on_error(self, msg: str) -> None:
|
|
self.errors.append(msg)
|
|
|
|
def resolve_approval(self, *_a: Any, **_kw: Any) -> None:
|
|
self._approval_event.set()
|
|
|
|
|
|
class _FakeSession:
|
|
def __init__(self, ws_id: str = "", user_id: str = "") -> None:
|
|
self.ws_id = ws_id
|
|
self.user_id = user_id
|
|
self.model = "test-model"
|
|
self.model_alias = ""
|
|
self.reasoning_effort = ""
|
|
self.context_window = 100000
|
|
self.messages: list[dict[str, Any]] = []
|
|
self._last_usage: dict[str, int] | None = None
|
|
self._pending_retry: str | None = None
|
|
# Real sessions always carry one; the /send route's cancel-drain
|
|
# poll reads it whenever a worker is live at send time.
|
|
self._cancel_event = threading.Event()
|
|
self.sends: list[tuple[str, Any, Any]] = []
|
|
self.commands: list[str] = []
|
|
self.compacts = 0
|
|
self.compact_raises: BaseException | None = None
|
|
self.send_raises: BaseException | None = None
|
|
self.exit_commands: set[str] = set()
|
|
self.queued_flushes = 0
|
|
self.queued_text = ""
|
|
# Gates let a test hold the worker slot open mid-command so a
|
|
# concurrent /send provably lands in the PARK path.
|
|
self.compact_gate: threading.Event | None = None
|
|
self.command_gate: threading.Event | None = None
|
|
self.command_raises: BaseException | None = None
|
|
# Gate for send() — lets a test hold a TURN in flight so a
|
|
# deferred entry provably lands in (or is refused by) the
|
|
# interjection fallback at drain time.
|
|
self.send_gate: threading.Event | None = None
|
|
# Every queue_message call — the defer invariant is that command
|
|
# windows NEVER reach the interjection queue.
|
|
self.queue_calls: list[str] = []
|
|
# When set, queue_message records the attempt and then raises it
|
|
# (e.g. CrossUserInterjectionError for the drain's re-park arm).
|
|
self.queue_raises: BaseException | None = None
|
|
# DELETE /send fall-through: ids the route asked this session to
|
|
# dequeue (the fake never holds interjections, so it returns
|
|
# False and the route falls through to the deferred-send list).
|
|
self.dequeues: list[str] = []
|
|
|
|
def send(self, text: str, *, attachments: Any = None, send_id: Any = None) -> None:
|
|
self.sends.append((text, attachments, send_id))
|
|
if self.send_gate is not None:
|
|
self.send_gate.wait(timeout=10)
|
|
if self.send_raises is not None:
|
|
raise self.send_raises
|
|
|
|
def queue_message(
|
|
self,
|
|
text: str,
|
|
attachment_ids: Any = None,
|
|
queue_msg_id: str | None = None,
|
|
interjector_user_id: str = "",
|
|
) -> tuple[str, str, str]:
|
|
self.queue_calls.append(text)
|
|
if self.queue_raises is not None:
|
|
raise self.queue_raises
|
|
cap = INTERJECTION_CAP_CHARS
|
|
cleaned = text[:cap] + "..." if len(text) > cap else text
|
|
return cleaned, "notice", queue_msg_id or "m1"
|
|
|
|
def dequeue_message(self, msg_id: str) -> bool:
|
|
self.dequeues.append(msg_id)
|
|
return False
|
|
|
|
def set_watch_runner(self, *_a: Any, **_kw: Any) -> None:
|
|
pass
|
|
|
|
def resume(self, _ws_id: str, *, fork: bool = False) -> bool:
|
|
return False
|
|
|
|
def cancel(self) -> None:
|
|
pass
|
|
|
|
def close(self) -> None:
|
|
pass
|
|
|
|
def handle_command(self, cmd: str) -> bool:
|
|
self.commands.append(cmd)
|
|
if self.command_gate is not None:
|
|
self.command_gate.wait(timeout=10)
|
|
if self.command_raises is not None:
|
|
raise self.command_raises
|
|
return cmd in self.exit_commands
|
|
|
|
def compact_now(self) -> bool:
|
|
self.compacts += 1
|
|
if self.compact_gate is not None:
|
|
self.compact_gate.wait(timeout=10)
|
|
if self.compact_raises is not None:
|
|
raise self.compact_raises
|
|
return True
|
|
|
|
def flush_queued_messages(self) -> bool:
|
|
self.queued_flushes += 1
|
|
had, self.queued_text = bool(self.queued_text), ""
|
|
return had
|
|
|
|
def request_title_refresh(self, _title: str) -> None:
|
|
pass
|
|
|
|
|
|
@pytest.fixture
|
|
def app_client(tmp_path, monkeypatch):
|
|
"""Full turnstone-server app with in-memory workstreams + fake sessions."""
|
|
from turnstone.core.adapters.interactive_adapter import InteractiveAdapter
|
|
from turnstone.core.metrics import MetricsCollector
|
|
from turnstone.core.session_manager import SessionManager
|
|
from turnstone.core.storage import get_storage, init_storage, reset_storage
|
|
from turnstone.server import WebUI, create_app
|
|
|
|
reset_storage()
|
|
init_storage("sqlite", path=str(tmp_path / "t.db"), run_migrations=False)
|
|
|
|
metrics = MetricsCollector()
|
|
metrics.model = "test-model"
|
|
monkeypatch.setattr("turnstone.server._metrics", metrics)
|
|
monkeypatch.setattr("turnstone.server.WebUI", _FakeUI)
|
|
|
|
def _factory(ui: Any, _model: Any, ws_id: str, **_kw: Any) -> _FakeSession:
|
|
uid = getattr(ui, "_user_id", "")
|
|
return _FakeSession(ws_id=ws_id, user_id=uid)
|
|
|
|
gq: queue.Queue[dict[str, Any]] = queue.Queue(maxsize=1000)
|
|
WebUI._global_queue = gq
|
|
adapter = InteractiveAdapter(
|
|
global_queue=gq,
|
|
ui_factory=lambda ws: _FakeUI(
|
|
ws_id=ws.id,
|
|
user_id=ws.user_id,
|
|
),
|
|
session_factory=_factory,
|
|
)
|
|
mgr = SessionManager(
|
|
adapter, storage=get_storage(), max_active=10, node_id="node-test", event_emitter=adapter
|
|
)
|
|
app = create_app(
|
|
workstreams=mgr,
|
|
global_queue=gq,
|
|
global_listeners=[],
|
|
global_listeners_lock=threading.Lock(),
|
|
skip_permissions=False,
|
|
jwt_secret=_TEST_JWT_SECRET,
|
|
auth_storage=get_storage(),
|
|
)
|
|
client = TestClient(app, raise_server_exceptions=False)
|
|
try:
|
|
yield client, mgr
|
|
finally:
|
|
client.close()
|
|
reset_storage()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# PR #1 HTTP-boundary kind validation (q-4) — previously untested
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestKindValidationOnCreate:
|
|
"""POST /v1/api/workstreams/new — kind field validation at the HTTP edge."""
|
|
|
|
def test_rejects_kind_coordinator(self, app_client):
|
|
client, _mgr = app_client
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"kind": "coordinator", "name": "x"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 400
|
|
assert "coordinator" in resp.json()["error"].lower()
|
|
|
|
def test_rejects_unknown_kind(self, app_client):
|
|
client, _mgr = app_client
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"kind": "interative", "name": "x"}, # typo
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 400
|
|
assert "unknown" in resp.json()["error"].lower()
|
|
|
|
def test_accepts_default_kind(self, app_client):
|
|
client, _mgr = app_client
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "x"}, # kind omitted
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
|
|
def test_rejects_cross_tenant_parent_ws_id(self, app_client, tmp_path):
|
|
"""parent_ws_id pointing at another user's coordinator → 403."""
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
# Victim creates a coordinator directly in storage (console path).
|
|
storage.register_workstream(
|
|
"victim-coord",
|
|
node_id="console",
|
|
name="victim",
|
|
kind="coordinator",
|
|
user_id="victim-user",
|
|
)
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "attacker", "parent_ws_id": "victim-coord"},
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 403
|
|
assert "coordinator you own" in resp.json()["error"]
|
|
|
|
|
|
class TestOpenKindGate:
|
|
"""POST /v1/api/workstreams/{ws_id}/open refuses coordinator rows.
|
|
|
|
Post-lift behavior change: the lifted ``open`` body delegates the
|
|
kind check to ``SessionManager.open()`` (which returns ``None``
|
|
for kind mismatch / missing row / tombstone — all the
|
|
"manager has no such ws_id" cases). The pre-lift handler had a
|
|
separate pre-mgr storage probe that returned a kind-specific
|
|
400 ("Workstream is not an interactive kind"); the lift
|
|
consolidates on a single 404 ("Workstream not found"). Security
|
|
boundary unchanged — caller still can't open a coord row from
|
|
the interactive node — but the error code + message converge
|
|
with the rest of the not-found paths.
|
|
"""
|
|
|
|
def test_refuses_to_open_coordinator(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
storage.register_workstream(
|
|
"coord-1",
|
|
node_id="console",
|
|
name="c",
|
|
kind="coordinator",
|
|
user_id="user-1",
|
|
)
|
|
resp = client.post(
|
|
"/v1/api/workstreams/coord-1/open",
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 404
|
|
assert "not found" in resp.json()["error"].lower()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# PR #2 authz cluster — cross-tenant gates on interactive-ws mutations
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _register_ws(storage: Any, ws_id: str, owner: str) -> None:
|
|
storage.register_workstream(ws_id, node_id="node-test", name=ws_id, user_id=owner)
|
|
|
|
|
|
class TestCrossTenantDelete:
|
|
def test_any_caller_can_delete(self, app_client):
|
|
# Trusted-team model: scope auth gates the endpoint, not
|
|
# row-level ownership. ``user_id`` stays on audit + storage
|
|
# metadata.
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-victim", "victim-user")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-victim/delete",
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 200
|
|
|
|
def test_owner_delete_records_audit(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-own", "user-1")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-own/delete",
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
events = storage.list_audit_events(action="workstream.deleted")
|
|
assert any(e["resource_id"] == "ws-own" for e in events)
|
|
|
|
|
|
class TestCrossTenantApprove:
|
|
def test_non_owner_cannot_approve(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-victim", "victim-user")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-victim/approve",
|
|
json={"approved": True},
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 404
|
|
|
|
|
|
class TestCrossTenantClose:
|
|
def test_non_owner_cannot_close(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-victim", "victim-user")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-victim/close",
|
|
json={},
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 404
|
|
|
|
|
|
class TestPermissionGatesOnLifecycle:
|
|
"""Gates that previously didn't exist — ``workstreams.create``,
|
|
``workstreams.close``, ``tools.approve`` were declared, seeded into
|
|
builtin-operator, surfaced in the admin Roles UI, and never wired
|
|
to a single ``require_permission`` site. PR added the gates; these
|
|
tests confirm a JWT without each perm gets 403."""
|
|
|
|
def test_create_without_perm_returns_403(self, app_client):
|
|
client, _mgr = app_client
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "no-perm"},
|
|
headers=_auth("user-1", permissions=frozenset()),
|
|
)
|
|
assert resp.status_code == 403
|
|
assert "workstreams.create" in resp.json()["error"]
|
|
|
|
def test_close_without_perm_returns_403(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-1", "user-1")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-1/close",
|
|
json={},
|
|
headers=_auth("user-1", permissions=frozenset()),
|
|
)
|
|
assert resp.status_code == 403
|
|
assert "workstreams.close" in resp.json()["error"]
|
|
|
|
def test_approve_without_perm_returns_403(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-1", "user-1")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-1/approve",
|
|
json={"approved": True},
|
|
headers=_auth("user-1", permissions=frozenset()),
|
|
)
|
|
assert resp.status_code == 403
|
|
assert "tools.approve" in resp.json()["error"]
|
|
|
|
def test_create_with_perm_passes_gate(self, app_client):
|
|
# Sanity: same call WITH the perm reaches the post-gate logic
|
|
# (whatever its outcome — a successful create or a non-403
|
|
# validation/state error is fine; only the gate behaviour is
|
|
# under test here).
|
|
client, _mgr = app_client
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "with-perm"},
|
|
headers=_auth("user-1", permissions=frozenset({"workstreams.create"})),
|
|
)
|
|
assert resp.status_code != 403, resp.json()
|
|
|
|
# Positive coverage for the admin.coordinator OR-fallback on each
|
|
# of the three lifted verbs. Without these, a future refactor
|
|
# that dropped admin.coordinator from the accepted_permissions
|
|
# tuple would regress coord-session children silently — the proxy
|
|
# tests only exercise the route_proxy verb dict, not the lift.
|
|
|
|
def test_create_with_admin_coordinator_passes_gate(self, app_client):
|
|
client, _mgr = app_client
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "coord-child"},
|
|
headers=_auth("user-1", permissions=frozenset({"admin.coordinator"})),
|
|
)
|
|
assert resp.status_code != 403, resp.json()
|
|
|
|
def test_close_with_admin_coordinator_passes_gate(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-1", "user-1")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-1/close",
|
|
json={},
|
|
headers=_auth("user-1", permissions=frozenset({"admin.coordinator"})),
|
|
)
|
|
assert resp.status_code != 403, resp.json()
|
|
|
|
def test_approve_with_admin_coordinator_passes_gate(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-1", "user-1")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-1/approve",
|
|
json={"approved": True},
|
|
headers=_auth("user-1", permissions=frozenset({"admin.coordinator"})),
|
|
)
|
|
assert resp.status_code != 403, resp.json()
|
|
|
|
|
|
class TestCrossTenantTitle:
|
|
def test_refresh_title_requires_live_session(self, app_client):
|
|
# Trusted-team model: scope-level auth is the gate; any caller
|
|
# can hit the endpoint. A not-currently-active workstream
|
|
# still 404s because the refresh needs the live session, not
|
|
# because of tenant mismatch.
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-victim", "victim-user")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-victim/refresh-title",
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 404
|
|
assert "not active" in resp.json().get("error", "") or "not found" in resp.json().get(
|
|
"error", ""
|
|
)
|
|
|
|
def test_any_caller_can_set_title(self, app_client):
|
|
# Trusted-team model: title is editable by any authenticated
|
|
# caller; ``user_id`` remains metadata.
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-victim", "victim-user")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-victim/title",
|
|
json={"title": "updated title"},
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 200
|
|
|
|
|
|
class TestCrossTenantOpen:
|
|
def test_any_caller_can_open_persisted(self, app_client):
|
|
# Trusted-team model: open is gated on scope auth, not on row
|
|
# ownership. The persisted ``user_id`` stays as metadata.
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-victim", "victim-user")
|
|
resp = client.post(
|
|
"/v1/api/workstreams/ws-victim/open",
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 200
|
|
|
|
|
|
class TestListWorkstreamsTrustedTeamVisibility:
|
|
"""Listing endpoints (/workstreams, /dashboard, /workstreams/saved)
|
|
return the cluster-wide set to any authenticated caller. Mutations
|
|
are gated independently on the per-workstream handlers — see
|
|
TestCrossTenant{Delete,Approve,Close,Title,Open} for those gates."""
|
|
|
|
def test_list_returns_all_owners(self, app_client):
|
|
client, _mgr = app_client
|
|
resp_a = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "a"},
|
|
headers=_auth("user-a"),
|
|
)
|
|
resp_b = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "b"},
|
|
headers=_auth("user-b"),
|
|
)
|
|
assert resp_a.status_code == 200 and resp_b.status_code == 200
|
|
ws_a, ws_b = resp_a.json()["ws_id"], resp_b.json()["ws_id"]
|
|
|
|
# user-a now sees both.
|
|
resp = client.get("/v1/api/workstreams", headers=_auth("user-a"))
|
|
assert resp.status_code == 200
|
|
# Row key renamed id → ws_id in the Stage 2 list-verb lift.
|
|
ids = {w["ws_id"] for w in resp.json()["workstreams"]}
|
|
assert {ws_a, ws_b}.issubset(ids), ids
|
|
|
|
def test_active_list_row_shape_includes_unified_fields(self, app_client):
|
|
"""Stage 2 list-verb-lift parity regression — interactive
|
|
active-list row carries the always-include fields (ws_id,
|
|
name, state, kind, parent_ws_id, user_id) that the lifted
|
|
``make_list_handler`` produces on every kind. Mirrors the
|
|
coord-side ``test_active_list_row_shape_includes_unified_fields``
|
|
in ``test_coordinator_endpoints.py`` so a future regression
|
|
that drops a field on either branch is caught."""
|
|
client, _mgr = app_client
|
|
create_resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "shape-check"},
|
|
headers=_auth("user-shape"),
|
|
)
|
|
assert create_resp.status_code == 200
|
|
ws_id = create_resp.json()["ws_id"]
|
|
|
|
resp = client.get("/v1/api/workstreams", headers=_auth("user-shape"))
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert "workstreams" in body
|
|
rows = [w for w in body["workstreams"] if w["ws_id"] == ws_id]
|
|
assert len(rows) == 1
|
|
row = rows[0]
|
|
# Always-include row shape — interactive populates kind=
|
|
# INTERACTIVE; user_id is post-lift parity (was coord-only).
|
|
assert set(row.keys()) == {
|
|
"ws_id",
|
|
"name",
|
|
"state",
|
|
"kind",
|
|
"parent_ws_id",
|
|
"user_id",
|
|
"project_id",
|
|
"persona",
|
|
}
|
|
assert row["kind"] == "interactive"
|
|
assert row["user_id"] == "user-shape"
|
|
# parent_ws_id is None for top-level interactive workstreams
|
|
# (only coord-spawned children carry it).
|
|
assert row["parent_ws_id"] is None
|
|
|
|
|
|
class TestDashboardTrustedTeamVisibility:
|
|
def test_dashboard_aggregate_includes_all_owners(self, app_client):
|
|
client, _mgr = app_client
|
|
client.post("/v1/api/workstreams/new", json={"name": "a"}, headers=_auth("user-a"))
|
|
client.post("/v1/api/workstreams/new", json={"name": "b"}, headers=_auth("user-b"))
|
|
client.post("/v1/api/workstreams/new", json={"name": "b2"}, headers=_auth("user-b"))
|
|
|
|
resp = client.get("/v1/api/dashboard", headers=_auth("user-b"))
|
|
assert resp.status_code == 200
|
|
data = resp.json()
|
|
# All three workstreams visible regardless of caller identity.
|
|
assert data["aggregate"]["total_count"] == 3
|
|
owners = {w["user_id"] for w in data["workstreams"]}
|
|
assert {"user-a", "user-b"}.issubset(owners)
|
|
|
|
def test_dashboard_pending_approval_details_default_empty(self, app_client):
|
|
"""No pending approval → the list field is explicitly empty on
|
|
the wire so consumers can distinguish "nothing pending" from
|
|
"absent key". Replaces 1.6's singular ``pending_approval_detail``
|
|
null (breaking, 1.7)."""
|
|
client, _mgr = app_client
|
|
client.post("/v1/api/workstreams/new", json={"name": "a"}, headers=_auth("user-a"))
|
|
resp = client.get("/v1/api/dashboard", headers=_auth("user-a"))
|
|
assert resp.status_code == 200
|
|
rows = resp.json()["workstreams"]
|
|
assert len(rows) == 1
|
|
assert "pending_approval_details" in rows[0]
|
|
assert rows[0]["pending_approval_details"] == []
|
|
# The 1.6 singular field is GONE, not null — a consumer still
|
|
# reading it should break loudly, not read None forever.
|
|
assert "pending_approval_detail" not in rows[0]
|
|
|
|
def test_dashboard_pending_approval_details_merge_judge_verdict(self, app_client):
|
|
"""When _pending_approval is set on a ws's UI, /dashboard
|
|
embeds one detail entry per live cycle with merged items +
|
|
judge_verdict so coord live-bulk callers can render inline
|
|
approve/deny buttons."""
|
|
client, mgr = app_client
|
|
client.post("/v1/api/workstreams/new", json={"name": "a"}, headers=_auth("user-a"))
|
|
ws_id = next(iter(mgr.list_all())).id
|
|
ui = mgr.get(ws_id).ui
|
|
ui._pending_approval = {
|
|
"type": "approve_request",
|
|
"cycle_id": "cyc-1",
|
|
"items": [
|
|
{
|
|
"call_id": "c-1",
|
|
"header": "bash",
|
|
"preview": "$ ls",
|
|
"func_name": "bash",
|
|
"approval_label": "bash",
|
|
"needs_approval": True,
|
|
}
|
|
],
|
|
"judge_pending": False,
|
|
}
|
|
ui._llm_verdicts["c-1"] = {
|
|
"recommendation": "deny",
|
|
"risk_level": "crit",
|
|
"confidence": 0.93,
|
|
"tier": "llm",
|
|
}
|
|
resp = client.get("/v1/api/dashboard", headers=_auth("user-a"))
|
|
assert resp.status_code == 200
|
|
row = next(w for w in resp.json()["workstreams"] if w["ws_id"] == ws_id)
|
|
details = row["pending_approval_details"]
|
|
assert len(details) == 1
|
|
detail = details[0]
|
|
assert detail["cycle_id"] == "cyc-1"
|
|
assert detail["call_id"] == "c-1"
|
|
assert detail["judge_pending"] is False
|
|
item = detail["items"][0]
|
|
assert item["func_name"] == "bash"
|
|
assert item["judge_verdict"]["recommendation"] == "deny"
|
|
assert item["judge_verdict"]["risk_level"] == "crit"
|
|
|
|
|
|
class TestSavedWorkstreamsTrustedTeamVisibility:
|
|
"""Listing returns the cluster-wide set across all owners. Resuming
|
|
an owned saved workstream goes through the per-workstream ownership
|
|
gate on /open (see TestCrossTenantOpen); ownerless persisted rows
|
|
are claimable by any authenticated caller via /open, consistent
|
|
with the same trusted-team model."""
|
|
|
|
def _seed(self, client):
|
|
"""Create two workstreams per user, each with a message so they
|
|
land in list_workstreams_with_history (the SQL gates on an
|
|
EXISTS conversation)."""
|
|
from turnstone.core.storage import get_storage
|
|
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "alice-saved", "alice")
|
|
storage.save_message("alice-saved", "user", "alice's plan")
|
|
_register_ws(storage, "bob-saved", "bob")
|
|
storage.save_message("bob-saved", "user", "bob's plan")
|
|
return storage
|
|
|
|
def test_any_caller_sees_all_rows(self, app_client):
|
|
client, _mgr = app_client
|
|
self._seed(client)
|
|
resp = client.get("/v1/api/workstreams/saved", headers=_auth("alice"))
|
|
assert resp.status_code == 200
|
|
ids = {r["ws_id"] for r in resp.json()["workstreams"]}
|
|
assert {"alice-saved", "bob-saved"}.issubset(ids), ids
|
|
|
|
def test_service_scope_sees_all_rows(self, app_client):
|
|
"""Service-scope still works — same set, different auth path."""
|
|
client, _mgr = app_client
|
|
self._seed(client)
|
|
resp = client.get(
|
|
"/v1/api/workstreams/saved",
|
|
headers=_auth("cluster-collector", scopes=frozenset({"read", "service"})),
|
|
)
|
|
assert resp.status_code == 200
|
|
ids = {r["ws_id"] for r in resp.json()["workstreams"]}
|
|
assert {"alice-saved", "bob-saved"}.issubset(ids)
|
|
|
|
def test_enriched_fields_in_response(self, app_client):
|
|
"""Saved-list rows carry the enrichment fields, incl. the
|
|
Python-computed context_ratio (latest usage prompt_tokens / model
|
|
context window)."""
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "rich-ws", "alice")
|
|
storage.save_message("rich-ws", "user", "do a thing")
|
|
storage.save_workstream_config("rich-ws", {"model_alias": "m1", "skill": "news"})
|
|
storage.record_usage_event("ev-rich", ws_id="rich-ws", prompt_tokens=500)
|
|
storage.create_model_definition(
|
|
"def-rich", alias="m1", model="m1-model", context_window=1000
|
|
)
|
|
|
|
resp = client.get("/v1/api/workstreams/saved", headers=_auth("alice"))
|
|
assert resp.status_code == 200
|
|
row = next(r for r in resp.json()["workstreams"] if r["ws_id"] == "rich-ws")
|
|
assert row["model_alias"] == "m1"
|
|
assert row["launch_skill"] == "news"
|
|
assert row["context_tokens"] == 500
|
|
assert row["context_ratio"] == 0.5 # 500 / 1000
|
|
assert row["child_count"] == 0
|
|
|
|
def test_context_ratio_zero_when_window_unknown(self, app_client):
|
|
"""A model_alias absent from model_definitions (e.g. config.toml-only)
|
|
leaves context_window NULL → context_ratio degrades to 0.0 instead of
|
|
erroring on the division."""
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "no-window-ws", "alice")
|
|
storage.save_message("no-window-ws", "user", "hi")
|
|
storage.save_workstream_config("no-window-ws", {"model_alias": "toml-only"})
|
|
storage.record_usage_event("ev-now", ws_id="no-window-ws", prompt_tokens=500)
|
|
|
|
resp = client.get("/v1/api/workstreams/saved", headers=_auth("alice"))
|
|
assert resp.status_code == 200
|
|
row = next(r for r in resp.json()["workstreams"] if r["ws_id"] == "no-window-ws")
|
|
assert row["context_tokens"] == 500
|
|
assert row["context_ratio"] == 0.0
|
|
assert row["model_alias"] == "toml-only"
|
|
|
|
def test_orphan_rows_visible(self, app_client):
|
|
"""Ownerless rows (empty user_id from migrations / startup
|
|
``name="default"``) appear in the cluster-wide listing alongside
|
|
owned rows. /open lets any authenticated caller claim them —
|
|
intentional under the trusted-team model — so the listing isn't
|
|
leaking anything the resume path wouldn't already grant."""
|
|
client, _mgr = app_client
|
|
storage = self._seed(client)
|
|
_register_ws(storage, "orphan-saved", "")
|
|
storage.save_message("orphan-saved", "user", "orphan content")
|
|
resp = client.get(
|
|
"/v1/api/workstreams/saved",
|
|
headers=_auth("alice", scopes=frozenset({"read"})),
|
|
)
|
|
assert resp.status_code == 200
|
|
ids = {r["ws_id"] for r in resp.json()["workstreams"]}
|
|
assert "orphan-saved" in ids
|
|
|
|
def test_coordinator_rows_excluded_even_for_service(self, app_client):
|
|
"""kind filter is orthogonal to the user_id filter — even a
|
|
service caller (cluster-wide) must not see coordinator rows on
|
|
the interactive 'saved workstreams' endpoint."""
|
|
from turnstone.core.storage import get_storage
|
|
from turnstone.core.workstream import WorkstreamKind
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
storage.register_workstream(
|
|
"coord-row",
|
|
node_id="console",
|
|
user_id="alice",
|
|
name="alice-coord",
|
|
kind=WorkstreamKind.COORDINATOR,
|
|
parent_ws_id=None,
|
|
)
|
|
storage.save_message("coord-row", "user", "planning")
|
|
_register_ws(storage, "alice-interactive", "alice")
|
|
storage.save_message("alice-interactive", "user", "interactive")
|
|
|
|
resp = client.get(
|
|
"/v1/api/workstreams/saved",
|
|
headers=_auth("alice", scopes=frozenset({"read", "service"})),
|
|
)
|
|
assert resp.status_code == 200
|
|
ids = {r["ws_id"] for r in resp.json()["workstreams"]}
|
|
assert "alice-interactive" in ids
|
|
assert "coord-row" not in ids
|
|
|
|
|
|
class TestGlobalEventsServiceGate:
|
|
def test_non_service_rejected(self, app_client):
|
|
client, _mgr = app_client
|
|
resp = client.get(
|
|
"/v1/api/events/global",
|
|
headers=_auth("user-a"), # no service scope
|
|
)
|
|
assert resp.status_code == 403
|
|
assert "service" in resp.json()["error"].lower()
|
|
|
|
def test_service_scope_accepted(self, app_client):
|
|
"""Regression for the console-collector 403 footgun: the
|
|
collector's ServiceTokenManager is configured in console/server.py
|
|
with scopes ``{"read", "service"}``. This gate must accept
|
|
exactly that scope set so the collector's SSE subscription
|
|
doesn't silently 403 out (#sev-0). Any future scope renaming
|
|
that would drop ``"service"`` from the node-side check breaks
|
|
this test before it breaks the dashboard.
|
|
|
|
Probe a deliberately-wrong ``expected_node_id`` — the handler
|
|
runs the scope gate first, then the node-identity check. A
|
|
409 response proves we made it past the scope gate (which is
|
|
what this test is asserting), while also avoiding an
|
|
indefinitely-open SSE stream the TestClient would never close.
|
|
"""
|
|
client, _mgr = app_client
|
|
# Exact scope set the collector uses today.
|
|
collector_scopes = frozenset({"read", "service"})
|
|
resp = client.get(
|
|
"/v1/api/events/global?expected_node_id=definitely-wrong-node-id",
|
|
headers=_auth("console-collector", scopes=collector_scopes),
|
|
)
|
|
# 409 = the scope gate passed and we hit the node-identity
|
|
# mismatch branch. Anything else (403 / 500 / 200 stream)
|
|
# is a failure for this contract.
|
|
assert resp.status_code == 409, (
|
|
f"service-scoped token did not reach node-id check: "
|
|
f"{resp.status_code} {resp.text[:120]}"
|
|
)
|
|
|
|
|
|
class TestPerWsSseGate:
|
|
def test_non_owner_rejected(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
_register_ws(storage, "ws-victim", "victim-user")
|
|
resp = client.get(
|
|
"/v1/api/workstreams/ws-victim/events",
|
|
headers=_auth("attacker-user"),
|
|
)
|
|
assert resp.status_code == 404
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Audit events on successful mutations (sec-11)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestAuditEventsOnMutations:
|
|
def test_workstream_created_emits_audit(self, app_client):
|
|
from turnstone.core.storage import get_storage
|
|
|
|
client, _mgr = app_client
|
|
storage = get_storage()
|
|
assert storage is not None
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "auditme"},
|
|
headers=_auth("user-audit"),
|
|
)
|
|
assert resp.status_code == 200
|
|
ws_id = resp.json()["ws_id"]
|
|
events = storage.list_audit_events(action="workstream.created")
|
|
matching = [e for e in events if e["resource_id"] == ws_id]
|
|
assert matching, "audit row absent for newly created workstream"
|
|
detail = json.loads(matching[0]["detail"])
|
|
assert detail["kind"] == "interactive"
|
|
|
|
|
|
class TestInteractiveCancelLifted:
|
|
"""HTTP-level coverage for the post-lift interactive ``cancel``
|
|
handler at ``POST /v1/api/workstreams/{ws_id}/cancel``. The lifted
|
|
``make_cancel_handler`` body is shared with coord. Pre-lift
|
|
``cancel_generation`` was untested at the HTTP layer; coord
|
|
exercised the lifted body via ``test_coordinator_endpoints.py``.
|
|
This class adds the missing interactive-side parity."""
|
|
|
|
def _create_ws(self, client) -> str:
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "cancel-target"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
return resp.json()["ws_id"]
|
|
|
|
def test_cancel_returns_dropped_shape(self, app_client):
|
|
"""Always-include shape: response carries ``dropped`` (the
|
|
forensic snapshot) regardless of whether anything was running."""
|
|
client, _mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
resp = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/cancel",
|
|
json={},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body["status"] == "ok"
|
|
assert "dropped" in body
|
|
assert body["dropped"]["was_running"] is False
|
|
|
|
def test_cancel_force_clears_worker_thread_and_running_flag(self, app_client):
|
|
"""Force-cancel parity with coord: clears ``worker_thread`` AND
|
|
``_worker_running`` so a follow-up send doesn't route through
|
|
``enqueue()`` to the abandoned worker's queue (bug-2 from the
|
|
cancel-lift /review). Mirrors
|
|
``test_cancel_force_flag_abandons_worker_thread_and_emits_stream_end``
|
|
on the coord side."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
# Simulate an in-flight worker the lifted cancel needs to
|
|
# abandon. The fake session's cancel() is a no-op, so the
|
|
# cancel flag side-effect doesn't matter — what matters is
|
|
# the (worker_thread, _worker_running) pair after force-cancel.
|
|
ws._worker_running = True
|
|
ws.worker_thread = threading.Thread(target=lambda: None, daemon=True)
|
|
|
|
resp = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/cancel",
|
|
json={"force": True},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
# Both fields cleared together — invariant from session_worker
|
|
# ("readers gating on either flag see a coherent
|
|
# (worker_thread, _worker_running) pair").
|
|
assert ws.worker_thread is None
|
|
assert ws._worker_running is False
|
|
|
|
def test_cancel_returns_400_when_session_missing(self, app_client):
|
|
"""Parity with coord: a placeholder workstream (session=None)
|
|
gets a 400 ``"No session"`` rather than a silent no-op 200.
|
|
Pre-lift interactive already returned 400 here; the lift
|
|
preserves the behaviour and propagates it to coord."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
ws.session = None # force the build-failed shape
|
|
|
|
resp = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/cancel",
|
|
json={},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 400
|
|
assert resp.json()["error"] == "No session"
|
|
|
|
|
|
from tests._replay_helpers import make_replay_mocks as _make_interactive_replay_mocks # noqa: E402
|
|
|
|
|
|
class TestInteractiveEventsLifted:
|
|
"""Unit + HTTP coverage for the lifted ``events`` SSE handler.
|
|
|
|
Substantive coverage targets the ``_interactive_events_replay``
|
|
callback (the kind-specific initial-replay generator the lifted
|
|
body iterates before the live loop) and the legacy URL shim.
|
|
The live SSE loop itself (``ws_closed`` exit + ``is_disconnected``
|
|
check) is hard to assert against ``TestClient`` because each
|
|
event arrives as a separate ``data:`` line and the stream runs
|
|
forever; the loop is the same shape used by every other lifted
|
|
SSE-shaped path (cancel / close / open / send), so a regression
|
|
in the loop body would surface across many test files. Live-loop
|
|
smoke coverage is a deferred follow-up tracked in
|
|
``1.5.0-stable-handoff.md``'s "Risk flags for the next session"
|
|
section.
|
|
"""
|
|
|
|
def test_events_replay_yields_connected_first(self):
|
|
"""Pre-lift ``events_sse`` yielded a ``connected`` event
|
|
first (model + skip_permissions). The lifted callback
|
|
preserves the order so client SSE handlers that key on
|
|
the connected event for state setup keep working."""
|
|
from turnstone.server import _interactive_events_replay
|
|
|
|
ws, ui, request = _make_interactive_replay_mocks()
|
|
out = list(_interactive_events_replay(ws, ui, request))
|
|
assert out[0]["type"] == "connected"
|
|
assert out[0]["model"] == "gpt-5"
|
|
assert out[0]["model_alias"] == "default"
|
|
assert out[0]["skip_permissions"] is False
|
|
|
|
def test_events_replay_includes_status_only_when_last_usage_present(self):
|
|
"""The ``status`` event populates the per-tab token-usage
|
|
bar on resume. Skipped when ``session._last_usage`` is None
|
|
(a freshly-created workstream that hasn't completed a turn)."""
|
|
from turnstone.server import _interactive_events_replay
|
|
|
|
ws, ui, request = _make_interactive_replay_mocks()
|
|
out = list(_interactive_events_replay(ws, ui, request))
|
|
assert "status" not in {ev["type"] for ev in out}
|
|
|
|
def test_events_replay_yields_pending_approval_then_verdicts(self):
|
|
"""When an approval is pending, the order is approval +
|
|
cached verdicts (so the client renders the prompt and then
|
|
the LLM-judge intent verdicts that fired during it). Pre-lift
|
|
ordering preserved."""
|
|
from turnstone.server import _interactive_events_replay
|
|
|
|
ws, ui, request = _make_interactive_replay_mocks(
|
|
_pending_approval={"type": "approve_request", "items": []},
|
|
_llm_verdicts={"v1": {"verdict_id": "v1", "tier": "judge"}},
|
|
)
|
|
|
|
out = list(_interactive_events_replay(ws, ui, request))
|
|
types = [ev["type"] for ev in out]
|
|
# The approve_request, then the intent_verdict.
|
|
approve_idx = types.index("approve_request")
|
|
verdict_idx = types.index("intent_verdict")
|
|
assert approve_idx < verdict_idx
|
|
|
|
def test_events_replay_skips_when_session_missing(self):
|
|
"""Defensive: a placeholder workstream whose session is
|
|
``None`` (close-then-reopen race) yields an empty replay
|
|
rather than NPE'ing on ``session.model``. The lifted body
|
|
already 409s for missing UI; this guards the rare case
|
|
where UI exists but session was detached."""
|
|
from turnstone.server import _interactive_events_replay
|
|
|
|
ws = MagicMock()
|
|
ws.session = None
|
|
ui = MagicMock()
|
|
request = MagicMock()
|
|
out = list(_interactive_events_replay(ws, ui, request))
|
|
assert out == []
|
|
|
|
def test_events_replay_omits_conversation_history(self):
|
|
"""PR A: conversation history is no longer replayed over SSE.
|
|
The frontend fetches it via ``GET /history`` (REST) on page
|
|
load and re-fetches on ``clear_ui``; the replay must not yield a
|
|
``history`` event (which previously shipped a multi-MB message
|
|
list on every (re)connect)."""
|
|
from turnstone.server import _interactive_events_replay
|
|
|
|
ws, ui, request = _make_interactive_replay_mocks(
|
|
_pending_approval={"type": "approve_request", "items": []},
|
|
)
|
|
out = list(_interactive_events_replay(ws, ui, request))
|
|
assert "history" not in {ev["type"] for ev in out}
|
|
|
|
def test_events_path_keyed_url_resolves_to_404_for_unknown_ws(self, app_client):
|
|
"""``GET /v1/api/workstreams/{ws_id}/events`` returns 404 for an
|
|
unknown ws_id. Pre-1.5 the same intent was tested against
|
|
``GET /api/events?ws_id=...`` via the legacy query-keyed
|
|
adapter; that URL family was removed in 1.5 along with the
|
|
adapter."""
|
|
client, _mgr = app_client
|
|
resp = client.get(
|
|
"/v1/api/workstreams/does-not-exist/events",
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 404
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# POST /v1/api/command — /compact worker dispatch (progress-events fix)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def test_command_backstop_sits_under_console_proxy_timeout() -> None:
|
|
"""The /command ``running`` answer only ever traverses a proxied pane
|
|
while the node-side completion backstop is STRICTLY under the console
|
|
proxy's client timeout — two constants in two processes whose
|
|
inequality used to live in a comment. This is the enforcement: edit
|
|
either constant past the other and this fails before a proxied
|
|
deployment silently loses the degraded answer again."""
|
|
from turnstone.console.server import _PROXY_CLIENT_TIMEOUT_S
|
|
from turnstone.server import _COMMAND_RESPONSE_BACKSTOP_S
|
|
|
|
assert _COMMAND_RESPONSE_BACKSTOP_S < _PROXY_CLIENT_TIMEOUT_S
|
|
|
|
|
|
class TestCompactCommandDispatch:
|
|
"""Manual /compact runs on the workstream's worker slot: the event loop
|
|
stays free to stream the compaction progress events, and a concurrent
|
|
send takes the queue path instead of racing the history swap."""
|
|
|
|
def _create_ws(self, client) -> str:
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "compact-me"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
return resp.json()["ws_id"]
|
|
|
|
def test_compact_dispatches_to_worker(self, app_client):
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
assert resp.json() == {"status": "ok"}
|
|
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
# The compaction runs on the spawned worker — join it (the runner
|
|
# clears _worker_running before the thread exits, so the join
|
|
# subsumes the flag).
|
|
ws.worker_thread.join(timeout=5)
|
|
assert not ws.worker_thread.is_alive()
|
|
assert ws.session.compacts == 1
|
|
assert ws._worker_running is False
|
|
# The slot was classified as a command window (what the /send
|
|
# route's defer keys on; stale after exit is harmless — every
|
|
# reader conjoins _worker_running).
|
|
assert ws.worker_kind == "command"
|
|
# Exit seam: stranded-text backstop flush only — no drain, no
|
|
# answering send (sends during the window defer in the /send
|
|
# route and dispatch as their own workers afterwards).
|
|
assert ws.session.queued_flushes == 1
|
|
assert ws.session.sends == []
|
|
# The worker wrapped the run in busy/idle state transitions.
|
|
assert ws.ui.states == ["thinking", "idle"]
|
|
|
|
# Drain outcomes are polled with tests._helpers.wait_until (flake-
|
|
# hardened final re-check; raises instead of returning False) at
|
|
# timeout=8.0 — the drain thread runs at a 0.25s cadence, so dispatch
|
|
# is asynchronous with respect to the released command worker and the
|
|
# helper's 5s default is too tight for CI descheduling stalls.
|
|
|
|
def _drain_idle(self, ws) -> bool:
|
|
"""True once the pending-send drain retired itself (list empty,
|
|
single-flight slot cleared) — the leaked-task guard for these
|
|
tests."""
|
|
with ws._lock:
|
|
return not ws._pending_sends and ws._pending_drain is None
|
|
|
|
def test_send_during_compact_window_defers_then_dispatches_fully(self, app_client):
|
|
"""A /send during a manual /compact is answered "queued"
|
|
IMMEDIATELY (no parked POST — a 30s-bounded caller like the
|
|
coordinator client or console proxy must never lose a message to
|
|
a multi-minute window) and then runs as an ordinary full-fidelity
|
|
send — never the interjection queue, whose INTERJECTION_CAP_CHARS
|
|
cap silently truncated pasted logs/code and whose cross-user guard
|
|
locked second participants out for the whole compaction."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.json() == {"status": "ok"}
|
|
big = "x" * 5000 # over the interjection cap — must survive intact
|
|
# The response is immediate even though the window is wedged open:
|
|
# this request runs on the same synchronous test client, so a
|
|
# parked POST would deadlock the test rather than pass it.
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": big},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.status_code == 200
|
|
body = r.json()
|
|
assert body["status"] == "queued"
|
|
# The SDK-facing discriminator between "interjection-queued into a
|
|
# live turn" and "parked on the deferred list" (SendResponse).
|
|
assert body["deferred"] is True
|
|
assert body["msg_id"]
|
|
assert body["priority"] == "notice"
|
|
# Registered, not dispatched: the window is still open.
|
|
assert ws.session.sends == []
|
|
assert ws.session.queue_calls == [] # the queue is unreachable
|
|
with ws._lock:
|
|
assert len(ws._pending_sends) == 1
|
|
gate.set() # compaction finishes; the drain dispatches
|
|
wait_until(lambda: ws.session.sends, timeout=8.0)
|
|
# Full fidelity: the exact 5000-char text, via a normal send (an
|
|
# oversized entry must take the fresh-spawn arm — the interjection
|
|
# fallback would truncate it).
|
|
assert [s[0] for s in ws.session.sends] == [big]
|
|
assert ws.session.queue_calls == []
|
|
# The dispatched send carries the deferred msg_id as its send_id,
|
|
# so the client's queued bubble reconciles against the turn.
|
|
assert ws.session.sends[0][2] == body["msg_id"]
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_cancelled_compact_flushes_queue_without_answering(self, app_client):
|
|
"""A user-stopped compaction must not auto-run a turn they may no
|
|
longer want — queued text lands in the transcript via the flush
|
|
drain instead (the cancel-seam precedent)."""
|
|
from turnstone.core.session import GenerationCancelled
|
|
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
ws.session.compact_raises = GenerationCancelled()
|
|
ws.session.queued_text = "queued mid-compact"
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
ws.worker_thread.join(timeout=5)
|
|
assert ws.session.compacts == 1
|
|
assert ws.session.queued_flushes == 1
|
|
assert ws.session.sends == []
|
|
assert ws.ui.states == ["thinking", "idle"]
|
|
|
|
def test_compact_refused_while_worker_running(self, app_client):
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
with ws._lock:
|
|
ws._worker_running = True # a turn is in flight
|
|
try:
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
finally:
|
|
with ws._lock:
|
|
ws._worker_running = False
|
|
assert resp.status_code == 409 # loud refusal for status-code-only callers
|
|
body = resp.json()
|
|
assert body["status"] == "busy"
|
|
assert "busy" in body["error"].lower()
|
|
# Never ran inline, never queued a phantom compaction.
|
|
assert ws.session.compacts == 0
|
|
assert ws.session.commands == []
|
|
|
|
def test_non_compact_commands_complete_before_response(self, app_client):
|
|
"""Quick commands dispatch through the same worker slot (mutual
|
|
exclusion vs sends / a running compaction / each other) but the
|
|
endpoint awaits completion, preserving the synchronous contract:
|
|
handle_command has finished by the time the response returns."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/skill", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
assert resp.json() == {"status": "ok"}
|
|
ws = mgr.get(ws_id)
|
|
# handle_command completed before the response (the done-Event
|
|
# gates it); join before asserting the slot state (the runner
|
|
# clears _worker_running before the thread exits, so the join
|
|
# subsumes the flag; without it this assert races the worker's
|
|
# last steps).
|
|
assert ws.session.commands == ["/skill"]
|
|
ws.worker_thread.join(timeout=5)
|
|
assert not ws.worker_thread.is_alive()
|
|
assert ws._worker_running is False
|
|
# No exit seam work at all for quick commands: no drain, no flush,
|
|
# no follow-up send, and no state chatter (they never left idle).
|
|
assert ws.session.queued_flushes == 0
|
|
assert ws.session.sends == []
|
|
assert ws.ui.states == []
|
|
|
|
def test_deferred_send_delivers_attachments_after_release(self, app_client, monkeypatch):
|
|
"""Attachments are peek-resolved BEFORE the defer (the bytes ride
|
|
the pending entry), so a send deferred through a long compaction
|
|
still delivers them on dispatch — the queue-path refusal
|
|
(attachments_busy) must never apply to a command window."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
|
|
def fake_resolve(requested_ids, _ws_id, _user_id):
|
|
assert list(requested_ids) == ["a1"]
|
|
return (["fake-attachment-bytes"], ["a1"], [])
|
|
|
|
monkeypatch.setattr("turnstone.core.attachments.resolve_staged_attachments", fake_resolve)
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "with attachment", "attachment_ids": ["a1"]},
|
|
headers=_auth("user-1"),
|
|
)
|
|
body = r.json()
|
|
assert body["status"] == "queued"
|
|
# The attachments are TAKEN by the deferred entry (the client
|
|
# consumes its chips off this list — a retract discards them).
|
|
assert body["attached_ids"] == ["a1"]
|
|
assert ws.session.sends == [] # still deferred
|
|
gate.set()
|
|
wait_until(lambda: ws.session.sends, timeout=8.0)
|
|
text, attachments, sid = ws.session.sends[0]
|
|
assert text == "with attachment"
|
|
assert attachments == ["fake-attachment-bytes"]
|
|
assert sid == body["msg_id"]
|
|
assert ws.session.queue_calls == []
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_send_during_quick_command_window_defers(self, app_client):
|
|
"""The defer applies to EVERY command window, not just /compact —
|
|
a send racing a quick command is answered "queued" and dispatches
|
|
after the window with full fidelity instead of entering the
|
|
interjection queue."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.command_gate = gate
|
|
cmd_result: dict = {}
|
|
|
|
def _cmd() -> None:
|
|
r = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/skill", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
cmd_result["body"] = r.json()
|
|
|
|
runner = threading.Thread(target=_cmd, daemon=True)
|
|
runner.start()
|
|
# Wait until the command worker actually holds the slot.
|
|
wait_until(lambda: ws._worker_running, timeout=8.0)
|
|
assert ws.worker_kind == "command"
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "mid-command send"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
send_body = r.json()
|
|
assert send_body["status"] == "queued"
|
|
assert ws.session.queue_calls == []
|
|
assert ws.session.sends == []
|
|
gate.set()
|
|
runner.join(timeout=10)
|
|
assert cmd_result["body"] == {"status": "ok"}
|
|
wait_until(lambda: ws.session.sends, timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == ["mid-command send"]
|
|
assert ws.session.queue_calls == []
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_clear_ui_rides_the_worker_not_the_endpoint(self, app_client):
|
|
"""The clear_ui follow-up runs on the worker after handle_command —
|
|
parked after the endpoint's 60s done-wait it was silently skipped
|
|
for any command that outlived the backstop, leaving every pane
|
|
rendering a transcript the server no longer holds."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/clear", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
ws = mgr.get(ws_id)
|
|
ws.worker_thread.join(timeout=5)
|
|
assert {"type": "clear_ui"} in ws.ui._enqueued
|
|
|
|
def _run_abandoned_command(self, client, mgr, ws_id, command="/clear", raises=None):
|
|
"""Dispatch a gated command, force-abandon its worker mid-run, then
|
|
release the gate so the abandoned thread finishes late. Returns
|
|
(endpoint response body, the abandoned worker thread)."""
|
|
ws = mgr.get(ws_id)
|
|
gate = threading.Event()
|
|
ws.session.command_gate = gate
|
|
ws.session.command_raises = raises
|
|
result: dict = {}
|
|
|
|
def _cmd() -> None:
|
|
r = client.post(
|
|
"/v1/api/command",
|
|
json={"command": command, "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
result["body"] = r.json()
|
|
|
|
runner = threading.Thread(target=_cmd, daemon=True)
|
|
runner.start()
|
|
wait_until(lambda: ws._worker_running, timeout=8.0)
|
|
worker = ws.worker_thread
|
|
# The cancel handler's force path shape: abandon the worker.
|
|
with ws._lock:
|
|
ws.worker_thread = None
|
|
ws._worker_running = False
|
|
gate.set() # the wedged command unwedges LATE
|
|
worker.join(timeout=5)
|
|
runner.join(timeout=10)
|
|
assert not worker.is_alive() and not runner.is_alive()
|
|
return result["body"], worker
|
|
|
|
def test_abandoned_command_worker_fires_no_followups(self, app_client):
|
|
"""A force-cancelled wedged command that unwedges minutes later
|
|
must not fire clear_ui (every pane would wipe its transcript
|
|
mid-successor-turn) — the owner guard every sibling worker closure
|
|
applies. done still fires, so the endpoint never hangs."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
body, _ = self._run_abandoned_command(client, mgr, ws_id, command="/clear")
|
|
assert not any(e.get("type") == "clear_ui" for e in ws.ui._enqueued)
|
|
# handle_command genuinely completed, so the (unhung) endpoint's
|
|
# answer reflects that.
|
|
assert body["status"] == "ok"
|
|
|
|
def test_abandoned_command_worker_swallows_late_error(self, app_client):
|
|
"""The except arm carries the same guard: a stray late 'Command
|
|
error:' from an abandoned worker must not land mid-successor-turn."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
self._run_abandoned_command(
|
|
client, mgr, ws_id, command="/skill", raises=RuntimeError("late boom")
|
|
)
|
|
assert ws.ui.errors == []
|
|
|
|
def test_abandoned_init_worker_still_fires_notify(self, app_client, monkeypatch):
|
|
"""The completion notify is WORKSTREAM-scoped, not slot-scoped:
|
|
_fire_notify_targets has exactly one call site (the init worker),
|
|
successor turns never notify, so an owner gate here has no
|
|
duplicate to prevent — it only converts force-cancel into
|
|
permanent notification loss for scheduled/unattended workstreams.
|
|
This pins the un-gated behavior against the tempting symmetry
|
|
'fix' (the finally's siblings ARE owner-gated, correctly — they
|
|
mutate live slot/UI state)."""
|
|
client, mgr = app_client
|
|
fired: list = []
|
|
monkeypatch.setattr(
|
|
"turnstone.server._fire_notify_targets",
|
|
lambda ws, content: fired.append(ws.id),
|
|
)
|
|
gate = threading.Event()
|
|
started = threading.Event()
|
|
orig_send = _FakeSession.send
|
|
|
|
def wedged_send(self, text, **kwargs):
|
|
started.set()
|
|
assert gate.wait(timeout=10)
|
|
orig_send(self, text, **kwargs)
|
|
|
|
monkeypatch.setattr(_FakeSession, "send", wedged_send)
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "sched", "initial_message": "do the thing"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
ws_id = resp.json()["ws_id"]
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
assert started.wait(timeout=10)
|
|
worker = ws.worker_thread
|
|
with ws._lock: # the force-cancel abandon shape
|
|
ws.worker_thread = None
|
|
ws._worker_running = False
|
|
gate.set() # the abandoned init completes LATE
|
|
worker.join(timeout=10)
|
|
assert not worker.is_alive()
|
|
assert fired == [ws_id] # exactly once, despite the abandonment
|
|
|
|
def test_compact_on_error_workstream_restores_error_badge(self, app_client):
|
|
"""/compact on an ERROR workstream must exit back to 'error', not
|
|
stamp 'idle' over the operator's investigatable badge — the
|
|
compaction neither retried nor resolved the failed turn (mirrors
|
|
the orphan reaper's ERROR carve-out). The existing happy-path
|
|
test pins ['thinking', 'idle'] for an IDLE workstream."""
|
|
from turnstone.core.workstream import WorkstreamState
|
|
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
ws.state = WorkstreamState.ERROR # a prior fatal turn's badge
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.json() == {"status": "ok"}
|
|
ws.worker_thread.join(timeout=5)
|
|
assert not ws.worker_thread.is_alive()
|
|
# THINKING during the window is honest (work is happening); the
|
|
# exit restores the badge instead of clearing it.
|
|
assert ws.ui.states == ["thinking", "error"]
|
|
|
|
def test_retracted_deferred_send_never_dispatches(self, app_client):
|
|
"""DELETE /send with a deferred msg_id retracts the entry before
|
|
dispatch — the drain must skip it entirely. The route falls
|
|
through the session's interjection dequeue (which misses) to the
|
|
workstream's pending list."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "changed my mind"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
msg_id = r.json()["msg_id"]
|
|
d = client.request(
|
|
"DELETE",
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"msg_id": msg_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert d.json() == {"status": "removed"}
|
|
assert ws.session.dequeues == [msg_id] # fall-through was exercised
|
|
gate.set()
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
assert ws.session.sends == []
|
|
assert ws.session.queue_calls == []
|
|
# Unknown ids still answer not_found after checking both holders.
|
|
d2 = client.request(
|
|
"DELETE",
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"msg_id": "nope"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert d2.json() == {"status": "not_found"}
|
|
|
|
def test_deferred_sends_dispatch_in_arrival_order(self, app_client):
|
|
"""Two sends deferred behind one window dispatch in arrival order,
|
|
each as its own full-fidelity turn (the fake's sends complete
|
|
instantly, so the second never needs the interjection fallback)."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
first = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "first"},
|
|
headers=_auth("user-1"),
|
|
).json()
|
|
second = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "second"},
|
|
headers=_auth("user-1"),
|
|
).json()
|
|
assert first["status"] == second["status"] == "queued"
|
|
assert first["msg_id"] != second["msg_id"]
|
|
gate.set()
|
|
wait_until(lambda: len(ws.session.sends) == 2, timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == ["first", "second"]
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_force_cancel_of_wedged_command_releases_drain(self, app_client):
|
|
"""Force-cancelling a wedged command clears the slot flags — the
|
|
drain polls the same (_worker_running, worker_kind) pair, so a
|
|
deferred message dispatches WITHOUT waiting for the abandoned
|
|
thread to unwedge (the operator's escape hatch keeps working
|
|
under defer exactly as it did under park)."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
zombie = ws.worker_thread
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "deferred behind wedge"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.json()["status"] == "queued"
|
|
resp = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/cancel",
|
|
json={"force": True},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
# Dispatch happens while the zombie is STILL wedged on the gate.
|
|
wait_until(lambda: ws.session.sends, timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == ["deferred behind wedge"]
|
|
gate.set() # let the abandoned thread finish and be joined
|
|
zombie.join(timeout=5)
|
|
assert not zombie.is_alive()
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_force_cancel_runs_abandon_machinery_before_idle_emission(self, app_client):
|
|
"""The force branch runs the session's abandon machinery BEFORE
|
|
clearing ownership and emitting idle. Subscribers on the IDLE
|
|
fan-out (the coordinator idle observer's operator-Stop gate,
|
|
the wake watcher's queue read) use the latch and the demote to
|
|
tell an operator-forced IDLE from a turn reaching idle under
|
|
its own power — and the thread force-cancel abandons is stuck
|
|
by definition, so its own exception handler cannot be relied
|
|
on to have run first."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
zombie = ws.worker_thread
|
|
calls: list[str] = []
|
|
ws.session._drain_pending_advisories = lambda: calls.append("abandon")
|
|
ws.ui.on_state_change = lambda s: calls.append(f"state:{s}")
|
|
resp = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/cancel",
|
|
json={"force": True},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
assert calls == ["abandon", "state:idle"]
|
|
gate.set()
|
|
zombie.join(timeout=5)
|
|
assert not zombie.is_alive()
|
|
|
|
def test_ws_close_mid_window_drops_pending_and_drain_exits(self, app_client):
|
|
"""A workstream closed with deferred sends outstanding drops them
|
|
(documented at-most-once contract) and the drain task retires
|
|
itself — no leaked task, no dispatch into a closed workstream."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "never delivered"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.json()["status"] == "queued"
|
|
with ws._lock:
|
|
ws._closed = True # the SessionManager.close tombstone shape
|
|
gate.set()
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
assert ws.session.sends == []
|
|
assert ws.session.queue_calls == []
|
|
|
|
def test_deferred_send_dispatches_into_post_swap_session(self, app_client):
|
|
"""The drain re-captures ws.session per attempt: a /resume-style
|
|
identity swap during the window routes the deferred message into
|
|
the POST-swap session — the same guarantee the park's
|
|
per-iteration re-capture provided."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
old_session = ws.session
|
|
gate = threading.Event()
|
|
old_session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "post-swap please"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.json()["status"] == "queued"
|
|
new_session = _FakeSession(ws_id=ws_id, user_id="user-1")
|
|
ws.session = new_session # the in-place identity swap
|
|
gate.set()
|
|
wait_until(lambda: new_session.sends, timeout=8.0)
|
|
assert [s[0] for s in new_session.sends] == ["post-swap please"]
|
|
assert old_session.sends == []
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_rejected_deferred_entry_waits_for_slot_then_fresh_spawns(self, app_client):
|
|
"""The drain's rejection arm: an entry the interjection fallback
|
|
refuses (cross-user here; attachments/oversized take the same
|
|
path) is NOT dropped — it waits for the slot to free and then
|
|
dispatches as its own fresh turn. The queued ack must never
|
|
become a silent drop while the workstream lives."""
|
|
from turnstone.core.session import CrossUserInterjectionError
|
|
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
send_gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
ws.session.send_gate = send_gate
|
|
ws.session.queue_raises = CrossUserInterjectionError("another participant's turn")
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
for msg in ("first", "second"):
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": msg},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.json()["status"] == "queued"
|
|
gate.set()
|
|
# Entry 1 spawns fresh and WEDGES in send() on send_gate — a live
|
|
# turn now holds the slot, so entry 2 takes the interjection
|
|
# fallback and is refused (the fake raises cross-user).
|
|
wait_until(lambda: ws.session.queue_calls == ["second"], timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == ["first"]
|
|
send_gate.set() # the turn ends; the slot frees
|
|
wait_until(lambda: len(ws.session.sends) == 2, timeout=8.0)
|
|
# Entry 2 dispatched as its own fresh turn — exactly one refused
|
|
# queue attempt, then the fresh-spawn arm.
|
|
assert [s[0] for s in ws.session.sends] == ["first", "second"]
|
|
assert ws.session.queue_calls == ["second"]
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_drain_crash_requeues_entry_and_retries(self, app_client):
|
|
"""A crash inside a claimed entry's dispatch attempt (the real-world
|
|
shape: Thread.start raising under thread exhaustion) must not eat
|
|
the acknowledged message — the per-iteration handler re-inserts it
|
|
at head and retries after a backoff, keeping the docstring's
|
|
no-silent-drop contract."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "survive the crash"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.json()["status"] == "queued"
|
|
with ws._lock:
|
|
entry = ws._pending_sends[0]
|
|
real_attempt = entry.attempt
|
|
calls = {"n": 0}
|
|
|
|
def crash_once(session):
|
|
calls["n"] += 1
|
|
if calls["n"] == 1:
|
|
raise RuntimeError("simulated Thread.start failure")
|
|
return real_attempt(session)
|
|
|
|
entry.attempt = crash_once
|
|
gate.set()
|
|
wait_until(lambda: ws.session.sends, timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == ["survive the crash"]
|
|
assert calls["n"] == 2 # crashed once, re-queued, dispatched on retry
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_persistently_crashing_entry_is_never_dropped_and_retract_frees_drain(self, app_client):
|
|
"""A persistently crashing attempt keeps the entry alive (retry
|
|
loop, never a silent drop) — and the user's DELETE still works
|
|
mid-crash-loop: the re-inserted entry is marked retracted, the
|
|
loop-top purge drops it, and the drain retires cleanly."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "doomed"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
msg_id = r.json()["msg_id"]
|
|
with ws._lock:
|
|
entry = ws._pending_sends[0]
|
|
calls = {"n": 0}
|
|
|
|
def always_crash(_session):
|
|
calls["n"] += 1
|
|
raise RuntimeError("persistent dispatch failure")
|
|
|
|
entry.attempt = always_crash
|
|
gate.set()
|
|
# Two attempts across the ~1s backoff prove the loop survived the
|
|
# first crash with the entry intact (a dropped entry can't be
|
|
# re-attempted) and nothing was dispatched behind the user's back.
|
|
wait_until(lambda: calls["n"] >= 2, timeout=8.0)
|
|
assert ws.session.sends == []
|
|
# ``calls["n"]`` increments as the FIRST statement of the attempt, so
|
|
# the counter crosses 2 while the drain still holds the entry CLAIMED
|
|
# (popped, dispatch in flight). Retract only scans ``_pending_sends``
|
|
# and correctly answers not_found for a claimed entry, so firing the
|
|
# DELETE off the counter alone races the crash path's re-insert and
|
|
# loses whenever the re-insert is slower than the request -- which it
|
|
# is under pytest, where the log handlers make the intervening
|
|
# ``log.exception`` roughly as expensive as wait_until's poll interval.
|
|
# Wait for the state this test is actually about: the entry BACK on
|
|
# the list, mid-crash-loop, retractable.
|
|
wait_until(
|
|
lambda: any(e.msg_id == msg_id for e in list(ws._pending_sends)),
|
|
timeout=8.0,
|
|
)
|
|
d = client.request(
|
|
"DELETE",
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"msg_id": msg_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert d.json() == {"status": "removed"}
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
assert ws.session.sends == []
|
|
|
|
def test_fresh_send_defers_behind_claimed_entry(self, app_client):
|
|
"""The order barrier's drain-alive term: a send arriving while the
|
|
drain has CLAIMED an entry (off the list, dispatch in flight) must
|
|
defer behind it, not overtake — the pre-fix route dispatched
|
|
immediately (the list looked empty) and inverted answer order
|
|
against the acknowledged send."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r1 = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "first"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r1.json()["deferred"] is True
|
|
with ws._lock:
|
|
entry = ws._pending_sends[0]
|
|
real_attempt = entry.attempt
|
|
claimed = threading.Event()
|
|
release = threading.Event()
|
|
|
|
def gated_attempt(session):
|
|
claimed.set()
|
|
assert release.wait(timeout=10)
|
|
return real_attempt(session)
|
|
|
|
entry.attempt = gated_attempt
|
|
gate.set() # window closes; the drain claims "first" and blocks
|
|
wait_until(claimed.is_set, timeout=8.0)
|
|
with ws._lock:
|
|
assert ws._pending_sends == [] # claimed — the list alone says "free"
|
|
big = "y" * 5000 # full-fidelity: can't fold into any live turn
|
|
r2 = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": big},
|
|
headers=_auth("user-1"),
|
|
)
|
|
# Deferred via the barrier (drain alive), never dispatched directly.
|
|
assert r2.json()["status"] == "queued"
|
|
assert r2.json()["deferred"] is True
|
|
assert ws.session.sends == []
|
|
release.set()
|
|
wait_until(lambda: len(ws.session.sends) == 2, timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == ["first", big]
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_drain_exit_re_arms_wake_gate_after_pure_retraction(self, app_client, monkeypatch):
|
|
"""A pending list that empties by RETRACTION never runs a deferred
|
|
turn, so no worker-exit backstop would re-run the wake gate that
|
|
yielded to the barrier — the drain's clean exit must re-arm it
|
|
itself (trigger="drain-exit"), or a nudge parked during the window
|
|
strands until some unrelated future dispatch."""
|
|
from turnstone.core import idle_nudge_watcher
|
|
|
|
wake_calls: list[str] = []
|
|
monkeypatch.setattr(
|
|
idle_nudge_watcher,
|
|
"wake_workstream_if_pending",
|
|
lambda ws, *, trigger="unspecified": (wake_calls.append(trigger), False)[1],
|
|
)
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "will be retracted"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
d = client.request(
|
|
"DELETE",
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"msg_id": r.json()["msg_id"]},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert d.json() == {"status": "removed"}
|
|
gate.set()
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
wait_until(lambda: "drain-exit" in wake_calls, timeout=8.0)
|
|
assert ws.session.sends == []
|
|
|
|
def test_deferred_dispatch_emits_message_dispatched(self, app_client):
|
|
"""Fresh-spawn arm of the settle protocol: dispatching a deferred
|
|
entry emits pane-tier ``message_dispatched {msg_id}`` (no
|
|
``folded`` key) so the queued chip promotes exactly when the
|
|
retraction window truly closes."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "settle me"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
msg_id = r.json()["msg_id"]
|
|
queued_events = [e for e in ws.ui._enqueued if e.get("type") == "message_queued"]
|
|
assert [e["msg_id"] for e in queued_events] == [msg_id]
|
|
gate.set()
|
|
wait_until(
|
|
lambda: any(e.get("type") == "message_dispatched" for e in ws.ui._enqueued),
|
|
timeout=8.0,
|
|
)
|
|
dispatched = [e for e in ws.ui._enqueued if e.get("type") == "message_dispatched"]
|
|
assert dispatched == [{"type": "message_dispatched", "msg_id": msg_id}]
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_folded_deferred_dispatch_emits_folded_flag(self, app_client):
|
|
"""Interjection fold-in arm: a deferred entry that dispatches INTO a
|
|
live turn's queue emits ``folded: true`` — the client must clear
|
|
only its deferred flag (DELETE still genuinely removes the message
|
|
from the interjection queue until the seam drains), not promote."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
send_gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
ws.session.send_gate = send_gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
first = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "first"},
|
|
headers=_auth("user-1"),
|
|
).json()
|
|
second = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "second"},
|
|
headers=_auth("user-1"),
|
|
).json()
|
|
gate.set()
|
|
# "first" spawns fresh and wedges in send() — a live turn holds the
|
|
# slot, so "second" (small, no attachments) folds into its
|
|
# interjection queue.
|
|
wait_until(lambda: ws.session.queue_calls == ["second"], timeout=8.0)
|
|
send_gate.set()
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
dispatched = [e for e in ws.ui._enqueued if e.get("type") == "message_dispatched"]
|
|
assert dispatched == [
|
|
{"type": "message_dispatched", "msg_id": first["msg_id"]},
|
|
{"type": "message_dispatched", "msg_id": second["msg_id"], "folded": True},
|
|
]
|
|
assert [s[0] for s in ws.session.sends] == ["first"]
|
|
|
|
def test_deferred_send_list_saturation_returns_queue_full(self, app_client):
|
|
"""The deferred list is bounded (workstream.PENDING_SENDS_MAX — the
|
|
shared constant ChatSession._QUEUE_MAX aliases): the 11th pending
|
|
send answers queue_full instead of silently pinning message +
|
|
attachment bytes per entry for a whole command window and then
|
|
costing one unattended turn each."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
for i in range(10):
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": f"m{i}"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.json()["status"] == "queued", i
|
|
r11 = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "one too many"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
body = r11.json()
|
|
assert body["status"] == "queue_full"
|
|
assert "msg_id" not in body # never acknowledged
|
|
with ws._lock:
|
|
assert len(ws._pending_sends) == 10
|
|
gate.set()
|
|
wait_until(lambda: len(ws.session.sends) == 10, timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == [f"m{i}" for i in range(10)]
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_drain_spawn_failure_rolls_back_and_answers_queue_full(self, app_client, monkeypatch):
|
|
"""Thread.start failing at the drain-spawn site must not leave a
|
|
phantom acknowledged-nowhere entry (the 500-after-registration
|
|
class): entry and slot roll back under the same lock, the client
|
|
gets the retryable queue_full, and a retry once resources return
|
|
dispatches exactly once — no duplicate turns."""
|
|
from turnstone.core import session_routes
|
|
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
|
|
class _ExhaustedThread(threading.Thread):
|
|
def start(self) -> None:
|
|
raise RuntimeError("can't start new thread")
|
|
|
|
real_factory = session_routes._make_drain_thread
|
|
monkeypatch.setattr(
|
|
session_routes,
|
|
"_make_drain_thread",
|
|
lambda _ws: _ExhaustedThread(target=lambda: None, daemon=True),
|
|
)
|
|
r = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "refused under exhaustion"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r.json()["status"] == "queue_full"
|
|
with ws._lock:
|
|
assert ws._pending_sends == [] # rolled back — no phantom
|
|
assert ws._pending_drain is None
|
|
# Resources recover: a plain retry defers and dispatches ONCE.
|
|
monkeypatch.setattr(session_routes, "_make_drain_thread", real_factory)
|
|
r2 = client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "the retry"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert r2.json()["status"] == "queued"
|
|
assert r2.json()["deferred"] is True
|
|
gate.set()
|
|
wait_until(lambda: ws.session.sends, timeout=8.0)
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
assert [s[0] for s in ws.session.sends] == ["the retry"]
|
|
|
|
def test_command_spawn_failure_answers_503_not_ok(self, app_client, monkeypatch):
|
|
"""A command worker that never spawned must not answer 200 ok — the
|
|
endpoint's generic catch-all used to swallow the dispatcher's
|
|
Thread.start re-raise, telling SDK callers their /clear ran while
|
|
the context stayed un-cleared."""
|
|
from turnstone.core import session_worker
|
|
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
|
|
def _exhausted(*_a, **_kw):
|
|
raise RuntimeError("can't start new thread")
|
|
|
|
monkeypatch.setattr(session_worker, "send", _exhausted)
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/clear", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 503
|
|
assert resp.json()["status"] == "error"
|
|
assert ws.session.commands == [] # the command never ran
|
|
|
|
def test_drain_exit_wake_raise_is_contained(self, app_client, monkeypatch, caplog):
|
|
"""A raising drain-exit wake (session_worker.send re-raises
|
|
Thread.start failures) must stay inside the wake's own guard —
|
|
it runs AFTER the drain retired its slot, so reaching the
|
|
last-resort handler would clear a slot this thread no longer
|
|
owns (a successor drain's live registration)."""
|
|
import logging
|
|
|
|
from turnstone.core import idle_nudge_watcher
|
|
|
|
def _raising_gate(ws, *, trigger="unspecified"):
|
|
if trigger == "drain-exit":
|
|
raise RuntimeError("can't start new thread")
|
|
return False
|
|
|
|
monkeypatch.setattr(idle_nudge_watcher, "wake_workstream_if_pending", _raising_gate)
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "dispatch me"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
with caplog.at_level(logging.WARNING, logger="turnstone.core.session_routes"):
|
|
gate.set()
|
|
wait_until(lambda: ws.session.sends, timeout=8.0)
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
wait_until(
|
|
lambda: any("drain_exit_wake_failed" in r.message for r in caplog.records),
|
|
timeout=8.0,
|
|
)
|
|
# The raise was contained by the wake's own guard — never the
|
|
# last-resort handler (whose log would be a false failure for a
|
|
# drain that exited cleanly).
|
|
assert not any("pending_drain_failed" in r.message for r in caplog.records)
|
|
# The seam stays fully serviceable: another defer cycle works.
|
|
gate2 = threading.Event()
|
|
ws.session.compact_gate = gate2
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "and me"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
gate2.set()
|
|
wait_until(lambda: len(ws.session.sends) == 2, timeout=8.0)
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
|
|
def test_closed_drain_exit_runs_no_wake(self, app_client, monkeypatch):
|
|
"""The drain-exit wake backstop belongs to the CLEAN exit only: a
|
|
closed workstream is torn down, and firing the gate there would
|
|
be work on a corpse. (The clean-exit case is pinned by
|
|
test_drain_exit_re_arms_wake_gate_after_pure_retraction.)"""
|
|
from turnstone.core import idle_nudge_watcher
|
|
|
|
triggers: list[str] = []
|
|
monkeypatch.setattr(
|
|
idle_nudge_watcher,
|
|
"wake_workstream_if_pending",
|
|
lambda ws, *, trigger="unspecified": (triggers.append(trigger), False)[1],
|
|
)
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "dies with the ws"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
with ws._lock:
|
|
ws._closed = True # tombstone, as SessionManager.close sets it
|
|
gate.set()
|
|
wait_until(lambda: self._drain_idle(ws), timeout=8.0)
|
|
assert "drain-exit" not in triggers
|
|
assert ws.session.sends == []
|
|
|
|
def test_drain_last_resort_clear_is_identity_guarded(self, app_client, monkeypatch, caplog):
|
|
"""Drive the drain's OUTER except (loop machinery failing — here
|
|
the closed-arm log call) and prove the last-resort slot-clear
|
|
executes its identity guard: the slot is this drain's own, so it
|
|
clears; a successor's registration would be left alone (the
|
|
guard compares thread identity, the sibling exit-seam pattern)."""
|
|
import logging
|
|
|
|
from turnstone.core import session_routes
|
|
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
gate = threading.Event()
|
|
ws.session.compact_gate = gate
|
|
client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/compact", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
client.post(
|
|
f"/v1/api/workstreams/{ws_id}/send",
|
|
json={"message": "stranded by the crash"},
|
|
headers=_auth("user-1"),
|
|
)
|
|
|
|
def _broken_warning(*_a, **_kw):
|
|
raise RuntimeError("log backend down")
|
|
|
|
monkeypatch.setattr(session_routes.log, "warning", _broken_warning)
|
|
with ws._lock:
|
|
ws._closed = True # closed arm with a pending entry → log.warning
|
|
with caplog.at_level(logging.ERROR, logger="turnstone.core.session_routes"):
|
|
gate.set()
|
|
wait_until(
|
|
lambda: any("pending_drain_failed" in r.message for r in caplog.records),
|
|
timeout=8.0,
|
|
)
|
|
# The identity-guarded clear ran for the drain's OWN slot.
|
|
wait_until(lambda: ws._pending_drain is None, timeout=8.0)
|
|
assert ws.session.sends == []
|
|
|
|
def test_exit_command_emits_ended_info_and_never_answers(self, app_client):
|
|
"""should_exit commands shut the session down — the worker emits
|
|
the ended notice and never launches an answering turn (sends
|
|
during the window defer in the /send route and belong to whatever
|
|
follows the shutdown, not to this worker)."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
ws.session.exit_commands = {"/exit"}
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/exit", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
assert resp.status_code == 200
|
|
ws.worker_thread.join(timeout=5)
|
|
assert ws.session.sends == []
|
|
assert any("Session ended" in m for m in ws.ui.infos)
|
|
|
|
def test_non_compact_command_refused_while_worker_running(self, app_client):
|
|
"""The old inline path was serialized by the event loop itself; the
|
|
worker-slot dispatch restores that mutual exclusion with an explicit
|
|
busy answer — /clear can no longer interleave with a live turn or a
|
|
running compaction."""
|
|
client, mgr = app_client
|
|
ws_id = self._create_ws(client)
|
|
ws = mgr.get(ws_id)
|
|
assert ws is not None
|
|
with ws._lock:
|
|
ws._worker_running = True # a turn / compaction is in flight
|
|
try:
|
|
resp = client.post(
|
|
"/v1/api/command",
|
|
json={"command": "/clear", "ws_id": ws_id},
|
|
headers=_auth("user-1"),
|
|
)
|
|
finally:
|
|
with ws._lock:
|
|
ws._worker_running = False
|
|
assert resp.status_code == 409
|
|
assert resp.json()["status"] == "busy"
|
|
assert ws.session.commands == []
|
|
|
|
|
|
class TestRequireProjectMountWiring:
|
|
"""server.require_project is wired on the REAL interactive create mount
|
|
(create_gate_require_project=True on interactive_endpoint_config). Synthetic-
|
|
cfg unit tests can't catch a mis-wire on the actual mount, so drive the
|
|
mounted endpoint end-to-end."""
|
|
|
|
def test_projectless_interactive_create_gated_when_on(self, app_client, make_config_store):
|
|
client, _mgr = app_client
|
|
client.app.state.config_store = make_config_store(**{"server.require_project": True})
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "no-project"},
|
|
headers=_auth("user-1", permissions=frozenset({"workstreams.create"})),
|
|
)
|
|
assert resp.status_code == 400, resp.json()
|
|
assert resp.json().get("code") == "require_project"
|
|
|
|
def test_projectless_interactive_create_allowed_when_off(self, app_client, make_config_store):
|
|
client, _mgr = app_client
|
|
client.app.state.config_store = make_config_store() # flag off (default)
|
|
resp = client.post(
|
|
"/v1/api/workstreams/new",
|
|
json={"name": "no-project"},
|
|
headers=_auth("user-1", permissions=frozenset({"workstreams.create"})),
|
|
)
|
|
body = resp.json()
|
|
assert not (resp.status_code == 400 and body.get("code") == "require_project"), body
|