refactor(kanban-cli): specify/decompose share triage-load + aux-call helpers; decompose_task split into routing/single/fanout phases

This commit is contained in:
Teknium
2026-09-02 20:25:28 -07:00
parent 9d896da1cc
commit 35876c9100
2 changed files with 217 additions and 266 deletions

View File

@@ -1,23 +1,18 @@
"""Kanban decomposer — fan a triage task out into a graph of child tasks.
Invoked by ``hermes kanban decompose [task_id | --all]`` and the
auto-decompose path in the gateway dispatcher loop. Reads the user's
profile roster (with descriptions) and asks the auxiliary LLM to
return a task graph in JSON. Then atomically creates the children,
links them under the root, and flips the root ``triage -> todo``.
Invoked by ``hermes kanban decompose [task_id | --all]`` and the gateway
dispatcher's auto-decompose path. Reads the profile roster (with
descriptions), asks the auxiliary LLM for a task graph in JSON, then
atomically creates the children, links them under the root, and flips the
root ``triage -> todo``. The root stays alive as parent of every leaf child so
it wakes back up when the graph completes and its assignee (the orchestrator
profile) can judge completion and add more work.
The root task stays alive and becomes the parent of every leaf child,
so when the whole graph completes the root wakes back up — its
assignee (the orchestrator profile) gets a chance to judge completion
and add more tasks if the work isn't done yet.
Design notes: mirrors ``kanban_specify`` (lazy aux import, lenient parse,
never raises on expected failures). The prompt sees the configured profile
roster; undescribed profiles are listed with a note so name-matching still
works. ``fanout=false`` collapses to the ``specify`` behaviour (tighten +
promote, no children), making ``decompose`` a strict superset. Unknown
assignees are rewritten to ``default_assignee`` — a child NEVER ends up with
``assignee=None``.
Mirrors ``kanban_specify`` (lazy aux import, lenient parse, never raises on
expected failures). ``fanout=false`` collapses to the ``specify`` behaviour
(tighten + promote, no children), making ``decompose`` a strict superset.
Unknown assignees are rewritten to ``default_assignee`` — a child NEVER ends
up with ``assignee=None``.
"""
from __future__ import annotations
@@ -29,7 +24,9 @@ from typing import Optional
from hermes_cli import kanban_db as kb
from hermes_cli import profiles as profiles_mod
from hermes_cli.kanban_specify import _extract_json_blob, _title_body, _truncate
from hermes_cli.kanban_specify import (
_call_aux, _extract_json_blob, _load_triage_task, _task_prompt_fields, _title_body,
)
from hermes_cli.kanban_specify import _profile_author as _specify_author
logger = logging.getLogger(__name__)
@@ -138,10 +135,8 @@ def _load_config() -> dict:
def _resolve_profile_from_cfg(cfg: dict, key: str) -> str:
"""``kanban.<key>`` if it names an existing profile, else the active
default profile — so a task is never stranded for lack of an owner.
``orchestrator_profile`` owns the root after fan-out; ``default_assignee``
catches children the decomposer can't route.
"""
catches children the decomposer can't route."""
kanban_cfg = cfg.get("kanban", {}) if isinstance(cfg, dict) else {}
explicit = (kanban_cfg.get(key) or "").strip()
if explicit:
@@ -159,13 +154,12 @@ def _resolve_profile_from_cfg(cfg: dict, key: str) -> str:
def _build_roster() -> tuple[list[dict], set[str]]:
"""``(roster_for_prompt, valid_assignee_names)``; entries are
``{name, description, has_description}``."""
roster: list[dict] = []
valid: set[str] = set()
try:
all_profiles = profiles_mod.list_profiles()
except Exception as exc:
logger.warning("decompose: failed to list profiles: %s", exc)
return roster, valid
return [], set()
roster = []
for p in all_profiles:
desc = (p.description or "").strip()
roster.append({
@@ -173,34 +167,131 @@ def _build_roster() -> tuple[list[dict], set[str]]:
"description": desc or f"(no description; profile named {p.name!r})",
"has_description": bool(desc),
})
valid.add(p.name)
return roster, valid
return roster, {p.name for p in all_profiles}
def _format_roster(roster: list[dict]) -> str:
if not roster:
return " (no profiles installed — decomposer cannot route work)"
lines = []
for entry in roster:
tag = "" if entry["has_description"] else " ⚠ undescribed"
lines.append(f" - {entry['name']}{tag}: {entry['description']}")
return "\n".join(lines)
return "\n".join(
f" - {entry['name']}{'' if entry['has_description'] else ' ⚠ undescribed'}: {entry['description']}"
for entry in roster
)
def _normalize_assignee_choice(
assignee: object,
*,
default_assignee: str,
valid_names: set[str],
) -> str:
def _normalize_assignee_choice(assignee: object, *, default_assignee: str, valid_names: set[str]) -> str:
"""A valid assignee, else ``default_assignee`` — promoted work is never
left unassigned."""
if not isinstance(assignee, str) or not assignee.strip():
return default_assignee
chosen = assignee.strip()
if chosen not in valid_names:
return default_assignee
return chosen
return chosen if chosen in valid_names else default_assignee
@dataclass
class _Routing:
"""Config-derived routing context for one decomposition."""
orchestrator: str
default_assignee: str
auto_promote: bool
roster: list[dict]
valid_names: set[str]
def _load_routing() -> _Routing:
cfg = _load_config()
kanban_cfg = cfg.get("kanban", {}) if isinstance(cfg, dict) else {}
roster, valid_names = _build_roster()
return _Routing(
orchestrator=_resolve_profile_from_cfg(cfg, "orchestrator_profile"),
default_assignee=_resolve_profile_from_cfg(cfg, "default_assignee"),
auto_promote=bool(kanban_cfg.get("auto_promote_children", True)),
roster=roster,
valid_names=valid_names,
)
def _apply_single(task: kb.Task, parsed: dict, routing: _Routing, author: str) -> DecomposeOutcome:
"""``fanout=false``: single-task spec promotion (same effect as specify)."""
title_val, body_val = _title_body(parsed)
assignee_val = None
if not task.assignee:
assignee_val = _normalize_assignee_choice(
parsed.get("assignee"), default_assignee=routing.default_assignee, valid_names=routing.valid_names,
)
if title_val is None and body_val is None:
return DecomposeOutcome(task.id, False, "decomposer returned fanout=false with no title/body")
with kb.connect_closing() as conn:
ok = kb.specify_triage_task(
conn, task.id, title=title_val, body=body_val, assignee=assignee_val, author=author,
)
if not ok:
return DecomposeOutcome(task.id, False, "task moved out of triage before promotion")
return DecomposeOutcome(task.id, True, "single task (no fanout)", fanout=False, new_title=title_val)
def _clean_children(task_id: str, raw_tasks: list, routing: _Routing) -> tuple[list[dict], str]:
"""Validate/normalise the LLM's ``tasks`` list; ``(children, "")`` or ``([], reason)``.
Unknown assignees route to the default; never assignee=None."""
children: list[dict] = []
for idx, entry in enumerate(raw_tasks):
if not isinstance(entry, dict):
return [], f"tasks[{idx}] is not an object"
title = entry.get("title")
if not isinstance(title, str) or not title.strip():
return [], f"tasks[{idx}].title is missing or empty"
body = entry.get("body")
assignee = entry.get("assignee")
chosen = _normalize_assignee_choice(
assignee, default_assignee=routing.default_assignee, valid_names=routing.valid_names,
)
if isinstance(assignee, str) and assignee.strip() and assignee.strip() not in routing.valid_names:
logger.info(
"decompose: task %s child %d picked unknown assignee %r — "
"routing to default_assignee %r",
task_id, idx, assignee, routing.default_assignee,
)
parents = entry.get("parents") or []
if not isinstance(parents, list):
parents = []
children.append({
"title": title.strip()[:200],
"body": body.strip() if isinstance(body, str) else "",
"assignee": chosen,
# Drop non-int, out-of-range and self parent indices.
"parents": [p for p in parents if isinstance(p, int) and 0 <= p < len(raw_tasks) and p != idx],
})
return children, ""
def _apply_fanout(task_id: str, parsed: dict, routing: _Routing, author: str) -> DecomposeOutcome:
raw_tasks = parsed.get("tasks") or []
if not isinstance(raw_tasks, list) or not raw_tasks:
return DecomposeOutcome(task_id, False, "decomposer returned fanout=true with empty tasks list")
children, reason = _clean_children(task_id, raw_tasks, routing)
if reason:
return DecomposeOutcome(task_id, False, reason)
try:
with kb.connect_closing() as conn:
child_ids = kb.decompose_triage_task(
conn,
task_id,
root_assignee=routing.orchestrator,
children=children,
author=author,
auto_promote=routing.auto_promote,
)
except ValueError as exc:
return DecomposeOutcome(task_id, False, f"DB rejected graph: {exc}")
except Exception as exc:
logger.exception("decompose: DB error on task %s", task_id)
return DecomposeOutcome(task_id, False, f"DB error: {type(exc).__name__}")
if child_ids is None:
return DecomposeOutcome(task_id, False, "task moved out of triage before decomposition")
return DecomposeOutcome(
task_id, True, f"decomposed into {len(child_ids)} children", fanout=True, child_ids=child_ids,
)
def decompose_task(
@@ -212,182 +303,35 @@ def decompose_task(
"""Decompose a triage task into a graph of child tasks. Expected failures
(not in triage, no aux client, API error, malformed/empty reply) surface
as ``ok=False``."""
with kb.connect_closing() as conn:
task = kb.get_task(conn, task_id)
task, reason = _load_triage_task(task_id)
if task is None:
return DecomposeOutcome(task_id, False, "unknown task id")
if task.status != "triage":
return DecomposeOutcome(
task_id, False, f"task is not in triage (status={task.status!r})"
)
return DecomposeOutcome(task_id, False, reason)
cfg = _load_config()
orchestrator = _resolve_profile_from_cfg(cfg, "orchestrator_profile")
default_assignee = _resolve_profile_from_cfg(cfg, "default_assignee")
kanban_cfg = cfg.get("kanban", {}) if isinstance(cfg, dict) else {}
auto_promote = bool(kanban_cfg.get("auto_promote_children", True))
roster, valid_names = _build_roster()
try:
from agent.auxiliary_client import call_llm # type: ignore
except Exception as exc:
logger.debug("decompose: auxiliary client import failed: %s", exc)
return DecomposeOutcome(task_id, False, "auxiliary client unavailable")
user_msg = _USER_TEMPLATE.format(
task_id=task.id,
title=_truncate(task.title or "", 400),
body=_truncate(task.body or "(no body)", 4000),
roster=_format_roster(roster),
default_assignee=default_assignee,
routing = _load_routing()
raw, reason = _call_aux(
"decompose", task_id, aux_task="kanban_decomposer", system=_SYSTEM_PROMPT,
user=_USER_TEMPLATE.format(
**_task_prompt_fields(task),
roster=_format_roster(routing.roster),
default_assignee=routing.default_assignee,
),
max_tokens=4000, timeout=timeout or 180, log=logger,
)
try:
# call_llm applies all auxiliary.kanban_decomposer.* config
# (provider/model/base_url, extra_body, reasoning_effort, retries).
resp = call_llm(
task="kanban_decomposer",
messages=[
{"role": "system", "content": _SYSTEM_PROMPT},
{"role": "user", "content": user_msg},
],
temperature=0.3,
max_tokens=4000,
timeout=timeout or 180,
)
except Exception as exc:
logger.info(
"decompose: API call failed for %s (%s)", task_id, exc,
)
return DecomposeOutcome(task_id, False, f"LLM error: {type(exc).__name__}")
try:
raw = resp.choices[0].message.content or ""
except Exception:
raw = ""
if raw is None:
return DecomposeOutcome(task_id, False, reason)
parsed = _extract_json_blob(raw, _FENCE_RE)
if parsed is None:
return DecomposeOutcome(task_id, False, "LLM returned malformed JSON")
fanout = bool(parsed.get("fanout"))
audit_author = author or _profile_author()
if not fanout:
# Fall back to single-task spec promotion (same effect as specify).
title_val, body_val = _title_body(parsed)
assignee_val = None
if not task.assignee:
assignee_val = _normalize_assignee_choice(
parsed.get("assignee"),
default_assignee=default_assignee,
valid_names=valid_names,
)
if title_val is None and body_val is None:
return DecomposeOutcome(
task_id, False, "decomposer returned fanout=false with no title/body",
)
with kb.connect_closing() as conn:
ok = kb.specify_triage_task(
conn,
task_id,
title=title_val,
body=body_val,
assignee=assignee_val,
author=audit_author,
)
if not ok:
return DecomposeOutcome(
task_id, False, "task moved out of triage before promotion",
)
return DecomposeOutcome(
task_id, True, "single task (no fanout)",
fanout=False, new_title=title_val,
)
raw_tasks = parsed.get("tasks") or []
if not isinstance(raw_tasks, list) or not raw_tasks:
return DecomposeOutcome(
task_id, False, "decomposer returned fanout=true with empty tasks list",
)
# Unknown assignees route to the default; never assignee=None.
children: list[dict] = []
for idx, entry in enumerate(raw_tasks):
if not isinstance(entry, dict):
return DecomposeOutcome(
task_id, False, f"tasks[{idx}] is not an object",
)
title = entry.get("title")
if not isinstance(title, str) or not title.strip():
return DecomposeOutcome(
task_id, False, f"tasks[{idx}].title is missing or empty",
)
body = entry.get("body")
if not isinstance(body, str):
body = ""
assignee = entry.get("assignee")
chosen = _normalize_assignee_choice(
assignee,
default_assignee=default_assignee,
valid_names=valid_names,
)
if (
isinstance(assignee, str)
and assignee.strip()
and assignee.strip() not in valid_names
):
logger.info(
"decompose: task %s child %d picked unknown assignee %r — "
"routing to default_assignee %r",
task_id, idx, assignee, default_assignee,
)
parents = entry.get("parents") or []
if not isinstance(parents, list):
parents = []
# Clean parent indices: drop non-int and out-of-range.
clean_parents = [p for p in parents if isinstance(p, int) and 0 <= p < len(raw_tasks) and p != idx]
children.append({
"title": title.strip()[:200],
"body": body.strip(),
"assignee": chosen,
"parents": clean_parents,
})
try:
with kb.connect_closing() as conn:
child_ids = kb.decompose_triage_task(
conn,
task_id,
root_assignee=orchestrator,
children=children,
author=audit_author,
auto_promote=auto_promote,
)
except ValueError as exc:
return DecomposeOutcome(task_id, False, f"DB rejected graph: {exc}")
except Exception as exc:
logger.exception("decompose: DB error on task %s", task_id)
return DecomposeOutcome(task_id, False, f"DB error: {type(exc).__name__}")
if child_ids is None:
return DecomposeOutcome(
task_id, False, "task moved out of triage before decomposition",
)
return DecomposeOutcome(
task_id, True, f"decomposed into {len(child_ids)} children",
fanout=True, child_ids=child_ids,
)
if not parsed.get("fanout"):
return _apply_single(task, parsed, routing, audit_author)
return _apply_fanout(task_id, parsed, routing, audit_author)
def list_triage_ids(*, tenant: Optional[str] = None) -> list[str]:
"""Return task ids currently in the triage column."""
with kb.connect_closing() as conn:
rows = kb.list_tasks(
conn,
status="triage",
tenant=tenant,
limit=1000,
)
rows = kb.list_tasks(conn, status="triage", tenant=tenant, limit=1000)
return [row.id for row in rows]

View File

@@ -99,9 +99,8 @@ def _extract_json_blob(raw: str, fence_re: re.Pattern = _FENCE_RE) -> Optional[d
last = stripped.rfind("}")
if first == -1 or last == -1 or last <= first:
return None
candidate = stripped[first : last + 1]
try:
val = json.loads(candidate)
val = json.loads(stripped[first : last + 1])
except (ValueError, json.JSONDecodeError):
return None
return val if isinstance(val, dict) else None
@@ -124,6 +123,60 @@ def _profile_author(default: str = "specifier") -> str:
return os.environ.get("HERMES_PROFILE") or os.environ.get("USER") or default
def _load_triage_task(task_id: str) -> tuple[Optional[kb.Task], str]:
"""``(task, "")`` when the task exists and is in triage, else ``(None, reason)``."""
with kb.connect_closing() as conn:
task = kb.get_task(conn, task_id)
if task is None:
return None, "unknown task id"
if task.status != "triage":
return None, f"task is not in triage (status={task.status!r})"
return task, ""
def _task_prompt_fields(task: kb.Task) -> dict[str, str]:
"""Bounded ``task_id``/``title``/``body`` for the user prompt templates."""
return {
"task_id": task.id,
"title": _truncate(task.title or "", 400),
"body": _truncate(task.body or "(no body)", 4000),
}
def _call_aux(verb: str, task_id: str, *, aux_task: str, system: str, user: str,
max_tokens: int, timeout: int, log: logging.Logger = logger) -> tuple[Optional[str], str]:
"""One auxiliary LLM call; ``(reply_text, "")`` or ``(None, reason)``.
``call_llm`` applies all ``auxiliary.<aux_task>.*`` config (provider/model/
base_url, extra_body, reasoning_effort, retries). Imported lazily so a
missing aux client degrades to a skip instead of an import-time crash.
"""
try:
from agent.auxiliary_client import call_llm
except Exception as exc: # pragma: no cover — import smoke test
log.debug("%s: auxiliary client import failed: %s", verb, exc)
return None, "auxiliary client unavailable"
try:
resp = call_llm(
task=aux_task,
messages=[
{"role": "system", "content": system},
{"role": "user", "content": user},
],
temperature=0.3,
max_tokens=max_tokens,
timeout=timeout,
)
except Exception as exc:
suffix = " — skipping" if verb == "specify" else ""
log.info("%s: API call failed for %s (%s)%s", verb, task_id, exc, suffix)
return None, f"LLM error: {type(exc).__name__}"
try:
return resp.choices[0].message.content or "", ""
except Exception:
return "", ""
def specify_task(
task_id: str,
*,
@@ -133,73 +186,29 @@ def specify_task(
"""Specify one triage task and promote it to ``todo``. Expected failures
(not in triage, no aux client, API error, malformed reply) surface as
``ok=False`` so an ``--all`` sweep continues."""
with kb.connect_closing() as conn:
task = kb.get_task(conn, task_id)
task, reason = _load_triage_task(task_id)
if task is None:
return SpecifyOutcome(task_id, False, "unknown task id")
if task.status != "triage":
return SpecifyOutcome(
task_id, False, f"task is not in triage (status={task.status!r})"
)
return SpecifyOutcome(task_id, False, reason)
try:
from agent.auxiliary_client import call_llm
except Exception as exc: # pragma: no cover — import smoke test
logger.debug("specify: auxiliary client import failed: %s", exc)
return SpecifyOutcome(task_id, False, "auxiliary client unavailable")
user_msg = _USER_TEMPLATE.format(
task_id=task.id,
title=_truncate(task.title or "", 400),
body=_truncate(task.body or "(no body)", 4000),
raw, reason = _call_aux(
"specify", task_id, aux_task="triage_specifier", system=_SYSTEM_PROMPT,
user=_USER_TEMPLATE.format(**_task_prompt_fields(task)),
max_tokens=HERMES_KANBAN_SPECIFY_MAX_TOKENS, timeout=timeout or 120,
)
try:
# call_llm applies all auxiliary.triage_specifier.* config
# (provider/model/base_url, extra_body, reasoning_effort, retries).
resp = call_llm(
task="triage_specifier",
messages=[
{"role": "system", "content": _SYSTEM_PROMPT},
{"role": "user", "content": user_msg},
],
temperature=0.3,
max_tokens=HERMES_KANBAN_SPECIFY_MAX_TOKENS,
timeout=timeout or 120,
)
except Exception as exc:
logger.info(
"specify: API call failed for %s (%s) — skipping",
task_id, exc,
)
return SpecifyOutcome(
task_id, False, f"LLM error: {type(exc).__name__}"
)
try:
raw = (resp.choices[0].message.content or "").strip()
except Exception:
raw = ""
if raw is None:
return SpecifyOutcome(task_id, False, reason)
raw = raw.strip()
parsed = _extract_json_blob(raw)
new_title: Optional[str]
new_body: Optional[str]
if parsed is None:
# Whole reply becomes the body; the user can edit afterward.
stripped_raw = raw.strip()
if not stripped_raw:
return SpecifyOutcome(
task_id, False, "LLM returned an empty response"
)
new_title = None
new_body = stripped_raw
if not raw:
return SpecifyOutcome(task_id, False, "LLM returned an empty response")
new_title, new_body = None, raw
else:
new_title, new_body = _title_body(parsed)
if new_body is None and new_title is None:
return SpecifyOutcome(
task_id, False, "LLM response missing title and body"
)
return SpecifyOutcome(task_id, False, "LLM response missing title and body")
with kb.connect_closing() as conn:
ok = kb.specify_triage_task(
@@ -211,9 +220,7 @@ def specify_task(
)
if not ok:
# Race: promoted/archived between our read and the write.
return SpecifyOutcome(
task_id, False, "task moved out of triage before promotion"
)
return SpecifyOutcome(task_id, False, "task moved out of triage before promotion")
return SpecifyOutcome(task_id, True, "specified", new_title=new_title)