feat(sse_client,cli,tui): implement issue #7 — empty-data skip + MalformedSseData
Bundles initial TDD impl + Volva-code-review F1/F3 amendments. sse_client.py: - New MalformedSseData(raw) exception; truncates raw to 200 chars at __init__ (mirrors MalformedSseId.raw[:64] precedent). - _iter_events gains `if sse.data == '': continue` BEFORE _parse_sse_id. Empty-data frames are silently skipped per issue #7 INV-001 (keepalive semantics). Empty-data + bad-id is still a keepalive; intentional ordering, don't reorder. - _iter_events json.loads(sse.data) now wrapped — JSONDecodeError → MalformedSseData(raw=sse.data). cli.py: - Imports MalformedSseData; _run_turn ERROR_ROUTING gains the case → stderr `[malformed_sse_data] raw={exc.raw!r}` + exit 22 (protocol- failure bucket, same as MalformedSseId/TurnIdFlip). tui.py: - Imports MalformedSseData; _stream_turn_worker ERROR_ROUTING gains the case → transcript label; finally block restores state→idle per INV-008 (mid-session errors don't exit the app). Tests (6 new): - test_sse_client.py: empty_data_skipped (tracer — 4 frames in, 3 events out), malformed_data_raises, whitespace_data_raises, malformed_data_truncation, AND empty_data_skip_preserves_last_seen_sse_id (F1 from Volva code-review — drop-after-empty probes internal last_sse_id non-advancement via SseConnectionDropped.last_seen_sse_id). - test_cli.py: malformed_sse_data (tightened to assert exact `[malformed_sse_data] raw='not-json'` shape per F3), malformed_sse_data_truncation (5000-char payload — verifies truncation carries through presenter rendering, F3). - test_tui.py: malformed_sse_data_returns_to_idle (state→idle per INV-008; app does NOT exit). Smoke validation (2026-05-22): the original crashing prompt ("what about system 1 and system 2 framing?") now completes cleanly end-to-end. mimir streamed 3193 tokens (50 seconds, 374980-token context), `[done] turn_id=96 duration_ms=50436`. Empty-data frames somewhere in the stream silently skipped; no crash. 172/172 tests GREEN; ruff clean; all 5 issue contracts (#1, #3, #4, #5, #7) drift-check clean. Persistent-memory updated per the commit-along rule: status reflects v0+#7 milestone; new dated decisions for #5/#6/#7 filing + #7 implementation; foot-gun entry for unguarded json.loads(sse.data).
This commit is contained in:
@@ -694,3 +694,145 @@ class TestCancelTurn:
|
||||
assert first.cancelled is True
|
||||
with pytest.raises(CancelAlreadyCompleted):
|
||||
await cancel_turn(client, "s1", 42)
|
||||
|
||||
|
||||
# ============================================================================
|
||||
# Issue #7: empty-data skip + MalformedSseData raise
|
||||
# ============================================================================
|
||||
|
||||
|
||||
def _sse_empty_chunk(sse_id: str) -> bytes:
|
||||
"""SSE frame with id but empty data (server-emitted keepalive shape)."""
|
||||
return f"id: {sse_id}\ndata:\n\n".encode()
|
||||
|
||||
|
||||
def _sse_raw_chunk(sse_id: str, raw_data: str) -> bytes:
|
||||
"""SSE frame with id + arbitrary raw data (for testing malformed JSON)."""
|
||||
return f"id: {sse_id}\ndata: {raw_data}\n\n".encode()
|
||||
|
||||
|
||||
class TestEmptyDataSkipped:
|
||||
@respx.mock
|
||||
async def test_empty_data_skipped(self) -> None:
|
||||
"""empty_data_skipped [trace]: 4 frames in, 3 events out; skip preserves last_sse_id."""
|
||||
from ratatoskr.sse_client import Done as _Done
|
||||
from ratatoskr.sse_client import Text as _Text
|
||||
|
||||
stream = (
|
||||
_sse_chunk("42:1", {"type": "text", "content": "first"})
|
||||
+ _sse_empty_chunk("42:2") # ← skipped silently
|
||||
+ _sse_chunk("42:3", {"type": "text", "content": "second"})
|
||||
+ _sse_chunk("42:4", _DONE_42_6)
|
||||
)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
events = [e async for e in stream_turn(client, "s1", "hi")]
|
||||
# Exactly 3 events: Text, Text, Done — empty-data frame at 42:2 is invisible
|
||||
assert len(events) == 3
|
||||
assert isinstance(events[0], _Text) and events[0].sse_id == SseId(42, 1)
|
||||
assert isinstance(events[1], _Text) and events[1].sse_id == SseId(42, 3)
|
||||
assert isinstance(events[2], _Done) and events[2].sse_id == SseId(42, 4)
|
||||
# Per INV-003: skip MUST NOT advance through 42:2. The second Text's sse_id
|
||||
# is (42, 3) — directly verifies the skip didn't bookkeep 42:2.
|
||||
assert events[1].sse_id.seq == 3, "skip advanced through 42:2"
|
||||
|
||||
@respx.mock
|
||||
async def test_empty_data_skip_preserves_last_seen_sse_id(self) -> None:
|
||||
"""empty_skip_does_not_advance [trace]: drop-after-empty → last_seen is last real event."""
|
||||
from ratatoskr.sse_client import SseConnectionDropped
|
||||
|
||||
# Stream: text(42:1), empty(42:2), then drop. Per INV-003, the skipped
|
||||
# 42:2 must NOT advance internal last_sse_id. If the consumer caught a
|
||||
# drop, SseConnectionDropped.last_seen_sse_id should be (42, 1) — the
|
||||
# last *real* event — NOT (42, 2).
|
||||
first = _sse_chunk("42:1", {"type": "text", "content": "x"})
|
||||
empty = _sse_empty_chunk("42:2")
|
||||
|
||||
class _DropAfterEmpty(httpx.AsyncByteStream):
|
||||
async def __aiter__(self): # type: ignore[no-untyped-def]
|
||||
yield first
|
||||
yield empty
|
||||
raise httpx.RemoteProtocolError("drop after skip")
|
||||
|
||||
async def aclose(self) -> None:
|
||||
return None
|
||||
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200,
|
||||
headers={"content-type": "text/event-stream"},
|
||||
stream=_DropAfterEmpty(),
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(SseConnectionDropped) as exc_info:
|
||||
async for _ in stream_turn(client, "s1", "hi"):
|
||||
pass
|
||||
assert exc_info.value.last_seen_sse_id == SseId(42, 1), (
|
||||
f"skip advanced last_sse_id through 42:2; got {exc_info.value.last_seen_sse_id}"
|
||||
)
|
||||
|
||||
@respx.mock
|
||||
async def test_malformed_data_raises(self) -> None:
|
||||
"""malformed_data_raises [error]: text + bad-JSON → yields Text then MalformedSseData."""
|
||||
from ratatoskr.sse_client import MalformedSseData
|
||||
from ratatoskr.sse_client import Text as _Text
|
||||
|
||||
stream = (
|
||||
_sse_chunk("42:1", {"type": "text", "content": "hi"})
|
||||
+ _sse_raw_chunk("42:2", "not-json")
|
||||
)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
yielded: list = []
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(MalformedSseData) as exc_info:
|
||||
async for e in stream_turn(client, "s1", "hi"):
|
||||
yielded.append(e)
|
||||
assert len(yielded) == 1
|
||||
assert isinstance(yielded[0], _Text)
|
||||
assert exc_info.value.raw == "not-json"
|
||||
|
||||
@respx.mock
|
||||
async def test_whitespace_data_raises(self) -> None:
|
||||
"""whitespace_data_raises [adv]: single-space data → MalformedSseData (NOT skipped)."""
|
||||
from ratatoskr.sse_client import MalformedSseData
|
||||
|
||||
stream = _sse_raw_chunk("42:1", " ") # single space — non-empty, JSON-invalid
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(MalformedSseData):
|
||||
async for _ in stream_turn(client, "s1", "hi"):
|
||||
pass
|
||||
|
||||
@respx.mock
|
||||
async def test_malformed_data_truncation(self) -> None:
|
||||
"""malformed_data_truncation [security]: 5000-char bad data → raw truncated to 200."""
|
||||
from ratatoskr.sse_client import MalformedSseData
|
||||
|
||||
huge_bad = "x" * 5000 # not JSON; very long
|
||||
stream = _sse_raw_chunk("42:1", huge_bad)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(MalformedSseData) as exc_info:
|
||||
async for _ in stream_turn(client, "s1", "hi"):
|
||||
pass
|
||||
assert len(exc_info.value.raw) == 200
|
||||
assert exc_info.value.raw == "x" * 200
|
||||
# Exception message also only contains the truncated form
|
||||
assert "x" * 5000 not in str(exc_info.value)
|
||||
|
||||
Reference in New Issue
Block a user