diff --git a/hermes_cli/kanban_decompose.py b/hermes_cli/kanban_decompose.py index fffe52850f..3af8ac9099 100644 --- a/hermes_cli/kanban_decompose.py +++ b/hermes_cli/kanban_decompose.py @@ -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.`` 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] diff --git a/hermes_cli/kanban_specify.py b/hermes_cli/kanban_specify.py index 24605d3fb5..0067c601d8 100644 --- a/hermes_cli/kanban_specify.py +++ b/hermes_cli/kanban_specify.py @@ -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..*`` 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)