fix(tui-gateway): report recently failed async delegations in subagent.list
This commit is contained in:
@@ -2062,10 +2062,10 @@ export interface SessionForeignImportResult {
|
||||
session_id: string
|
||||
already_imported?: boolean
|
||||
}
|
||||
/** ``delegations`` is reserved for async delegation records and is currently always empty. */
|
||||
/** ``delegations``: recently failed async delegation tasks for the session (durable store), newest first. */
|
||||
export interface SubagentListResult {
|
||||
subagents?: SubagentSnapshot[]
|
||||
delegations?: Record<string, unknown>[]
|
||||
delegations?: FailedDelegation[]
|
||||
}
|
||||
/** ``methods_subagents._SUBAGENT_SNAPSHOT_FIELDS`` projection of one live child record. */
|
||||
export interface SubagentSnapshot {
|
||||
@@ -2083,6 +2083,16 @@ export interface SubagentSnapshot {
|
||||
}
|
||||
/** Lifecycle of one delegated child (``tools/delegate_tool_child_run.py``); ``failed`` / ``error`` / ``timeout`` / ``interrupted`` / ``completed`` are terminal. */
|
||||
export type SubagentStatus = 'queued' | 'running' | 'completed' | 'failed' | 'error' | 'timeout' | 'interrupted'
|
||||
/** ``async_delegation.failed_delegations_for_session`` row: one failed task of an async delegation. */
|
||||
export interface FailedDelegation {
|
||||
delegation_id: string
|
||||
task_index?: number
|
||||
status: string
|
||||
goal?: string
|
||||
error?: string | null
|
||||
dispatched_at?: number | null
|
||||
completed_at?: number | null
|
||||
}
|
||||
export interface SubagentIdParams {
|
||||
session_id: string
|
||||
profile?: string | null
|
||||
|
||||
@@ -13092,6 +13092,72 @@
|
||||
"title": "ErrorSurface",
|
||||
"type": "object"
|
||||
},
|
||||
"FailedDelegation": {
|
||||
"additionalProperties": false,
|
||||
"description": "``async_delegation.failed_delegations_for_session`` row: one failed task of an async delegation.",
|
||||
"properties": {
|
||||
"delegation_id": {
|
||||
"title": "Delegation Id",
|
||||
"type": "string"
|
||||
},
|
||||
"task_index": {
|
||||
"default": 0,
|
||||
"title": "Task Index",
|
||||
"type": "integer"
|
||||
},
|
||||
"status": {
|
||||
"title": "Status",
|
||||
"type": "string"
|
||||
},
|
||||
"goal": {
|
||||
"default": "",
|
||||
"title": "Goal",
|
||||
"type": "string"
|
||||
},
|
||||
"error": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "string"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"default": null,
|
||||
"title": "Error"
|
||||
},
|
||||
"dispatched_at": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "number"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"default": null,
|
||||
"title": "Dispatched At"
|
||||
},
|
||||
"completed_at": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "number"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"default": null,
|
||||
"title": "Completed At"
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
"delegation_id",
|
||||
"status"
|
||||
],
|
||||
"title": "FailedDelegation",
|
||||
"type": "object"
|
||||
},
|
||||
"FileAttachParams": {
|
||||
"additionalProperties": false,
|
||||
"description": "``path`` when the file is gateway-visible, else ``data_url`` carries the bytes; ``name`` labels\nan uploaded file.",
|
||||
@@ -32863,7 +32929,7 @@
|
||||
},
|
||||
"SubagentListResult": {
|
||||
"additionalProperties": false,
|
||||
"description": "``delegations`` is reserved for async delegation records and is currently always empty.",
|
||||
"description": "``delegations``: recently failed async delegation tasks for the session (durable store), newest first.",
|
||||
"properties": {
|
||||
"subagents": {
|
||||
"items": {
|
||||
@@ -32874,8 +32940,7 @@
|
||||
},
|
||||
"delegations": {
|
||||
"items": {
|
||||
"additionalProperties": true,
|
||||
"type": "object"
|
||||
"$ref": "#/components/schemas/FailedDelegation"
|
||||
},
|
||||
"title": "Delegations",
|
||||
"type": "array"
|
||||
|
||||
63
tests/tools/test_async_delegation_failed_surface.py
Normal file
63
tests/tools/test_async_delegation_failed_surface.py
Normal file
@@ -0,0 +1,63 @@
|
||||
"""Failed async delegations stay visible after the live roster forgets them (#97202).
|
||||
|
||||
A renderer reload drops the in-memory subagent roster, and an ended child leaves it anyway, so a
|
||||
failed delegation had nowhere to show up. ``failed_delegations_for_session`` reads the durable row.
|
||||
These tests write real rows into a temp ``state.db`` through the module's own persistence calls.
|
||||
"""
|
||||
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from hermes_constants import reset_hermes_home_override, set_hermes_home_override
|
||||
from tools import async_delegation as ad
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def home(tmp_path):
|
||||
ad._reset_for_tests()
|
||||
token = set_hermes_home_override(str(tmp_path))
|
||||
yield tmp_path
|
||||
reset_hermes_home_override(token)
|
||||
ad._reset_for_tests()
|
||||
|
||||
|
||||
def _delegation(delegation_id, *, status, result, ui="ui-1", parent="agent-1", completed_at=None, **task):
|
||||
ad._persist_dispatch({"delegation_id": delegation_id, "session_key": "", "origin_ui_session_id": ui,
|
||||
"parent_session_id": parent, "dispatched_at": time.time() - 60, **task})
|
||||
ad._persist_completion({"delegation_id": delegation_id, "status": status,
|
||||
"completed_at": completed_at or time.time()}, result)
|
||||
|
||||
|
||||
def test_a_failed_single_delegation_is_listed_for_its_session():
|
||||
_delegation("d-fail", status="error", goal="audit billing",
|
||||
result={"status": "error", "error": "interrupted: waiting for model response"})
|
||||
_delegation("d-ok", status="completed", goal="write docs", result={"status": "completed", "summary": "done"})
|
||||
|
||||
assert ad.failed_delegations_for_session("ui-1") == [{
|
||||
"delegation_id": "d-fail", "task_index": 0, "status": "error", "goal": "audit billing",
|
||||
"error": "interrupted: waiting for model response", "dispatched_at": pytest.approx(time.time() - 60, abs=30),
|
||||
"completed_at": pytest.approx(time.time(), abs=30)}]
|
||||
|
||||
|
||||
def test_a_batch_that_completed_still_surfaces_its_failed_task():
|
||||
_delegation("d-batch", status="completed", is_batch=True, goals=["scan api", "scan web"], task_indexes=[3, 4],
|
||||
result={"results": [{"task_index": 3, "status": "completed", "summary": "ok"},
|
||||
{"task_index": 4, "status": "timeout", "error": "no progress"}]})
|
||||
|
||||
[row] = ad.failed_delegations_for_session("ui-1")
|
||||
assert (row["task_index"], row["goal"], row["status"], row["error"]) == (4, "scan web", "timeout", "no progress")
|
||||
|
||||
|
||||
def test_the_durable_parent_id_claims_rows_after_a_reload_remints_the_ui_session():
|
||||
_delegation("d-fail", status="stalled", goal="refactor", result={"status": "error", "error": "stalled"})
|
||||
|
||||
assert [r["delegation_id"] for r in ad.failed_delegations_for_session("ui-reminted", "agent-1")] == ["d-fail"]
|
||||
|
||||
|
||||
def test_other_sessions_and_old_failures_stay_out():
|
||||
_delegation("d-other", status="error", goal="x", ui="ui-2", parent="agent-2", result={"status": "error"})
|
||||
_delegation("d-old", status="error", goal="y", completed_at=time.time() - 3 * 86400, result={"status": "error"})
|
||||
|
||||
assert ad.failed_delegations_for_session("ui-1", "agent-1") == []
|
||||
assert ad.failed_delegations_for_session() == []
|
||||
@@ -303,3 +303,23 @@ def test_list_follows_the_conversation_across_ui_sid_and_compression_rotation(ru
|
||||
finally:
|
||||
_unregister_subagent("child")
|
||||
db.close()
|
||||
|
||||
|
||||
def test_list_surfaces_failed_delegations_that_outlived_the_live_roster(runtime):
|
||||
"""A failed child is gone from the live roster (ended, or a renderer reload dropped it); the
|
||||
durable row still reaches ``delegations`` for its own session only (#97202)."""
|
||||
import time
|
||||
|
||||
from tools import async_delegation as bg
|
||||
|
||||
_server, _owner, _transport, call = runtime
|
||||
for did, ui in (("d-mine", "ui-owner"), ("d-foreign", "other")):
|
||||
bg._persist_dispatch({"delegation_id": did, "session_key": "", "origin_ui_session_id": ui,
|
||||
"parent_session_id": None, "dispatched_at": time.time(), "goal": f"{did} goal"})
|
||||
bg._persist_completion({"delegation_id": did, "status": "error", "completed_at": time.time()},
|
||||
{"status": "error", "error": "interrupted: waiting for model response"})
|
||||
|
||||
snapshot = call("subagent.list")["result"]
|
||||
assert snapshot["subagents"] == []
|
||||
assert [(d["delegation_id"], d["goal"], d["status"]) for d in snapshot["delegations"]] == [
|
||||
("d-mine", "d-mine goal", "error")]
|
||||
|
||||
@@ -559,6 +559,61 @@ def get_durable_delegation(delegation_id: str) -> Optional[Dict[str, Any]]:
|
||||
"delivery_attempts": row[6], "origin_session_id": row[7] or ""}
|
||||
|
||||
|
||||
_FAILED_TASK_STATES = frozenset({"error", "failed", "failure", "timeout", "stalled", "unknown", "interrupted"})
|
||||
_FAILURE_SURFACE_WINDOW_S = 24 * 3600.0
|
||||
|
||||
|
||||
def _json_object(raw: Optional[str]) -> Dict[str, Any]:
|
||||
try:
|
||||
value = json.loads(raw) if raw else {}
|
||||
except ValueError:
|
||||
return {}
|
||||
return value if isinstance(value, dict) else {}
|
||||
|
||||
|
||||
def failed_delegations_for_session(
|
||||
origin_ui_session_id: str = "", parent_session_id: str = "", *, limit: int = 20, now: Optional[float] = None,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Recently failed async delegation tasks owned by a session, newest first.
|
||||
|
||||
The live roster forgets a child once it ends and does not survive a renderer reload, so a failed
|
||||
delegation had nowhere to show (#97202). This reads the durable row instead: one entry per failed
|
||||
task (a batch unit that "completed" can still carry failed tasks) with ``delegation_id``,
|
||||
``task_index``, ``goal``, ``status``, ``error``, ``dispatched_at`` and ``completed_at``. Either selector claims a row:
|
||||
the UI session id at dispatch, or the spawner's durable session id (survives a reload re-mint)."""
|
||||
selectors = [(col, val) for col, val in (
|
||||
("origin_ui_session_id", origin_ui_session_id), ("parent_session_id", parent_session_id)) if val]
|
||||
if not selectors:
|
||||
return []
|
||||
cutoff = (now if now is not None else time.time()) - _FAILURE_SURFACE_WINDOW_S
|
||||
owner_sql = " OR ".join(f"{col}=?" for col, _ in selectors)
|
||||
with _DB_LOCK, _transaction() as conn:
|
||||
rows = conn.execute(
|
||||
f"""SELECT delegation_id, state, dispatched_at, completed_at, task_json, result_json FROM async_delegations
|
||||
WHERE ({owner_sql}) AND state NOT IN ('running','finalizing') AND completed_at >= ?
|
||||
ORDER BY completed_at DESC LIMIT ?""",
|
||||
(*(val for _, val in selectors), cutoff, limit)).fetchall()
|
||||
failed: List[Dict[str, Any]] = []
|
||||
for delegation_id, state, dispatched_at, completed_at, task_json, result_json in rows:
|
||||
task, result = _json_object(task_json), _json_object(result_json)
|
||||
goals = task.get("goals") if isinstance(task.get("goals"), list) and task["goals"] else [task.get("goal") or ""]
|
||||
goal_for = dict(zip(task.get("task_indexes") or range(len(goals)), goals))
|
||||
tasks = result["results"] if isinstance(result.get("results"), list) else [] if task.get("is_batch") else [result]
|
||||
if not tasks and str(state).lower() in _FAILED_TASK_STATES:
|
||||
tasks = [{"task_index": 0, "error": result.get("error")}]
|
||||
for entry in tasks:
|
||||
status = str(entry.get("status") or state or "").lower()
|
||||
if status not in _FAILED_TASK_STATES:
|
||||
continue
|
||||
index = entry.get("task_index") if isinstance(entry.get("task_index"), int) else 0
|
||||
error = entry.get("error") or result.get("error")
|
||||
failed.append({
|
||||
"delegation_id": delegation_id, "task_index": index, "status": status,
|
||||
"goal": str(goal_for.get(index, goals[0]) or ""), "error": str(error) if error else None,
|
||||
"dispatched_at": dispatched_at, "completed_at": completed_at})
|
||||
return failed[:limit]
|
||||
|
||||
|
||||
# ── In-memory registry queries ──────────────────────────────────────────────
|
||||
def _get_executor(max_workers: int) -> ThreadPoolExecutor:
|
||||
"""Lazily create (or grow in place, never shrink) the shared daemon executor. Raising
|
||||
|
||||
@@ -614,11 +614,23 @@ class SubagentSnapshot(Result):
|
||||
accepting_steer: bool | None = None
|
||||
|
||||
|
||||
class FailedDelegation(Result):
|
||||
"""``async_delegation.failed_delegations_for_session`` row: one failed task of an async delegation."""
|
||||
|
||||
delegation_id: str
|
||||
task_index: int = 0
|
||||
status: str
|
||||
goal: str = ""
|
||||
error: str | None = None
|
||||
dispatched_at: float | None = None
|
||||
completed_at: float | None = None
|
||||
|
||||
|
||||
class SubagentListResult(Result):
|
||||
"""``delegations`` is reserved for async delegation records and is currently always empty."""
|
||||
"""``delegations``: recently failed async delegation tasks for the session (durable store), newest first."""
|
||||
|
||||
subagents: list[SubagentSnapshot] = Field(default_factory=list)
|
||||
delegations: list[dict[str, JsonValue]] = Field(default_factory=list)
|
||||
delegations: list[FailedDelegation] = Field(default_factory=list)
|
||||
|
||||
|
||||
method("subagent.list", params=SessionParams, result=SubagentListResult,
|
||||
|
||||
@@ -59,10 +59,25 @@ def _(rid, params):
|
||||
live = _visible_subagent_records(session_id, transport, owner)
|
||||
return _ok(rid, {
|
||||
"subagents": [{key: r.get(key) for key in _SUBAGENT_SNAPSHOT_FIELDS} for r in live],
|
||||
"delegations": [],
|
||||
"delegations": _failed_delegations(session_id, owner),
|
||||
})
|
||||
|
||||
|
||||
def _failed_delegations(session_id, owner):
|
||||
"""Recently failed async delegation tasks for this session from the durable store (the live
|
||||
roster forgets ended children and dies with a renderer reload, #97202). Read under the session's
|
||||
profile home, where its delegations were persisted; a store error degrades to no rows."""
|
||||
from tools.async_delegation import failed_delegations_for_session
|
||||
|
||||
agent_session_id = str(getattr(owner.get("agent"), "session_id", "") or "")
|
||||
try:
|
||||
with _session_home_scope(owner):
|
||||
return failed_delegations_for_session(session_id, agent_session_id)
|
||||
except Exception:
|
||||
logger.debug("subagent.list: failed-delegation read failed for %s", session_id, exc_info=True)
|
||||
return []
|
||||
|
||||
|
||||
@method("subagent.interrupt")
|
||||
def _(rid, params):
|
||||
from agent.interrupt_compat import request_hard_interrupt
|
||||
|
||||
Reference in New Issue
Block a user