diff --git a/EvoScientist/cli/tui_interactive.py b/EvoScientist/cli/tui_interactive.py index 1b06955..3054cc6 100644 --- a/EvoScientist/cli/tui_interactive.py +++ b/EvoScientist/cli/tui_interactive.py @@ -525,6 +525,7 @@ def run_textual_interactive( CompactingWidget, LoadingWidget, MCPLoaderWidget, + PanelWidget, SubAgentWidget, SummarizationWidget, SystemMessage, @@ -1590,6 +1591,7 @@ def run_textual_interactive( todo_w: TodoWidget | None = None tool_widgets: dict[str, ToolCallWidget] = {} subagent_widgets: dict[str, SubAgentWidget] = {} + panel_widgets: dict[str, PanelWidget] = {} @dataclass class _ResponseDisplayState: @@ -2100,6 +2102,40 @@ def run_textual_interactive( if sa_w is not None: sa_w.finalize() + elif event_type == "panel_dispatch_start": + eval_id = event.get("eval_id", "") or "_unbatched" + panel_w = panel_widgets.get(eval_id) + if panel_w is None: + panel_w = PanelWidget(eval_id) + # Register before awaiting mount: a cancel + # during the await would otherwise orphan a + # ticking panel outside the cleanup loop. + panel_widgets[eval_id] = panel_w + await container.mount(panel_w) + await panel_w.start_dispatch( + event["id"], + event.get("subagent_type", ""), + event.get("label", "") or event.get("description", ""), + ) + + elif event_type == "panel_dispatch_complete": + eval_id = event.get("eval_id", "") or "_unbatched" + panel_w = panel_widgets.get(eval_id) + if panel_w is not None: + panel_w.complete_dispatch( + event["id"], int(event.get("duration_ms", 0)) + ) + + elif event_type == "panel_dispatch_error": + eval_id = event.get("eval_id", "") or "_unbatched" + panel_w = panel_widgets.get(eval_id) + if panel_w is not None: + panel_w.fail_dispatch( + event["id"], + int(event.get("duration_ms", 0)), + event.get("error", ""), + ) + elif event_type == "ask_user": questions = event.get("questions", []) if questions: @@ -2320,6 +2356,13 @@ def run_textual_interactive( sa_w.finalize() except Exception: pass + # Finalize any still-running panel dispatches so their + # per-row spinner timers stop instead of ticking forever. + for panel_w in panel_widgets.values(): + try: + panel_w.finalize_running() + except Exception: + pass # Finalize thinking widget if thinking_w is not None and thinking_w._is_active: try: diff --git a/EvoScientist/cli/widgets/__init__.py b/EvoScientist/cli/widgets/__init__.py index 7bd817d..cb50eff 100644 --- a/EvoScientist/cli/widgets/__init__.py +++ b/EvoScientist/cli/widgets/__init__.py @@ -7,6 +7,7 @@ from .compact_summary_widget import CompactSummaryWidget from .compacting_widget import CompactingWidget from .loading_widget import LoadingWidget from .mcp_loader_widget import MCPLoaderWidget +from .panel_widget import PanelWidget from .subagent_widget import SubAgentWidget from .summarization_widget import SummarizationWidget from .system_message import SystemMessage @@ -25,6 +26,7 @@ __all__ = [ "CompactingWidget", "LoadingWidget", "MCPLoaderWidget", + "PanelWidget", "SubAgentWidget", "SummarizationWidget", "SystemMessage", diff --git a/EvoScientist/cli/widgets/panel_widget.py b/EvoScientist/cli/widgets/panel_widget.py new file mode 100644 index 0000000..a837cf1 --- /dev/null +++ b/EvoScientist/cli/widgets/panel_widget.py @@ -0,0 +1,257 @@ +"""Panel widget — in-eval ``task()`` fan-out live view. + +Groups all sub-agent dispatches from a single ``code_interpreter`` eval into +one bordered container, one row per dispatch. Each row shows the expert / +subagent type, the label, a running elapsed timer, and a status dot that +flips to ``ok``/``err`` on completion. + +Sourced from the ``custom`` stream events emitted by ``langchain_quickjs`` +(see ``.stream.emitter.panel_dispatch_start`` etc.). Keyed by ``eval_id`` +so parallel dispatches from the same eval appear stacked; distinct evals +get distinct panels. +""" + +from __future__ import annotations + +import time + +from rich.text import Text +from textual.containers import Vertical +from textual.widget import Widget +from textual.widgets import Static + +from ..status_bar import SPINNER_FRAMES + +_ROW_LABEL_MAX_CHARS = 48 + + +class _DispatchRow(Widget): + """One row inside a PanelWidget — a single ``task()`` dispatch. + + Subclasses ``Widget`` directly and overrides ``render()`` to build the + row's ``Text`` on demand. Earlier revisions subclassed ``Static`` (visual + stayed ``None`` past the first paint) and ``Vertical`` with an inner + ``Static`` (container-layout race); rendering from ``render()`` is the + standard pattern for single-line widgets and dodges both issues by + letting Textual manage the visual lifecycle itself. + """ + + DEFAULT_CSS = """ + _DispatchRow { + height: 1; + width: 100%; + } + """ + + def __init__(self, subagent_type: str, label: str) -> None: + super().__init__() + self._subagent_type = subagent_type + self._label = label + self._started_at = time.monotonic() + self._status: str = "running" # "running" | "ok" | "err" + self._duration_ms: int | None = None + self._error: str = "" + self._frame = 0 + + def render(self) -> Text: + line = Text() + if self._status == "running": + line.append(f" {SPINNER_FRAMES[self._frame]} ", style="cyan") + elif self._status == "ok": + line.append(" \u2713 ", style="green") + else: + line.append(" \u2717 ", style="red") + line.append(f"{self._subagent_type} ", style="bold") + if self._label: + trimmed = self._label + if len(trimmed) > _ROW_LABEL_MAX_CHARS: + trimmed = trimmed[: _ROW_LABEL_MAX_CHARS - 1] + "\u2026" + line.append(f"\u2014 {trimmed} ", style="dim") + line.append(self._elapsed_display(), style="dim") + if self._status == "err" and self._error: + err = self._error.split("\n", 1)[0] + if len(err) > 60: + err = err[:59] + "\u2026" + line.append(f" {err}", style="red") + return line + + def tick(self) -> None: + if self._status == "running": + self._frame = (self._frame + 1) % len(SPINNER_FRAMES) + self.refresh() + + def complete(self, duration_ms: int) -> None: + self._status = "ok" + self._duration_ms = duration_ms + self.refresh() + + def fail(self, duration_ms: int, error: str) -> None: + self._status = "err" + self._duration_ms = duration_ms + self._error = error + self.refresh() + + def _elapsed_display(self) -> str: + if self._duration_ms is not None: + secs = self._duration_ms / 1000.0 + else: + secs = time.monotonic() - self._started_at + return f"{secs:5.1f}s" + + +class PanelWidget(Vertical): + """Container for one eval's fan-out — bordered box, one row per dispatch.""" + + DEFAULT_CSS = """ + PanelWidget { + height: auto; + margin: 0 0; + } + PanelWidget .panel-header { + height: auto; + color: #22d3ee; + } + PanelWidget .panel-rows { + height: auto; + padding: 0 0 0 2; + } + PanelWidget .panel-footer { + height: auto; + color: #22d3ee; + } + PanelWidget.--completed .panel-header { + color: #4ade80; + } + PanelWidget.--completed .panel-footer { + color: #4ade80; + } + """ + + def __init__(self, eval_id: str) -> None: + super().__init__() + self._eval_id = eval_id + self._rows: dict[str, _DispatchRow] = {} + self._timer_handle = None + self._is_active = True + + @property + def eval_id(self) -> str: + return self._eval_id + + @property + def dispatch_count(self) -> int: + return len(self._rows) + + def compose(self): + yield Static("", classes="panel-header") + yield Vertical(classes="panel-rows") + yield Static("", classes="panel-footer") + + def on_mount(self) -> None: + self._timer_handle = self.set_interval(0.1, self._tick) + self._render_header() + self._render_footer() + + def _tick(self) -> None: + for row in self._rows.values(): + row.tick() + self._render_header() + + async def start_dispatch( + self, dispatch_id: str, subagent_type: str, label: str + ) -> None: + if dispatch_id in self._rows: + return + # Re-arm if the panel already finalized: a Promise.allSettled retry + # of a failed subset dispatches under the same eval_id, so a new + # start after _maybe_finalize() has stopped the timer must undo + # the three effects of finalize (latch, class, timer) or the new + # row's spinner/elapsed stay frozen and future completions never + # refresh the header. + if not self._is_active: + self._is_active = True + self.remove_class("--completed") + self._timer_handle = self.set_interval(0.1, self._tick) + self._render_footer() + row = _DispatchRow(subagent_type, label) + rows_container = self.query_one(".panel-rows", Vertical) + await rows_container.mount(row) + self._rows[dispatch_id] = row + self._render_header() + + def complete_dispatch(self, dispatch_id: str, duration_ms: int) -> None: + row = self._rows.get(dispatch_id) + if row is not None: + row.complete(duration_ms) + self._maybe_finalize() + + def fail_dispatch(self, dispatch_id: str, duration_ms: int, error: str) -> None: + row = self._rows.get(dispatch_id) + if row is not None: + row.fail(duration_ms, error) + self._maybe_finalize() + + def finalize_running(self, reason: str = "interrupted") -> None: + """Fail all still-running rows with their measured elapsed time. + + Called from the TUI turn-cleanup ``finally`` so the 100ms interval + timer stops when a turn is cancelled mid-dispatch — without this, + no reference remains to stop the timer once the outer scope exits. + """ + if not self._is_active: + return + now = time.monotonic() + for dispatch_id, row in list(self._rows.items()): + if row._status == "running": + elapsed_ms = int((now - row._started_at) * 1000) + self.fail_dispatch(dispatch_id, elapsed_ms, reason) + # Cover the zero-running-rows case: cancel-during-mount can leave + # the timer armed with either no rows registered (first dispatch) + # or every registered row already terminal (allSettled retry). + # The loop skips both, so call _maybe_finalize unconditionally — + # it is a no-op once _is_active has flipped. + self._maybe_finalize() + + def _maybe_finalize(self) -> None: + if not self._is_active: + return + if all(row._status != "running" for row in self._rows.values()): + self._is_active = False + if self._timer_handle is not None: + self._timer_handle.stop() + self._timer_handle = None + self.add_class("--completed") + self._render_header() + self._render_footer() + + def _summary_counts(self) -> tuple[int, int, int]: + running = ok = err = 0 + for row in self._rows.values(): + if row._status == "running": + running += 1 + elif row._status == "ok": + ok += 1 + else: + err += 1 + return running, ok, err + + def _render_header(self) -> None: + header = self.query_one(".panel-header", Static) + running, ok, err = self._summary_counts() + line = Text() + if self._is_active: + line.append("\u250c \u25b6 Expert panel ", style="bold cyan") + line.append( + f"({running} running, {ok} done, {err} failed)", style="dim cyan" + ) + else: + line.append("\u2713 Expert panel ", style="bold green") + line.append(f"({ok} done, {err} failed)", style="dim green") + header.update(line) + + def _render_footer(self) -> None: + footer = self.query_one(".panel-footer", Static) + if self._is_active: + footer.update(Text("\u2514 running...", style="dim cyan")) + else: + footer.update(Text("")) diff --git a/EvoScientist/prompts.py b/EvoScientist/prompts.py index d9cc784..1eb45ff 100644 --- a/EvoScientist/prompts.py +++ b/EvoScientist/prompts.py @@ -338,6 +338,15 @@ Launch multiple sub-agents only when experiments are independent: - Debug → fix → re-run — must observe the outcome before proceeding - Ablation design — requires knowing which components matter first +## Dispatch Mechanisms +Three ways to reach a sub-agent — pick based on what you need: + +- **Sequential `task`** — the default. Emit one `task({subagent_type: ..., description: ...})` tool call, wait for the result, integrate, continue. Use for a single-shot consult, including an **expert consult** when the user has invited an expert to the thread (see any `` cue in the system prompt). + +- **In-eval `task()` fan-out via `code_interpreter`** — write a short JS script that dispatches N `task()` calls concurrently and synthesises results in the same eval. Use for independent parallel work: **expert panels** (dispatch to multiple invited experts and synthesise), ELO-style tournaments, N-way method / dataset comparisons where results are independent. Prefer `Promise.allSettled` over `Promise.all` so one failed dispatch does not fail the whole eval — inspect each entry's `status` and retry only the failed subset. + +- **`start_async_task`** — spawn a long-running background job that returns a task ID immediately; poll with `check_async_task` or continue when the async notification arrives. Use for work that will take minutes to hours (long training runs, exhaustive experiments, whole pipelines). The user can keep working in the main conversation while it runs. + ## When to Stop Iterating After each stage, ask: "Would a critical reviewer accept this evidence?" @@ -359,6 +368,7 @@ After each stage, ask: "Would a critical reviewer accept this evidence?" - Bias towards a single sub-agent — add concurrency only when the workload is genuinely independent. - Avoid premature decomposition — one focused task per sub-agent. - Each sub-agent returns self-contained findings with concrete artifacts. +- For parallel fan-outs in `code_interpreter`, use `Promise.allSettled` — a single failed dispatch must not fail the whole eval. ## When a sub-agent reports a blocked command An async sub-agent cannot ask the user anything — it runs on its own thread. If it diff --git a/EvoScientist/stream/emitter.py b/EvoScientist/stream/emitter.py index cfa268a..9fcdebe 100644 --- a/EvoScientist/stream/emitter.py +++ b/EvoScientist/stream/emitter.py @@ -139,6 +139,73 @@ class StreamEventEmitter: {"type": "subagent_end", "name": name, "instance_id": instance_id}, ) + @staticmethod + def panel_dispatch_start( + *, + eval_id: str, + dispatch_id: str, + subagent_type: str, + label: str, + description: str, + ) -> StreamEvent: + """A ``task()`` dispatch inside a ``code_interpreter`` eval started. + + Sourced from ``langchain_quickjs``'s custom-stream subagent events. + Distinct family from ``subagent_start`` because the semantics differ: + short-lived expert dispatch inside a QuickJS eval, grouped by + ``eval_id`` (parent ``code_interpreter`` tool_call_id), not a full + nested subagent turn. + """ + return StreamEvent( + "panel_dispatch_start", + { + "type": "panel_dispatch_start", + "eval_id": eval_id, + "id": dispatch_id, + "subagent_type": subagent_type, + "label": label, + "description": description, + }, + ) + + @staticmethod + def panel_dispatch_complete( + *, + eval_id: str, + dispatch_id: str, + duration_ms: int, + ) -> StreamEvent: + """A ``task()`` dispatch inside a ``code_interpreter`` eval finished OK.""" + return StreamEvent( + "panel_dispatch_complete", + { + "type": "panel_dispatch_complete", + "eval_id": eval_id, + "id": dispatch_id, + "duration_ms": duration_ms, + }, + ) + + @staticmethod + def panel_dispatch_error( + *, + eval_id: str, + dispatch_id: str, + duration_ms: int, + error: str, + ) -> StreamEvent: + """A ``task()`` dispatch inside a ``code_interpreter`` eval raised.""" + return StreamEvent( + "panel_dispatch_error", + { + "type": "panel_dispatch_error", + "eval_id": eval_id, + "id": dispatch_id, + "duration_ms": duration_ms, + "error": error, + }, + ) + @staticmethod def done(response: str = "") -> StreamEvent: """Done event.""" diff --git a/EvoScientist/stream/events.py b/EvoScientist/stream/events.py index 5143c10..6a69803 100644 --- a/EvoScientist/stream/events.py +++ b/EvoScientist/stream/events.py @@ -255,6 +255,62 @@ class _V3EventProcessor: return events if method == "input.requested": return self._process_input_requested(event.get("params")) + if method == "custom": + return self._process_custom_event(_event_data(event)) + return [] + + def _process_custom_event(self, data: object) -> list[dict[str, Any]]: + """Translate a ``custom`` stream payload into UI events. + + Currently handles the ``subagent`` lifecycle emitted by + ``langchain_quickjs`` for in-eval ``task()`` fan-out: ``start`` / + ``complete`` / ``error`` become ``panel_dispatch_start`` / + ``panel_dispatch_complete`` / ``panel_dispatch_error`` UI events. + Unrecognised event types are ignored so future producers can extend + the custom channel without breaking existing consumers. + """ + payload = _as_raw_map(data) + if payload is None or payload.get("type") != "subagent": + return [] + phase = payload.get("phase") + eval_id = payload.get("eval_id") + dispatch_id = payload.get("id") + if not isinstance(dispatch_id, str): + return [] + eval_id_str = eval_id if isinstance(eval_id, str) else "" + if phase == "start": + subagent_type = payload.get("subagent_type") + label = payload.get("label") + description = payload.get("description") + return [ + self.emitter.panel_dispatch_start( + eval_id=eval_id_str, + dispatch_id=dispatch_id, + subagent_type=subagent_type + if isinstance(subagent_type, str) + else "", + label=label if isinstance(label, str) else "", + description=description if isinstance(description, str) else "", + ).data + ] + if phase == "complete": + return [ + self.emitter.panel_dispatch_complete( + eval_id=eval_id_str, + dispatch_id=dispatch_id, + duration_ms=payload.get("duration_ms", 0), + ).data + ] + if phase == "error": + error = payload.get("error") + return [ + self.emitter.panel_dispatch_error( + eval_id=eval_id_str, + dispatch_id=dispatch_id, + duration_ms=payload.get("duration_ms", 0), + error=error if isinstance(error, str) else "", + ).data + ] return [] @classmethod @@ -851,7 +907,7 @@ async def stream_agent_events( _run_raised: bool = False event_sink_token = None try: - from langgraph.stream.transformers import UpdatesTransformer + from langgraph.stream.transformers import CustomTransformer, UpdatesTransformer from ..middleware.events import ( MiddlewareEventSink, @@ -867,7 +923,7 @@ async def stream_agent_events( astream_input, config=config, version="v3", - transformers=[UpdatesTransformer], + transformers=[UpdatesTransformer, CustomTransformer], ) except AttributeError as exc: raise RuntimeError( diff --git a/tests/stream_v3_fakes.py b/tests/stream_v3_fakes.py index 4269fc8..681e486 100644 --- a/tests/stream_v3_fakes.py +++ b/tests/stream_v3_fakes.py @@ -113,6 +113,18 @@ def message_tool_call_block( ) +def custom_subagent_event( + payload: dict[str, Any], + namespace: Iterable[Any] = (), +) -> dict[str, Any]: + """Build a ``custom``-method v3 event carrying a subagent-lifecycle payload. + + Mirrors the shape ``langchain_quickjs._subagent`` emits via + ``stream_writer(event)`` for in-eval ``task()`` fan-out. + """ + return protocol_event("custom", payload, namespace) + + def tool_started( name: str, args: dict[str, Any] | None = None, diff --git a/tests/test_panel_widget.py b/tests/test_panel_widget.py new file mode 100644 index 0000000..acde658 --- /dev/null +++ b/tests/test_panel_widget.py @@ -0,0 +1,194 @@ +"""Unit tests for PanelWidget (TUI in-eval fan-out live view). + +Widget-mount / compose paths need a Textual App context and are exercised +through the interactive TUI end-to-end. Compositor-free tests cover +construction and state-transition helpers; ``TestPanelWidgetReentry`` runs +against a real Textual harness to probe the ``_is_active`` latch behaviour +after finalize (the ``Promise.allSettled`` → retry-failed-subset scenario). +""" + +from __future__ import annotations + +import unittest + +import pytest + +pytest.importorskip("textual") + + +class TestPanelWidgetState: + """PanelWidget state transitions independent of the compositor.""" + + def test_construction(self): + from EvoScientist.cli.widgets.panel_widget import PanelWidget + + w = PanelWidget("ci_eval_1") + assert w.eval_id == "ci_eval_1" + assert w.dispatch_count == 0 + assert w._is_active is True + + def test_summary_counts_all_running(self): + from EvoScientist.cli.widgets.panel_widget import PanelWidget, _DispatchRow + + w = PanelWidget("e1") + w._rows["d1"] = _DispatchRow("innovator", "a") + w._rows["d2"] = _DispatchRow("pragmatist", "b") + running, ok, err = w._summary_counts() + assert (running, ok, err) == (2, 0, 0) + + def test_summary_counts_mixed(self): + from EvoScientist.cli.widgets.panel_widget import PanelWidget, _DispatchRow + + w = PanelWidget("e1") + r1 = _DispatchRow("innovator", "a") + r2 = _DispatchRow("pragmatist", "b") + r3 = _DispatchRow("critic", "c") + r1._status = "ok" + r2._status = "err" + w._rows["d1"] = r1 + w._rows["d2"] = r2 + w._rows["d3"] = r3 + running, ok, err = w._summary_counts() + assert (running, ok, err) == (1, 1, 1) + + +class TestDispatchRow(unittest.TestCase): + """_DispatchRow state transitions.""" + + def test_construction_sets_running(self): + from EvoScientist.cli.widgets.panel_widget import _DispatchRow + + row = _DispatchRow("idea-brainstorm", "innovator voice") + assert row._subagent_type == "idea-brainstorm" + assert row._label == "innovator voice" + assert row._status == "running" + assert row._duration_ms is None + + def test_elapsed_display_uses_recorded_duration(self): + from EvoScientist.cli.widgets.panel_widget import _DispatchRow + + row = _DispatchRow("x", "y") + row._duration_ms = 2500 + display = row._elapsed_display() + assert "2.5" in display + + +class TestPanelWidgetReentry: + """Regression: retry-failed-subset under the same eval_id re-arms the panel. + + The ``Promise.allSettled`` → retry-failed-subset pattern (encouraged by + the DELEGATION_STRATEGY guardrail) can send a second wave of + ``panel_dispatch_start`` events under the same ``eval_id`` after the + panel has already finalized. ``start_dispatch`` must detect the + finalized state and undo the three effects of ``_maybe_finalize`` + (latch, ``--completed`` class, interval timer) so the new row's + spinner + elapsed keep ticking and the header refreshes on completion. + """ + + async def _boot_panel(self): + from textual.app import App, ComposeResult + + from EvoScientist.cli.widgets.panel_widget import PanelWidget + + class _PanelApp(App[None]): + def compose(self) -> ComposeResult: + yield PanelWidget("eval-42") + + app = _PanelApp() + return app + + def _header_text(self, panel) -> str: + from textual.widgets import Static + + header = panel.query_one(".panel-header", Static) + if header.size.height == 0: + return "" + return "".join(seg.text for seg in header.render_line(0)).strip() + + async def test_reentry_after_finalize_rearms_panel(self): + from EvoScientist.cli.widgets.panel_widget import PanelWidget + + app = await self._boot_panel() + async with app.run_test(size=(120, 20)) as pilot: + panel = app.query_one(PanelWidget) + await pilot.pause() + + # Wave 1: two dispatches, one settles ok, one settles err. + await panel.start_dispatch("d1", "innovator", "brainstorm a") + await panel.start_dispatch("d2", "pragmatist", "brainstorm b") + await pilot.pause() + panel.complete_dispatch("d1", 1200) + panel.fail_dispatch("d2", 800, "boom") + await pilot.pause() + + # Baseline: finalize fired. + assert panel._is_active is False + assert panel._timer_handle is None + assert panel.has_class("--completed") + header_after_wave1 = self._header_text(panel) + assert "1 done" in header_after_wave1 + assert "1 failed" in header_after_wave1 + assert "running" not in header_after_wave1 + + # Wave 2: eval retries the failed subset under the same eval_id. + await panel.start_dispatch("d3", "pragmatist", "brainstorm b retry") + await pilot.pause() + + # Panel re-armed: latch flipped, timer restarted, class dropped. + assert panel._is_active is True + assert panel._timer_handle is not None + assert not panel.has_class("--completed") + + # Header shows the running count again. + header_during_wave2 = self._header_text(panel) + assert "1 running" in header_during_wave2 + assert "1 done" in header_during_wave2 + assert "1 failed" in header_during_wave2 + + # d3 completes: header refreshes to 2 done, panel re-finalizes. + panel.complete_dispatch("d3", 900) + await pilot.pause() + header_after_wave2 = self._header_text(panel) + assert "2 done" in header_after_wave2 + assert "1 failed" in header_after_wave2 + assert "running" not in header_after_wave2 + assert panel._is_active is False + assert panel._timer_handle is None + assert panel.has_class("--completed") + + +class TestPanelWidgetFinalizeEmpty: + """Regression: ``finalize_running`` stops the timer even with zero + running rows. + + The timer is armed by ``on_mount`` (first-dispatch case) or re-armed + by ``start_dispatch`` (retry-failed-subset case) before any row lands + in the "running" state. If cancel arrives before the first row is + registered, ``finalize_running``'s loop body never runs, so + ``_maybe_finalize`` must be called unconditionally after the loop. + """ + + async def test_finalize_running_stops_timer_with_no_rows(self): + from textual.app import App, ComposeResult + + from EvoScientist.cli.widgets.panel_widget import PanelWidget + + class _PanelApp(App[None]): + def compose(self) -> ComposeResult: + yield PanelWidget("eval-empty") + + app = _PanelApp() + async with app.run_test(size=(120, 20)) as pilot: + panel = app.query_one(PanelWidget) + await pilot.pause() + + # Timer armed by on_mount, no dispatches yet. + assert panel._is_active is True + assert panel._timer_handle is not None + assert panel._rows == {} + + panel.finalize_running() + + assert panel._is_active is False + assert panel._timer_handle is None + assert panel.has_class("--completed") diff --git a/tests/test_stream_events.py b/tests/test_stream_events.py index a2f6525..828f8b3 100644 --- a/tests/test_stream_events.py +++ b/tests/test_stream_events.py @@ -29,6 +29,7 @@ from tests.stream_v3_fakes import ( SubscriptionSensitiveV3Agent, async_iter, collect_events, + custom_subagent_event, message_delta, message_finish, message_tool_call_block, @@ -1251,6 +1252,157 @@ class TestUsageStatsExtraction: assert len(usage_events) == 0 +class TestPanelDispatchEvents: + """Custom-stream subagent lifecycle from in-eval task() fan-out.""" + + async def test_start_event_becomes_panel_dispatch_start(self): + agent = FakeV3Agent( + [ + custom_subagent_event( + { + "type": "subagent", + "phase": "start", + "id": "ptc_task_abc12345", + "eval_id": "ci_eval_1", + "subagent_type": "idea-brainstorm", + "label": "innovator voice", + "description": "generate one bold candidate", + } + ), + ] + ) + events = await collect_events(agent) + starts = [e for e in events if e.get("type") == "panel_dispatch_start"] + assert len(starts) == 1 + start = starts[0] + assert start["id"] == "ptc_task_abc12345" + assert start["eval_id"] == "ci_eval_1" + assert start["subagent_type"] == "idea-brainstorm" + assert start["label"] == "innovator voice" + assert start["description"] == "generate one bold candidate" + + async def test_complete_event_becomes_panel_dispatch_complete(self): + agent = FakeV3Agent( + [ + custom_subagent_event( + { + "type": "subagent", + "phase": "complete", + "id": "ptc_task_abc12345", + "eval_id": "ci_eval_1", + "duration_ms": 1234, + } + ), + ] + ) + events = await collect_events(agent) + completes = [e for e in events if e.get("type") == "panel_dispatch_complete"] + assert len(completes) == 1 + assert completes[0]["id"] == "ptc_task_abc12345" + assert completes[0]["duration_ms"] == 1234 + + async def test_error_event_becomes_panel_dispatch_error(self): + agent = FakeV3Agent( + [ + custom_subagent_event( + { + "type": "subagent", + "phase": "error", + "id": "ptc_task_abc12345", + "eval_id": "ci_eval_1", + "duration_ms": 42, + "error": "boom", + } + ), + ] + ) + events = await collect_events(agent) + errors = [e for e in events if e.get("type") == "panel_dispatch_error"] + assert len(errors) == 1 + assert errors[0]["error"] == "boom" + assert errors[0]["duration_ms"] == 42 + + async def test_unknown_custom_type_ignored(self): + """Custom payloads whose ``type`` is not ``subagent`` don't emit events.""" + agent = FakeV3Agent( + [ + custom_subagent_event({"type": "something-else", "value": 1}), + ] + ) + events = await collect_events(agent) + assert not any(e.get("type", "").startswith("panel_dispatch") for e in events) + + async def test_missing_id_dropped(self): + """A malformed subagent event without ``id`` is silently dropped.""" + agent = FakeV3Agent( + [ + custom_subagent_event( + { + "type": "subagent", + "phase": "start", + "subagent_type": "x", + "label": "y", + } + ), + ] + ) + events = await collect_events(agent) + assert not any(e.get("type", "").startswith("panel_dispatch") for e in events) + + async def test_grouped_fanout_shares_eval_id(self): + """Parallel dispatches from one eval share the same ``eval_id``.""" + agent = FakeV3Agent( + [ + custom_subagent_event( + { + "type": "subagent", + "phase": "start", + "id": "d1", + "eval_id": "e1", + "subagent_type": "innovator", + "label": "a", + "description": "d", + } + ), + custom_subagent_event( + { + "type": "subagent", + "phase": "start", + "id": "d2", + "eval_id": "e1", + "subagent_type": "pragmatist", + "label": "b", + "description": "d", + } + ), + custom_subagent_event( + { + "type": "subagent", + "phase": "complete", + "id": "d1", + "eval_id": "e1", + "duration_ms": 100, + } + ), + custom_subagent_event( + { + "type": "subagent", + "phase": "complete", + "id": "d2", + "eval_id": "e1", + "duration_ms": 200, + } + ), + ] + ) + events = await collect_events(agent) + panel_events = [ + e for e in events if e.get("type", "").startswith("panel_dispatch") + ] + assert len(panel_events) == 4 + assert {e["eval_id"] for e in panel_events} == {"e1"} + + class TestSummarizationHelpers: """Summarization extraction helpers."""