"""Tool-progress frame emission (M2): the tool.start / tool.end lifecycle. Mixin for ``adapter.IrisAdapter``. The gateway accumulates tool lines in one editable bubble; on an edit the full buffer is re-sent, so new lines are diffed against ``seen_tool_lines`` and each new tool closes the previously open one (attaching the output/duration captured by the ``post_tool_call`` hook). """ from typing import Any from gateway.platforms.base import SendResult from . import protocol from .classify import ( _extract_code_block, _extract_verbose_args, _mint_message_id, _parse_tool_line, _short_preview_from_args, _TurnState, ) from .hooks import _reset_tool_results, _tool_emoji, _tool_end_fields from .mixin_base import IrisAdapterBase class ToolProgressHandlers(IrisAdapterBase): """Tool-progress lifecycle (see module docstring).""" async def _emit_tool_lines( self, chat_id: str, content: str, state: _TurnState, thread_id: str | None, *, is_edit: bool, ) -> SendResult: """Emit ``tool.start`` for each NEW tool line in *content*. The gateway accumulates tool lines in one editable bubble; on an edit the full buffer is re-sent, so we diff against ``seen_tool_lines`` to emit only the new ones. A new tool closes the previously-open tool. """ 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: key = line.strip() if key in state.seen_tool_lines: continue state.seen_tool_lines.add(key) parsed = self._parse_tool_line_or_block(line, content) if parsed is None: continue name, preview, args = parsed # A new tool begins: close the previously-open one, attaching the # output/duration/ok captured by the post_tool_call hook. if state.open_tool_index is not None: extra = _tool_end_fields(state.open_tool_name or "") await self._broadcast_or_log( chat_id, protocol.tool_end( chat_id, state.open_tool_index, state.open_tool_name or "", ok=extra.get("ok", True), duration=extra.get("duration"), output_preview=extra.get("output_preview"), thread_id=thread_id, ), ) state.tool_index += 1 state.open_tool_index = state.tool_index state.open_tool_name = name await self._broadcast_or_log( chat_id, protocol.tool_start( chat_id, state.tool_index, name, preview=preview, args=args, emoji=_tool_emoji(name), thread_id=thread_id, ), ) return SendResult(success=True, message_id=message_id) @staticmethod def _parse_tool_line_or_block( line: str, content: str ) -> tuple[str, str | None, dict[str, Any] | None] | None: """Parse a tool line into ``(name, preview, args)``. Expands a terminal code block to its command, and a verbose header (`` (keys)``) to its full args JSON (the JSON sits on the following line). ``args`` is ``None`` unless the line is a verbose header with a parseable JSON body. """ parsed = _parse_tool_line(line) if parsed is None: return None name, preview = parsed # Terminal code block: the command lives in the fenced lines that # follow the " terminal" head line. if name == "terminal" and preview is None and "```" in content: cmd = _extract_code_block(content) if cmd: return name, cmd, None # Verbose mode: recover the full args from the following JSON line. args = _extract_verbose_args(line, content) if args is not None and preview is None: preview = _short_preview_from_args(args) return name, preview, args async def _close_open_tool( self, chat_id: str, state: _TurnState, thread_id: str | None ) -> None: """Emit ``tool.end`` for the currently-open tool, if any. A tool is considered complete when the next tool starts OR a new content segment begins (the model only produces content after the tool it was waiting on has returned). """ if state.open_tool_index is not None: extra = _tool_end_fields(state.open_tool_name or "") await self._broadcast_or_log( chat_id, protocol.tool_end( chat_id, state.open_tool_index, state.open_tool_name or "", ok=extra.get("ok", True), duration=extra.get("duration"), output_preview=extra.get("output_preview"), thread_id=thread_id, ), ) state.open_tool_index = None state.open_tool_name = None def _reset_tool_state(self, state: _TurnState) -> None: """Clear per-turn tool bookkeeping (called at turn finalization).""" state.tool_msg_id = None state.seen_tool_lines = set() state.tool_index = 0 state.open_tool_index = None state.open_tool_name = None # Drop any captured tool results not consumed by a tool.end this turn # (e.g. tool_progress off) so they can't leak into the next turn. _reset_tool_results()