Files
turnstone/tests/test_server_authz.py
2026-08-11 21:58:21 -07:00

3613 lines
144 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._auto_approve_tools_source: dict[str, str] = {}
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[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_principals: list[str | None] = []
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
# Principals for whom the fake models a retained queue item owned by
# somebody else. The real session checks this at fresh-slot admission.
self.foreign_queue_principals: set[str] = set()
self.fork_calls: list[tuple[str, str, bool]] = []
# 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 = "",
turn_principal_id: str | None = None,
) -> tuple[str, str, str]:
del turn_principal_id
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 has_foreign_queued_messages(self, principal_id: str) -> bool:
return principal_id in self.foreign_queue_principals
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 fork_from_storage(
self,
source_ws_id: str,
*,
principal_id: str,
source_reservation_token: str,
trusted_internal: bool = False,
) -> Any:
from turnstone.core.storage import get_storage
self.fork_calls.append((source_ws_id, principal_id, trusted_internal))
storage = get_storage()
assert storage is not None
assert source_reservation_token
snapshot = storage.clone_workstream(
source_ws_id,
self.ws_id,
principal_id=principal_id,
trusted_internal=trusted_internal,
)
self.messages = list(snapshot.turns)
return snapshot
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, *, principal_id: str | None = None) -> bool:
self.compacts += 1
self.compact_principals.append(principal_id)
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, *, principal_id: str = "") -> None:
del principal_id
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"]
def test_rejects_creating_parent_ws_id(self, app_client):
"""A child cannot attach to a coordinator before lifecycle birth."""
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
assert storage.register_workstream(
"pending-coord",
node_id="console",
name="pending",
state="creating",
kind="coordinator",
user_id="user-1",
)
resp = client.post(
"/v1/api/workstreams/new",
json={"name": "too-early", "parent_ws_id": "pending-coord"},
headers=_auth("user-1"),
)
assert resp.status_code == 400
assert "known workstream" in resp.json()["error"]
assert mgr.count == 0
class TestAutoApproveToolsOnCreate:
"""The mounted node handler applies CSV/list per-tool approval policy."""
@pytest.mark.parametrize(
"raw_tools",
[
" read_file,write_file,read_file,, ",
[" read_file ", "write_file", "read_file", ""],
],
)
def test_accepts_canonicalizes_and_tags_tools_with_blanket_off(
self,
app_client: Any,
raw_tools: Any,
) -> None:
from turnstone.core.session_ui_base import AutoApproveReason
client, mgr = app_client
response = client.post(
"/v1/api/workstreams/new",
json={"name": "per-tool", "auto_approve_tools": raw_tools},
headers=_auth("user-1"),
)
assert response.status_code == 200, response.text
ws = mgr.get(response.json()["ws_id"])
assert ws is not None
assert ws.ui.auto_approve is False
assert ws.ui.auto_approve_tools == {"read_file", "write_file"}
assert ws.ui._auto_approve_tools_source == {
"read_file": AutoApproveReason.AUTO_APPROVE_TOOLS,
"write_file": AutoApproveReason.AUTO_APPROVE_TOOLS,
}
@pytest.mark.parametrize(
("raw_tools", "error"),
[
({"read_file": True}, "comma-separated string or array"),
(["read_file", 7], "auto_approve_tools[1] must be a string"),
],
)
def test_rejects_malformed_tools_before_create(
self,
app_client: Any,
raw_tools: Any,
error: str,
) -> None:
client, mgr = app_client
response = client.post(
"/v1/api/workstreams/new",
json={"auto_approve_tools": raw_tools},
headers=_auth("user-1"),
)
assert response.status_code == 400
assert error in response.json()["error"]
assert mgr.count == 0
def test_schema_exposes_live_create_fields(self) -> None:
from turnstone.api.server_schemas import CreateWorkstreamRequest
request = CreateWorkstreamRequest(
auto_approve_tools=["read_file"],
judge_model="judge-fast",
user_id="forwarded-user",
)
assert request.auto_approve_tools == ["read_file"]
assert request.judge_model == "judge-fast"
assert request.user_id == "forwarded-user"
class TestMountedCreateUserIdOverride:
"""The real interactive mount trusts only the console service identity."""
@staticmethod
def _service_headers(*, include_service_scope: bool) -> dict[str, str]:
from turnstone.core.auth import JWT_AUD_SERVER, create_jwt
scopes = {"read", "write", "approve"}
if include_service_scope:
scopes.add("service")
token = create_jwt(
user_id="console-service",
scopes=frozenset(scopes),
source="console",
secret=_TEST_JWT_SECRET,
audience=JWT_AUD_SERVER,
permissions=frozenset({"workstreams.create"}),
)
return {"Authorization": f"Bearer {token}"}
def test_ordinary_caller_cannot_override_owner(self, app_client: Any) -> None:
client, mgr = app_client
response = client.post(
"/v1/api/workstreams/new",
json={"name": "ordinary", "user_id": "impersonated"},
headers=_auth("ordinary-user"),
)
assert response.status_code == 200, response.text
ws = mgr.get(response.json()["ws_id"])
assert ws is not None
assert ws.user_id == "ordinary-user"
def test_console_without_service_scope_cannot_override_owner(self, app_client: Any) -> None:
client, mgr = app_client
response = client.post(
"/v1/api/workstreams/new",
json={"name": "unscoped-console", "user_id": "impersonated"},
headers=self._service_headers(include_service_scope=False),
)
assert response.status_code == 200, response.text
ws = mgr.get(response.json()["ws_id"])
assert ws is not None
assert ws.user_id == "console-service"
def test_console_service_can_forward_owner(self, app_client: Any) -> None:
client, mgr = app_client
response = client.post(
"/v1/api/workstreams/new",
json={"name": "trusted-console", "user_id": "forwarded-user"},
headers=self._service_headers(include_service_scope=True),
)
assert response.status_code == 200, response.text
ws = mgr.get(response.json()["ws_id"])
assert ws is not None
assert ws.user_id == "forwarded-user"
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",
"persistence_state",
}
assert row["kind"] == "interactive"
assert row["user_id"] == "user-shape"
assert row["persistence_state"] == "healthy"
# 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_projects_only_sanitized_persistence_state(self, app_client):
client, mgr = app_client
created = client.post(
"/v1/api/workstreams/new",
json={"name": "needs-history-repair"},
headers=_auth("user-a"),
)
ws = mgr.get(created.json()["ws_id"])
ws.session.conversation_persistence_status = lambda: {
"state": "retrying",
"attempts": 2,
"last_failure_at": "not-public",
}
row = client.get("/v1/api/dashboard", headers=_auth("user-a")).json()["workstreams"][0]
assert row["persistence_state"] == "retrying"
assert "attempts" not in row
assert "last_failure_at" not in row
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?user_turn=1",
headers=_auth("attacker-user"),
)
assert resp.status_code == 404
assert "user_turn" not in resp.text
assert "projection" not in resp.text
# ---------------------------------------------------------------------------
# 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_force_cancel_never_abandons_successor_that_replaces_pinned_target(
self,
app_client,
):
"""A exits and B spawns after Stop starts but before force cleanup."""
client, mgr = app_client
ws_id = self._create_ws(client)
ws = mgr.get(ws_id)
assert ws is not None and ws.session is not None and ws.ui is not None
predecessor = threading.Thread(target=lambda: None, name="cancel-target-a")
successor = threading.Thread(target=lambda: None, name="fresh-successor-b")
with ws._lock:
ws._worker_running = True
ws.worker_thread = predecessor
ws._worker_principal_id = "alice"
def _cancel_and_replace_target() -> None:
# Deterministically occupy the exact window between the handler's
# target snapshot and its later force-ownership decision.
with ws._lock:
assert ws.worker_thread is predecessor
ws._worker_running = False
ws.worker_thread = None
ws._worker_principal_id = ""
ws._worker_running = True
ws.worker_thread = successor
ws._worker_principal_id = "bob"
ws.session.cancel = _cancel_and_replace_target # type: ignore[method-assign]
resp = client.post(
f"/v1/api/workstreams/{ws_id}/cancel",
json={"force": True},
headers=_auth("user-1"),
)
assert resp.status_code == 200
assert ws.worker_thread is successor
assert ws._worker_running is True
assert ws._worker_principal_id == "bob"
event_types = [event.get("type") for event in ws.ui._enqueued]
assert "cancelled" in event_types
assert "stream_end" not in event_types
assert "idle" not in ws.ui.states
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_shared_preamble_yields_connected_first(self):
"""The shared handler's preamble preserves connected-event shape."""
from turnstone.core.session_replay import session_replay_preamble
ws, ui, request = _make_interactive_replay_mocks()
out = list(session_replay_preamble(ws.session, ui, project_name="Visible Project"))
assert out[0]["type"] == "connected"
assert out[0]["model"] == "gpt-5"
assert out[0]["model_alias"] == "default"
assert out[0]["project_name"] == "Visible Project"
assert out[0]["skip_permissions"] is False
def test_shared_preamble_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.core.session_replay import session_replay_preamble
ws, ui, _request = _make_interactive_replay_mocks()
out = list(session_replay_preamble(ws.session, ui))
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.session.compact_principals == ["user-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_foreign_retained_queue_refuses_fresh_send_before_worker_spawn(self, app_client):
"""A persistence-retained interjection keeps its original owner.
The foreign-owner check runs inside the atomic fresh-slot admission,
but ``session_worker.send`` reports that refusal as ``False``. The
route must preserve the typed conflict outcome instead of falling
through to its generic queue-full response, and no worker may bind or
append the later participant's turn.
"""
client, mgr = app_client
ws_id = self._create_ws(client)
ws = mgr.get(ws_id)
assert ws is not None
ws.session.foreign_queue_principals.add("user-1")
response = client.post(
f"/v1/api/workstreams/{ws_id}/send",
json={"message": "must wait for the retained owner"},
headers=_auth("user-1"),
)
assert response.status_code == 409
assert response.json()["status"] == "cross_user_interjection"
assert ws.session.sends == []
assert ws.session.queue_calls == []
assert ws.worker_thread is None
assert ws._worker_running is False
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]
assert [e["sender"] for e in queued_events] == ["user-1"]
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 TestRemoteLifecycleCommandBoundary:
"""The HTTP command bridge never reaches storage-global REPL lifecycle code."""
@pytest.mark.parametrize(
"command",
[
"/workstreams",
"/resume secret-alias",
"/resume missing-alias",
"/delete secret-alias",
"/delete missing-alias",
],
)
def test_cli_only_command_is_rejected_before_dispatch_without_oracle(
self,
app_client,
command: str,
) -> None:
"""Known/private and missing targets have one inert, non-disclosing result."""
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
storage.create_project("secret-project", "Secret Project", "victim-user")
storage.register_workstream(
"secret-ws",
node_id="node-test",
name="secret-name",
user_id="victim-user",
project_id="secret-project",
)
assert storage.set_workstream_alias("secret-ws", "secret-alias") is True
storage.save_message("secret-ws", "user", "private history")
created = client.post(
"/v1/api/workstreams/new",
json={"name": "caller-workstream"},
headers=_auth("caller-user"),
)
assert created.status_code == 200
caller_ws_id = created.json()["ws_id"]
caller_ws = mgr.get(caller_ws_id)
assert caller_ws is not None
assert caller_ws.session is not None
original_session_id = caller_ws.session.ws_id
original_messages = list(caller_ws.session.messages)
response = client.post(
"/v1/api/command",
json={"command": command, "ws_id": caller_ws_id},
headers=_auth("caller-user"),
)
assert response.status_code == 400
assert response.json() == {
"error": "This workstream command is only available in the local CLI."
}
assert caller_ws.worker_thread is None
assert caller_ws.session.commands == []
assert caller_ws.session.ws_id == original_session_id
assert caller_ws.session.messages == original_messages
assert storage.get_workstream("secret-ws") is not None
assert [turn.text for turn in storage.load_message_turns("secret-ws")] == [
"private history"
]
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
def test_flag_off_forwarded_service_cannot_fork_private_nonmember(
self, app_client, make_config_store, monkeypatch
):
"""The resume read gate runs before persona/history reads or destination create.
A console service token may forward the real end-user id, but its cross-tenant
service scope belongs to the console identity. It must not make that different
effective user omniscient.
"""
from turnstone.core.auth import JWT_AUD_SERVER, create_jwt
from turnstone.core.storage import get_storage
client, mgr = app_client
client.app.state.config_store = make_config_store() # invariant is flag-independent
storage = get_storage()
assert storage is not None
storage.create_project("private-project", "Secret", "victim-user")
source_id = "a" * 32
storage.register_workstream(
source_id,
node_id="node-test",
name="secret-source",
user_id="victim-user",
project_id="private-project",
)
assert storage.set_workstream_alias(source_id, "secret-alias") is True
source_config_read = MagicMock(side_effect=AssertionError("private source config read"))
monkeypatch.setattr(storage, "load_workstream_config", source_config_read)
token = create_jwt(
user_id="console-service",
scopes=frozenset({"read", "write", "approve", "service"}),
source="console",
secret=_TEST_JWT_SECRET,
audience=JWT_AUD_SERVER,
permissions=frozenset({"workstreams.create"}),
)
resp = client.post(
"/v1/api/workstreams/new",
json={
"name": "forbidden-fork",
"resume_ws": "secret-alias",
"user_id": "forwarded-user",
},
headers={"Authorization": f"Bearer {token}"},
)
assert resp.status_code == 404
assert resp.json() == {"error": "Workstream not found"}
assert mgr.list_all() == []
source_config_read.assert_not_called()
class TestCreateForkRollback:
"""A requested fork is one create transaction at the mounted HTTP edge."""
@staticmethod
def _register_source(storage: Any, ws_id: str, *, with_history: bool) -> None:
storage.register_workstream(
ws_id,
node_id="node-test",
name="fork-source",
user_id="user-1",
)
if with_history:
storage.save_message(ws_id, "user", "durable source history")
@staticmethod
def _global_events(client: Any) -> list[dict[str, Any]]:
events: list[dict[str, Any]] = []
global_queue = client.app.state.global_queue
while True:
try:
events.append(global_queue.get_nowait())
except queue.Empty:
return events
def test_source_disappears_after_validation_rolls_back_destination(
self, app_client, monkeypatch
) -> None:
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
source_id = "b" * 32
destination_id = "c" * 32
self._register_source(storage, source_id, with_history=False)
original_fork = _FakeSession.fork_from_storage
def _disappear_at_pre_commit(
session: _FakeSession,
fork_source_id: str,
*,
principal_id: str,
source_reservation_token: str,
trusted_internal: bool = False,
) -> Any:
assert fork_source_id == source_id
assert storage.delete_workstream(source_id) is True
return original_fork(
session,
fork_source_id,
principal_id=principal_id,
source_reservation_token=source_reservation_token,
trusted_internal=trusted_internal,
)
monkeypatch.setattr(_FakeSession, "fork_from_storage", _disappear_at_pre_commit)
resp = client.post(
"/v1/api/workstreams/new",
json={
"ws_id": destination_id,
"name": "vanished-source-fork",
"resume_ws": source_id,
},
headers=_auth("user-1"),
)
assert resp.status_code == 409
assert resp.json() == {"error": "Fork source is no longer available"}
assert mgr.get(destination_id) is None
assert storage.get_workstream(destination_id) is None
assert not [
event
for event in storage.list_audit_events(action="workstream.created")
if event["resource_id"] == destination_id
]
assert not {
event["type"]
for event in self._global_events(client)
if event.get("type") in {"ws_created", "ws_rename"}
}
def test_source_replacement_after_preflight_cannot_inherit_fork(
self,
app_client,
monkeypatch,
) -> None:
from turnstone.core.storage import ForkCloneExpectation, get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
source_id = "8" * 32
destination_id = "9" * 32
self._register_source(storage, source_id, with_history=True)
replacement_token = "replacement-source-incarnation"
def _replace_at_pre_commit(
session: _FakeSession,
fork_source_id: str,
*,
principal_id: str,
source_reservation_token: str,
trusted_internal: bool = False,
) -> Any:
assert fork_source_id == source_id
assert source_reservation_token
assert storage.get_workstream_reservation_token(source_id) == (source_reservation_token)
assert storage.delete_workstream(source_id) is True
assert storage.register_workstream(
source_id,
user_id="user-1",
name="replacement-source",
state="idle",
kind="interactive",
fork_reservation_token=replacement_token,
)
storage.save_message(source_id, "user", "replacement must not fork")
destination_token = str(getattr(session, "_fork_reservation_token", ""))
assert destination_token
return storage.clone_workstream(
source_id,
session.ws_id,
principal_id=principal_id,
trusted_internal=trusted_internal,
expected_session=ForkCloneExpectation(
persona_config=(),
project_id="",
project_name="",
project_writable=False,
destination_reservation_token=destination_token,
source_reservation_token=source_reservation_token,
),
)
monkeypatch.setattr(_FakeSession, "fork_from_storage", _replace_at_pre_commit)
response = client.post(
"/v1/api/workstreams/new",
json={
"ws_id": destination_id,
"name": "replaced-source-fork",
"resume_ws": source_id,
},
headers=_auth("user-1"),
)
assert response.status_code == 409
assert response.json() == {"error": "Fork source is no longer available"}
assert mgr.get(destination_id) is None
assert storage.get_workstream(destination_id) is None
replacement = storage.get_workstream(source_id)
assert replacement is not None
assert replacement["name"] == "replacement-source"
assert "fork_reservation_token" not in replacement
assert storage.get_workstream_reservation_token(source_id) == replacement_token
assert [turn.text for turn in storage.load_message_turns(source_id)] == [
"replacement must not fork"
]
def test_destination_storage_failure_rolls_back_destination(
self, app_client, monkeypatch
) -> None:
from turnstone.core.storage import ForkDestinationConflictError, get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
source_id = "d" * 32
destination_id = "e" * 32
self._register_source(storage, source_id, with_history=True)
assert storage.load_message_turns(source_id)
fork_calls: list[tuple[str, str, bool]] = []
def _fail_fork(
_session: _FakeSession,
fork_source_id: str,
*,
principal_id: str,
source_reservation_token: str,
trusted_internal: bool = False,
) -> Any:
fork_calls.append((fork_source_id, principal_id, trusted_internal))
raise ForkDestinationConflictError("destination raced")
monkeypatch.setattr(_FakeSession, "fork_from_storage", _fail_fork)
resp = client.post(
"/v1/api/workstreams/new",
json={
"ws_id": destination_id,
"name": "failed-history-fork",
"resume_ws": source_id,
},
headers=_auth("user-1"),
)
assert resp.status_code == 409
assert resp.json() == {"error": "Workstream creation was superseded"}
assert fork_calls == [(source_id, "user-1", False)]
assert mgr.get(destination_id) is None
assert storage.get_workstream(destination_id) is None
assert storage.load_message_turns(source_id)
assert not [
event
for event in storage.list_audit_events(action="workstream.created")
if event["resource_id"] == destination_id
]
assert not {
event["type"]
for event in self._global_events(client)
if event.get("type") in {"ws_created", "ws_rename"}
}
def test_caller_chosen_destination_collision_preserves_existing_row(self, app_client) -> None:
"""An unloaded durable row is not a blank reservation for a fork.
``register_workstream`` historically ignored a duplicate id. That
made a caller-chosen id collide with an existing empty workstream,
after which the clone could replace its config/history and the create
rollback could delete it outright. Reject the collision before either
mutation, even when owner and project happen to match.
"""
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
project_id = "collision-project"
source_id = "8" * 32
destination_id = "9" * 32
storage.create_project(project_id, "Collision", "user-1", visibility="private")
storage.register_workstream(
source_id,
node_id="source-node",
name="source",
state="idle",
user_id="user-1",
project_id=project_id,
)
storage.save_message(source_id, "user", "source history")
storage.register_workstream(
destination_id,
node_id="original-node",
name="existing empty destination",
state="idle",
user_id="user-1",
project_id=project_id,
)
storage.save_workstream_config(destination_id, {"keep": "untouched"})
before = storage.get_workstream(destination_id)
assert before is not None
assert storage.load_message_turns(destination_id) == []
resp = client.post(
"/v1/api/workstreams/new",
json={
"ws_id": destination_id,
"name": "colliding fork",
"resume_ws": source_id,
},
headers=_auth("user-1"),
)
assert resp.status_code == 409, resp.json()
assert mgr.get(destination_id) is None
assert storage.get_workstream(destination_id) == before
assert storage.load_message_turns(destination_id) == []
assert storage.load_workstream_config(destination_id) == {"keep": "untouched"}
assert [turn.text for turn in storage.load_message_turns(source_id)] == ["source history"]
assert not [
event
for event in storage.list_audit_events(action="workstream.created")
if event["resource_id"] == destination_id
]
assert not {
event["type"]
for event in self._global_events(client)
if event.get("type") in {"ws_created", "ws_rename"}
}
@pytest.mark.parametrize("lifecycle", ["close", "delete"])
def test_destination_retired_while_clone_blocked_cannot_publish_or_resurrect(
self,
app_client,
monkeypatch,
lifecycle: str,
) -> None:
"""A pending create loses ownership when its exact slot is retired."""
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
source_id = "a" * 32
destination_id = "b" * 32
self._register_source(storage, source_id, with_history=True)
entered_clone = threading.Event()
release_clone = threading.Event()
original_fork = _FakeSession.fork_from_storage
def _blocked_fork(
session: _FakeSession,
fork_source_id: str,
*,
principal_id: str,
source_reservation_token: str,
trusted_internal: bool = False,
) -> Any:
entered_clone.set()
assert release_clone.wait(timeout=10), "test did not release clone"
return original_fork(
session,
fork_source_id,
principal_id=principal_id,
source_reservation_token=source_reservation_token,
trusted_internal=trusted_internal,
)
monkeypatch.setattr(_FakeSession, "fork_from_storage", _blocked_fork)
responses: list[Any] = []
request_errors: list[BaseException] = []
def _create_fork() -> None:
try:
responses.append(
client.post(
"/v1/api/workstreams/new",
json={"ws_id": destination_id, "resume_ws": source_id},
headers=_auth("user-1"),
)
)
except BaseException as exc: # pragma: no cover - diagnostic capture
request_errors.append(exc)
request_thread = threading.Thread(target=_create_fork, daemon=True)
request_thread.start()
assert entered_clone.wait(timeout=5), "fork request never reached clone"
# Deferred creates are deliberately hidden from public manager lookup
# until lifecycle birth; inspect the exact internal reservation to
# drive the terminal-race seam.
pending = mgr._workstreams.get(destination_id)
assert pending is not None
try:
if lifecycle == "close":
assert mgr.close(destination_id) is True
else:
assert storage.delete_workstream(destination_id) is True
assert mgr.delete(destination_id) is True
finally:
release_clone.set()
request_thread.join(timeout=10)
assert not request_thread.is_alive()
assert request_errors == []
assert len(responses) == 1
resp = responses[0]
assert resp.status_code == 409, resp.text
assert mgr.get(destination_id) is None
destination = storage.get_workstream(destination_id)
if lifecycle == "delete":
assert destination is None
else:
# The create rollback may delete this never-advertised row. If it
# elects to preserve the concurrent close, it must stay retired.
assert destination is None or destination["state"] == "closed"
if destination is not None:
assert storage.load_message_turns(destination_id) == []
assert [turn.text for turn in storage.load_message_turns(source_id)] == [
"durable source history"
]
assert not [
event
for event in storage.list_audit_events(action="workstream.created")
if event["resource_id"] == destination_id
]
events = self._global_events(client)
assert not {
event["type"] for event in events if event.get("type") in {"ws_created", "ws_rename"}
}
# The reservation never crossed lifecycle birth, so terminal cleanup
# is silent: neither a phantom create nor a close-without-create is
# observable on the global stream.
assert not [event for event in events if event.get("ws_id") == destination_id]
@pytest.mark.anyio
@pytest.mark.parametrize("anyio_backend", ["asyncio"])
async def test_cancelled_request_waits_for_clone_then_rolls_back(
self,
app_client,
monkeypatch,
anyio_backend: str,
) -> None:
"""Request cancellation cannot outrun the non-cancellable clone."""
import asyncio
import httpx
assert anyio_backend == "asyncio"
sync_client, mgr = app_client
storage = sync_client.app.state.auth_storage
assert storage is not None
source_id = "c" * 32
destination_id = "d" * 32
attachment_id = "e" * 64
self._register_source(storage, source_id, with_history=False)
source_row = storage.save_message(source_id, "user", "attached source")
storage.save_attachment(
attachment_id,
"source.txt",
"text/plain",
6,
"text",
b"source",
)
storage.set_message_attachments(source_id, source_row, [attachment_id])
source_attachment = storage.get_attachment(attachment_id)
assert source_attachment is not None and source_attachment["refcount"] == 1
entered_clone = threading.Event()
release_clone = threading.Event()
clone_finished = threading.Event()
original_fork = _FakeSession.fork_from_storage
def _blocked_fork(
session: _FakeSession,
fork_source_id: str,
*,
principal_id: str,
source_reservation_token: str,
trusted_internal: bool = False,
) -> Any:
entered_clone.set()
assert release_clone.wait(timeout=10), "test did not release clone"
try:
return original_fork(
session,
fork_source_id,
principal_id=principal_id,
source_reservation_token=source_reservation_token,
trusted_internal=trusted_internal,
)
finally:
clone_finished.set()
monkeypatch.setattr(_FakeSession, "fork_from_storage", _blocked_fork)
transport = httpx.ASGITransport(app=sync_client.app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
request_task = asyncio.create_task(
client.post(
"/v1/api/workstreams/new",
json={"ws_id": destination_id, "resume_ws": source_id},
headers=_auth("user-1"),
)
)
for _ in range(500):
if entered_clone.is_set():
break
await asyncio.sleep(0.01)
assert entered_clone.is_set(), "fork request never reached clone"
request_task.cancel()
try:
await asyncio.sleep(0.05)
assert not request_task.done()
assert not clone_finished.is_set()
finally:
release_clone.set()
with pytest.raises(asyncio.CancelledError):
await request_task
assert clone_finished.is_set()
assert mgr.get(destination_id) is None
assert storage.get_workstream(destination_id) is None
assert [turn.text for turn in storage.load_message_turns(source_id)] == ["attached source"]
attachment = storage.get_attachment(attachment_id)
assert attachment is not None and attachment["refcount"] == 1
assert not [
event
for event in storage.list_audit_events(action="workstream.created")
if event["resource_id"] == destination_id
]
assert not {
event["type"]
for event in self._global_events(sync_client)
if event.get("type") in {"ws_created", "ws_rename"}
}
def test_source_project_deleted_after_preflight_is_uniform_409(
self, app_client, monkeypatch
) -> None:
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
storage.create_project("fork-project", "Fork", "user-1", visibility="public")
source_id = "6" * 32
destination_id = "7" * 32
storage.register_workstream(
source_id,
node_id="node-test",
name="project-source",
user_id="user-1",
project_id="fork-project",
)
original_fork = _FakeSession.fork_from_storage
def _delete_project_at_pre_commit(
session: _FakeSession,
fork_source_id: str,
*,
principal_id: str,
source_reservation_token: str,
trusted_internal: bool = False,
) -> Any:
assert storage.delete_project("fork-project") is True
return original_fork(
session,
fork_source_id,
principal_id=principal_id,
source_reservation_token=source_reservation_token,
trusted_internal=trusted_internal,
)
monkeypatch.setattr(_FakeSession, "fork_from_storage", _delete_project_at_pre_commit)
resp = client.post(
"/v1/api/workstreams/new",
json={"ws_id": destination_id, "resume_ws": source_id},
headers=_auth("user-1"),
)
assert resp.status_code == 409
assert resp.json() == {"error": "Fork source is no longer available"}
assert mgr.get(destination_id) is None
assert storage.get_workstream(destination_id) is None
assert storage.get_workstream(source_id) is not None
assert not [
event
for event in storage.list_audit_events(action="workstream.created")
if event["resource_id"] == destination_id
]
assert not {
event["type"]
for event in self._global_events(client)
if event.get("type") in {"ws_created", "ws_rename"}
}
def test_empty_source_is_successful_fork(self, app_client) -> None:
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
source_id = "f" * 32
destination_id = "1" * 32
self._register_source(storage, source_id, with_history=False)
assert storage.load_message_turns(source_id) == []
resp = client.post(
"/v1/api/workstreams/new",
json={
"ws_id": destination_id,
"name": "empty-source-fork",
"resume_ws": source_id,
},
headers=_auth("user-1"),
)
assert resp.status_code == 200
assert resp.json()["resumed"] is True
assert resp.json()["message_count"] == 0
destination = mgr.get(destination_id)
assert destination is not None
assert destination.session is not None
assert destination.session.fork_calls == [(source_id, "user-1", False)]
assert {event["type"] for event in destination.ui._enqueued} >= {"clear_ui"}
assert storage.get_workstream(destination_id) is not None
events = self._global_events(client)
assert [event["type"] for event in events if event.get("type") == "ws_created"] == [
"ws_created"
]
assert [event["type"] for event in events if event.get("type") == "ws_rename"] == [
"ws_rename"
]
def test_nonempty_source_commits_before_publication(self, app_client) -> None:
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
source_id = "2" * 32
destination_id = "3" * 32
self._register_source(storage, source_id, with_history=True)
resp = client.post(
"/v1/api/workstreams/new",
json={
"ws_id": destination_id,
"name": "copied-history-fork",
"resume_ws": source_id,
},
headers=_auth("user-1"),
)
assert resp.status_code == 200, resp.json()
assert resp.json()["resumed"] is True
assert resp.json()["message_count"] == 1
assert [turn.text for turn in storage.load_message_turns(destination_id)] == [
"durable source history"
]
destination = mgr.get(destination_id)
assert destination is not None and destination.session is not None
assert [turn.text for turn in destination.session.messages] == ["durable source history"]
assert destination.session.fork_calls == [(source_id, "user-1", False)]
created = [
event
for event in storage.list_audit_events(action="workstream.created")
if event["resource_id"] == destination_id
]
assert len(created) == 1
def test_successful_fork_still_dispatches_initial_message(self, app_client) -> None:
from turnstone.core.storage import get_storage
client, mgr = app_client
storage = get_storage()
assert storage is not None
source_id = "4" * 32
destination_id = "5" * 32
self._register_source(storage, source_id, with_history=False)
resp = client.post(
"/v1/api/workstreams/new",
json={
"ws_id": destination_id,
"resume_ws": source_id,
"initial_message": "continue from the fork",
},
headers=_auth("user-1"),
)
assert resp.status_code == 200, resp.json()
assert resp.json()["resumed"] is True
destination = mgr.get(destination_id)
assert destination is not None and destination.session is not None
assert destination.worker_thread is not None
destination.worker_thread.join(timeout=5)
assert not destination.worker_thread.is_alive()
assert destination.session.sends == [("continue from the fork", None, None)]