From 5dfea72f751a7fb9b424577df444ccdb2bfdb35e Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Mon, 7 Sep 2026 02:58:11 -0700 Subject: [PATCH] refactor: extract Kanban graph persistence into topical sibling --- cron/AGENTS.md | 2 +- hermes_cli/kanban_db.py | 173 +--------------- hermes_cli/kanban_db_graph.py | 189 ++++++++++++++++++ hermes_cli/kanban_decompose.py | 3 +- tests/hermes_cli/test_kanban_decompose_db.py | 7 +- .../test_kanban_worktree_isolation.py | 3 +- 6 files changed, 200 insertions(+), 177 deletions(-) create mode 100644 hermes_cli/kanban_db_graph.py diff --git a/cron/AGENTS.md b/cron/AGENTS.md index 756f1eed0e..78a13580bf 100644 --- a/cron/AGENTS.md +++ b/cron/AGENTS.md @@ -33,7 +33,7 @@ Durable SQLite-backed board letting multiple profiles/workers collaborate. Users zero outside a kanban task (footprint ladder rung 3). - **CLI:** `hermes_cli/kanban.py` facade + 14 `kanban_*.py` siblings (`boards`, `db`, `db_connect`, - `db_dispatch`, `db_notify`, `workspace`, ...). Verbs: `init, create, list (ls), show, assign, link, + `db_dispatch`, `db_notify`, `db_graph` (task initialization and decomposition), `workspace`, ...). Verbs: `init, create, list (ls), show, assign, link, unlink, comment, attach, attachments, attach-rm, complete, request-review, request-changes, reopen-review, block, unblock, archive, tail`, plus `watch, stats, runs, log, assignees, heartbeat, notify-*, dispatch, daemon, gc`. Argparse alias dispatch must accept both `list` and `ls` (root). diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 5337525d13..2dfe348858 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -1242,7 +1242,9 @@ def create_task( ``project_source_task_id``: cross-profile fallback when ``project_id`` is not in the active profile's projects.db — see ``_resolve_project_link``. """ + from hermes_cli.kanban_db_graph import _initial_task_status from hermes_cli.kanban_pr_acceptance import validate_contract + completion_contract = validate_contract(completion_contract) model_override, provider_override = _validate_model_override(model_override, provider_override) reasoning_effort = normalize_reasoning_effort(reasoning_effort) @@ -1367,30 +1369,6 @@ def _board_meta_for(board: Optional[str]) -> dict: return read_board_metadata(board if board else get_current_board()) -def _initial_task_status( - conn: sqlite3.Connection, parents: tuple[str, ...], initial_status: str, triage: bool, -) -> str: - """Status for a new task: ``blocked``/``triage`` when parked by the caller, - else ``ready`` unless a parent is not yet ``done`` (-> ``todo``). Parent ids - are validated in every mode (even triage) so link rows never dangle.""" - if parents: - missing = _missing_task_ids(conn, parents) - if missing: - raise ValueError(f"unknown parent task(s): {', '.join(missing)}") - if initial_status == "blocked": - return "blocked" - if triage: - return "triage" - if parents: - rows = conn.execute( - "SELECT status FROM tasks WHERE id IN " - "(" + ",".join("?" * len(parents)) + ")", parents, - ).fetchall() - if any(r["status"] != "done" for r in rows): - return "todo" - return "ready" - - def _project_branch_name(project_obj: Any, task_id: str, title: Optional[str]) -> Optional[str]: from hermes_cli import projects_db as _pdb @@ -3496,153 +3474,6 @@ def specify_triage_task( return True -def _validate_children_graph(children: list) -> None: - """DB-free shape check + Kahn's cycle check on the sibling graph (a cycle - would deadlock every involved child in ``todo`` forever).""" - for idx, child in enumerate(children): - if not isinstance(child, dict): - raise ValueError(f"child[{idx}] is not a dict") - title = child.get("title") - if not isinstance(title, str) or not title.strip(): - raise ValueError(f"child[{idx}].title is required") - parents_idx = child.get("parents") or [] - if not isinstance(parents_idx, list): - raise ValueError(f"child[{idx}].parents must be a list") - for p in parents_idx: - if not isinstance(p, int) or p < 0 or p >= len(children): - raise ValueError(f"child[{idx}].parents[{p}] is not a valid index into children") - if p == idx: - raise ValueError(f"child[{idx}] cannot list itself as a parent") - - in_deg = [0] * len(children) - adj: list[list[int]] = [[] for _ in children] - for i, c in enumerate(children): - for p in (c.get("parents") or []): - adj[p].append(i) - in_deg[i] += 1 - queue = [i for i in range(len(children)) if in_deg[i] == 0] - seen = 0 - while queue: - seen += 1 - for nb in adj[queue.pop()]: - in_deg[nb] -= 1 - if in_deg[nb] == 0: - queue.append(nb) - if seen != len(children): - raise ValueError("cyclic dependency detected in decomposed children list") - - -def decompose_triage_task( - conn: sqlite3.Connection, task_id: str, *, root_assignee: Optional[str], children: list[dict], - author: Optional[str] = None, auto_promote: bool = True, -) -> Optional[list[str]]: - """Fan a triage task out into children and move the root to ``todo``; the root - waits on every child and wakes (``ready``) when all are done. - - ``children``: dicts of ``title`` (required), ``body``, ``assignee``, - ``parents`` (indices into this list), optional workspace overrides. - Returns child ids in input order, or None when the root is missing / not - in triage. Atomic: a malformed entry aborts the whole fan-out. - """ - if not children: - return None - if root_assignee is not None: - root_assignee = _canonical_assignee(root_assignee) - _validate_children_graph(children) - - # ONE txn so the fan-out is atomic; helpers that open their own write_txn - # (create_task, link_tasks, add_comment) must not be called in here. - now = int(time.time()) - with write_txn(conn): - root_row = conn.execute( - "SELECT id, status, tenant, workspace_kind, workspace_path " - "FROM tasks WHERE id = ?", (task_id,), - ).fetchone() - if root_row is None or root_row["status"] != "triage": - return None - child_ids = [ - _insert_decomposed_child(conn, task_id, root_row, child, author, now) - for child in children - ] - # Sibling edges within the decomposed graph. - for idx, child in enumerate(children): - for p_idx in child.get("parents") or []: - parent_id, child_id = child_ids[p_idx], child_ids[idx] - _link(conn, parent_id, child_id) - _append_event(conn, child_id, "linked", {"parent": parent_id, "child": child_id}) - # Root waits for the whole graph: link it under EVERY child (simpler - # than computing leaves; cycle-free since the root is only ever a child). - for cid in child_ids: - _link(conn, cid, task_id) - # Flip the root triage -> todo, assignee -> orchestrator. - sets = ["status = 'todo'"] - params: list[Any] = [] - if root_assignee is not None: - sets.append("assignee = ?") - params.append(root_assignee) - params.append(task_id) - conn.execute(f"UPDATE tasks SET {', '.join(sets)} WHERE id = ?", tuple(params)) - if author and author.strip(): - _insert_comment( - conn, task_id, author.strip(), - "Decomposed into " + ", ".join(child_ids) - + ". Root will wake when all children complete.", - now, - ) - _append_event( - conn, task_id, "decomposed", {"child_ids": child_ids, "root_assignee": root_assignee}, - ) - # Outside the txn (own IMMEDIATE txn). ``auto_promote=False`` leaves the - # children in ``todo`` for manual-review-first workflows. - if auto_promote: - recompute_ready(conn) - return child_ids - - -def _insert_decomposed_child( - conn: sqlite3.Connection, root_id: str, root_row: sqlite3.Row, child: dict, - author: Optional[str], now: int, -) -> str: - """Insert one decomposed child as ``todo`` (linked under the root later so - the dispatcher only ever sees a coherent graph); returns its id. - - Workspace: per-child override wins, else inherit the root's kind. Path - inherits only when kinds match (a 'dir' child must not point at the - root's worktree) and NEVER for worktrees — siblings dispatch concurrently - and one shared checkout would put them all on the first sibling's branch - with no lock; leaving it unset makes dispatch materialize a fresh - ``/.worktrees/`` per child from the board anchor. - """ - root_ws_kind = root_row["workspace_kind"] or "scratch" - child_ws_kind = child.get("workspace_kind") or root_ws_kind - if child.get("workspace_path"): - child_ws_path = child.get("workspace_path") - elif child_ws_kind == "worktree": - child_ws_path = None - elif child_ws_kind == root_ws_kind: - child_ws_path = root_row["workspace_path"] - else: - child_ws_path = None - new_id = _new_task_id() - body = child.get("body") - conn.execute( - "INSERT INTO tasks " - "(id, title, body, assignee, status, workspace_kind, " - " workspace_path, tenant, created_at, created_by) " - "VALUES (?, ?, ?, ?, 'todo', ?, ?, ?, ?, ?)", - ( - new_id, child["title"].strip(), body if isinstance(body, str) else None, - _canonical_assignee(child.get("assignee")), child_ws_kind, child_ws_path, - root_row["tenant"], now, (author or "decomposer"), - ), - ) - _append_event( - conn, new_id, "created", {"by": author or "decomposer", "from_decompose_of": root_id}, - ) - _inherit_notify_subs(conn, new_id, (root_id,), created_at=now) - return new_id - - def archive_task(conn: sqlite3.Connection, task_id: str) -> bool: with write_txn(conn): cur = conn.execute( diff --git a/hermes_cli/kanban_db_graph.py b/hermes_cli/kanban_db_graph.py new file mode 100644 index 0000000000..d15ab359ea --- /dev/null +++ b/hermes_cli/kanban_db_graph.py @@ -0,0 +1,189 @@ +"""Task graph initialization and atomic decomposition persistence.""" +from __future__ import annotations + +import sqlite3 +import time +from typing import Any, Optional + +def _initial_task_status( + conn: sqlite3.Connection, parents: tuple[str, ...], initial_status: str, triage: bool, +) -> str: + """Status for a new task: ``blocked``/``triage`` when parked by the caller, + else ``ready`` unless a parent is not yet ``done`` (-> ``todo``). Parent ids + are validated in every mode (even triage) so link rows never dangle.""" + from hermes_cli.kanban_db import _missing_task_ids + + if parents: + missing = _missing_task_ids(conn, parents) + if missing: + raise ValueError(f"unknown parent task(s): {', '.join(missing)}") + if initial_status == "blocked": + return "blocked" + if triage: + return "triage" + if parents: + rows = conn.execute( + "SELECT status FROM tasks WHERE id IN " + "(" + ",".join("?" * len(parents)) + ")", parents, + ).fetchall() + if any(r["status"] != "done" for r in rows): + return "todo" + return "ready" + + +def _validate_children_graph(children: list) -> None: + """DB-free shape check + Kahn's cycle check on the sibling graph (a cycle + would deadlock every involved child in ``todo`` forever).""" + for idx, child in enumerate(children): + if not isinstance(child, dict): + raise ValueError(f"child[{idx}] is not a dict") + title = child.get("title") + if not isinstance(title, str) or not title.strip(): + raise ValueError(f"child[{idx}].title is required") + parents_idx = child.get("parents") or [] + if not isinstance(parents_idx, list): + raise ValueError(f"child[{idx}].parents must be a list") + for p in parents_idx: + if not isinstance(p, int) or p < 0 or p >= len(children): + raise ValueError(f"child[{idx}].parents[{p}] is not a valid index into children") + if p == idx: + raise ValueError(f"child[{idx}] cannot list itself as a parent") + + in_deg = [0] * len(children) + adj: list[list[int]] = [[] for _ in children] + for i, c in enumerate(children): + for p in (c.get("parents") or []): + adj[p].append(i) + in_deg[i] += 1 + queue = [i for i in range(len(children)) if in_deg[i] == 0] + seen = 0 + while queue: + seen += 1 + for nb in adj[queue.pop()]: + in_deg[nb] -= 1 + if in_deg[nb] == 0: + queue.append(nb) + if seen != len(children): + raise ValueError("cyclic dependency detected in decomposed children list") + + +def decompose_triage_task( + conn: sqlite3.Connection, task_id: str, *, root_assignee: Optional[str], children: list[dict], + author: Optional[str] = None, auto_promote: bool = True, +) -> Optional[list[str]]: + """Fan a triage task out into children and move the root to ``todo``; the root + waits on every child and wakes (``ready``) when all are done. + + ``children``: dicts of ``title`` (required), ``body``, ``assignee``, + ``parents`` (indices into this list), optional workspace overrides. + Returns child ids in input order, or None when the root is missing / not + in triage. Atomic: a malformed entry aborts the whole fan-out. + """ + from hermes_cli.kanban_db import ( + _canonical_assignee, _link, _append_event, _insert_comment, + write_txn, recompute_ready, + ) + + if not children: + return None + if root_assignee is not None: + root_assignee = _canonical_assignee(root_assignee) + _validate_children_graph(children) + + # ONE txn so the fan-out is atomic; helpers that open their own write_txn + # (create_task, link_tasks, add_comment) must not be called in here. + now = int(time.time()) + with write_txn(conn): + root_row = conn.execute( + "SELECT id, status, tenant, workspace_kind, workspace_path " + "FROM tasks WHERE id = ?", (task_id,), + ).fetchone() + if root_row is None or root_row["status"] != "triage": + return None + child_ids = [ + _insert_decomposed_child(conn, task_id, root_row, child, author, now) + for child in children + ] + # Sibling edges within the decomposed graph. + for idx, child in enumerate(children): + for p_idx in child.get("parents") or []: + parent_id, child_id = child_ids[p_idx], child_ids[idx] + _link(conn, parent_id, child_id) + _append_event(conn, child_id, "linked", {"parent": parent_id, "child": child_id}) + # Root waits for the whole graph: link it under EVERY child (simpler + # than computing leaves; cycle-free since the root is only ever a child). + for cid in child_ids: + _link(conn, cid, task_id) + # Flip the root triage -> todo, assignee -> orchestrator. + sets = ["status = 'todo'"] + params: list[Any] = [] + if root_assignee is not None: + sets.append("assignee = ?") + params.append(root_assignee) + params.append(task_id) + conn.execute(f"UPDATE tasks SET {', '.join(sets)} WHERE id = ?", tuple(params)) + if author and author.strip(): + _insert_comment( + conn, task_id, author.strip(), + "Decomposed into " + ", ".join(child_ids) + + ". Root will wake when all children complete.", + now, + ) + _append_event( + conn, task_id, "decomposed", {"child_ids": child_ids, "root_assignee": root_assignee}, + ) + # Outside the txn (own IMMEDIATE txn). ``auto_promote=False`` leaves the + # children in ``todo`` for manual-review-first workflows. + if auto_promote: + recompute_ready(conn) + return child_ids + + +def _insert_decomposed_child( + conn: sqlite3.Connection, root_id: str, root_row: sqlite3.Row, child: dict, + author: Optional[str], now: int, +) -> str: + """Insert one decomposed child as ``todo`` (linked under the root later so + the dispatcher only ever sees a coherent graph); returns its id. + + Workspace: per-child override wins, else inherit the root's kind. Path + inherits only when kinds match (a 'dir' child must not point at the + root's worktree) and NEVER for worktrees — siblings dispatch concurrently + and one shared checkout would put them all on the first sibling's branch + with no lock; leaving it unset makes dispatch materialize a fresh + ``/.worktrees/`` per child from the board anchor. + """ + from hermes_cli.kanban_db import ( + _new_task_id, _canonical_assignee, _append_event, _inherit_notify_subs, + ) + + root_ws_kind = root_row["workspace_kind"] or "scratch" + child_ws_kind = child.get("workspace_kind") or root_ws_kind + if child.get("workspace_path"): + child_ws_path = child.get("workspace_path") + elif child_ws_kind == "worktree": + child_ws_path = None + elif child_ws_kind == root_ws_kind: + child_ws_path = root_row["workspace_path"] + else: + child_ws_path = None + new_id = _new_task_id() + body = child.get("body") + conn.execute( + "INSERT INTO tasks " + "(id, title, body, assignee, status, workspace_kind, " + " workspace_path, tenant, created_at, created_by) " + "VALUES (?, ?, ?, ?, 'todo', ?, ?, ?, ?, ?)", + ( + new_id, child["title"].strip(), body if isinstance(body, str) else None, + _canonical_assignee(child.get("assignee")), child_ws_kind, child_ws_path, + root_row["tenant"], now, (author or "decomposer"), + ), + ) + _append_event( + conn, new_id, "created", {"by": author or "decomposer", "from_decompose_of": root_id}, + ) + _inherit_notify_subs(conn, new_id, (root_id,), created_at=now) + return new_id + + diff --git a/hermes_cli/kanban_decompose.py b/hermes_cli/kanban_decompose.py index d2a93ecacf..b17723a9e2 100644 --- a/hermes_cli/kanban_decompose.py +++ b/hermes_cli/kanban_decompose.py @@ -23,6 +23,7 @@ from dataclasses import dataclass from typing import Optional from hermes_cli import kanban_db as kb +from hermes_cli.kanban_db_graph import decompose_triage_task from hermes_cli import kanban_db_connect as kbc from hermes_cli import profiles as profiles_mod from hermes_cli.kanban_specify import ( @@ -275,7 +276,7 @@ def _apply_fanout(task_id: str, parsed: dict, routing: _Routing, author: str) -> return DecomposeOutcome(task_id, False, reason) try: with kbc.connect_closing() as conn: - child_ids = kb.decompose_triage_task( + child_ids = decompose_triage_task( conn, task_id, root_assignee=routing.orchestrator, diff --git a/tests/hermes_cli/test_kanban_decompose_db.py b/tests/hermes_cli/test_kanban_decompose_db.py index 73109419e6..d7eec08821 100644 --- a/tests/hermes_cli/test_kanban_decompose_db.py +++ b/tests/hermes_cli/test_kanban_decompose_db.py @@ -1,4 +1,4 @@ -"""Tests for kb.decompose_triage_task — the DB-layer atomic fan-out +"""Tests for decompose_triage_task — the DB-layer atomic fan-out from the triage column. LLM-free by design. """ @@ -9,6 +9,7 @@ from pathlib import Path import pytest from hermes_cli import kanban_db as kb +from hermes_cli.kanban_db_graph import decompose_triage_task from hermes_cli import kanban_db_connect as kbc @@ -43,7 +44,7 @@ def test_decompose_creates_children_and_promotes_root(kanban_home): {"title": "build it", "body": "write code", "assignee": "engineer", "parents": [0]}, ] with kbc.connect() as conn: - child_ids = kb.decompose_triage_task( + child_ids = decompose_triage_task( conn, tid, root_assignee="orchestrator", @@ -72,7 +73,7 @@ def test_decompose_creates_children_and_promotes_root(kanban_home): def test_decompose_records_audit_comment_and_event(kanban_home): with kbc.connect() as conn: tid = _create_triage(conn) - child_ids = kb.decompose_triage_task( + child_ids = decompose_triage_task( conn, tid, root_assignee="orch", diff --git a/tests/hermes_cli/test_kanban_worktree_isolation.py b/tests/hermes_cli/test_kanban_worktree_isolation.py index 5932161fb9..10a9ab1aff 100644 --- a/tests/hermes_cli/test_kanban_worktree_isolation.py +++ b/tests/hermes_cli/test_kanban_worktree_isolation.py @@ -22,6 +22,7 @@ from pathlib import Path import pytest from hermes_cli import kanban_db as kb +from hermes_cli.kanban_db_graph import decompose_triage_task from hermes_cli import kanban_db_connect as kbc from hermes_cli import kanban_db_workspace as kbw @@ -78,7 +79,7 @@ def test_decompose_worktree_children_get_own_workspace(kanban_home): ) conn.commit() - child_ids = kb.decompose_triage_task( + child_ids = decompose_triage_task( conn, root, root_assignee="orchestrator",