Live agent todo list: compact scrollable strip above the composer
CI / Gateway plugin tests (push) Successful in 5m5s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m47s

Add a todo.update frame (server->app) carrying the agent's full current
todo list. The gateway emits it whenever the hermes todo tool completes
(the tool result is authoritative even for merge writes) and re-sends a
snapshot right after hello so a reconnecting device re-learns the plan.
Ephemeral: never outboxed.

The app renders it as a compact strip above the composer (max 3 lines,
the rest scrollable) mirroring the hermes desktop composer status stack:
pending = hollow ring, in_progress = spinner, completed = green check,
cancelled = struck through. It auto-scrolls to the current task whenever
the active task changes, and hides itself once the list is empty or fully
resolved.
This commit is contained in:
ARIA committed 2026-08-23 00:23:02 +02:00
1 parent 1ff2ef380c
commit 742916903b
10 files changed
+447 -3

No files matched your search

+109
View File
@@ -289,6 +289,17 @@ def _on_post_tool_call(**kwargs: Any) -> None:
_tool_results.append(record)
while len(_tool_results) > _MAX_TOOL_RESULTS:
_tool_results.popleft()
# Live todo list: the todo tool's result is the authoritative full list,
# so emit it the moment the call completes — the tool.end frame only
# arrives when the NEXT tool starts or the turn ends, which would lag the
# app's strip behind the agent's actual progress. Best-effort: a parse
# failure (truncated preview) or a missing live adapter just skips it.
if kwargs.get("tool_name") == "todo" and record["status"] == "ok":
todos = _parse_todo_result(record["result"])
adapter = _live_adapter
if todos is not None and adapter is not None and adapter._loop is not None:
with contextlib.suppress(Exception):
asyncio.run_coroutine_threadsafe(adapter._emit_todo_update(todos), adapter._loop)
def _take_tool_result(tool_name: str) -> dict[str, Any] | None:
@@ -325,6 +336,43 @@ def _tool_end_fields(tool_name: str) -> dict[str, Any]:
return fields
# The live adapter instance (module-level so the synchronous plugin hooks
# below can reach it). A personal iris gateway runs exactly one adapter;
# set on connect, cleared on disconnect.
_live_adapter: "IrisAdapter | None" = None
# Valid todo item statuses (hermes tools/todo_tool.py VALID_STATUSES).
_TODO_STATUSES = ("pending", "in_progress", "completed", "cancelled")
def _parse_todo_result(result: Any) -> list[dict[str, str]] | None:
"""Parse the ``todo`` tool's result into a clean item list.
The tool returns ``{"todos": [...], "summary": {...}}`` — the FULL current
list, which is authoritative even for ``merge`` writes (whose args carry
only the changed items) and read-only calls. Returns ``None`` when the
result is not a parseable todo list (error string, truncated preview, …)
so the caller skips the emission instead of broadcasting garbage.
"""
try:
data = json.loads(result)
except (TypeError, ValueError):
return None
items = data.get("todos") if isinstance(data, dict) else None
if not isinstance(items, list):
return None
todos: list[dict[str, str]] = []
for it in items:
if not isinstance(it, dict):
continue
content = str(it.get("content") or "").strip()
status = str(it.get("status") or "")
if not content or status not in _TODO_STATUSES:
continue
todos.append({"id": str(it.get("id") or ""), "content": content, "status": status})
return todos
def _tool_emoji(tool_name: str) -> str | None:
"""Cosmetic per-tool emoji for the ``tool.start`` frame.
@@ -1279,6 +1327,19 @@ class IrisAdapter(BasePlatformAdapter):
# M5: per-chat push throttle (epoch seconds of the last successful
# push). The coalesced frame still reaches the app via sync.
self._last_push_at: dict[str, float] = {}
# Live todo lists (lane key -> items): the agent's planning state,
# re-served as a snapshot when a device opens its event stream.
# In-memory only — a gateway restart drops it (the agent re-emits on
# the next todo call).
self._todo_lists: dict[str, list[dict[str, str]]] = {}
# The lane (chat_id, thread_id) of the turn currently in flight — the
# post_tool_call hook is global (no chat id), so the todo emission
# routes to the lane the adapter last saw tool activity for. A
# personal gateway serves one active turn at a time.
self._active_lane: tuple[str, str | None] | None = None
# The adapter's event loop (captured on connect): the synchronous
# plugin hooks schedule broadcasts onto it via run_coroutine_threadsafe.
self._loop: asyncio.AbstractEventLoop | None = None
def _turn_state(self, chat_id: str) -> _TurnState:
st = self._turns.get(chat_id)
@@ -1295,6 +1356,7 @@ class IrisAdapter(BasePlatformAdapter):
async def connect(self, *, is_reconnect: bool = False) -> bool:
"""Bring the platform up: bind the HTTP server on host:http_port."""
global _live_adapter
if not self.token:
logger.error("iris: IRIS_TOKEN must be set")
self._set_fatal_error(
@@ -1343,11 +1405,17 @@ class IrisAdapter(BasePlatformAdapter):
self._connected = True
self._mark_connected()
# The synchronous plugin hooks (post_tool_call) schedule frame
# broadcasts onto this loop; register the live adapter so they can
# reach it (single-adapter personal gateway).
self._loop = asyncio.get_running_loop()
_live_adapter = self
logger.info("iris: connected; HTTP server on %s:%s", self.host, self.http_port)
return True
async def disconnect(self) -> None:
"""Tear down the platform: stop the server, close device streams."""
global _live_adapter
# Tell live clients the gateway is going away (restart/shutdown) so
# the app can distinguish a clean gateway teardown from a plain
# network drop: the "Gateway restarting" chat notice is shown only
@@ -1367,6 +1435,9 @@ class IrisAdapter(BasePlatformAdapter):
self._outbox.close()
self._connected = False
self._mark_disconnected()
if _live_adapter is self:
_live_adapter = None
self._loop = None
logger.info("iris: disconnected")
# ── Outbound (agent -> app) ───────────────────────────────────────────
@@ -1636,6 +1707,9 @@ class IrisAdapter(BasePlatformAdapter):
message_id = state.tool_msg_id or _mint_message_id()
state.tool_msg_id = message_id
state.active = True
# Tool activity marks this lane as the turn in flight: the global
# post_tool_call hook (no chat id) routes todo emissions here.
self._active_lane = (chat_id, thread_id)
lines = [ln for ln in content.splitlines() if ln.strip()]
for line in lines:
@@ -1764,6 +1838,41 @@ class IrisAdapter(BasePlatformAdapter):
return None
return self._http_reply_sinks.pop(device_id, None)
# ── Live todo list (todo.update) ─────────────────────────────────────
@staticmethod
def _lane_key(chat_id: str, thread_id: str | None) -> str:
"""Lane key for a (chat, thread) pair (mirrors the app's ChatStore)."""
return chat_id if not thread_id else f"{chat_id}::{thread_id}"
async def _emit_todo_update(self, todos: list[dict[str, str]]) -> None:
"""Store + broadcast the agent's current todo list for the active lane.
Scheduled from the synchronous post_tool_call hook (the todo tool
just completed); *todos* is the tool result's full list.
"""
chat_id, thread_id = self._active_lane or (self.home_channel, None)
self._todo_lists[self._lane_key(chat_id, thread_id)] = todos
await self._broadcast_or_log(
chat_id, protocol.todo_update(chat_id, todos, thread_id=thread_id)
)
def todo_snapshot_frames(self) -> list["protocol.Frame"]:
"""``todo.update`` frames for every lane with a non-empty list.
Written by the HTTP server right after the hello frame when a device
opens its event stream, so a (re)connecting app re-learns the agent's
current plan (the frames are ephemeral — never outboxed).
"""
frames: list[protocol.Frame] = []
for lane, todos in self._todo_lists.items():
if not todos:
continue
parts = lane.split("::", 1)
thread_id = parts[1] if len(parts) > 1 else None
frames.append(protocol.todo_update(parts[0], todos, thread_id=thread_id))
return frames
async def _broadcast_both(self, frame: "protocol.Frame") -> None:
"""Bare (non-outbox) broadcast to live subscribers (docs/19): the
frame reaches every live SSE/long-poll subscriber."""
+5 -3
View File
@@ -223,9 +223,7 @@ class HttpServer:
return
self._httpd = httpd
self._thread = threading.Thread(
target=httpd.serve_forever, name="iris-http", daemon=True
)
self._thread = threading.Thread(target=httpd.serve_forever, name="iris-http", daemon=True)
self._thread.start()
self.enabled = True
scheme = "https" if (self._adapter.http_cert and self._adapter.http_key) else "http"
@@ -513,6 +511,10 @@ class HttpServer:
self._write_sse(
handler, "frame", None, protocol.status(self._adapter.gateway_status()).to_json()
)
# 2b. Live todo-list snapshot (ephemeral state, never outboxed):
# a (re)connecting device re-learns the agent's current plan here.
for snap in self._adapter.todo_snapshot_frames():
self._write_sse(handler, "frame", None, snap.to_json())
# 3. Live frames (cursor=None frames have no id).
while not sub.closed.is_set():
try:
+32
View File
@@ -48,6 +48,9 @@ TYPE_TOOL_START = "tool.start"
TYPE_TOOL_PROGRESS = "tool.progress"
TYPE_TOOL_END = "tool.end"
# Agent todo list (live planning state; ephemeral, never outboxed)
TYPE_TODO_UPDATE = "todo.update"
# Intermediate assistant beat (M2)
TYPE_COMMENTARY = "commentary"
@@ -486,6 +489,35 @@ def tool_end(
)
# ---------------------------------------------------------------------------
# Todo-update frame (agent planning state)
# ---------------------------------------------------------------------------
def todo_update(
chat_id: str,
todos: list[dict[str, str]],
*,
thread_id: str | None = None,
) -> Frame:
"""The agent's current todo list for a chat/thread lane.
Emitted whenever the ``todo`` tool completes (the tool result is the
authoritative full list — it also covers ``merge`` writes, whose args
carry only the changed items) and, as a snapshot, when a device opens
its event stream. Ephemeral state: never outboxed, so a reconnecting
device learns the current list from the snapshot instead of a replay.
Each item is ``{id, content, status}`` with status one of
``pending | in_progress | completed | cancelled``.
"""
return Frame(
type=TYPE_TODO_UPDATE,
chat_id=chat_id,
thread_id=thread_id,
payload={"todos": todos},
)
# ---------------------------------------------------------------------------
# Commentary frame (M2)
# ---------------------------------------------------------------------------