From aba17304bd0f03096821c36bb3592bf036947354 Mon Sep 17 00:00:00 2001 From: Vuong Hoang Date: Sun, 19 Jul 2026 07:09:26 -0700 Subject: [PATCH] =?UTF-8?q?fix(#20):=20heid-bug-hunt=20fixups=20=E2=80=94?= =?UTF-8?q?=20cutover=20edge-path=20robustness=20(slice-2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Triaged the heid-bug-hunt panel (Gróa 8 / Hulda 6 / Regin 6; Heid source-checked + refuted 2 Regin FPs). The lens pulled real weight — confirmed bugs the conformance review structurally could not see. Confirmed bugs fixed: - SessionRetired (410) stream-open maps to wt.SessionApiFailed, but neither cli _run_turn nor web gen() caught it → crash / dropped SSE stream. Both presenters now catch it (cli → exit 20; web → labeled `event: error`). (Gróa#2) + cli regression test. - cli forwarded consumer_key unconditionally; an UNBOUND create with the env key set would auth as the Bifrost consumer, not the default bearer. Guarded in the adapter (consumer_key only when bifrost is set). (Gróa#4 + Regin#4) + test. - cli _turn_id_from_sse_id crashed on a None/non-str sse_id (web guarded, cli didn't) → now tolerant. (Gróa#1 + Hulda#2) + test. - _cancel_and_log broadened to `except Exception` — after the code-review's ApiError default, a cancel could raise SessionApiFailed it didn't catch, breaking INV-009 (never-raise). (Gróa#3, Heid-endorsed over Regin's refuted mechanism). Open-world degrade-not-crash (contract posture): render hardened — float duration_ms (_format_duration_safe), non-mapping usage/snapshot guards, unknown event type degrades instead of asserting (Gróa#5/#6 + Hulda#3); web _event_to_browser_payload guards a non-mapping `raw` (Hulda#4); web _wt_client bearer extraction is now case-insensitive + whitespace-robust (Hulda#5 + Regin#5). + render-degrade test. Rejected (verified): Regin#1 (httpx IS caught), Regin#2 (wtsdk IS worldtree_sdk), Regin#3 (sse_client.AgentNotAvailable IS caught by SseConnectFailed) — all FPs; Hulda#1 (deleted funcs "break callers") — grep-verified zero callers pre-deletion. Accepted-known-risk: lenient sse_id parse, CancelFailed status=0, async-gen aclose (pre-existing pattern, not a cutover regression). Suite 497 green; wt/cli/web ruff + wt mypy clean. Patch. --- pyproject.toml | 2 +- src/ratatoskr/cli.py | 61 ++++++++++++++++++++++++------------- src/ratatoskr/web/server.py | 21 +++++++++---- src/ratatoskr/wt.py | 8 ++++- tests/test_cli.py | 48 +++++++++++++++++++++++++++++ tests/test_wt.py | 8 +++++ uv.lock | 2 +- 7 files changed, 120 insertions(+), 30 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 7b60ae3..b9fe3f5 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "ratatoskr" -version = "0.21.9" +version = "0.21.10" description = "Worldtree Conversation API debug console (web + headless CLI) — multi-pane observability" readme = "README.md" requires-python = ">=3.12" diff --git a/src/ratatoskr/cli.py b/src/ratatoskr/cli.py index 3082ba2..88b85ae 100644 --- a/src/ratatoskr/cli.py +++ b/src/ratatoskr/cli.py @@ -62,9 +62,6 @@ from ratatoskr.sessions import ( # (--whoami / --characters / --set-persona / --seed-first-message) stay on the # `sessions` wrappers until their own slices. from ratatoskr.sse_client import ( - CancelAlreadyCompleted, - CancelFailed, - CancelTurnNotFound, MalformedSseData, MalformedSseId, SseConnectFailed, @@ -347,21 +344,34 @@ def _format_usage(usage: dict[str, int], *, arrow: str) -> str: return f"{p} in {arrow} {c} out ({t} total, {ci} cached)" -def _format_usage_safe(usage: Mapping[str, int] | None) -> str: +def _format_usage_safe(usage: object) -> str: """Tolerant wrapper over `_format_usage` for the SDK's open-world - `DoneEvent.usage` (typed optional): the canonical four-key usage formats; - anything absent or malformed degrades to `(n/a)` rather than crashing the - presenter (same posture as `_format_whoami`).""" + `DoneEvent.usage`: the canonical four-key mapping formats; anything absent or + malformed (None, a non-mapping like `5`, a partial dict) degrades to `(n/a)` + rather than crashing the presenter (same posture as `_format_whoami`).""" keys = ("prompt_tokens", "completion_tokens", "total_tokens", "cached_input_tokens") - if usage is not None and all(k in usage for k in keys): + if isinstance(usage, Mapping) and all(k in usage for k in keys): return _format_usage(dict(usage), arrow="->") return "(n/a)" -def _turn_id_from_sse_id(sse_id: str) -> int | None: +def _format_duration_safe(ms: object) -> str: + """Tolerant wrapper over `_format_duration_ms` for the open-world + `DoneEvent.duration_ms`: a finite non-negative number formats (a float wire + value is floored to int); anything else degrades to `n/a` rather than tripping + `_format_duration_ms`'s int assertion.""" + if isinstance(ms, (int, float)) and not isinstance(ms, bool) and ms >= 0: + return _format_duration_ms(int(ms)) + return "n/a" + + +def _turn_id_from_sse_id(sse_id: object) -> int | None: """The turn component of the SDK's composite sse_id (`"{turn}:{seq}"`). This is the mid-stream cancel target: it is present on EVERY frame, unlike the SDK's - top-level `turn_id`, which is the body field (absent on text/thinking events).""" + top-level `turn_id`, which is the body field (absent on text/thinking events). + Tolerant of a malformed/absent sse_id (open-world) — mirrors the web helper.""" + if not isinstance(sse_id, str): + return None head, _, _ = sse_id.partition(":") try: turn = int(head) @@ -389,14 +399,19 @@ class CliPresenterState: malformed/partial event degrades to a placeholder rather than crashing the presenter — the same posture as `_format_whoami`. """ - assert isinstance( + if not isinstance( event, ( WorkerPhaseEvent, ThinkingEvent, TextEvent, TextBoundaryEvent, ToolStartEvent, ToolResultEvent, DoneEvent, ErrorEvent, CancelledEvent, AffectUpdateEvent, AwaitingLlmFirstTokenEvent, ), - ) + ): + # Open-world: an unknown / future SDK event type degrades to a one-line + # note rather than aborting the presenter. (The SDK skips unknown wire + # types today, so this is belt-and-suspenders for a future SDK event set.) + stderr.write(f". unknown_event: {type(event).__name__}\n") + return # Thinking events accumulate into the open run. if isinstance(event, ThinkingEvent): content = event.content or "" @@ -430,7 +445,7 @@ class CliPresenterState: if isinstance(event, DoneEvent): stderr.write( f"[done] turn_id={event.turn_id} model={event.model} " - f"duration={_format_duration_ms(event.duration_ms or 0)} " + f"duration={_format_duration_safe(event.duration_ms)} " f"usage {_format_usage_safe(event.usage)}\n" ) return @@ -470,7 +485,8 @@ class CliPresenterState: if isinstance(event, AffectUpdateEvent): # Worldtree #204 / v0.28.0. CLI surface is debug telemetry — # one line to stderr with status + (for current) dominant_emotion. - if event.snapshot is not None: + # isinstance(Mapping) guards an open-world non-mapping snapshot. + if isinstance(event.snapshot, Mapping): dom = event.snapshot.get("dominant_emotion") stderr.write( f". affect_update: status={event.status} turn_id={event.turn_id} " @@ -505,13 +521,11 @@ async def _cancel_and_log( assert isinstance(turn_id, int) and turn_id > 0 try: await wt.cancel_turn(client, session_id, turn_id) - except ( - CancelFailed, - CancelTurnNotFound, - CancelAlreadyCompleted, - ConnectFailed, # SDK normalizes a transport drop to ConnectFailed(status=0) - httpx.RequestError, - ) as exc: + except Exception as exc: + # Any cancel failure (mapped ratatoskr cancel exceptions, an adapter-defaulted + # SessionApiFailed, an SDK ConnectFailed, a transport error, or anything the + # SDK doesn't normalize) is logged and swallowed — the fire-and-forget cancel + # must never propagate into _run_turn's finally. stderr.write(f"[cancel_failed] {type(exc).__name__}: {exc}\n") @@ -574,6 +588,11 @@ async def _run_turn( except StopAsyncIteration: stderr.write("[connection_dropped] last_seen=\n") return 21 + except wt.SessionApiFailed as exc: + # The adapter maps a stream-open SessionRetired (410) here; without + # this the retired-session stream would crash out of _run_turn. + stderr.write(f"[session_api_failed] status={exc.status} body={exc.body!r}\n") + return 20 except SseConnectFailed as exc: stderr.write(f"[sse_connect_failed] status={exc.status} body={exc.body!r}\n") return 20 diff --git a/src/ratatoskr/web/server.py b/src/ratatoskr/web/server.py index e96841d..0d6581a 100644 --- a/src/ratatoskr/web/server.py +++ b/src/ratatoskr/web/server.py @@ -11,7 +11,7 @@ from __future__ import annotations import asyncio import itertools import json -from collections.abc import AsyncIterator, Callable +from collections.abc import AsyncIterator, Callable, Mapping from dataclasses import asdict, dataclass, is_dataclass from importlib.metadata import version as _pkg_version @@ -72,7 +72,10 @@ def _wt_client(client: httpx.AsyncClient, *, max_reconnects: int = 5) -> Worldtr a placeholder key (respx ignores auth).""" base_url = str(client.base_url) or "http://localhost" header = client.headers.get("Authorization", "") - api_key = header[len("Bearer "):].strip() if header.startswith("Bearer ") else "" + # Case-insensitive scheme + tolerant of extra whitespace, so a valid bearer is + # not silently dropped to the placeholder key (which would misauthenticate). + parts = header.split(None, 1) + api_key = parts[1].strip() if len(parts) == 2 and parts[0].lower() == "bearer" else "" return wt.build_client( base_url, api_key=api_key or "ratatoskr", transport=client, max_reconnects=max_reconnects ) @@ -276,8 +279,11 @@ def _event_to_browser_payload(event: object) -> tuple[str, dict]: …), NOT the SDK class name. Open-world: additive server fields pass through. """ browser_type = getattr(event, "type", "") or "" - raw = getattr(event, "raw", None) or {} - data = {k: v for k, v in dict(raw).items() if k != "type"} + raw = getattr(event, "raw", None) + # Open-world: degrade a non-mapping `raw` to an empty payload rather than letting + # dict(raw) raise (which would abort the SSE stream mid-response). + src = raw if isinstance(raw, Mapping) else {} + data = {k: v for k, v in src.items() if k != "type"} data["sse_id"] = getattr(event, "sse_id", None) return browser_type, data @@ -342,8 +348,11 @@ async def _stream_turn_endpoint(request: Request) -> StreamingResponse: if isinstance(event, (DoneEvent, ErrorEvent, CancelledEvent)): handle.status = event.type or "done" break - except (SseConnectFailed, SseConnectionDropped, MalformedSseId, - MalformedSseData, TurnIdFlip) as exc: + except (wt.SessionApiFailed, SseConnectFailed, SseConnectionDropped, + MalformedSseId, MalformedSseData, TurnIdFlip) as exc: + # wt.SessionApiFailed covers the adapter's SessionRetired (410) mapping; + # without it a retired-session stream would escape gen() after partial + # frames as an uncaught 500, not a labeled `event: error`. yield _format_sse( "error", {"exception": type(exc).__name__, "message": str(exc)}, diff --git a/src/ratatoskr/wt.py b/src/ratatoskr/wt.py index c7e307b..8f09c4f 100644 --- a/src/ratatoskr/wt.py +++ b/src/ratatoskr/wt.py @@ -193,7 +193,13 @@ async def create_session( body["bifrost"] = {"endpoint_url": bifrost.endpoint_url, "scope": bifrost.scope} try: - return await client.sessions.create(body, consumer_key=consumer_key) + # consumer_key is a BOUND-create credential only — never forward it on an + # unbound create, or the SDK's credential precedence (consumer_key > default) + # would authenticate as the Bifrost consumer instead of the default bearer. + # Centralized here so both surfaces are guarded (the web endpoint already is). + return await client.sessions.create( + body, consumer_key=consumer_key if bifrost is not None else None + ) except ApiError as exc: if exc.status == 404: raise AgentNotFound(agent_id=agent_id) from exc diff --git a/tests/test_cli.py b/tests/test_cli.py index 14e211a..991ded8 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -647,6 +647,39 @@ class TestCliPresenterState: state.render(_make_done(duration_ms=72000), stdout=io.StringIO(), stderr=stderr) assert "duration=1.2m" in stderr.getvalue() + def test_render_degrades_on_malformed_open_world_fields(self) -> None: + """Open-world hardening (heid-bug-hunt Gróa#5 / Hulda#3): a DoneEvent with a + float duration_ms + a non-mapping usage, and an AffectUpdate with a non-mapping + snapshot, DEGRADE rather than crash the presenter.""" + from ratatoskr.cli import CliPresenterState + + stderr = io.StringIO() + state = CliPresenterState() + done = build_event( + "done", "42:9", 42, + {"type": "done", "duration_ms": 1234.0, "usage": 5, "model": "m"}, + ) + state.render(done, stdout=io.StringIO(), stderr=stderr) # must not raise + out = stderr.getvalue() + # float duration floored to int (1234ms → "1.2s"); non-mapping usage → "(n/a)". + assert "[done]" in out and "duration=1.2s" in out and "usage (n/a)" in out + # AffectUpdate with a list snapshot → no AttributeError on .get. + affect = build_event( + "affect_update", "42:1", 42, + {"type": "affect_update", "status": "current", "snapshot": []}, + ) + CliPresenterState().render(affect, stdout=io.StringIO(), stderr=io.StringIO()) + + def test_turn_id_from_sse_id_tolerates_non_str(self) -> None: + """Open-world hardening (heid-bug-hunt Gróa#1 / Hulda#2): a None/non-str sse_id + yields None instead of crashing on .partition.""" + from ratatoskr.cli import _turn_id_from_sse_id + + assert _turn_id_from_sse_id(None) is None + assert _turn_id_from_sse_id(42) is None + assert _turn_id_from_sse_id("42:1") == 42 + assert _turn_id_from_sse_id("0:1") is None + def test_usage_format_ascii_arrow(self) -> None: """usage_format_ascii_arrow [trace]: stderr label contains the natural-language usage shape with ASCII arrow (-> not →) for CLI scriptability. @@ -940,6 +973,21 @@ class TestRunTurn: assert "[sse_connect_failed]" in out assert "status=404" in out + @respx.mock + async def test_session_retired_410_maps_to_session_api_failed(self) -> None: + """session_retired [error]: 410 stream-open → SessionRetired → SessionApiFailed + → exit 20. Without the presenter catch this crashed _run_turn (heid-bug-hunt Gróa#2).""" + respx.post("https://w.example/sessions/s-1/messages").mock( + return_value=httpx.Response(410, json={"error_code": "session_retired"}) + ) + sigint = asyncio.Event() + stdout, stderr = io.StringIO(), io.StringIO() + async with httpx.AsyncClient(base_url="https://w.example") as _tp: + client = _wtc(_tp) + exit_code = await _run_turn(client, "s-1", "hi", sigint, stdout=stdout, stderr=stderr) + assert exit_code == 20 + assert "[session_api_failed]" in stderr.getvalue() + @respx.mock async def test_connection_dropped(self) -> None: """connection_dropped [error]: RemoteProtocolError mid-stream → exit 21.""" diff --git a/tests/test_wt.py b/tests/test_wt.py index 6c797fc..2abe79b 100644 --- a/tests/test_wt.py +++ b/tests/test_wt.py @@ -218,6 +218,14 @@ class TestCreateSession: # INV-CUT: the consumer key rides the SDK's per-request auth, NOT a header. assert kwargs["consumer_key"] == "ck-real" + async def test_unbound_create_drops_consumer_key(self) -> None: + # A consumer_key must NOT reach the SDK on an UNBOUND create — the SDK's + # credential precedence would otherwise auth as the Bifrost consumer instead + # of the default bearer (heid-bug-hunt Gróa#4 / Regin#4). + fake = _FakeSessions(result={"session_id": "s"}) + await create_session(_wt(fake), "mimir", consumer_key="ck-should-be-dropped") + assert fake.calls[-1][2]["consumer_key"] is None + async def test_bifrost_without_consumer_key_rejected_pre_http(self) -> None: fake = _FakeSessions(result={"session_id": "s"}) binding = BifrostBinding(endpoint_url="http://h:8391", scope=None) diff --git a/uv.lock b/uv.lock index 51de766..e0895c2 100644 --- a/uv.lock +++ b/uv.lock @@ -472,7 +472,7 @@ wheels = [ [[package]] name = "ratatoskr" -version = "0.21.9" +version = "0.21.10" source = { editable = "." } dependencies = [ { name = "httpx" },