refactor(plugins/kanban): table-driven orchestration resolution, tighter estimate parsing / profile validation, drop no-op re-raise

This commit is contained in:
Teknium
2026-09-02 22:19:30 -07:00
parent eccfcb5f8e
commit b60caea6b2
+49 -72
View File
@@ -180,16 +180,15 @@ def _compute_task_diagnostics(conn: sqlite3.Connection, task_ids: Optional[list[
dashboard's typical working set, paginate if profiling shows a hotspot."""
from hermes_cli.config import load_config
if task_ids is not None and not task_ids:
return {}
diag_config = kd.config_from_runtime_config(load_config())
if task_ids is not None:
if not task_ids:
return {}
rows = conn.execute(f"SELECT * FROM tasks WHERE id IN ({_placeholders(task_ids)})", tuple(task_ids)).fetchall()
else:
rows = conn.execute("SELECT * FROM tasks WHERE status != 'archived'").fetchall()
if not rows:
return {}
row_ids = [r["id"] for r in rows]
def _rows_by_task(table: str) -> dict[str, list]:
@@ -206,8 +205,7 @@ def _compute_task_diagnostics(conn: sqlite3.Connection, task_ids: Optional[list[
for r in rows:
tid = r["id"]
diags = kd.compute_task_diagnostics(
r, events_by_task.get(tid, []), runs_by_task.get(tid, []), config=diag_config, graph=graph_by_task.get(tid),
)
r, events_by_task[tid], runs_by_task[tid], config=diag_config, graph=graph_by_task.get(tid))
if diags:
out[tid] = [d.to_dict() for d in diags]
return out
@@ -318,7 +316,7 @@ def get_task(
with _board_conn(board) as (board, conn):
if (run_state_type is None) ^ (run_state_name is None):
raise HTTPException(status_code=400, detail="run_state_type and run_state_name must be passed together or omitted")
if run_state_type is not None and run_state_type not in ("status", "outcome"):
if run_state_type not in (None, "status", "outcome"):
raise HTTPException(status_code=400, detail="run_state_type must be 'status' or 'outcome'")
task = _require_task(conn, task_id)
# Drawer returns the FULL summary (cards on /board carry a 200-char preview).
@@ -360,10 +358,8 @@ class CreateTaskBody(BaseModel):
goal_max_turns: Optional[int] = None
model_override: Optional[str] = None
provider_override: Optional[str] = None
# Thinking depth (none|minimal|…|ultra); None inherits the profile's own level.
reasoning_effort: Optional[str] = None
# When omitted, create_task inherits the board's scoped project (if any).
project_id: Optional[str] = None
reasoning_effort: Optional[str] = None # none|minimal|…|ultra; None inherits the profile's level
project_id: Optional[str] = None # None inherits the board's scoped project (if any)
@router.post("/tasks")
@@ -374,15 +370,14 @@ def create_task(payload: CreateTaskBody, board: Optional[str] = Query(None)):
task = kanban_db.get_task(conn, task_id)
body: dict[str, Any] = {"task": _task_dict(task) if task else None}
# Dispatcher-presence warning so the UI can banner a ready+assigned task that would
# otherwise sit idle (no gateway / dispatch_in_gateway=false). triage/todo are
# expected to wait; unassigned tasks can't dispatch anyway.
# otherwise sit idle (no gateway / dispatch_in_gateway=false); triage/todo are expected
# to wait, unassigned tasks can't dispatch anyway. Probe the request's active home: the
# dashboard backend may run under a different HERMES_HOME than the board's profile.
if task and task.status == "ready" and task.assignee:
try:
from hermes_cli.kanban import _check_dispatcher_presence
from hermes_constants import get_hermes_home
# Probe the request's active home: the dashboard backend may run under a
# different HERMES_HOME than the board's profile.
running, message = _check_dispatcher_presence(hermes_home=get_hermes_home())
if not running and message:
body["warning"] = message
@@ -431,8 +426,6 @@ async def upload_task_attachment(
status_code=413,
detail=f"attachment exceeds {KANBAN_ATTACHMENT_MAX_BYTES // (1024 * 1024)} MB limit")
out.write(chunk)
except HTTPException:
raise
except OSError as exc:
raise HTTPException(status_code=500, detail=f"failed to store attachment: {exc}")
@@ -479,17 +472,15 @@ class UpdateTaskBody(BaseModel):
body: Optional[str] = None
result: Optional[str] = None
block_reason: Optional[str] = None
# Structured handoff fields forwarded to complete_task on -> 'done'
# (parity with ``hermes kanban complete --summary/--metadata``).
# Handoff fields forwarded to complete_task on -> 'done' (parity with ``hermes kanban complete``).
summary: Optional[str] = None
metadata: Optional[dict] = None
# Model/provider override. In a PATCH ``None`` means "field not sent", so
# ``clear_model_override=True`` is the explicit clear signal.
# In a PATCH ``None`` means "field not sent", so ``clear_*=True`` is the explicit clear signal.
# ``reasoning_effort="none"`` is a VALUE (thinking off); it is cleared separately so
# dropping a model override doesn't silently reset the depth.
model_override: Optional[str] = None
provider_override: Optional[str] = None
clear_model_override: bool = False
# Thinking depth: ``"none"`` is a VALUE (thinking off); clear separately so
# dropping a model override doesn't silently reset the depth.
reasoning_effort: Optional[str] = None
clear_reasoning_effort: bool = False
@@ -861,11 +852,8 @@ def list_diagnostics(
"task_id": tid, "task_title": r["title"] if r else None, "task_status": r["status"] if r else None,
"task_assignee": r["assignee"] if r else None, "diagnostics": dl})
sev_idx = {s: i for i, s in enumerate(kd.SEVERITY_ORDER)}
def _sort_key(row):
top = row["diagnostics"][0]
return (-sev_idx.get(top.get("severity"), -1), -(top.get("last_seen_at") or 0))
out.sort(key=_sort_key)
out.sort(key=lambda row: (
-sev_idx.get(row["diagnostics"][0].get("severity"), -1), -(row["diagnostics"][0].get("last_seen_at") or 0)))
return {"diagnostics": out, "count": sum(len(d["diagnostics"]) for d in out)}
@@ -913,9 +901,9 @@ def inspect_run_endpoint(run_id: int, board: Optional[str] = _BOARD_Q):
if r.ended_at is not None:
return _dead("run already ended")
if r.worker_pid is None:
return _dead("no worker_pid recorded")
pid = r.worker_pid
if pid is None:
return _dead("no worker_pid recorded")
if _psutil is None:
return _dead("psutil not available", pid=pid)
try:
@@ -1082,11 +1070,8 @@ def _run_estimate(title: str, body: Optional[str]) -> dict:
# Same tolerant JSON-blob extraction the specifier uses.
try:
blob = raw
if not blob.lstrip().startswith("{"):
m = re.search(r"\{.*\}", blob, re.DOTALL)
blob = m.group(0) if m else blob
obj = json.loads(blob)
m = None if raw.lstrip().startswith("{") else re.search(r"\{.*\}", raw, re.DOTALL)
obj = json.loads(m.group(0) if m else raw)
parsed = obj if isinstance(obj, dict) else None
except Exception:
parsed = None
@@ -1097,10 +1082,9 @@ def _run_estimate(title: str, body: Optional[str]) -> dict:
except (TypeError, ValueError):
est_tokens = 0
complexity = str(parsed.get("complexity") or "").strip().upper()
if complexity not in {"S", "M", "L"}:
complexity = None
rationale = str(parsed.get("rationale") or "").strip() or None
return {"ok": True, "est_tokens": est_tokens, "complexity": complexity, "rationale": rationale, "model": model}
return {
"ok": True, "est_tokens": est_tokens, "complexity": complexity if complexity in {"S", "M", "L"} else None,
"rationale": str(parsed.get("rationale") or "").strip() or None, "model": model}
# --- Plugin config ----------------------------------------------------------
@@ -1175,12 +1159,9 @@ def get_home_channels(task_id: Optional[str] = Query(None), board: Optional[str]
with _board_conn(board) as (board, conn):
subs = kanban_db.list_notify_subs(conn, task_id)
subscribed_homes = {
(str(sub.get("platform") or ""), str(sub.get("chat_id") or ""), str(sub.get("thread_id") or "")) for sub in subs
}
return {
"home_channels": [
{**home, "subscribed": (home["platform"], home["chat_id"], home["thread_id"]) in subscribed_homes}
for home in homes]}
(str(sub.get("platform") or ""), str(sub.get("chat_id") or ""), str(sub.get("thread_id") or "")) for sub in subs}
return {"home_channels": [
{**home, "subscribed": (home["platform"], home["chat_id"], home["thread_id"]) in subscribed_homes} for home in homes]}
@router.post("/tasks/{task_id}/home-subscribe/{platform}")
@@ -1283,8 +1264,7 @@ class CreateBoardBody(BaseModel):
icon: Optional[str] = None
color: Optional[str] = None
default_workdir: Optional[str] = None
# Project (id or slug) scoping the board: default_workdir mirrors the
# project's primary repo and new tasks inherit the project.
# Project (id or slug) scoping the board: default_workdir mirrors its primary repo, tasks inherit it.
project_id: Optional[str] = None
switch: bool = False
@@ -1401,8 +1381,7 @@ def list_boards(include_archived: bool = Query(False)):
# so counting them in the switcher badge would visibly disagree.
b["total"] = sum(n for status, n in b["counts"].items() if status != "archived")
b["default_workspace_kind"] = _default_workspace_kind(b)
pid = b.get("project_id") or None
b["project_id"] = pid
pid = b["project_id"] = b.get("project_id") or None
proj = proj_map.get(pid) if pid else None
b["project_name"] = proj.name if proj else None
return {"boards": boards, "current": current}
@@ -1619,28 +1598,24 @@ def get_orchestration_settings():
values (fallbacks filled the same way the decomposer does)."""
cfg = _load_config_or_empty()
kanban_cfg = (cfg.get("kanban") or {}) if isinstance(cfg, dict) else {}
explicit_orch = (kanban_cfg.get("orchestrator_profile") or "").strip()
explicit_default = (kanban_cfg.get("default_assignee") or "").strip()
resolved_orch = explicit_orch
resolved_default = explicit_default
explicit = {k: (kanban_cfg.get(k) or "").strip() for k in _PROFILE_SETTINGS}
resolved = dict(explicit)
try:
from hermes_cli import profiles as profiles_mod
active_default = profiles_mod.get_active_profile_name() or "default"
if not resolved_orch or not profiles_mod.profile_exists(resolved_orch):
resolved_orch = active_default
if not resolved_default or not profiles_mod.profile_exists(resolved_default):
resolved_default = active_default
for k, v in explicit.items():
if not v or not profiles_mod.profile_exists(v):
resolved[k] = active_default
except Exception:
active_default = "default"
resolved_orch = resolved_orch or active_default
resolved_default = resolved_default or active_default
resolved = {k: v or active_default for k, v in resolved.items()}
return {
"orchestrator_profile": explicit_orch,
"default_assignee": explicit_default,
"orchestrator_profile": explicit["orchestrator_profile"],
"default_assignee": explicit["default_assignee"],
"auto_decompose": bool(kanban_cfg.get("auto_decompose", True)),
"auto_promote_children": bool(kanban_cfg.get("auto_promote_children", True)),
"resolved_orchestrator_profile": resolved_orch,
"resolved_default_assignee": resolved_default,
"resolved_orchestrator_profile": resolved["orchestrator_profile"],
"resolved_default_assignee": resolved["default_assignee"],
"active_profile": active_default}
@@ -1649,12 +1624,11 @@ def _validated_profile_name(raw: Optional[str], profiles_mod) -> str:
name = (raw or "").strip()
if name and profiles_mod is not None:
try:
if not profiles_mod.profile_exists(name):
raise HTTPException(status_code=400, detail=f"profile '{name}' does not exist")
except HTTPException:
raise
exists = profiles_mod.profile_exists(name)
except Exception:
pass
exists = True
if not exists:
raise HTTPException(status_code=400, detail=f"profile '{name}' does not exist")
return name
@@ -1694,6 +1668,13 @@ def _int_param(ws: WebSocket, name: str) -> int:
return 0
def _ws_board(raw: Optional[str]) -> Optional[str]:
try:
return kanban_db._normalize_board_slug(raw) if raw else None
except ValueError:
return None
@router.websocket("/events")
async def stream_events(ws: WebSocket):
if not _ws_upgrade_authorized(ws):
@@ -1717,11 +1698,7 @@ async def stream_events(ws: WebSocket):
cursor = _int_param(ws, "since")
# Board is pinned at the handshake; the UI opens a new WS on board change
# rather than reconciling two cursors mid-stream.
ws_board_raw = ws.query_params.get("board")
try:
ws_board = kanban_db._normalize_board_slug(ws_board_raw) if ws_board_raw else None
except ValueError:
ws_board = None
ws_board = _ws_board(ws.query_params.get("board"))
def _fetch_new(cursor_val: int) -> tuple[int, list[dict]]:
nonlocal event_conn