Files
hermes-agent/tui_gateway/tool_progress.py
Siddharth Balyan 70f5dc5f46 feat(connectors): the backend API for the desktop Connectors page; connect an app without a chat session (#115191)
* feat(connectors): the backend serves a connector's tool list, cached for 24 hours

The Connectors page opens one app and shows every tool it has. The backend
had no way to read that list.

- `tools/connectors/portal/`: a client for the portal's tool-list route and a
  JSON cache under the Hermes home, one file per portal origin and connector.
  An entry is fresh for 24 hours. After that the read revalidates with the
  stored ETag: 304 keeps the list, 404 deletes the entry, an upstream failure
  serves the stored list marked stale, and a 401 never serves the cache.
- `connectors.tools {slug, refresh}`: account-level, routed by `profile`, no
  chat session. Errors carry a fixed `reason` from one closed set on the rail.
- Every connector model that is not operation state moves into
  `tui_gateway/contracts/connectors.py`. Handlers that no chat session owns
  live in `tui_gateway/methods_connectors_account.py`.

The wire model is tolerant: an unknown facet reads as unclassified and one odd
tool never blanks a connector.

* feat(connectors): catalog, accounts and member tool rules by RPC

The Connectors page needs the app catalog, the connected account of one app,
a way to disconnect it, and the member's own on/off rules. None had an RPC.

- `connectors.catalog`: name, description, category and logo of each app.
- `connectors.accounts`, `connectors.accounts.remove`: read the accounts at
  the tool gateway and remove one by id.
- `connectors.policy.get`: the rule layers that apply to the member, widest
  first. The body is a union on `mode`, so a reader can name who turned a
  tool off.
- `connectors.policy.set`: one change, a union on `type` (the tools of one
  connector, or one connector on or off), with the revision the user saw. A
  stale revision answers `POLICY_CONFLICT`. The backend composes the upstream
  write in one pure function, so no renderer learns the upstream rules.
- Bundled MCP manifests can name their hosted twin with `connector:`, so the
  page can show one card per app.

* feat(connectors): connect an app without a chat session

Every connector RPC took a `session_id`, and a connect that did not come from
the model's tool call minted a link with no watcher. The Connectors page has
no chat session, and its card must flip to connected by itself.

- `connectors.list`, `connectors.connect`, `connectors.operation.status`,
  `connectors.operation.wake` and `connection.respond` take `owner`, a union
  on `type`: `session` (today's behaviour and authorization) or `account`
  (routed by `profile`, authorized by the live transport like `mcp.*`).
  `session_id` is gone from these params; every desktop caller sends `owner`.
- An account connect runs the same operation lifecycle on a background
  thread, under the profile's scope, so the watcher reads the account and
  settles the operation. A second connect for an app that is already
  connecting returns the open operation and mints nothing.
- `connection.update` carries `owner`. An account operation has no session to
  address, so its updates go out on the session-less broadcast path.

* feat(mcp-catalog): eighteen more bundled entries name their hosted connector

A bundled MCP entry and a hosted connector for the same app are one card
on the Connectors page only when the manifest names its hosted twin.
Linear and Notion had the field. These entries get it too: airtable,
asana, attio, calendly, dropbox, figma, railway, supabase, todoist,
betterstack, canva, cloudflare, datadog, intercom, neon, sentry, stripe
and vercel. Atlassian maps to two hosted connectors and Prisma Postgres
is not clearly the same app, so both stay without one.

* refactor(connectors): the account handlers share one gate, one params model and one write table

The six account-level handlers each repeated the availability gate, the
auth catch and the catch-all reply. One decorator now owns that, and each
handler validates its params with its contract model instead of a ladder
of isinstance checks. The five connection RPCs share one guard for the
unexpected-failure reply.

The four write composers for the member rules were the same function
with a different list key and polarity. They are one table now.

The owner union lives in contracts/common.py, so the params side and the
event side stop declaring it twice and the import cycle is gone.

An account operation start carries one event and a flag, so the wait for
the sign-in link blocks instead of polling every 50 ms. run_operation
loses its two account-only parameters; drive_operation is the second
entry point.

Tests: four deleted (they exercised pydantic or the mock), three merged
into tables, two added (a client that still sends the old top-level
session_id is refused; all six account RPCs run off the server loop).
The shared reply helper and the HTTP and managed-client fakes move to
one place each. Comments are one line or gone.

* fix(connectors): a missing tool-list route reads as "unavailable", not "connector gone"

The tool-list read treated every 404 as the portal's "this connector is
not in the catalog" answer. It deleted the cache entry and answered
CONNECTOR_NOT_FOUND, so a page would offer to remove an app that is
connected and works. A portal that does not serve the route yet answers
a bare 404 for every app.

Only the portal's own {"error": "connector_not_found"} means the
connector is gone. Any other 404 is now a tool-list outage: the cached
list is served as stale, or the RPC answers TOOLS_UNAVAILABLE.

* fix(connectors): a connect from the page returns to the app after sign-in

The sign-in link carries a return target only when the session's surface
is the desktop. A chat session binds that surface. An account-owned call
has no chat session, so nothing bound it: the link was minted without a
return target and the browser ended on the portal's done page instead of
coming back to Hermes.

Every account-owned call now runs with the process's own surface bound,
next to its profile scope. The operation thread copies that context, so
the first link and every reissued link carry the return target and the
operation id.

* test(connectors): defer the new connector RPC coverage

The tests for the new account RPCs, the portal client, the tool-list cache
and the rule composer leave this PR and come back in one later change, after
the API is settled. The same was done for #111008.

Kept: the edits that existing tests need because the five connection RPCs
now take `owner` instead of `session_id`, and the rename of the managed
client seam.

Removed: six new test files, their two fakes and the gateway conftest, and
the new cases in test_mcp_catalog.py, test_connectors_gateway_client.py,
gateway-rpc.test.ts and notifications.test.ts. Reverting this commit restores
all of them.

* fix(cli): the connection panel hands the tool thread back at once

The classic CLI's connection callback waited on a queue for the user's first
decision. The operation's watcher starts only after the callback returns, and
the watcher is what polls a hosted account, runs the 300-second deadline and
sees Ctrl+C.

For a hosted connector the panel opens on the sign-in link, where the only
key that filled the queue was Cancel. The account was never polled: the user
signed in, the panel never changed, and Esc reported the app as skipped.
Ctrl+C set the interrupt flag but left the thread parked on the queue, so the
turn never ended.

The callback now opens the panel and returns, as the gateway's callback does
for the desktop and the Ink TUI. The panel's actions already reach the
operation through apply_answer on the UI thread, so the queue is removed. An
install with a form still waits for Connect, because the backend starts no
work for a pending row. Ctrl+C now settles the operation as `interrupt`, and
open rows become `not_connected`.

Checked on the e2e rig with the fake tool gateway: hosted connect completes on
the third status read; Ctrl+C ends the turn and the polling stops; an MCP
install with a plain and a secret field still saves config and both values.

* fix(connectors): "run it again" lives in the library, so the classic CLI can use it

Making a new sign-in link for a failed or expired hosted connector was
implemented only in the JSON-RPC layer (`_reissue`). The classic CLI does not
go through JSON-RPC: its Connect button on a failed row called apply_answer,
which does nothing for a hosted operation because it has no MCP runner. The
panel showed "Waiting…" until the deadline.

`tools.connectors.run.reissue(operation, names)` now holds the checks and the
per-kind action, and returns a refusal reason or None. The gateway maps each
reason to the same JSON-RPC error as before. The CLI calls it for a hosted
row; a refusal is shown on the row. MCP rows keep their path, because Connect
on a failed MCP row re-sends the form values.

Checked on the e2e rig: a scripted failed sign-in, then Connect: a second mint
with `reinitiate: true`, a new link with a new connection id, then connected.

* feat(connectors): the account list and disconnect go through the portal

`connectors.accounts` and `connectors.accounts.remove` called the tool
gateway. They now call the portal's account-management routes
(`GET /api/v1/connectors/accounts`, `DELETE /api/v1/connectors/accounts/{id}`),
which apply the organisation membership checks and write the disconnect audit
row. There is no fallback to the gateway when the portal is unavailable, and a
removal is never retried.

The read of ONE account stays on the gateway (`GET v1/connectors/accounts/{id}`):
the portal has no such route, and the operation watcher polls it once per second.

`ConnectorClient.list_accounts` and `delete_account` are removed. The removed
account's reply model carries `connector`, which both services send.

* fix(connectors): the account RPCs answer what the portal really sends

Checked against the portal source and against the staging and production
services.

- Errors are read from the upstream error code, not the HTTP status. A rule
  write answered 409 for a stale revision and for a user with no organisation;
  both read as "the policy changed". `org_required` is now `ORG_REQUIRED` and
  403 `no_access` is `ORG_ACCESS_DENIED` on every account RPC; only a rejected
  sign-in is `NEEDS_NOUS_AUTH`. `connectors.list` and `connectors.connect` with
  the account owner map these too.
- `connectors.policy.get` and `connectors.policy.set` carry `effective`: the
  portal's own result for this user, with its stamp and without provider or
  subject ids. Nothing is recomputed locally.
- A rule write needs the revision the user saw: `expected_revision` is required
  and must be a revision string; a bad one is refused before any HTTP call.
- A tool row carries `no_auth`; a list without the upstream flag is an invalid
  answer, not `false`.
- `connectors.accounts.remove` returns the app of the removed account. An
  invalid id is `INVALID_PARAMS`.
- The tool-list cache is per signed-in member (a hash of the token's `sub`),
  so two Nous accounts on one profile do not share entries.
- A malformed slug is a local error, not a 404 from a server nobody called.

Live, staging: no revision and a malformed revision refused locally; a good
revision wrote one disabled Gmail tool and returned it in `effective`; the
same revision again answered `POLICY_CONFLICT`; the list row showed the tool;
the restore brought the member rules back to the start. Live, staging and
production, read-only: all 60 tool lists (5483 tools) parse.

* fix(connectors): the operation RPCs match their contract; a settled card cannot start a new link

Found by two adversarial reviews of the RPC layer and its types.

- `connectors.connect` from a chat session with no open operation is refused
  (`UNKNOWN_OPERATION`). It used to call `manage_connections` through the tool
  registry with no card: it made a link nobody watched, returned a reply
  without the required `settled` field, and named an operation that was never
  registered. There is one way into an operation: the agent's call, or the
  account owner's `connectors.connect`. "Run it again" inside an open
  operation is unchanged.
- `connection.update` for a session is routed by session key AND profile; two
  profiles with the same key no longer cross-deliver a sign-in link. The event
  payload gets the same redaction as the RPC replies.
- `connection.respond` runs on the long-handler pool: an approval can start MCP
  OAuth discovery, which blocked every RPC of the gateway while it ran.
- `connectors.list` rows are a closed snake_case model: `connector`, `enabled`,
  `connected`, `connection_status`, `status_reason`, `gateway_disabled_tools`.
  The last one is display data: the gateway enforces the rules, the backend
  only passes the list on. The phantom `name` and `description` are gone, and
  the desktop uses the generated types instead of hand-written copies.
- `tools_listing` (model-only data) no longer rides on `connectors.operation.status`.
- `unavailable` is removed from the target states and settle reasons: nothing
  produces it. The contract generator now fails when a contract enum and its
  domain enum differ.
- `ConnectorErrorReason` is part of the generated TypeScript and OpenRPC.
- The desktop sends `connection.respond` on the socket that holds the session,
  as wake and reissue already did.
- Contract violations are logged every time, at error level.
- An account connect whose prepare step is slow returns the live operation
  instead of an error while the operation keeps running.
- The MCP-manifest `connector` field leaves this PR (it moves to a later one
  on top of the catalog-reader change). `hermes_cli/mcp_catalog.py` and
  `optional-mcps/` are untouched by this PR again.

anti-slop: no net-new findings (15 touched files).

* fix(connectors): the model gets no sign-in link wherever a card exists; side agents cannot connect

The flag that tells the model "a connection card exists" was the session
platform (`== "desktop"`). The Ink TUI and the classic CLI also draw a card,
so there a connector call on an unconnected app handed the model the raw
`connect_url` and told it to pass the link to the user.

- The agent turn now declares how a link can reach the user
  (`tools/connectors/turn.py`): CARD when the agent was built with a
  connection callback, SIDE for a subagent or a background turn, LINK for a
  headless run (`-q`, cron, ACP, api_server, messaging). It is set once per
  tool batch in the agent loop and read by the connector dispatch path, which
  never sees the agent. The session platform decides return-to-app only.
- CARD: the result carries `connect_card_available` and our hint, never the
  link and never the gateway's own hint.
- SIDE: subagents (`delegate_tool`), gateway background turns and the classic
  CLI `/bg` are built with `side_agent=True`. They hold no `manage_connections`
  tool on any path that derives the tool list, and a connector call on an
  unconnected app gets no link, only "report this to the main agent".
- LINK is unchanged.
- The hosted path with no card builds a detached operation, as the MCP path
  does, so no `connection.update` is emitted for an operation no client asked
  for. Names and docstrings that said "off desktop" now say "no card".
- A settled card is dead on the desktop: `reissueConnectionTarget` and
  `respondToConnectionRequest` share one guard and send nothing for a settled
  or unknown operation.
- The model-facing settled result no longer carries `connection_id`; the model
  repeated it to the user.

Shown on the real clients with a real model (rig, fake tool gateway): Ink TUI
and classic CLI get `connect_card_available` and no link, the model opens the
card, the account connects, the retried call succeeds; `-q` still gets the
link; a subagent and a background turn have no `manage_connections` and get
the no-link hint; on the desktop a card settled with Continue has no enabled
control and sends no RPC.

* feat(tools): every call made through tool_search + tool_call shows a real label on all three clients

A bridged call showed as a generic `tool_call` row in the Ink TUI and as
`⚡ tool_call` in the classic CLI, because the display looked the name up in
the tool registry and bridged names are made at run time. The desktop labelled
only batches that were all hosted connector calls, by parsing names itself.

- `tools/tool_labels.py` is the one place that turns a bridged call into a
  label: kind, app, action, emoji and text. Hosted: `connectors__gmail__GMAIL_SEND_EMAIL`
  → "Gmail · send email". MCP: "Linear · list issues". A local deferred tool
  keeps its own emoji, verb and primary-argument preview. A batch gets exactly
  one label per entry, always; an entry with no name gets a generic label.
- Classic CLI: one row per inner call; the duration on the last row; the
  failure text on the row of the call that failed. With friendly labels off
  it prints what it printed before.
- Gateway: tool start, progress and complete events and stored transcript rows
  carry a typed `labels` field. It does not depend on the classic CLI's
  display setting. Clients no longer parse tool names.
- Ink TUI: rows from the labels; the verbose trail keeps Args and Result.
- Desktop: `ConnectorExecution` renders hosted, MCP and mixed turns from the
  labels, one row per call. The labels reach the row under a key no tool
  argument can use. The connect card it drew under a failed tool result is
  gone: after `CONNECTION_REQUIRED` the one way in is the agent's own
  `manage_connections` call.
- `tool_search` and `tool_describe` rows read "Searching tools · <query>" and
  "Reading tool details · N tools".

Shown on the real desktop (video and screenshots), the Ink TUI and the classic
CLI with the rig: hosted rows, MCP rows, a two-entry batch, a failed entry, a
`CONNECTION_REQUIRED` row with no card under it, labels after a reload, and the
desktop rows with the classic CLI setting off.

* fix(connectors): the model can tell "hosted tools unavailable" from "no such tool"; manage_connections routes MCP names correctly

- A failed hosted search or describe used to return nothing, by design, so the
  model saw only local tools and told the user that a connected app was
  missing. The local results are unchanged; when the hosted leg failed, the
  `tool_search` and `tool_describe` results carry
  `connectors: {status: "unavailable", reason: "unreachable" | "sign_in_expired"}`
  and one hint line. A rejected token is `sign_in_expired`; an entitlement
  refusal or a shut gate adds nothing. `tool_describe` no longer lists those
  names under `not_found` next to "search again".
- NS-932. The description now says which side a name belongs to: a bare name
  is a hosted connector account; `mcp: true` only when the user asks for an MCP
  server, a local server or an install, or when the name exists only in the
  catalog; connect and reconnect are hosted verbs, install, enable and
  authorize are MCP verbs. It names the three clients that draw a card.
- A misrouted target is refused with the call that works. Only when the
  gateway does not know the connector (confirmed on that failure path) and the
  name is a catalog entry does the target fail with "X is a local MCP server.
  Call manage_connections with action install ...". It is a per-target
  outcome: other targets of the same call keep their links and their card. A
  vendor failure on a name both sides know stays an ordinary failed row. The
  MCP side mirrors it, and never for an entry that is only not installed.
- "Do not re-ask after a skip or a timeout" no longer stops the model when the
  USER asks for that app again; the description and the settled-result notes
  say so. A builder saw the model refuse a direct user request.

Shown on the Ink TUI and the classic CLI with a real model: a dead gateway and
a 401; "connect fxmail" goes hosted; "install the fx-noauth MCP server" goes
MCP; "connect fx-noauth" reaches the MCP install card in one corrective round
with no hosted mint; a two-target call where one is misrouted still connects
the other with exactly one mint.

* fix(tui): the connection card answers every key, shows what is happening, and is dead once settled

Reproduced on the real Ink TUI with the rig, then fixed:

- The keyboard was dead during the sign-in wait: the card kept a `submitting`
  flag that the normal OAuth path never cleared, and Esc went through the same
  guard. The in-flight state now belongs to the answered row and clears when
  that row moves, when any later frame of the operation arrives, or after
  five seconds. Esc skips the row in every phase; Ctrl+C interrupts the turn
  (the input handler had no branch for this overlay); Shift+arrows scroll the
  transcript and the card ignores them; arrow keys no longer move the text
  cursor and the field focus at once.
- The card was lost at turn idle: the overlay flag was cleared while the
  operation stayed in the store, and a resume dropped the pending card. The
  flag survives idle, a resume shows the pending card again, a session switch
  clears it.
- States with no branch: `not_connected` and a row with no link fell into the
  credential form; `expired` vanished with no note. The title and the row text
  now name the action (connect, reconnect, install, enable, authorize); a
  failed or expired row with no fields offers Try again / Skip; a failed row
  WITH fields reopens the form over the typed draft, with the failure above it.
- A settled card is dead: at settle the overlay closes and one transcript line
  per app states the outcome. A settled or dismissed operation id is
  remembered, so no replay or resume can reopen its card. Esc in the last
  "Finishing…" moment hides the card and still writes the outcome lines.
- A failed `connection.respond` and a browser that did not open are shown on
  the card in one sentence.

Also: `tui_gateway/connector_payload.py` redacted the BOOLEAN `secret` flag of
a credential field to the string "[REDACTED]". On the desktop every credential
field therefore rendered as a password and lost its prefilled default. A
boolean is no longer redacted.

* chore(connectors): remove the comments and docstrings this branch added

Deletions only. Kept: tool directives (`# noqa`, `// eslint-disable`, ...),
`// SAFETY:` lines, and the docstrings of the contract models under
`tui_gateway/contracts/`, which become the descriptions in the generated
OpenRPC and TypeScript.

Checked that no code changed: every Python file has the same AST as before
once docstrings and `pass` are ignored (62 files), and every TypeScript file
prints the same with comments stripped by the TypeScript printer (32 files).
The generated contract files are unchanged.

* fix(connectors): a card restored after a reload answers again; every account RPC names auth and org failures

Found by the end-to-end runs on the pushed head.

- Desktop: after a window reload, Continue on the restored card sent nothing.
  The answer looked up the backend that holds the session with the runtime
  session id, the lookup wants the stored id, and a failed lookup returned
  silently. When the lookup gives no owner the answer now goes out on the
  window's active socket, which is what main does.
- `connectors.policy.get` answered `POLICY_UNAVAILABLE` for a rejected sign-in,
  a refused scope, a non-member and a missing organisation alike: the handler
  runs with the gateway's globals and did not import the reason enum, so its
  own error mapping raised. `connectors.accounts.remove` caught auth failures
  in its generic branch. `org_required` was mapped on `policy.set` only. All
  six account RPCs now answer `NEEDS_NOUS_AUTH`, `FORBIDDEN_SCOPE`,
  `ORG_ACCESS_DENIED` and `ORG_REQUIRED` for those four upstream answers.
2026-09-22 13:57:51 +05:30

433 lines
19 KiB
Python

"""Tool lifecycle callbacks (tool.start/complete/progress events), verbose-text capping/redaction, todo-state
projection. Bodies are rebound onto server.py's globals (method_ctx.bind_module) and reference them bare."""
from __future__ import annotations
from .method_ctx import bind_module
# Verbose tool text is capped to the Ink render budget (a hair more, so the "[omitted …]" label
# stays informative): unbounded output fed a render-tree blowup that OOM-killed the TUI parent.
# Full output stays in the agent context and the SQLite session, untouched.
# Tool Args/Result text shipped to the TUI for the verbose trail line. The TUI renders only a small
# persisted preview (ui-tui VERBOSE_TRAIL_MAX_CHARS), kept all session and expanded by default — so shipping
# more than that is pure pipe waste AND feeds the Ink render-tree blowup that silently OOM-killed the TUI
# parent (#34095).
_TUI_VERBOSE_TEXT_MAX_CHARS = 1_000
_TUI_VERBOSE_TEXT_MAX_LINES = 16
_TODO_TOOL_NAMES = ("todo_list", "todo") # legacy alias: pre-rename replays
def _cap_tui_verbose_text(text: str) -> str:
if len(text) <= _TUI_VERBOSE_TEXT_MAX_CHARS and text.count("\n") < _TUI_VERBOSE_TEXT_MAX_LINES:
return text
# Start of the last MAX_LINES lines, then pull forward to the char budget (never mid-line).
line_start = len(text) - len("\n".join(text.split("\n")[-_TUI_VERBOSE_TEXT_MAX_LINES:]))
start = max(line_start, len(text) - _TUI_VERBOSE_TEXT_MAX_CHARS)
if start > line_start:
next_break = text.find("\n", start)
if 0 <= next_break < len(text) - 1:
start = next_break + 1
tail = text[start:].lstrip()
omitted_chars = max(0, len(text) - len(tail))
omitted_lines = text[:start].count("\n")
omitted = f"{omitted_lines} lines / {omitted_chars} chars" if omitted_lines else f"{omitted_chars} chars"
return f"[showing verbose tail; omitted {omitted}]\n{tail}"
def _redact_tui_verbose_text(text: str) -> str:
try:
from agent.redact import redact_sensitive_text
redacted = redact_sensitive_text(str(text), force=True)
except Exception:
return ""
return _cap_tui_verbose_text(redacted)
def _verbose_text(render, fallback) -> str:
"""Redacted+capped ``render()``; ``fallback()`` when rendering raises."""
try:
raw = render()
except Exception:
raw = fallback()
return _redact_tui_verbose_text(raw)
def _tool_args_text(args: dict) -> str:
return _verbose_text(lambda: json.dumps(args or {}, indent=2, ensure_ascii=False, default=str), lambda: str(args or {}))
def _tool_result_text(result: object) -> str:
def render():
from agent.tool_dispatch_helpers import _multimodal_text_summary
return _multimodal_text_summary(result)
return _verbose_text(render, lambda: str(result))
def _fmt_tool_duration(seconds: float | None) -> str:
if seconds is None:
return ""
if seconds < 10:
return f"{seconds:.1f}s"
if seconds < 60:
return f"{round(seconds)}s"
mins, secs = divmod(int(round(seconds)), 60)
return f"{mins}m {secs}s" if secs else f"{mins}m"
def _count_list(obj: object, *path: str) -> int | None:
cur = obj
for key in path:
if not isinstance(cur, dict):
return None
cur = cur.get(key)
return len(cur) if isinstance(cur, list) else None
# tool name -> (count-from-result, verb, singular, plural) for the tool.complete summary line.
_SUMMARY_COUNTERS = {
"web_search": (lambda d: _count_list(d, "data", "web"), "Did", "search", "searches"),
"web_extract": (lambda d: _count_list(d, "results") or _count_list(d, "data", "results"), "Extracted", "page", "pages"),
}
def _tool_summary(name: str, result: str, duration_s: float | None) -> str | None:
try:
data = json.loads(result)
except Exception:
data = None
if not isinstance(data, dict):
return None
dur = _fmt_tool_duration(duration_s)
suffix = f" in {dur}" if dur else ""
warning = str(data.get("fallback_warning") or "").strip()
if warning:
return f"{warning}{suffix}"
entry = _SUMMARY_COUNTERS.get(name)
n = entry[0](data) if entry else None
return f"{entry[1]} {n} {entry[2] if n == 1 else entry[3]}{suffix}" if n is not None else None
def _normalize_todo_state(value: object) -> dict | None:
"""Return a client-safe full todo snapshot or ``None`` when malformed."""
if not isinstance(value, dict) or not isinstance(value.get("todos"), list):
return None
try:
revision = max(0, int(value.get("revision") or 0))
except (TypeError, ValueError):
return None
todos = list(value["todos"])
# Unused TodoStore snapshot() is {todos: [], revision: 0}: attaching it on resume stamps a client
# watermark and blocks unversioned tool.start merges. Empty at revision >= 1 is a real clear.
if not todos and revision == 0:
return None
return {"todos": todos, "revision": revision}
def _cache_todo_state(session: dict, state: dict | None) -> None:
"""Keep the newest snapshot on the session (revision-monotonic)."""
cached = _normalize_todo_state(session.get("todo_state")) if state is not None else None
if state is not None and (cached is None or state["revision"] >= cached["revision"]):
session["todo_state"] = state
def _session_todo_state(session: dict) -> dict | None:
"""Return the newest live/cached todo snapshot for a runtime session."""
cached = _normalize_todo_state(session.get("todo_state"))
live = None
snapshot = getattr(getattr(session.get("agent"), "_todo_store", None), "snapshot", None)
if callable(snapshot):
try:
live = _normalize_todo_state(snapshot())
except Exception:
logger.debug("failed to read live todo state", exc_info=True)
if live is not None and (cached is None or live["revision"] >= cached["revision"]):
cached = live
if cached is not None:
session["todo_state"] = cached
return cached
def _attach_todo_state(payload: dict, session: dict) -> dict:
"""Attach the authoritative todo snapshot to a session response."""
state = _session_todo_state(session)
if state is not None:
payload["todo_state"] = state
return payload
def _todo_state_from_history(history) -> dict | None:
"""Latest todo snapshot from a loaded transcript, for resume paths that answer before an AIAgent (and
its live TodoStore) exists: the newest tool result paired with an assistant ``todo`` call IS it."""
if not isinstance(history, list) or not history:
return None
try:
from tools.todo_tool import MAX_TODO_RESULT_CHARS
todo_call_ids = {
call.get("id")
for msg in history if isinstance(msg, dict)
for call in msg.get("tool_calls") or []
if (call.get("function") or {}).get("name") in _TODO_TOOL_NAMES and call.get("id")
}
if not todo_call_ids:
return None
for msg in reversed(history):
if not isinstance(msg, dict) or msg.get("role") != "tool" or msg.get("tool_call_id") not in todo_call_ids:
continue
content = msg.get("content", "")
if not isinstance(content, str) or len(content) > MAX_TODO_RESULT_CHARS or '"todos"' not in content:
continue
try:
return _normalize_todo_state(json.loads(content))
except Exception:
continue
return None
except Exception:
logger.debug("failed to derive todo state from history", exc_info=True)
return None
def _tool_labels(name: str, args: dict) -> list[dict] | None:
from agent.display import tool_labels_for_call
return [label.as_payload() for label in tool_labels_for_call(name, args)] or None
def _connector_tool_lifecycle(name: str, args: dict) -> bool:
from tools.connectors import is_connector_name
if name == "manage_connections" or is_connector_name(name):
return True
if name != "tool_call" or not isinstance(args, dict):
return False
calls = args.get("calls") if isinstance(args.get("calls"), list) else [args]
return any(isinstance(call, dict) and (call.get("name") == "manage_connections"
or is_connector_name(call.get("name"))) for call in calls)
def _connector_lifecycle_is_stale(sid: str, name: str, args: dict) -> bool:
if not _connector_tool_lifecycle(name, args):
return False
owner = _current_runtime_session_record.get()
return owner is not None and (_sessions.get(sid) is not owner or owner.get("_finalized", False))
def _emit_tool_lifecycle(event, sid, name, args, payload):
if not _connector_tool_lifecycle(name, args):
return _emit(event, sid, payload)
from tui_gateway.connector_payload import connector_ui_payload
from tui_gateway.event_replay import _stamp_event
payload = connector_ui_payload(payload)
# Capture the owner after projection so id reuse cannot redirect its link,
# then release the session lock before potentially blocking transport I/O.
with _sessions_lock:
if _connector_lifecycle_is_stale(sid, name, args):
return
transport = (_sessions.get(sid) or {}).get("transport")
if transport is None:
transport = current_transport() or _stdio_transport
frame = _event_frame(event, sid, payload)
_stamp_event(frame)
from tui_gateway.hosted_room_member_activity import project_room_member_activity
project_room_member_activity(frame, _sessions)
transport.write(frame)
def _on_tool_start(sid: str, tool_call_id: str, name: str, args: dict):
if _connector_lifecycle_is_stale(sid, name, args):
return
session = _sessions.get(sid)
if session is not None:
with contextlib.suppress(Exception):
from agent.display import capture_local_edit_snapshot
snapshot = capture_local_edit_snapshot(name, args)
if snapshot is not None:
session.setdefault("edit_snapshots", {})[tool_call_id] = snapshot
session.setdefault("tool_started_at", {})[tool_call_id] = time.time()
if (_tool_progress_enabled(sid) or _tool_lifecycle_required_for_ui(name)
or _connector_tool_lifecycle(name, args)):
payload: dict[str, object] = {"tool_id": tool_call_id, "name": name, "context": _tool_ctx(name, args)}
if (labels := _tool_labels(name, args)) is not None:
payload["labels"] = labels
# Full args (not just the 80-char `context` preview) so the desktop's expanded tool row is complete
# while the tool runs. args.todos may be a partial merge — tool.complete is the truth.
if args:
payload["args"] = args
if _session_verbose(sid) and (args_text := _tool_args_text(args)):
payload["args_text"] = args_text
_emit_tool_lifecycle("tool.start", sid, name, args, payload)
def _on_tool_complete(sid: str, tool_call_id: str, name: str, args: dict, result: str):
if _connector_lifecycle_is_stale(sid, name, args):
return
payload = {"tool_id": tool_call_id, "name": name, "args": args}
if (labels := _tool_labels(name, args)) is not None:
payload["labels"] = labels
session = _sessions.get(sid)
snapshot = session.setdefault("edit_snapshots", {}).pop(tool_call_id, None) if session is not None else None
started_at = session.setdefault("tool_started_at", {}).pop(tool_call_id, None) if session is not None else None
duration_s = time.time() - started_at if started_at else None
if duration_s is not None:
payload["duration_s"] = duration_s
try:
payload["result"] = json.loads(result)
except Exception:
payload["result"] = result
summary = _tool_summary(name, result, duration_s)
if summary:
payload["summary"] = summary
if _session_verbose(sid) and (result_text := _tool_result_text(result)):
payload["result_text"] = result_text
todo_state = _normalize_todo_state(payload.get("result")) if name in _TODO_TOOL_NAMES else None
if todo_state is not None:
payload.update(todo_state)
if session is not None:
_cache_todo_state(session, todo_state)
with contextlib.suppress(Exception):
from agent.display import render_edit_diff_with_delta
rendered: list[str] = []
if render_edit_diff_with_delta(name, result, function_args=args, snapshot=snapshot, print_fn=rendered.append):
payload["inline_diff"] = "\n".join(rendered)
if (_tool_progress_enabled(sid) or payload.get("inline_diff") or _tool_lifecycle_required_for_ui(name)
or name in _TODO_TOOL_NAMES or _connector_tool_lifecycle(name, args)):
_emit_tool_lifecycle("tool.complete", sid, name, args, payload)
# Task state is application data, not tool-progress chrome: a dedicated full-snapshot event lets
# every client reconcile without parsing tool args.
if todo_state is not None:
_emit("todo.updated", sid, todo_state)
# ── _on_tool_progress dispatch: each handler takes (sid, name, preview, kw) ─────────────────────
# `tool.started` is dropped on purpose: _on_tool_start already emits the authoritative tool.start with
# the stable id and args; an id-less duplicate row makes the desktop live view diverge from history.
def _progress_output_risk(sid, name, preview, kw):
metadata = kw.get("risk_metadata")
if isinstance(metadata, dict):
_emit("tool.output_risk", sid, {
"tool_id": str(kw.get("tool_call_id") or ""), "name": str(name), "risk": str(metadata.get("risk") or "low"),
"findings": [str(item) for item in metadata.get("findings", [])], "redacted": bool(metadata.get("redacted", False)),
})
def _progress_reasoning(sid, name, preview, kw):
_emit("reasoning.available", sid, {"text": str(preview), **({"verbose": True} if _session_verbose(sid) else {})})
def _progress_moa_reference(sid, name, preview, kw):
# MoA reference-model output, rendered as a labelled block before the aggregator's response.
# `name` is the slot label, `preview` the text.
ref_payload: dict[str, object] = {"label": str(name), "text": str(preview or "")}
for key, out in (("moa_index", "index"), ("moa_count", "count")):
if kw.get(key) is not None:
ref_payload[out] = kw[key]
_emit("moa.reference", sid, ref_payload)
def _progress_moa_progress(sid, name, preview, kw):
# Drives the status-bar `MOA: 2/3 refs done`; both counters required for deterministic rendering.
refs_done, refs_total = kw.get("moa_refs_done"), kw.get("moa_refs_total")
# Per-reference completion — drives the status-bar progress indicator (`MOA: 2/3 refs done`) requested
# in issue #59546. Only emitted when both counters are present so the client can render
# deterministically.
if refs_done is None or refs_total is None:
return
_emit("moa.progress", sid, {"label": str(name or ""), "refs_done": int(refs_done), "refs_total": int(refs_total)})
def _progress_moa_phase(sid, name, preview, kw):
# Currently only phase="aggregator" fires, once fan-out completes.
phase = kw.get("moa_phase")
if not phase:
return
phase_payload: dict[str, object] = {"phase": str(phase)}
for key, out in (("moa_refs_done", "refs_done"), ("moa_refs_total", "refs_total")):
if kw.get(key) is not None:
phase_payload[out] = int(kw[key])
if name:
phase_payload["aggregator"] = str(name)
_emit("moa.phase", sid, phase_payload)
def _not_none(v):
return v is not None
def _str_list(v):
return [str(x) for x in v]
def _int_or_skip(v):
"""Per-branch token/api rollups tolerate junk from older emitters: unparsable -> field omitted."""
try:
return int(v)
except (TypeError, ValueError):
return None
# Optional subagent.* payload fields in WIRE ORDER: (source key, present-when, coerce). Identity fields
# are all optional: older emitters omit them and the TUI spawn tree falls back to flat rendering.
# `tool_name`/`text` are fed from the positional name/preview; `output_tail` is a list of dicts.
_SUBAGENT_FIELDS = (
("subagent_id", bool, str), ("parent_id", bool, str), ("child_session_id", bool, str),
("delegation_id", bool, str), ("depth", _not_none, int), ("model", bool, str), ("tool_count", _not_none, int),
("toolsets", bool, _str_list), ("input_tokens", _not_none, _int_or_skip), ("output_tokens", _not_none, _int_or_skip),
("reasoning_tokens", _not_none, _int_or_skip), ("api_calls", _not_none, _int_or_skip),
("files_read", bool, _str_list), ("files_written", bool, _str_list), ("output_tail", bool, list),
("tool_name", bool, str), ("text", bool, str), ("status", bool, str), ("summary", bool, str),
("duration_seconds", _not_none, float),
)
def _progress_subagent(sid, name, preview, kw, event_type):
payload = {"goal": str(kw.get("goal") or ""), "task_count": int(kw.get("task_count") or 1), "task_index": int(kw.get("task_index") or 0)}
source = {**kw, "tool_name": name, "text": preview}
for key, present, coerce in _SUBAGENT_FIELDS:
if present(source.get(key)):
val = coerce(source[key])
if val is not None:
payload[key] = val
if preview and event_type == "subagent.tool":
payload["tool_preview"] = str(preview)
payload["text"] = str(preview)
# subagent.text is the child's per-token reply, relayed solely to feed a watch window's live mirror
# (keyed off the child sid); on the parent it's hundreds of ignored frames, so skip it.
if event_type != "subagent.text":
_emit(event_type, sid, payload)
_mirror_subagent_to_child(event_type, payload)
# event_type -> (handler, requires): `requires` names the arg that must be truthy for the row to be
# emitted at all ("name" / "preview" / None).
_PROGRESS_HANDLERS = {
"tool.output_risk": (_progress_output_risk, "name"), "reasoning.available": (_progress_reasoning, "preview"),
"moa.reference": (_progress_moa_reference, "name"),
"moa.aggregating": (lambda sid, name, preview, kw: _emit("moa.aggregating", sid, {"aggregator": str(name or "")}), None),
"moa.progress": (_progress_moa_progress, None), "moa.phase": (_progress_moa_phase, None),
}
def _on_tool_progress(
sid: str, event_type: str, name: str | None = None, preview: str | None = None,
_args: dict | None = None, **_kwargs,
):
if event_type == "tool.started" and name:
return
# Subagent lifecycle is application state (Desktop status stack, TUI spawn tree), not
# tool-progress chrome: it must survive display.tool_progress=off like todo.updated does.
if event_type.startswith("subagent."):
return _progress_subagent(sid, name, preview, _kwargs, event_type)
if not _tool_progress_enabled(sid):
return
handler, requires = _PROGRESS_HANDLERS.get(event_type, (None, None))
if handler is not None and (requires is None or {"name": name, "preview": preview}[requires]):
handler(sid, name, preview, _kwargs)
def register(server) -> None:
"""Publish this module's helpers + handlers onto ``server``, rebound to its globals."""
bind_module(globals(), server, skip=("_",))