mirror of
https://github.com/turnstonelabs/turnstone.git
synced 2026-08-12 23:12:23 -06:00
480a1426b3
* fix(session): fail-closed history-commit handoff (#981) The deleted-workstream discovery is now a terminal, ws_id-keyed latch: keyed conversation commits refuse admission once the durable parent is gone (convergence finalizers and force-abandon are exempt), history handoff refuses to mint a proof token so /history fails closed with a 503 instead of silently wiping the pane, and the SSE stream carries a workstream_gone resync reason. Discarded commits leave a forensic log of commit keys and roles, never content. Conversation rows gain a commit_key (migration 071): keyed saves are idempotent under retry, validated against the full commit identity, and refused when they would cross a workstream deletion. The prune orphan category now requires a NULL alias plus a two-hour updated grace, with cutoffs computed at discovery time and carried into both dialects' rechecks. The mid-turn interjection queue is owner-partitioned with no per-site mode flags: pops take the acting principal's and unowned rows, other participants' rows are structurally retained, and enforcement lives at queue admission plus the shared before_spawn gates. The retraction ledger is bounded by open pop windows: pops open a window atomically with the queue delete, restores close their ids atomically with the ledger consume, every other exit closes through one helper, and misses for unheld ids record nothing. The workstream-gone latch refuses unattended wakes at all three gates (watcher spawn, claim, delivery pre-pop), and the retry dispatcher regained its pre-envelope cancel/error convergence net. Persistence-state reporting derives through the session bound to each UI instead of a registry lookup by id that failed open to healthy during tombstone retention. The dashboard roster no longer re-inserts ghost entries from trailing activity events, the history tool-outcome scan tolerates interleaved non-turn rows, and the shared handoff-deadline handle owns its own retirement. Single-sourced across call sites: keyed-commit row values, attachment save wrappers, tail-truncation and conflict-resolution bodies for both storage dialects; worker-slot lifecycle field sets; the direct-commit admission frame; queued-row layout accessors; the string-aware comment stripper shared by every JS harness suite. Refs #981 #964 * fix(session): sweep handoff fixes to their sibling surfaces The interactive replay loop treated a system row as a tool-batch boundary, so every tool result after an interleaved row vanished from that pane while the coordinator rendered the same history correctly. Only a conversational turn ends the batch window now, matching the shared outcome index. Accepted user turns clear the composer's attachment chips on the same viewer policy that settles optimistic bubbles rather than on having matched a local bubble, so a workstream created with an upload no longer keeps a chip for an attachment the create dispatch already consumed. The coordinator's raced-Stop arm emits the stream-end hook it inherits alongside the idle state, leaving no unfinalized bubble or unflushed tool output. Ending a session surfaces a failure toast when the request never lands or answers with a non-JSON body. The per-second persistence reconcile now probes each session without blocking: a workstream whose generation and handoff locks are held is skipped until the next pass instead of contending the locks every commit needs. The one-shot repair that gates workstream creation at capacity keeps a definite probe — it has no next pass, and the sessions likeliest to be contended are the ones whose unresolved journals emptied its candidate list. Single-sourced: the attachment lane builds its conversation row through the shared commit-identity builder; the ordinary worker exit releases its slot through the lifecycle owner; both operator surfaces snapshot their counters through one non-consuming helper; the replay preamble loses its per-kind wrappers and its config hook; the browser harness suites share one brace walker; and each in-flight history attempt is one record carrying both its abort controller and its deadline. Refs #981 #964
401 lines
13 KiB
Python
401 lines
13 KiB
Python
"""Tests for the shared SSE reconnect helper in turnstone.channels._sse."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import json
|
|
from types import SimpleNamespace
|
|
from unittest.mock import AsyncMock, MagicMock
|
|
|
|
import httpx
|
|
import pytest
|
|
|
|
|
|
def _run(coro): # type: ignore[no-untyped-def]
|
|
return asyncio.run(coro)
|
|
|
|
|
|
class _FakeSSEEvent:
|
|
"""A fake ``httpx_sse.ServerSentEvent`` with the subset we read."""
|
|
|
|
def __init__(self, event: str, data: str) -> None:
|
|
self.event = event
|
|
self.data = data
|
|
|
|
|
|
class _FakeEventSource:
|
|
"""Context manager returned by our fake ``aconnect_sse``.
|
|
|
|
Captures the (status_code, events) the test wants to deliver.
|
|
``aiter_sse`` yields the events then returns; the caller then hits
|
|
the outer ``while True`` loop again, which will pick up the next
|
|
queued response via the shared iterator state on _FakeConnect.
|
|
"""
|
|
|
|
def __init__(self, *, status_code: int, events: list[_FakeSSEEvent]) -> None:
|
|
self.response = SimpleNamespace(
|
|
status_code=status_code,
|
|
request=MagicMock(),
|
|
)
|
|
self._events = events
|
|
|
|
async def __aenter__(self) -> _FakeEventSource:
|
|
return self
|
|
|
|
async def __aexit__(self, exc_type, exc, tb) -> None: # noqa: ANN001
|
|
return None
|
|
|
|
async def aiter_sse(self): # type: ignore[no-untyped-def]
|
|
for event in self._events:
|
|
yield event
|
|
|
|
|
|
class _FakeConnect:
|
|
"""Drop-in replacement for ``httpx_sse.aconnect_sse``.
|
|
|
|
On each call, pops the next ``_FakeEventSource`` from *queue*. When
|
|
the queue is empty, raises ``asyncio.CancelledError`` so the loop
|
|
terminates cleanly in tests.
|
|
"""
|
|
|
|
def __init__(self, queue: list[_FakeEventSource]) -> None:
|
|
self._queue = queue
|
|
self.call_count = 0
|
|
self.calls: list[tuple[tuple[object, ...], dict[str, object]]] = []
|
|
|
|
def __call__(self, *args, **kwargs): # noqa: ANN001, ANN204
|
|
self.call_count += 1
|
|
self.calls.append((args, kwargs))
|
|
if not self._queue:
|
|
raise asyncio.CancelledError
|
|
return self._queue.pop(0)
|
|
|
|
|
|
@pytest.fixture
|
|
def _fast_sleep(monkeypatch):
|
|
"""Patch asyncio.sleep so backoff doesn't actually wait; record calls."""
|
|
sleeps: list[float] = []
|
|
|
|
async def fake_sleep(delay: float) -> None:
|
|
sleeps.append(delay)
|
|
|
|
monkeypatch.setattr("turnstone.channels._sse.asyncio.sleep", fake_sleep)
|
|
return sleeps
|
|
|
|
|
|
def _valid_event_data(ws_id: str = "ws-1") -> str:
|
|
"""A payload ``ServerEvent.from_dict`` will accept (a ContentEvent)."""
|
|
return json.dumps(
|
|
{
|
|
"type": "content",
|
|
"ws_id": ws_id,
|
|
"text": "hello",
|
|
}
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 404 → on_stale + exit
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestStaleRoute:
|
|
def test_404_calls_on_stale_and_returns(self, monkeypatch, _fast_sleep):
|
|
from turnstone.channels import _sse
|
|
|
|
queue = [_FakeEventSource(status_code=404, events=[])]
|
|
fake_connect = _FakeConnect(queue)
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
on_stale = AsyncMock()
|
|
on_event = AsyncMock()
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=None,
|
|
on_event=on_event,
|
|
on_stale=on_stale,
|
|
)
|
|
)
|
|
|
|
on_stale.assert_awaited_once()
|
|
on_event.assert_not_awaited()
|
|
# No reconnect after 404.
|
|
assert fake_connect.call_count == 1
|
|
assert fake_connect.calls[0][1]["params"] == {"user_turn": 1}
|
|
assert _fast_sleep == []
|
|
|
|
def test_on_stale_exception_still_exits(self, monkeypatch, _fast_sleep):
|
|
"""If on_stale raises, the loop must not reconnect."""
|
|
from turnstone.channels import _sse
|
|
|
|
queue = [_FakeEventSource(status_code=404, events=[])]
|
|
fake_connect = _FakeConnect(queue)
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
on_stale = AsyncMock(side_effect=RuntimeError("storage down"))
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=None,
|
|
on_event=AsyncMock(),
|
|
on_stale=on_stale,
|
|
)
|
|
)
|
|
|
|
on_stale.assert_awaited_once()
|
|
# Still a single connect — no livelock.
|
|
assert fake_connect.call_count == 1
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 500+ → exponential backoff
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestBackoff:
|
|
def test_500_triggers_backoff_and_retries(self, monkeypatch, _fast_sleep):
|
|
from turnstone.channels import _sse
|
|
|
|
queue = [
|
|
_FakeEventSource(status_code=503, events=[]),
|
|
_FakeEventSource(status_code=503, events=[]),
|
|
_FakeEventSource(status_code=503, events=[]),
|
|
]
|
|
fake_connect = _FakeConnect(queue)
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
with contextlib.suppress(asyncio.CancelledError):
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=None,
|
|
on_event=AsyncMock(),
|
|
on_stale=AsyncMock(),
|
|
)
|
|
)
|
|
|
|
assert fake_connect.call_count >= 3
|
|
# First three recorded sleeps are 2s, 4s, 8s (starts at
|
|
# SSE_RECONNECT_DELAY, doubles each time, capped at
|
|
# SSE_MAX_RECONNECT_DELAY).
|
|
assert _fast_sleep[0] == _sse.SSE_RECONNECT_DELAY
|
|
assert _fast_sleep[1] == _sse.SSE_RECONNECT_DELAY * 2
|
|
assert _fast_sleep[2] == _sse.SSE_RECONNECT_DELAY * 4
|
|
|
|
def test_backoff_resets_after_successful_dispatch(self, monkeypatch, _fast_sleep):
|
|
"""After a 200 + successful event dispatch, the next error
|
|
restarts backoff at the initial delay."""
|
|
from turnstone.channels import _sse
|
|
|
|
good_event = _FakeSSEEvent(event="message", data=_valid_event_data())
|
|
queue = [
|
|
_FakeEventSource(status_code=503, events=[]),
|
|
_FakeEventSource(status_code=200, events=[good_event]),
|
|
_FakeEventSource(status_code=503, events=[]),
|
|
]
|
|
fake_connect = _FakeConnect(queue)
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
on_event = AsyncMock()
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
with contextlib.suppress(asyncio.CancelledError):
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=None,
|
|
on_event=on_event,
|
|
on_stale=AsyncMock(),
|
|
)
|
|
)
|
|
|
|
on_event.assert_awaited()
|
|
# Sleep sequence: 2 (after first 503), 2 (reset after 200/event),
|
|
# then CancelledError exits. First two sleeps are both the base
|
|
# delay — the reset did its job.
|
|
assert len(_fast_sleep) >= 2
|
|
assert _fast_sleep[0] == _sse.SSE_RECONNECT_DELAY
|
|
assert _fast_sleep[1] == _sse.SSE_RECONNECT_DELAY
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Event dispatch
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestEventDispatch:
|
|
def test_invalid_json_is_skipped(self, monkeypatch, _fast_sleep):
|
|
from turnstone.channels import _sse
|
|
|
|
bad = _FakeSSEEvent(event="message", data="{not json")
|
|
good = _FakeSSEEvent(event="message", data=_valid_event_data())
|
|
queue = [_FakeEventSource(status_code=200, events=[bad, good])]
|
|
fake_connect = _FakeConnect(queue)
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
on_event = AsyncMock()
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
with contextlib.suppress(asyncio.CancelledError):
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=None,
|
|
on_event=on_event,
|
|
on_stale=AsyncMock(),
|
|
)
|
|
)
|
|
|
|
# Good event delivered, bad one silently dropped.
|
|
assert on_event.await_count == 1
|
|
|
|
def test_on_event_exception_does_not_kill_stream(self, monkeypatch, _fast_sleep):
|
|
from turnstone.channels import _sse
|
|
|
|
e1 = _FakeSSEEvent(event="message", data=_valid_event_data())
|
|
e2 = _FakeSSEEvent(event="message", data=_valid_event_data())
|
|
queue = [_FakeEventSource(status_code=200, events=[e1, e2])]
|
|
fake_connect = _FakeConnect(queue)
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
on_event = AsyncMock(side_effect=[RuntimeError("boom"), None])
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
with contextlib.suppress(asyncio.CancelledError):
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=None,
|
|
on_event=on_event,
|
|
on_stale=AsyncMock(),
|
|
)
|
|
)
|
|
|
|
# Both events attempted — first raised but second still delivered.
|
|
assert on_event.await_count == 2
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Token factory
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestTokenFactory:
|
|
def test_header_refreshed_per_connection(self, monkeypatch, _fast_sleep):
|
|
"""token_factory is called once per reconnect so rotating service
|
|
JWTs stay fresh."""
|
|
from turnstone.channels import _sse
|
|
|
|
# Two reconnects followed by CancelledError to exit.
|
|
queue = [
|
|
_FakeEventSource(status_code=503, events=[]),
|
|
_FakeEventSource(status_code=503, events=[]),
|
|
]
|
|
fake_connect = _FakeConnect(queue)
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
tokens: list[str] = []
|
|
|
|
def factory() -> str:
|
|
tok = f"tok-{len(tokens)}"
|
|
tokens.append(tok)
|
|
return tok
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
with contextlib.suppress(asyncio.CancelledError):
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=factory,
|
|
on_event=AsyncMock(),
|
|
on_stale=AsyncMock(),
|
|
)
|
|
)
|
|
|
|
assert len(tokens) >= 2
|
|
assert tokens[0] != tokens[1]
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# httpx errors
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestTransportErrors:
|
|
def test_connect_error_falls_through_to_backoff(self, monkeypatch, _fast_sleep):
|
|
"""ConnectError is caught and treated as retryable."""
|
|
from turnstone.channels import _sse
|
|
|
|
call_order = {"n": 0}
|
|
|
|
def fake_connect(*args, **kwargs): # noqa: ANN001, ANN003
|
|
call_order["n"] += 1
|
|
if call_order["n"] == 1:
|
|
raise httpx.ConnectError("boom")
|
|
# Second attempt: signal the loop to exit.
|
|
raise asyncio.CancelledError
|
|
|
|
monkeypatch.setattr(_sse.httpx_sse, "aconnect_sse", fake_connect)
|
|
|
|
async def node_url_fn(ws_id: str) -> str:
|
|
return "http://node"
|
|
|
|
with contextlib.suppress(asyncio.CancelledError):
|
|
_run(
|
|
_sse.run_sse_stream(
|
|
http_client=MagicMock(),
|
|
log_prefix="test",
|
|
ws_id="ws-1",
|
|
node_url_fn=node_url_fn,
|
|
token_factory=None,
|
|
on_event=AsyncMock(),
|
|
on_stale=AsyncMock(),
|
|
)
|
|
)
|
|
|
|
assert call_order["n"] == 2
|
|
# Backoff ran once after the ConnectError.
|
|
assert _fast_sleep == [_sse.SSE_RECONNECT_DELAY]
|