"""Tests for ratatoskr.tui per docs/contracts/issues/4.contract.md.""" from unittest.mock import MagicMock import httpx import pytest import respx from textual.widgets import RichLog from ratatoskr.cli import ParsedArgs from ratatoskr.sse_client import ( Cancelled, Done, Error, SseId, Text, TextBoundary, Thinking, ToolResult, ToolStart, WorkerPhase, ) from ratatoskr.tui import RatatoskrApp, _cancel_via_sse, _render_event_to_log _CANCEL_OK_RESP = {"turn_id": 42, "cancelled": True, "reason": None, "partial_message_id": None} _CREATE_OK_RESP = { "session_id": "s-new12345", "agent_id": "mimir", "message_count": 0, "created_at": "2026-05-21T00:00:00+00:00", "last_active": "2026-05-21T00:00:00+00:00", "metadata": {}, } def _args_new(**overrides) -> ParsedArgs: base = dict( send_content=None, session_id=None, new=True, agent_id="mimir", api_key="k", server_url="https://w.example", raw=False, ) base.update(overrides) return ParsedArgs(**base) def _args_existing(session_id: str = "s-1existing", **overrides) -> ParsedArgs: base = dict( send_content=None, session_id=session_id, new=False, agent_id=None, api_key="k", server_url="https://w.example", raw=False, ) base.update(overrides) return ParsedArgs(**base) def _spy_writes(monkeypatch) -> list: """Patch RichLog.write to record every arg into a list (returned).""" writes: list = [] original = RichLog.write def spy(self, content, **kw): writes.append(content) return original(self, content, **kw) monkeypatch.setattr(RichLog, "write", spy) return writes def _resolved_app( args: ParsedArgs, *, session_id: str | None = None, agent_id: str | None = None, client: httpx.AsyncClient | None = None, ) -> RatatoskrApp: """Construct RatatoskrApp with pre-resolved state (issue #6 lifecycle). Production path: `run_tui` → `_resolve_then_run` opens AsyncClient, mints or attaches session, then constructs the App with the resolved tuple. This helper inlines that shape so tests bypass the pre-flight without re-implementing it. The client is opened here (and leaks at test teardown — acceptable; respx mocks all network calls and pytest exits cleanly). """ sid = session_id if session_id is not None else (args.session_id or "s-default") aid = args.agent_id if agent_id is None else agent_id if client is None: client = httpx.AsyncClient( base_url=args.server_url, headers={"Authorization": f"Bearer {args.api_key}"}, timeout=httpx.Timeout(connect=10.0, read=None, write=10.0, pool=10.0), ) return RatatoskrApp(args, session_id=sid, agent_id=aid, client=client) SID = SseId(42, 5) class TestRenderEventToLog: def test_text_renders_raw_delta(self) -> None: """text_renders_raw_delta [happy,tracer]: Text → log.write('hello').""" log = MagicMock() _render_event_to_log(Text(sse_id=SID, content="hello"), log=log, raw=False) log.write.assert_called_once_with("hello") def test_done_renders_label_only(self) -> None: """done_renders_label_only: …""" log = MagicMock() evt = Done( sse_id=SID, phase="completed", response="hi there", model="glm5-turbo", duration_ms=1234, usage={"prompt": 1, "completion": 2}, ) _render_event_to_log(evt, log=log, raw=False) log.write.assert_called_once() line = log.write.call_args[0][0] assert line.startswith("[done]") assert "turn_id=42" in line assert "model=glm5-turbo" in line # POST-003: the Done line is labels only; markdown render is the caller's job assert "hi there" not in line def test_error_renders_label(self) -> None: """error_renders_label: Error → log line starts with [error].""" log = MagicMock() evt = Error(sse_id=SID, phase="failed", message="boom", error_code="llm_output_invalid") _render_event_to_log(evt, log=log, raw=False) line = log.write.call_args[0][0] assert line.startswith("[error]") assert "turn_id=42" in line assert "code=llm_output_invalid" in line def test_cancelled_renders_label(self) -> None: """cancelled_renders_label: Cancelled → log line starts with [cancelled].""" log = MagicMock() evt = Cancelled( sse_id=SID, phase="cancelled", turn_id=42, reason="user", partial_message_id=7 ) _render_event_to_log(evt, log=log, raw=False) line = log.write.call_args[0][0] assert line.startswith("[cancelled]") assert "reason='user'" in line assert "partial_message_id=7" in line def test_worker_phase_renders_label(self) -> None: """worker_phase_renders_label: WorkerPhase → log line starts with [worker_phase].""" log = MagicMock() evt = WorkerPhase(sse_id=SID, phase="streaming", turn_id=42) _render_event_to_log(evt, log=log, raw=False) line = log.write.call_args[0][0] assert line.startswith("[worker_phase]") assert "phase=streaming" in line def test_thinking_truncated(self) -> None: """thinking_truncated [trace]: …""" log = MagicMock() _render_event_to_log(Thinking(sse_id=SID, content="a" * 500), log=log, raw=False) line = log.write.call_args[0][0] assert line.startswith("[thinking]") assert "a" * 500 not in line assert "a" * 200 in line def test_tool_start_renders_label(self) -> None: """tool_start_renders_label: ToolStart → [tool_start] name=... args=...""" log = MagicMock() evt = ToolStart(sse_id=SID, name="read_file", arguments={"path": "/x"}) _render_event_to_log(evt, log=log, raw=False) line = log.write.call_args[0][0] assert line.startswith("[tool_start] name=read_file args=") def test_tool_result_truncated(self) -> None: """tool_result_truncated [trace]: …""" log = MagicMock() evt = ToolResult(sse_id=SID, name="x", result="b" * 500, duration_ms=42) _render_event_to_log(evt, log=log, raw=False) line = log.write.call_args[0][0] assert line.startswith("[tool_result]") # The whole repr-portion of the result is truncated to 200; the full 500-b # string can never fit in line whole. assert "b" * 500 not in line def test_text_boundary_renders_label(self) -> None: """text_boundary_renders_label: TextBoundary → [text_boundary] kind=... char_offset=...""" log = MagicMock() evt = TextBoundary(sse_id=SID, kind="sentence", char_offset=128, ts="2026-05-21T00:00:00Z") _render_event_to_log(evt, log=log, raw=False) line = log.write.call_args[0][0] assert line.startswith("[text_boundary]") assert "kind=sentence" in line assert "char_offset=128" in line class TestCancelViaSse: @respx.mock async def test_happy_cancel(self) -> None: """happy_cancel [happy,tracer]: 200 OK → returns None; log has no [cancel_failed].""" respx.post("https://w.example/sessions/s-1/turns/42/cancel").mock( return_value=httpx.Response(200, json=_CANCEL_OK_RESP) ) log = MagicMock() async with httpx.AsyncClient(base_url="https://w.example") as client: result = await _cancel_via_sse(client, "s-1", 42, log=log) assert result is None log.write.assert_not_called() @respx.mock async def test_cancel_failed_500(self) -> None: """cancel_failed_500 [error]: …""" respx.post("https://w.example/sessions/s-1/turns/42/cancel").mock( return_value=httpx.Response(500, content=b"boom") ) log = MagicMock() async with httpx.AsyncClient(base_url="https://w.example") as client: await _cancel_via_sse(client, "s-1", 42, log=log) line = log.write.call_args[0][0] assert "[cancel_failed]" in line assert "CancelFailed" in line @respx.mock async def test_cancel_already_completed(self) -> None: """cancel_already_completed [scenario]: 409 → '[cancel_failed] CancelAlreadyCompleted:'.""" respx.post("https://w.example/sessions/s-1/turns/42/cancel").mock( return_value=httpx.Response(409) ) log = MagicMock() async with httpx.AsyncClient(base_url="https://w.example") as client: await _cancel_via_sse(client, "s-1", 42, log=log) line = log.write.call_args[0][0] assert "[cancel_failed]" in line assert "CancelAlreadyCompleted" in line @respx.mock async def test_transport_error_swallowed(self) -> None: """transport_error_swallowed [error]: …""" respx.post("https://w.example/sessions/s-1/turns/42/cancel").mock( side_effect=httpx.ConnectError("network down") ) log = MagicMock() async with httpx.AsyncClient(base_url="https://w.example") as client: await _cancel_via_sse(client, "s-1", 42, log=log) line = log.write.call_args[0][0] assert "[cancel_failed]" in line assert "ConnectError" in line class TestAppMount: """on_mount narrows per issue #6: only identity-widget population. Session resolution + AsyncClient open + error-on-resolve are exercised at the `_resolve_then_run` layer (see TestResolveThenRun); only happy mount paths remain here, exercised with pre-resolved state via _resolved_app. """ async def test_happy_new_session_mount(self) -> None: """happy_new_session_mount [happy,tracer]: identity populated from pre-resolved state.""" app = _resolved_app(_args_new(), session_id="s-new12345", agent_id="mimir") async with app.run_test() as pilot: await pilot.pause() assert app.session_id == "s-new12345" assert app.agent_id == "mimir" assert app.state == "idle" # INV-002: session identity visible — agent_id + last 8 of session_id assert "mimir" in (app.sub_title or "") assert app.session_id[-8:] in (app.sub_title or "") async def test_happy_existing_session_mount(self) -> None: """happy_existing_session_mount: identity shows when agent_id is None.""" app = _resolved_app(_args_existing(session_id="s-existing-tail8x")) async with app.run_test() as pilot: await pilot.pause() assert app.session_id == "s-existing-tail8x" assert app.state == "idle" # INV-002 carve-out: agent unknown → · … assert "" in (app.sub_title or "") assert app.session_id[-8:] in (app.sub_title or "") async def test_footer_identity_visible_first_frame(self) -> None: """footer_identity_visible_first_frame [trace]: identity widget rendered first frame.""" from textual.widgets import Static app = _resolved_app(_args_new(), session_id="s-new12345", agent_id="mimir") async with app.run_test() as pilot: await pilot.pause() identity_widget = app.query_one("#identity", Static) rendered = str(identity_widget.render()) assert "mimir" in rendered assert "·" in rendered assert app.session_id[-8:] in rendered import asyncio # noqa: E402 from textual.widgets import Input # noqa: E402 async def _noop_worker(self, content: str) -> None: """Fake _stream_turn_worker that never completes (lets state stay 'streaming').""" await asyncio.Future() # await forever; cancelled when test exits class TestOnInputSubmitted: @respx.mock async def test_happy_submit_echoes_and_spawns( self, monkeypatch: pytest.MonkeyPatch ) -> None: """happy_submit_echoes_and_spawns [happy,tracer]: …""" monkeypatch.setattr(RatatoskrApp, "_stream_turn_worker", _noop_worker) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) inp.value = "hello" await inp.action_submit() await pilot.pause() assert any("❯ hello" in str(w) for w in writes) # noqa: RUF001 assert inp.value == "" assert app.state == "streaming" assert app.stream_worker is not None @respx.mock async def test_empty_submit_no_op( self, monkeypatch: pytest.MonkeyPatch ) -> None: """empty_submit_no_op [trace]: '' + Enter → no change; no worker spawned.""" monkeypatch.setattr(RatatoskrApp, "_stream_turn_worker", _noop_worker) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() # Spy AFTER mount so identity-widget writes (if any) aren't counted. writes = _spy_writes(monkeypatch) inp = app.query_one("#prompt", Input) inp.value = "" await inp.action_submit() await pilot.pause() assert app.state == "idle" assert app.stream_worker is None # POST: no RichLog write fires on empty submit assert writes == [] @respx.mock async def test_submit_during_streaming_shows_busy_notice( self, monkeypatch: pytest.MonkeyPatch ) -> None: """submit_during_streaming_shows_busy_notice [adversarial]: …""" monkeypatch.setattr(RatatoskrApp, "_stream_turn_worker", _noop_worker) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) # First submit: enters streaming inp.value = "first" await inp.action_submit() await pilot.pause() first_worker = app.stream_worker assert app.state == "streaming" # Second submit while streaming → busy notice; no new worker writes.clear() inp.value = "second" await inp.action_submit() await pilot.pause() assert any("[busy] turn in flight; input ignored" in str(w) for w in writes) assert app.stream_worker is first_worker # unchanged assert app.state == "streaming" assert inp.value == "" @respx.mock async def test_submit_during_cancelling_shows_busy_notice( self, monkeypatch: pytest.MonkeyPatch ) -> None: """submit_during_cancelling_shows_busy_notice [adversarial]: …""" monkeypatch.setattr(RatatoskrApp, "_stream_turn_worker", _noop_worker) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() app.state = "cancelling" # bypass the natural transition for the test assert app.stream_worker is None # no live worker before non-idle submit inp = app.query_one("#prompt", Input) inp.value = "x" await inp.action_submit() await pilot.pause() assert any("[busy]" in str(w) for w in writes) assert app.state == "cancelling" # POST-005 (from issue #4 on_input_submitted contract): # input cleared; NO new worker spawned during non-idle submit. assert inp.value == "" assert app.stream_worker is None @respx.mock async def test_footer_hint_flips_to_cancel( self, monkeypatch: pytest.MonkeyPatch ) -> None: """footer_hint_flips_to_cancel [trace]: hint widget shows 'Ctrl-C to cancel'.""" from textual.widgets import Static monkeypatch.setattr(RatatoskrApp, "_stream_turn_worker", _noop_worker) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() hint_widget = app.query_one("#hint", Static) assert str(hint_widget.render()) == RatatoskrApp.HINT_IDLE inp = app.query_one("#prompt", Input) inp.value = "hi" await inp.action_submit() await pilot.pause() assert str(hint_widget.render()) == RatatoskrApp.HINT_STREAMING import json # noqa: E402 def _sse_chunk(sse_id: str, body: dict) -> bytes: return f"id: {sse_id}\ndata: {json.dumps(body)}\n\n".encode() _DONE_BODY = { "type": "done", "phase": "succeeded", "response": "hello", "model": "m", "duration_ms": 1, "usage": { "prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0, "cached_input_tokens": 0, }, } _CANCELLED_STREAM_BODY = { "type": "cancelled", "phase": "cancelled", "turn_id": 42, "reason": "user_cancel", "partial_message_id": None, } def _sse_resp(body: bytes | httpx.AsyncByteStream) -> httpx.Response: headers = {"content-type": "text/event-stream"} if isinstance(body, bytes): return httpx.Response(200, headers=headers, content=body) return httpx.Response(200, headers=headers, stream=body) async def _submit_and_wait(app: RatatoskrApp, pilot, content: str) -> None: """Type content into the input and submit; wait for worker to finish.""" inp = app.query_one("#prompt", Input) inp.value = content await inp.action_submit() await pilot.pause() # let the Input.Submitted message dispatch # Poll until the worker resolves (state returns to idle) for _ in range(100): if app.state == "idle" and app.stream_worker is not None: return await pilot.pause(0.02) class TestStreamTurnWorker: @respx.mock async def test_happy_text_done_renders_markdown( self, monkeypatch: pytest.MonkeyPatch ) -> None: """happy_text_done_renders_markdown [happy,tracer]: …""" stream = ( _sse_chunk("42:1", {"type": "text", "content": "hello"}) + _sse_chunk("42:2", _DONE_BODY) ) respx.post( "https://w.example/sessions/s-1existing/messages" ).mock(return_value=_sse_resp(stream)) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "hi") assert app.state == "idle" # Streamed delta + done label + rule + markdown render assert any(w == "hello" for w in writes) assert any("[done]" in str(w) for w in writes) # The post-Done markdown render uses rich Rule + Markdown — non-string writes. # INV-005: BOTH separator (Rule) AND markdown render must be present in non-raw. from rich.markdown import Markdown from rich.rule import Rule assert any(isinstance(w, Markdown) for w in writes) assert any(isinstance(w, Rule) for w in writes) @respx.mock async def test_raw_flag_skips_markdown_render( self, monkeypatch: pytest.MonkeyPatch ) -> None: """raw_flag_skips_markdown_render [trace]: …""" stream = ( _sse_chunk("42:1", {"type": "text", "content": "hi"}) + _sse_chunk("42:2", _DONE_BODY) ) respx.post( "https://w.example/sessions/s-1existing/messages" ).mock(return_value=_sse_resp(stream)) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing(raw=True)) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "x") # INV-005: with --raw, NEITHER Rule separator NOR Markdown render appears. from rich.markdown import Markdown from rich.rule import Rule assert not any(isinstance(w, Markdown) for w in writes) assert not any(isinstance(w, Rule) for w in writes) @respx.mock async def test_error_terminal_returns_to_idle( self, monkeypatch: pytest.MonkeyPatch ) -> None: """error_terminal_returns_to_idle [happy]: …""" stream = ( _sse_chunk("42:1", {"type": "text", "content": "x"}) + _sse_chunk( "42:2", { "type": "error", "phase": "failed", "error_code": "llm_output_invalid", "message": "boom", }, ) ) respx.post( "https://w.example/sessions/s-1existing/messages" ).mock(return_value=_sse_resp(stream)) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "x") assert app.state == "idle" assert any("[error]" in str(w) for w in writes) @respx.mock async def test_cancelled_terminal_returns_to_idle( self, monkeypatch: pytest.MonkeyPatch ) -> None: """cancelled_terminal_returns_to_idle [happy]: …""" stream = ( _sse_chunk("42:1", {"type": "text", "content": "x"}) + _sse_chunk("42:2", _CANCELLED_STREAM_BODY) ) respx.post( "https://w.example/sessions/s-1existing/messages" ).mock(return_value=_sse_resp(stream)) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "x") assert app.state == "idle" assert any("[cancelled]" in str(w) for w in writes) @respx.mock async def test_active_turn_id_set_on_first_event( self, monkeypatch: pytest.MonkeyPatch ) -> None: """active_turn_id_set_on_first_event [trace]: …""" # Use a gated stream: yield first event, then hold, so we can inspect mid-stream first = _sse_chunk("42:1", {"type": "text", "content": "x"}) gate = asyncio.Event() class _GatedAfterFirst(httpx.AsyncByteStream): async def __aiter__(self): yield first await gate.wait() yield _sse_chunk("42:2", _DONE_BODY) async def aclose(self) -> None: return None respx.post("https://w.example/sessions/s-1existing/messages").mock( return_value=_sse_resp(_GatedAfterFirst()) ) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) inp.value = "x" await inp.action_submit() # Wait for first event to be processed (active_turn_id set) for _ in range(50): if app.active_turn_id is not None: break await pilot.pause(0.02) assert app.active_turn_id == 42 # Release the gate so the worker can finish and the app can shut down cleanly gate.set() for _ in range(50): if app.state == "idle": break await pilot.pause(0.02) @respx.mock async def test_sse_connect_failed_returns_to_idle( self, monkeypatch: pytest.MonkeyPatch ) -> None: """sse_connect_failed_returns_to_idle [error]: …""" respx.post("https://w.example/sessions/s-1existing/messages").mock( return_value=httpx.Response(404, json={"error": "session_not_found"}) ) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "x") assert app.state == "idle" assert any("[sse_connect_failed]" in str(w) for w in writes) assert app.return_value is None # app NOT exited per INV-008 @respx.mock async def test_connection_dropped_returns_to_idle( self, monkeypatch: pytest.MonkeyPatch ) -> None: """connection_dropped_returns_to_idle [error]: …""" class _DropAfter(httpx.AsyncByteStream): async def __aiter__(self): yield _sse_chunk("42:1", {"type": "text", "content": "x"}) raise httpx.RemoteProtocolError("drop") async def aclose(self) -> None: return None respx.post("https://w.example/sessions/s-1existing/messages").mock( return_value=_sse_resp(_DropAfter()) ) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "x") assert app.state == "idle" assert any("[connection_dropped]" in str(w) for w in writes) @respx.mock async def test_malformed_sse_data_returns_to_idle( self, monkeypatch: pytest.MonkeyPatch ) -> None: """malformed_sse_data_returns_to_idle [error]: bad-JSON → [malformed_sse_data]; idle.""" stream = ( _sse_chunk("42:1", {"type": "text", "content": "x"}) + b"id: 42:2\ndata: not-json\n\n" ) respx.post( "https://w.example/sessions/s-1existing/messages" ).mock(return_value=_sse_resp(stream)) writes = _spy_writes(monkeypatch) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "x") assert app.state == "idle" assert any("[malformed_sse_data]" in str(w) for w in writes) assert any("not-json" in str(w) for w in writes) # INV-008: mid-session error does NOT exit the app assert app.return_value is None @respx.mock async def test_rendered_event_per_event( self, monkeypatch: pytest.MonkeyPatch ) -> None: """rendered_event_per_event [trace]: …""" chunks = ( _sse_chunk("42:1", {"type": "worker_phase", "phase": "streaming", "turn_id": 42}) + _sse_chunk("42:2", {"type": "text", "content": "hi"}) + _sse_chunk("42:3", _DONE_BODY) ) respx.post( "https://w.example/sessions/s-1existing/messages" ).mock(return_value=_sse_resp(chunks)) from ratatoskr import tui as tui_mod call_count = 0 original = tui_mod._render_event_to_log def spy(event, *, log, raw): nonlocal call_count call_count += 1 return original(event, log=log, raw=raw) monkeypatch.setattr(tui_mod, "_render_event_to_log", spy) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await _submit_and_wait(app, pilot, "x") assert call_count == 3 class TestActionInterrupt: @respx.mock async def test_idle_ctrl_c_exits_zero(self) -> None: """idle_ctrl_c_exits_zero [happy,tracer]: state=idle; ctrl+c → exit(0).""" app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() assert app.state == "idle" await pilot.press("ctrl+c") await pilot.pause() assert app.return_value == 0 @respx.mock async def test_streaming_first_ctrl_c_cancels( self, monkeypatch: pytest.MonkeyPatch ) -> None: """streaming_first_ctrl_c_cancels [scenario,tracer]: …""" # Stream that yields one text event (sets active_turn_id) then waits forever first_chunk = _sse_chunk("42:1", {"type": "text", "content": "x"}) gate = asyncio.Event() class _GatedAfterFirst(httpx.AsyncByteStream): async def __aiter__(self): yield first_chunk await gate.wait() async def aclose(self) -> None: return None respx.post("https://w.example/sessions/s-1existing/messages").mock( return_value=_sse_resp(_GatedAfterFirst()) ) cancel_observed = asyncio.Event() def cancel_handler(req: httpx.Request) -> httpx.Response: cancel_observed.set() return httpx.Response(200, json=_CANCEL_OK_RESP) cancel_route = respx.post("https://w.example/sessions/s-1existing/turns/42/cancel").mock( side_effect=cancel_handler ) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) inp.value = "go" await inp.action_submit() await pilot.pause() # Wait for active_turn_id to be set (first event consumed) for _ in range(50): if app.active_turn_id == 42: break await pilot.pause(0.02) assert app.active_turn_id == 42 assert app.state == "streaming" await pilot.press("ctrl+c") # Wait for cancel POST to land for _ in range(50): if cancel_observed.is_set(): break await pilot.pause(0.02) assert cancel_route.call_count == 1 assert app.state == "cancelling" from textual.widgets import Static hint_widget = app.query_one("#hint", Static) assert str(hint_widget.render()) == RatatoskrApp.HINT_CANCELLING # Release the gate so the stream worker can finish cleanly during teardown gate.set() @respx.mock async def test_streaming_no_turn_id_force_exits( self, monkeypatch: pytest.MonkeyPatch ) -> None: """streaming_no_turn_id_force_exits [scenario]: …""" cancel_route = respx.post("https://w.example/sessions/s-1existing/turns/0/cancel").mock( return_value=httpx.Response(200, json=_CANCEL_OK_RESP) ) # Stream that hangs forever (no events to set active_turn_id) gate = asyncio.Event() class _NeverYields(httpx.AsyncByteStream): async def __aiter__(self): await gate.wait() if False: yield b"" async def aclose(self) -> None: return None respx.post("https://w.example/sessions/s-1existing/messages").mock( return_value=_sse_resp(_NeverYields()) ) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) inp.value = "go" await inp.action_submit() await pilot.pause() assert app.state == "streaming" assert app.active_turn_id is None # Capture worker reference + spy on its .cancel() before ctrl+c worker_ref = app.stream_worker assert worker_ref is not None cancel_calls: list = [] original_cancel = type(worker_ref).cancel monkeypatch.setattr( type(worker_ref), "cancel", lambda self: (cancel_calls.append(self), original_cancel(self))[-1], ) await pilot.press("ctrl+c") await pilot.pause() gate.set() # let the gated stream resolve so teardown is clean assert app.return_value == 3 assert cancel_route.call_count == 0 # action_interrupt MUST cancel the stream worker on the no-active_turn_id force-exit path assert worker_ref in cancel_calls @respx.mock async def test_cancelling_second_ctrl_c_force_exits( self, monkeypatch: pytest.MonkeyPatch ) -> None: """cancelling_second_ctrl_c_force_exits [scenario]: …""" # Set up a real live stream worker (gated, hangs forever) so we can # observe action_interrupt's cancel() call on the second-Ctrl-C path. monkeypatch.setattr(RatatoskrApp, "_stream_turn_worker", _noop_worker) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) inp.value = "go" await inp.action_submit() await pilot.pause() assert app.stream_worker is not None app.state = "cancelling" # bypass the natural transition for the test worker_ref = app.stream_worker cancel_calls: list = [] original_cancel = type(worker_ref).cancel monkeypatch.setattr( type(worker_ref), "cancel", lambda self: (cancel_calls.append(self), original_cancel(self))[-1], ) await pilot.press("ctrl+c") await pilot.pause() assert app.return_value == 3 # Second-Ctrl-C in cancelling state MUST cancel the in-flight worker assert worker_ref in cancel_calls @respx.mock async def test_cancel_failed_swallowed(self) -> None: """cancel_failed_swallowed [scenario]: …""" first_chunk = _sse_chunk("42:1", {"type": "text", "content": "x"}) stream_gate = asyncio.Event() class _GatedAfterFirst(httpx.AsyncByteStream): async def __aiter__(self): yield first_chunk await stream_gate.wait() async def aclose(self) -> None: return None respx.post("https://w.example/sessions/s-1existing/messages").mock( return_value=_sse_resp(_GatedAfterFirst()) ) cancel_observed = asyncio.Event() def cancel_handler(req: httpx.Request) -> httpx.Response: cancel_observed.set() return httpx.Response(500, content=b"boom") respx.post("https://w.example/sessions/s-1existing/turns/42/cancel").mock( side_effect=cancel_handler ) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) inp.value = "go" await inp.action_submit() await pilot.pause() for _ in range(50): if app.active_turn_id == 42: break await pilot.pause(0.02) await pilot.press("ctrl+c") for _ in range(50): if cancel_observed.is_set(): break await pilot.pause(0.02) # Give _cancel_via_sse time to write the [cancel_failed] line await pilot.pause(0.05) log = app.query_one("#transcript", RichLog) rendered = "\n".join(str(strip.text) for strip in log.lines) assert "[cancel_failed]" in rendered assert app.state == "cancelling" stream_gate.set() # let stream finish for teardown class TestActionQuit: @respx.mock async def test_idle_ctrl_d_exits_zero(self) -> None: """idle_ctrl_d_exits_zero [happy,tracer]: state=idle; ctrl+d → exit(0).""" app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() await pilot.press("ctrl+d") await pilot.pause() assert app.return_value == 0 @respx.mock async def test_streaming_ctrl_d_force_exits( self, monkeypatch: pytest.MonkeyPatch ) -> None: """streaming_ctrl_d_force_exits [scenario]: …""" cancel_route = respx.post("https://w.example/sessions/s-1existing/turns/42/cancel").mock( return_value=httpx.Response(200, json=_CANCEL_OK_RESP) ) first_chunk = _sse_chunk("42:1", {"type": "text", "content": "x"}) gate = asyncio.Event() class _GatedAfterFirst(httpx.AsyncByteStream): async def __aiter__(self): yield first_chunk await gate.wait() async def aclose(self) -> None: return None respx.post("https://w.example/sessions/s-1existing/messages").mock( return_value=_sse_resp(_GatedAfterFirst()) ) app = _resolved_app(_args_existing()) async with app.run_test() as pilot: await pilot.pause() inp = app.query_one("#prompt", Input) inp.value = "go" await inp.action_submit() await pilot.pause() for _ in range(50): if app.active_turn_id == 42: break await pilot.pause(0.02) # Capture worker + spy on cancel before ctrl+d worker_ref = app.stream_worker assert worker_ref is not None cancel_calls: list = [] original_cancel = type(worker_ref).cancel monkeypatch.setattr( type(worker_ref), "cancel", lambda self: (cancel_calls.append(self), original_cancel(self))[-1], ) await pilot.press("ctrl+d") await pilot.pause() gate.set() assert app.return_value == 0 assert cancel_route.call_count == 0 # POST-002: Ctrl-D MUST cancel the in-flight stream worker (abandon-and-exit) assert worker_ref in cancel_calls from ratatoskr.tui import run_tui # noqa: E402 class TestResolveThenRun: """Tests at the `_resolve_then_run` layer — pre-`App.run()` session resolution + AsyncClient ownership + stderr error routing per issue #6. """ @respx.mock def test_happy_new_session_resolve(self, monkeypatch: pytest.MonkeyPatch) -> None: """happy_new_session_resolve [happy]: --new path through _resolve_then_run. Verifies POST /sessions count + SessionInfo propagation to RatatoskrApp's pre-resolved state (session_id / agent_id / client). The corresponding TestAppMount.test_happy_new_session_mount uses _resolved_app and bypasses _resolve_then_run entirely; this test exercises the production resolve path with a real POST /sessions mock. """ sessions_route = respx.post("https://w.example/sessions").mock( return_value=httpx.Response(201, json=_CREATE_OK_RESP) ) snapshot: dict = {} async def capture_run_async(self, *a, **kw): snapshot["session_id"] = self.session_id snapshot["agent_id"] = self.agent_id snapshot["client"] = self.client snapshot["client_open"] = not self.client.is_closed return 0 monkeypatch.setattr(RatatoskrApp, "run_async", capture_run_async) rc = run_tui(_args_new()) assert rc == 0 # Exactly one POST /sessions invocation by _resolve_then_run assert sessions_route.call_count == 1 # SessionInfo fields propagated into the constructed App assert snapshot["session_id"] == "s-new12345" assert snapshot["agent_id"] == "mimir" assert snapshot["client"] is not None assert snapshot["client_open"] is True @respx.mock def test_happy_new_with_end_user_id_resolve( self, monkeypatch: pytest.MonkeyPatch ) -> None: """happy_new_with_end_user_id_resolve [happy]: args.end_user_id threads into POST body. Issue #5 amends #4: _resolve_then_run's create_session call now forwards args.end_user_id (renamed from the contract's _mount target, since #6 moved session resolution out of on_mount into _resolve_then_run). """ import json as _json sessions_route = respx.post("https://w.example/sessions").mock( return_value=httpx.Response(201, json=_CREATE_OK_RESP) ) async def fake_run_async(self, *a, **kw): return 0 monkeypatch.setattr(RatatoskrApp, "run_async", fake_run_async) rc = run_tui(_args_new(end_user_id="alice")) assert rc == 0 assert sessions_route.call_count == 1 body = _json.loads(sessions_route.calls[0].request.content) assert body == {"agent_id": "mimir", "end_user_id": "alice"} @respx.mock def test_user_agent_header_sent(self, monkeypatch: pytest.MonkeyPatch) -> None: """user_agent_header_sent [trace]: outbound requests carry the ratatoskr User-Agent. Worldtree-dev (althing 2026-05-23) requested consumers send `User-Agent: ratatoskr/ ()` so server logs can distinguish ratatoskr traffic from other consumers. """ sessions_route = respx.post("https://w.example/sessions").mock( return_value=httpx.Response(201, json=_CREATE_OK_RESP) ) async def fake_run_async(self, *a, **kw): return 0 monkeypatch.setattr(RatatoskrApp, "run_async", fake_run_async) rc = run_tui(_args_new()) assert rc == 0 ua = sessions_route.calls[0].request.headers["User-Agent"] assert ua.startswith("ratatoskr/") assert "vh@phasefinal.com" in ua @respx.mock def test_alt_screen_never_opens_on_resolve_error( self, monkeypatch: pytest.MonkeyPatch ) -> None: """alt_screen_never_opens_on_resolve_error [trace]: 404 → run_tui=12; run_async unhit. Directly probes INV-001: session resolution failures MUST short-circuit BEFORE the alt-screen opens. """ respx.post("https://w.example/sessions").mock( return_value=httpx.Response(404, json={"error": "unknown_agent_id"}) ) sentinel_called = False async def sentinel(self, *a, **kw): nonlocal sentinel_called sentinel_called = True return 0 monkeypatch.setattr(RatatoskrApp, "run_async", sentinel) rc = run_tui(_args_new()) assert rc == 12 assert not sentinel_called @respx.mock def test_agent_not_found_on_resolve( self, capsys: pytest.CaptureFixture[str] ) -> None: """agent_not_found_on_resolve [error]: --new + 404 → stderr [agent_not_found]; exit 12.""" respx.post("https://w.example/sessions").mock( return_value=httpx.Response(404, json={"error": "unknown_agent_id"}) ) rc = run_tui(_args_new()) err = capsys.readouterr().err assert rc == 12 assert "[agent_not_found]" in err assert "agent_id=mimir" in err @respx.mock def test_session_api_failed_on_resolve( self, capsys: pytest.CaptureFixture[str] ) -> None: """session_api_failed_on_resolve [error]: --new + 500 → [session_api_failed] stderr.""" respx.post("https://w.example/sessions").mock( return_value=httpx.Response(500, content=b"server error") ) rc = run_tui(_args_new()) err = capsys.readouterr().err assert rc == 20 assert "[session_api_failed]" in err assert "status=500" in err @respx.mock def test_network_error_on_resolve( self, capsys: pytest.CaptureFixture[str] ) -> None: """network_error_on_resolve [error]: --new + ConnectError → [network_error] stderr.""" respx.post("https://w.example/sessions").mock(side_effect=httpx.ConnectError("down")) rc = run_tui(_args_new()) err = capsys.readouterr().err assert rc == 21 assert "[network_error]" in err assert "ConnectError" in err @respx.mock def test_stderr_label_format_matches_cli( self, capsys: pytest.CaptureFixture[str] ) -> None: """stderr_label_format_matches_cli [trace]: cli._amain and _resolve_then_run produce identical stderr lines for AgentNotFound (INV-006). """ # Re-fetch cli's ParsedArgs from the current module state — test_cli's # `importlib.reload(ratatoskr.cli)` rebinds the class, so the top-of-file # `from ratatoskr.cli import ParsedArgs` may now refer to a stale class. from ratatoskr import cli as cli_mod respx.post("https://w.example/sessions").mock( return_value=httpx.Response(404, json={"error": "unknown_agent_id"}) ) cli_args = cli_mod.ParsedArgs( send_content="x", session_id=None, new=True, agent_id="mimir", api_key="k", server_url="https://w.example", raw=False, ) # Drive cli._amain's error path (--send mode) cli_rc = asyncio.run(cli_mod._amain(cli_args)) cli_err = capsys.readouterr().err # Drive _resolve_then_run's error path (TUI mode); _args_new() uses the # pre-reload ParsedArgs which still matches tui.run_tui's isinstance check. tui_rc = run_tui(_args_new()) tui_err = capsys.readouterr().err # Same exit code, same verbatim stderr line. assert cli_rc == 12 assert tui_rc == 12 assert cli_err == tui_err assert cli_err == "[agent_not_found] agent_id=mimir\n" def test_client_open_after_resolve(self, monkeypatch: pytest.MonkeyPatch) -> None: """client_open_after_resolve [trace]: app.client is open at the time run_async runs.""" snapshot: dict = {} async def capture_run_async(self, *a, **kw): snapshot["client_is"] = self.client snapshot["closed_during_run"] = self.client.is_closed return 0 monkeypatch.setattr(RatatoskrApp, "run_async", capture_run_async) rc = run_tui(_args_existing()) assert rc == 0 assert snapshot["client_is"] is not None assert snapshot["closed_during_run"] is False def test_client_lifetime_owned_by_run_tui(self, monkeypatch: pytest.MonkeyPatch) -> None: """client_lifetime_owned_by_run_tui [trace]: open during run_async, closed after run_tui. Probes INV-002: the App is a consumer of an externally-owned client; the async-with in run_tui closes it, not on_unmount. """ snapshot: dict = {} async def capture_run_async(self, *a, **kw): # During run_async (the alt-screen lifetime) the client is open. snapshot["client"] = self.client snapshot["closed_during_run"] = self.client.is_closed return 0 monkeypatch.setattr(RatatoskrApp, "run_async", capture_run_async) rc = run_tui(_args_existing()) assert rc == 0 client = snapshot["client"] assert client is not None # Open while the app was running; closed by run_tui's async-with after. assert snapshot["closed_during_run"] is False assert client.is_closed is True def test_run_tui_closes_client_on_app_exit( self, monkeypatch: pytest.MonkeyPatch ) -> None: """run_tui_closes_client_on_app_exit: async-with closes client after app.run_async ret.""" seen_clients: list[httpx.AsyncClient] = [] async def fake_run_async(self, *a, **kw): seen_clients.append(self.client) return 0 monkeypatch.setattr(RatatoskrApp, "run_async", fake_run_async) rc = run_tui(_args_existing()) assert rc == 0 assert len(seen_clients) == 1 # After run_tui returns, the client should be closed by the async-with assert seen_clients[0].is_closed async def test_on_unmount_does_not_close_client(self) -> None: """on_unmount narrowed [trace]: probes INV-002 from the on_unmount side. The complementary check to test_client_lifetime_owned_by_run_tui (which patches run_async and so never exercises on_unmount). Here we DO run the real on_unmount via Pilot ctrl+d → app teardown, and assert the client is still open afterward (close site is run_tui's async-with, which is NOT entered in this Pilot-driven test). """ client = httpx.AsyncClient( base_url="https://w.example", headers={"Authorization": "Bearer k"}, timeout=httpx.Timeout(connect=10.0, read=None, write=10.0, pool=10.0), ) app = RatatoskrApp( _args_existing(), session_id="s-1existing", agent_id=None, client=client, ) async with app.run_test() as pilot: await pilot.pause() assert client.is_closed is False await pilot.press("ctrl+d") await pilot.pause() # After app.run_test() teardown, on_unmount has fired. Per INV-002 the # client MUST still be open — only run_tui's async-with closes it. assert client.is_closed is False await client.aclose() # test-side cleanup class TestRunTui: def test_happy_returns_zero_on_quit(self, monkeypatch: pytest.MonkeyPatch) -> None: """happy_returns_zero_on_quit [happy,tracer]: run_tui propagates app.run_async exit code.""" captured: list[ParsedArgs] = [] async def fake_run_async(self, *a, **kw): captured.append(self.args) return 0 monkeypatch.setattr(RatatoskrApp, "run_async", fake_run_async) rc = run_tui(_args_existing()) assert rc == 0 assert len(captured) == 1 assert captured[0].send_content is None def test_precondition_send_content_none(self) -> None: """precondition_send_content_none [adversarial]: …""" bad_args = ParsedArgs( send_content="x", # PRE-001 violation session_id="s-1", new=False, agent_id=None, api_key="k", server_url="https://w.example", raw=False, ) with pytest.raises(AssertionError): run_tui(bad_args)