diff --git a/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt b/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt index e499bd2..de11efb 100644 --- a/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt +++ b/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt @@ -20,9 +20,12 @@ import iris.protocol.TYPE_MESSAGE_START import iris.protocol.TYPE_MESSAGE_STOP import iris.protocol.TYPE_MESSAGE_UPDATE import iris.protocol.TYPE_PICKER_CHOICE +import iris.protocol.TYPE_TODO_UPDATE import iris.protocol.TYPE_TOOL_END import iris.protocol.TYPE_TOOL_PROGRESS import iris.protocol.TYPE_TOOL_START +import iris.protocol.TodoItem +import iris.protocol.TodoUpdatePayload import iris.protocol.ToolEndPayload import iris.protocol.ToolProgressPayload import iris.protocol.ToolStartPayload @@ -140,6 +143,12 @@ class ChatStore { private val _currentLane = MutableStateFlow(DEFAULT_LANE) val currentLane: StateFlow = _currentLane.asStateFlow() + /** The agent's live todo list per lane (todo.update; last-write-wins). + * Ephemeral: not persisted — the gateway re-sends a snapshot when the + * app reconnects, and the next `todo` tool call refreshes it. */ + private val _todos = MutableStateFlow>>(emptyMap()) + val todos: StateFlow>> = _todos.asStateFlow() + private var localSeq = 0 /** When false, `message.start`/`message.update` frames are ignored and each @@ -234,6 +243,7 @@ class ChatStore { TYPE_TOOL_START -> onToolStart(lane, frame) TYPE_TOOL_PROGRESS -> onToolProgress(lane, frame) TYPE_TOOL_END -> onToolEnd(lane, frame) + TYPE_TODO_UPDATE -> onTodoUpdate(lane, frame) TYPE_COMMENTARY -> onCommentary(lane, frame) TYPE_MEDIA_OFFER -> onMediaOffer(lane, frame) TYPE_PICKER_CHOICE -> onPickerChoice(lane, frame) @@ -516,6 +526,18 @@ class ChatStore { } } + // ── todo.update (the agent's live todo list, last-write-wins) ───────── + + private fun onTodoUpdate( + lane: String, + frame: Frame, + ) { + val p = frame.payloadAs() ?: return + val map = _todos.value.toMutableMap() + if (p.todos.isEmpty()) map.remove(lane) else map[lane] = p.todos + _todos.value = map + } + // ── commentary (dimmed interim beat) ────────────────────────────────── private fun onCommentary( diff --git a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt index 950903d..352be80 100644 --- a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt +++ b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt @@ -47,6 +47,11 @@ const val TYPE_MESSAGE_DELETED = "message.deleted" const val TYPE_TOOL_START = "tool.start" const val TYPE_TOOL_PROGRESS = "tool.progress" const val TYPE_TOOL_END = "tool.end" + +// Agent todo list (live planning state; ephemeral, never outboxed — a +// reconnecting device re-learns it from the snapshot after hello.ack). +const val TYPE_TODO_UPDATE = "todo.update" + const val TYPE_COMMENTARY = "commentary" // M4 — media (offer; upload/pull are HTTP, docs/19 §19.15) @@ -296,6 +301,30 @@ data class ToolEndPayload( @SerialName("output_preview") val outputPreview: String? = null, ) +// ── todo.update (server -> app): the agent's live todo list ──────────────── + +/** One item of the agent's todo list (hermes `todo` tool). [status] is one + * of [TODO_PENDING], [TODO_IN_PROGRESS], [TODO_COMPLETED], [TODO_CANCELLED]. */ +@Serializable +data class TodoItem( + val id: String, + val content: String, + val status: String, +) + +const val TODO_PENDING = "pending" +const val TODO_IN_PROGRESS = "in_progress" +const val TODO_COMPLETED = "completed" +const val TODO_CANCELLED = "cancelled" + +/** The agent's FULL current todo list for a lane (last-write-wins; the + * gateway emits it whenever the `todo` tool completes — the tool result is + * authoritative even for merge writes — and as a snapshot on connect). */ +@Serializable +data class TodoUpdatePayload( + val todos: List = emptyList(), +) + // ── M2: commentary frame (server -> app) ──────────────────────────────── @Serializable diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index 67b30ff..69543f1 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -56,6 +56,7 @@ import iris.protocol.TYPE_READ_RECEIPT import iris.protocol.TYPE_SEARCH_RESULTS import iris.protocol.TYPE_STATUS import iris.protocol.TYPE_SYNC_DONE +import iris.protocol.TYPE_TODO_UPDATE import iris.protocol.TYPE_TOOL_END import iris.protocol.TYPE_TOOL_PROGRESS import iris.protocol.TYPE_TOOL_START @@ -482,6 +483,13 @@ class IrisController( if (frame.cursor == null) chat.onFrame(frame) } + TYPE_TODO_UPDATE -> { + // The agent's live todo list (last-write-wins). Ephemeral: + // never outboxed, so it never carries a cursor and is + // always applied (a reconnect snapshot just re-sets it). + chat.onFrame(frame) + } + TYPE_MESSAGE, TYPE_MESSAGE_START, TYPE_MESSAGE_UPDATE, diff --git a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt index 7ad4d6e..5f6a25a 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt @@ -44,6 +44,7 @@ import androidx.compose.foundation.text.KeyboardActions import androidx.compose.foundation.text.KeyboardOptions import androidx.compose.foundation.verticalScroll import androidx.compose.material.icons.Icons +import androidx.compose.material.icons.filled.Check import androidx.compose.material.icons.filled.Favorite import androidx.compose.material.icons.filled.KeyboardArrowDown import androidx.compose.material.icons.filled.Settings @@ -101,6 +102,7 @@ import androidx.compose.ui.text.font.FontFamily import androidx.compose.ui.text.font.FontWeight import androidx.compose.ui.text.input.ImeAction import androidx.compose.ui.text.style.TextAlign +import androidx.compose.ui.text.style.TextDecoration import androidx.compose.ui.text.style.TextOverflow import androidx.compose.ui.unit.Dp import androidx.compose.ui.unit.dp @@ -129,6 +131,11 @@ import iris.protocol.ROLE_ASSISTANT import iris.protocol.ROLE_USER import iris.protocol.SearchHit import iris.protocol.SlashCommand +import iris.protocol.TODO_CANCELLED +import iris.protocol.TODO_COMPLETED +import iris.protocol.TODO_IN_PROGRESS +import iris.protocol.TODO_PENDING +import iris.protocol.TodoItem import iris.state.IrisController import iris.state.ToolDetail import iris.ui.MarkdownText @@ -675,6 +682,18 @@ fun ChatScreen(controller: IrisController) { } } + // The agent's live todo list (todo.update): a compact, scrollable + // strip above the composer so the plan stays visible without + // eating the chat. Hidden while the list is empty or fully + // resolved (mirrors the hermes desktop composer status stack). + if (!isAutomation) { + val todosMap by controller.chat.todos.collectAsState() + val todoList = todosMap[currentLane].orEmpty() + if (todoList.any { it.status == TODO_PENDING || it.status == TODO_IN_PROGRESS }) { + TodoStrip(todos = todoList) + } + } + // Pending attachments (M4; hidden in automation channels — no composer) if (attachments.isNotEmpty() && !isAutomation) { Row( @@ -3019,3 +3038,158 @@ private fun bestSlashScore( } return best } + +// ── Agent todo strip (todo.update) ──────────────────────────────────────── + +/** Rows visible in the strip before it scrolls (a phone screen is precious). */ +private const val TODO_STRIP_VISIBLE_ROWS = 3 + +/** + * The agent's live todo list: a compact strip above the composer (max 3 + * lines, the rest scrollable by tap-drag) so the user stays aware of the + * plan and the current task without the chat being disrupted. Style mirrors + * the hermes desktop composer status stack: pending = hollow ring, + * in_progress = spinner, completed = green check, cancelled = struck + * through. Whenever an item's status changes (or a new item appears) the + * list auto-scrolls to it. + */ +@Composable +private fun TodoStrip(todos: List) { + val rowHeight = 22.dp + val scrollState = rememberScrollState() + val rowHeightPx = with(LocalDensity.current) { rowHeight.toPx().roundToInt() } + // The list the strip last rendered — the diff against the new list drives + // the auto-scroll (first item whose status changed, or that is new). + var prevTodos by remember { mutableStateOf?>(null) } + + LaunchedEffect(todos) { + val prev = prevTodos + prevTodos = todos + val target = + if (prev == null) { + // First render: land on the current task (or the top). + todos.indexOfFirst { it.status == TODO_IN_PROGRESS }.takeIf { it >= 0 } ?: 0 + } else { + // Prefer the CURRENT task (in_progress) — that's what the user + // is tracking. When a task completes the NEXT one activates in + // the same update, so the first *changed* item is the one that + // just finished (still visible) while the new active task sits + // just below the fold; targeting it is what keeps the current + // task in view. Fall back to the first changed/new item when + // nothing is in flight (e.g. a freshly appended task). + todos + .indexOfFirst { it.status == TODO_IN_PROGRESS } + .takeIf { it >= 0 } + ?: todos.indices.indexOfFirst { i -> i >= prev.size || prev[i].status != todos[i].status } + } + if (target in todos.indices) { + val top = target * rowHeightPx + val viewport = TODO_STRIP_VISIBLE_ROWS * rowHeightPx + val maxScroll = (todos.size * rowHeightPx - viewport).coerceAtLeast(0) + if (top < scrollState.value || top + rowHeightPx > scrollState.value + viewport) { + scrollState.animateScrollTo(top.coerceIn(0, maxScroll)) + } + } + } + + val completed = todos.count { it.status == TODO_COMPLETED } + Column( + modifier = + Modifier + .fillMaxWidth() + .padding(horizontal = 12.dp, vertical = 4.dp), + ) { + Box( + modifier = + Modifier + .fillMaxWidth() + .clip(RoundedCornerShape(12.dp)) + .background(IrisColors.surface) + .border(1.dp, IrisColors.divider, RoundedCornerShape(12.dp)), + ) { + Column(modifier = Modifier.padding(horizontal = 10.dp, vertical = 6.dp)) { + Row(verticalAlignment = Alignment.CenterVertically) { + Text("📋", fontSize = 12.sp) + Spacer(Modifier.width(6.dp)) + Text( + "Tasks", + fontSize = 12.sp, + fontWeight = FontWeight.SemiBold, + color = IrisColors.textSecondary, + ) + Spacer(Modifier.weight(1f)) + Text( + "$completed/${todos.size}", + fontSize = 12.sp, + color = IrisColors.textSecondary, + ) + } + Spacer(Modifier.height(4.dp)) + Column( + modifier = + Modifier + .fillMaxWidth() + .height(rowHeight * TODO_STRIP_VISIBLE_ROWS) + .verticalScroll(scrollState), + ) { + todos.forEach { item -> TodoRow(item, rowHeight) } + } + } + } + } +} + +/** One todo row: status glyph + single-line content (truncated). */ +@Composable +private fun TodoRow( + item: TodoItem, + rowHeight: Dp, +) { + val resolved = item.status == TODO_COMPLETED || item.status == TODO_CANCELLED + Row( + modifier = + Modifier + .fillMaxWidth() + .height(rowHeight), + verticalAlignment = Alignment.CenterVertically, + horizontalArrangement = Arrangement.spacedBy(6.dp), + ) { + when (item.status) { + TODO_IN_PROGRESS -> { + CircularProgressIndicator( + modifier = Modifier.size(12.dp), + strokeWidth = 1.5.dp, + color = IrisColors.primary, + ) + } + + TODO_COMPLETED -> { + Icon( + imageVector = Icons.Filled.Check, + contentDescription = null, + tint = IrisColors.statusGreen, + modifier = Modifier.size(12.dp), + ) + } + + else -> { + Box( + modifier = + Modifier + .size(12.dp) + .clip(CircleShape) + .border(1.dp, IrisColors.textDim, CircleShape), + ) + } + } + Text( + text = item.content, + fontSize = 12.sp, + maxLines = 1, + overflow = TextOverflow.Ellipsis, + textDecoration = if (item.status == TODO_CANCELLED) TextDecoration.LineThrough else null, + color = if (resolved) IrisColors.textDim else IrisColors.textBright, + modifier = Modifier.weight(1f), + ) + } +} diff --git a/app/shared/src/commonTest/kotlin/iris/data/ChatStoreCacheTest.kt b/app/shared/src/commonTest/kotlin/iris/data/ChatStoreCacheTest.kt index 8cc58fc..8448d2c 100644 --- a/app/shared/src/commonTest/kotlin/iris/data/ChatStoreCacheTest.kt +++ b/app/shared/src/commonTest/kotlin/iris/data/ChatStoreCacheTest.kt @@ -2,10 +2,14 @@ package iris.data import iris.protocol.Frame import iris.protocol.IrisJson +import iris.protocol.TYPE_TODO_UPDATE import iris.protocol.TYPE_TOOL_START +import iris.protocol.TodoItem +import iris.protocol.TodoUpdatePayload import iris.protocol.ToolStartPayload import kotlin.test.Test import kotlin.test.assertEquals +import kotlin.test.assertNull class ChatStoreCacheTest { @Test @@ -110,4 +114,42 @@ class ChatStoreCacheTest { val ids = store.lanes.value["default"]!!.map { it.id } assertEquals(ids.size, ids.toSet().size) } + + @Test + fun todoUpdateSetsLaneListLastWriteWins() { + val store = ChatStore() + assertNull(store.todos.value["default"]) + store.onFrame(todoFrame("default", listOf(TodoItem("1", "a", "in_progress")))) + assertEquals(listOf("1"), store.todos.value["default"]!!.map { it.id }) + // A second update REPLACES the list (the gateway always sends the + // full list), and other lanes are untouched. + store.onFrame( + todoFrame( + "default", + listOf(TodoItem("1", "a", "completed"), TodoItem("2", "b", "pending")), + ), + ) + store.onFrame(todoFrame("chan_7", listOf(TodoItem("9", "z", "pending")))) + assertEquals(listOf("1", "2"), store.todos.value["default"]!!.map { it.id }) + assertEquals("completed", store.todos.value["default"]!![0].status) + assertEquals(listOf("9"), store.todos.value["chan_7"]!!.map { it.id }) + // An empty list clears the lane (the agent dropped its plan). + store.onFrame(todoFrame("default", emptyList())) + assertNull(store.todos.value["default"]) + assertEquals(listOf("9"), store.todos.value["chan_7"]!!.map { it.id }) + } + + private fun todoFrame( + chatId: String, + todos: List, + ): Frame = + Frame( + type = TYPE_TODO_UPDATE, + chatId = chatId, + payload = + IrisJson.instance.encodeToJsonElement( + TodoUpdatePayload.serializer(), + TodoUpdatePayload(todos = todos), + ), + ) } diff --git a/docs/04-wire-protocol.md b/docs/04-wire-protocol.md index 67ade76..119650a 100644 --- a/docs/04-wire-protocol.md +++ b/docs/04-wire-protocol.md @@ -140,6 +140,31 @@ short tail (full output is not streamed — it lives in agent history). `emoji` — e.g. `read_file` 📖, `write_file` ✍️, `terminal` 💻). Omitted when the tool is unknown, so the app falls back to its own default glyph. +### `todo.update` + +The agent's **full current todo list** for a chat/thread lane (last-write-wins). +Emitted whenever the `todo` tool completes — the tool *result* is the +authoritative list, which also covers `merge` writes (whose args carry only +the changed items) and read-only calls — and, as a **snapshot**, right after +`hello` when a device opens its event stream. + +```json +{"type":"todo.update","chat_id":"default","thread_id":null,"payload":{ + "todos":[ + {"id":"1","content":"Scaffold the module","status":"completed"}, + {"id":"2","content":"Wire the store","status":"in_progress"}, + {"id":"3","content":"Add tests","status":"pending"} + ]}} +``` + +`status` is one of `pending | in_progress | completed | cancelled`. The frame +is **ephemeral**: it is never outboxed, so a reconnecting device learns the +current list from the snapshot (not a replay), and a gateway restart drops it +(the agent re-emits on the next `todo` call). The app renders it as a compact +scrollable strip above the composer (max 3 lines; auto-scrolls to the item +whose status just changed) and hides it once the list is empty or fully +resolved. + ### `typing` / `typing.stop` ```json diff --git a/docs/protocol/frames.schema.json b/docs/protocol/frames.schema.json index 967b669..952d694 100644 --- a/docs/protocol/frames.schema.json +++ b/docs/protocol/frames.schema.json @@ -50,6 +50,7 @@ "tool.start": { "description": "Cosmetic per-tool emoji (resolved server-side via hermes' get_tool_emoji: active-skin overrides, then the tool registry's per-tool emoji); omitted when the tool is unknown so the app falls back to its own default glyph.", "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "preview": { "type": "string" }, "args": { "type": "object" }, "emoji": { "type": "string" } } }, "tool.progress": { "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "note": { "type": "string" } } }, "tool.end": { "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "ok": { "type": "boolean" }, "duration": { "type": "number" }, "output_preview": { "type": "string" } } }, + "todo.update": { "description": "The agent's FULL current todo list for a chat/thread lane (last-write-wins). Emitted whenever the `todo` tool completes (the tool result is authoritative even for merge writes, whose args carry only the changed items) and, as a snapshot, right after hello when a device opens its event stream. Ephemeral: never outboxed, so a reconnecting device learns the list from the snapshot, not a replay. The app renders it as a compact scrollable strip above the composer.", "payload": { "todos": { "type": "array", "items": { "type": "object", "properties": { "id": { "type": "string" }, "content": { "type": "string" }, "status": { "type": "string", "enum": ["pending", "in_progress", "completed", "cancelled"] } } } } } }, "typing": { "payload": { "on": { "type": "boolean" } } }, "notification": { "payload": { "kind": { "type": "string", "enum": ["channel_renamed", "channel_created", "channel_deleted", "cron", "approval", "clarify", "generic"] }, "title": { "type": "string" }, "body": { "type": "string" }, "ts": { "type": "integer" } } }, "channel.list": { "description": "Full channel directory (response to a channel.list request).", "payload": { "channels": { "type": "array", "items": { "$ref": "#/definitions/channel" } } } }, diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index 6943075..6fff61a 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -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.""" diff --git a/gateway-plugin/http_server.py b/gateway-plugin/http_server.py index 621c5c9..e93181e 100644 --- a/gateway-plugin/http_server.py +++ b/gateway-plugin/http_server.py @@ -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: diff --git a/gateway-plugin/protocol.py b/gateway-plugin/protocol.py index c985354..b038843 100644 --- a/gateway-plugin/protocol.py +++ b/gateway-plugin/protocol.py @@ -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) # ---------------------------------------------------------------------------