diff --git a/evals/kanban_graph_identity.py b/evals/kanban_graph_identity.py new file mode 100644 index 0000000000..507e11b9cb --- /dev/null +++ b/evals/kanban_graph_identity.py @@ -0,0 +1,86 @@ +"""Credential-free SQLite integration probe; run with --repo PATH.""" +import argparse +from concurrent.futures import ThreadPoolExecutor +import importlib +import json +import os +from pathlib import Path +import subprocess +import sys +import tempfile +import threading + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--repo", type=Path, required=True) + parser.add_argument("--child", action="store_true") + args = parser.parse_args() + if not args.child: + with tempfile.TemporaryDirectory(prefix="kanban-identity-") as home: + env = {"HOME": home, "HERMES_HOME": home + "/hermes", "PATH": os.defpath, + "LANG": "C.UTF-8", "TZ": "UTC", "PYTHONDONTWRITEBYTECODE": "1"} + return subprocess.run( + [sys.executable, str(Path(__file__).resolve()), "--repo", str(args.repo), "--child"], + cwd=args.repo, env=env, check=False, + ).returncode + sys.path.insert(0, str(args.repo)) + from hermes_cli import kanban_db as kb, kanban_db_connect as kbc + from hermes_cli.kanban_db_dispatch import dispatch_once + from tools.kanban_tools import _handle_create + from hermes_cli.kanban_decompose import _apply_fanout, _Routing + graph = importlib.import_module("hermes_cli.kanban_db_graph") if (args.repo / "hermes_cli/kanban_db_graph.py").exists() else kb + decompose = graph.decompose_triage_task + specs = [{"title": "work", "assignee": "default"}] + with kbc.connect_closing() as conn: + parent = kb.create_task(conn, title="prerequisite", tenant="business-a") + root = kb.create_task(conn, title="root", triage=True, tenant="business-a", parents=[parent]) + downstream = kb.create_task(conn, title="downstream", parents=[root], tenant="business-a") + child = kb.create_task(conn, title="manual child", parents=[parent]) + tool = json.loads(_handle_create({"title": "tool child", "assignee": "default", "parents": [parent]})) + assert tool["ok"] + out = {"db_tenant": kb.get_task(conn, child).tenant, + "tool_tenant": kb.get_task(conn, tool["task_id"]).tenant} + first = decompose(conn, root, root_assignee="default", children=specs) + assert first and kb.get_task(conn, first[0]).tenant == "business-a" + assert conn.execute("SELECT 1 FROM task_links WHERE parent_id=? AND child_id=?", (parent, root)).fetchone() + assert conn.execute("SELECT 1 FROM task_links WHERE parent_id=? AND child_id=?", (root, downstream)).fetchone() + out["legitimate_prelinked_first_fanout"] = True + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status='triage' WHERE id=?", (root,)) + out["repeat_refused"] = decompose(conn, root, root_assignee="default", children=specs) is None + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status='done' WHERE id=?", (root,)) + conn.execute("UPDATE task_events SET created_at=0 WHERE task_id=?", (root,)) + kb.gc_events(conn, older_than_seconds=1) + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status='triage' WHERE id=?", (root,)) + out["retention_repeat_refused"] = decompose(conn, root, root_assignee="default", children=specs) is None + ids = list(conn.execute("SELECT id FROM tasks ORDER BY id")) + for _ in range(3): + dispatch_once(conn, max_spawn=0) + out["three_dispatch_ticks_stable"] = ids == list(conn.execute("SELECT id FROM tasks ORDER BY id")) + ingress_root = kb.create_task(conn, title="ingress", tenant="business-a", triage=True) + routing = _Routing("default", "default", True, [], {"default"}) + assert _apply_fanout(ingress_root, {"tasks": specs}, routing, "probe").ok + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status='triage' WHERE id=?", (ingress_root,)) + out["fanout_ingress_repeat_refused"] = not _apply_fanout(ingress_root, {"tasks": specs}, routing, "probe").ok + race_root = kb.create_task(conn, title="race", tenant="business-a", triage=True) + barrier = threading.Barrier(2) + def race(): + with kbc.connect_closing() as conn: + barrier.wait(timeout=10) + return decompose(conn, race_root, root_assignee="default", children=specs) + with ThreadPoolExecutor(max_workers=2) as pool: + outcomes = list(pool.map(lambda _: race(), range(2))) + out["concurrent_winners"] = sum(value is not None for value in outcomes) + with kbc.connect_closing() as conn: + out["integrity"] = conn.execute("PRAGMA integrity_check").fetchone()[0] + assert out["three_dispatch_ticks_stable"] and out["concurrent_winners"] == 1 and out["integrity"] == "ok" + print(json.dumps(out, sort_keys=True)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 2dfe348858..046681515f 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -1242,7 +1242,7 @@ 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_db_graph import initial_task_state from hermes_cli.kanban_pr_acceptance import validate_contract completion_contract = validate_contract(completion_contract) @@ -1304,7 +1304,7 @@ def create_task( # allow_nested: graph builders compose create_task under one outer # commit so the dispatcher never sees a half-built graph. with write_txn(conn, allow_nested=True): - task_status = _initial_task_status(conn, parents, initial_status, triage) + task_status, tenant = initial_task_state(conn, parents, initial_status, triage, tenant) # Project worktree: fresh dir under the repo + deterministic # branch, instead of the random ``wt/`` worker fallback. if project_obj is not None and workspace_kind == "worktree": @@ -3848,11 +3848,11 @@ def task_age(task: Task) -> dict: # --- Retention + garbage collection --- def gc_events(conn: sqlite3.Connection, *, older_than_seconds: int = 30 * 24 * 3600) -> int: - """Delete events older than the cutoff on done/archived tasks only; returns the count.""" + """Prune old done/archived events, retaining decomposition identity until task deletion.""" cutoff = int(time.time()) - int(older_than_seconds) with write_txn(conn): cur = conn.execute( - "DELETE FROM task_events WHERE created_at < ? AND task_id IN " + "DELETE FROM task_events WHERE created_at < ? AND kind != 'decomposed' AND task_id IN " "(SELECT id FROM tasks WHERE status IN ('done', 'archived'))", (cutoff,), ) return int(cur.rowcount or 0) diff --git a/hermes_cli/kanban_db_graph.py b/hermes_cli/kanban_db_graph.py index d15ab359ea..ea329ed7b2 100644 --- a/hermes_cli/kanban_db_graph.py +++ b/hermes_cli/kanban_db_graph.py @@ -5,30 +5,33 @@ 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 +def initial_task_state( + conn: sqlite3.Connection, parents: tuple[str, ...], initial_status: str, + triage: bool, tenant: Optional[str], +) -> tuple[str, Optional[str]]: + """Resolve state and tenant under the creator's write transaction. + Parent order breaks ties in this soft namespace; explicit tenant wins. + Validate parents even for parked tasks so links never dangle. + """ + rows = {} if parents: - missing = _missing_task_ids(conn, parents) + rows = {row["id"]: row for row in conn.execute( + "SELECT id, status, tenant FROM tasks WHERE id IN " + "(" + ",".join("?" * len(parents)) + ")", parents, + )} + missing = [pid for pid in parents if pid not in rows] if missing: raise ValueError(f"unknown parent task(s): {', '.join(missing)}") + if tenant is None: + tenant = next((rows[pid]["tenant"] for pid in parents if rows[pid]["tenant"]), None) if initial_status == "blocked": - return "blocked" + return "blocked", tenant 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" + return "triage", tenant + if any(row["status"] != "done" for row in rows.values()): + return "todo", tenant + return "ready", tenant def _validate_children_graph(children: list) -> None: @@ -77,7 +80,7 @@ def decompose_triage_task( ``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. + in triage, or has already decomposed. Atomic: malformed entries abort fan-out. """ from hermes_cli.kanban_db import ( _canonical_assignee, _link, _append_event, _insert_comment, @@ -100,6 +103,13 @@ def decompose_triage_task( ).fetchone() if root_row is None or root_row["status"] != "triage": return None + # Dependency links alone do not imply lineage. The completion event is + # committed with the graph, and survives re-triage or unlinking. + if conn.execute( + "SELECT 1 FROM task_events WHERE task_id = ? AND kind = 'decomposed' LIMIT 1", + (task_id,), + ).fetchone(): + return None child_ids = [ _insert_decomposed_child(conn, task_id, root_row, child, author, now) for child in children @@ -185,5 +195,3 @@ def _insert_decomposed_child( ) _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 b17723a9e2..3bdb89927e 100644 --- a/hermes_cli/kanban_decompose.py +++ b/hermes_cli/kanban_decompose.py @@ -290,7 +290,7 @@ def _apply_fanout(task_id: str, parsed: dict, routing: _Routing, author: str) -> 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, False, "task already decomposed or moved out of triage") return DecomposeOutcome( task_id, True, f"decomposed into {len(child_ids)} children", fanout=True, child_ids=child_ids, ) diff --git a/tests/hermes_cli/test_kanban_graph_identity.py b/tests/hermes_cli/test_kanban_graph_identity.py new file mode 100644 index 0000000000..29c434caf5 --- /dev/null +++ b/tests/hermes_cli/test_kanban_graph_identity.py @@ -0,0 +1,55 @@ +"""Graph identity is completion history, not prerequisite edges.""" +import json +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 + + +def test_completed_decomposition_survives_retriage(tmp_path, monkeypatch): + monkeypatch.setattr(Path, "home", lambda: tmp_path) + monkeypatch.setenv("HERMES_HOME", str(tmp_path / "hermes")) + with kbc.connect_closing() as conn: + prerequisite = kb.create_task(conn, title="prerequisite", tenant="business-a") + root = kb.create_task(conn, title="root", triage=True, tenant="business-a", parents=[prerequisite]) + downstream = kb.create_task(conn, title="downstream", parents=[root], tenant="business-a") + specs = [{"title": "work", "assignee": "default"}] + first = decompose_triage_task(conn, root, root_assignee="default", children=specs) + assert first and kb.get_task(conn, first[0]).tenant == "business-a" + assert conn.execute("SELECT 1 FROM task_links WHERE parent_id=? AND child_id=?", (prerequisite, root)).fetchone() + assert conn.execute("SELECT 1 FROM task_links WHERE parent_id=? AND child_id=?", (root, downstream)).fetchone() + # Retention must not erase the identity of a completed fan-out. + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status='done' WHERE id=?", (root,)) + conn.execute("UPDATE task_events SET created_at=0 WHERE task_id=?", (root,)) + assert kb.gc_events(conn, older_than_seconds=1) > 0 + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status='triage' WHERE id=?", (root,)) + before = list(conn.execute("SELECT id FROM tasks ORDER BY id")) + assert decompose_triage_task(conn, root, root_assignee="default", children=specs) is None + assert list(conn.execute("SELECT id FROM tasks ORDER BY id")) == before + assert conn.execute("SELECT count(*) FROM task_events WHERE task_id=? AND kind='decomposed'", (root,)).fetchone()[0] == 1 + assert conn.execute("PRAGMA integrity_check").fetchone()[0] == "ok" + + +def test_parent_tenant_is_inherited_at_creation_boundary(tmp_path, monkeypatch): + monkeypatch.setattr(Path, "home", lambda: tmp_path) + monkeypatch.setenv("HERMES_HOME", str(tmp_path / "hermes")) + monkeypatch.delenv("HERMES_TENANT", raising=False) + monkeypatch.delenv("HERMES_KANBAN_TASK", raising=False) + from tools.kanban_tools import _handle_create + with kbc.connect_closing() as conn: + unscoped = kb.create_task(conn, title="unscoped") + parent = kb.create_task(conn, title="parent", tenant="business-a") + for explicit, expected in [(None, "business-a"), ("business-b", "business-b")]: + child = kb.create_task(conn, title="child", parents=[unscoped, parent], tenant=explicit) + assert kb.get_task(conn, child).tenant == expected + result = json.loads(_handle_create({"title": "tool child", "assignee": "default", "parents": [parent]})) + assert result["ok"] + assert kb.get_task(conn, result["task_id"]).tenant == "business-a" + with pytest.raises(ValueError, match="unknown parent"): + kb.create_task(conn, title="invalid", parents=["missing"]) + assert kb.get_task(conn, unscoped).tenant is None diff --git a/website/docs/user-guide/features/kanban.md b/website/docs/user-guide/features/kanban.md index f337128dab..47ccb1cb74 100644 --- a/website/docs/user-guide/features/kanban.md +++ b/website/docs/user-guide/features/kanban.md @@ -632,6 +632,16 @@ The kanban board has two ways to handle a task you drop into the Triage column: **Auto (default)** — `kanban.auto_decompose: true`. The gateway-embedded dispatcher runs the **decomposer** on each tick, capped by `kanban.auto_decompose_per_tick` (default 3 tasks per tick) so a bulk-load of triage tasks doesn't burst-spend the auxiliary LLM. The decomposer uses the built-in decomposition prompt plus the `auxiliary.kanban_decomposer` model path, reads your installed profiles + their descriptions, and asks the LLM to produce a JSON task graph: which tasks to spawn, who they go to, and which depend on which. The original triage task becomes the parent of every leaf in the graph, so it stays alive until the whole graph completes - and then promotes back to `ready` so its assignee (`kanban.orchestrator_profile`, or the active default profile when unset) can judge completion and add more tasks if the work isn't done. This is the "drop a one-liner, walk away" flow. +A completed built-in fan-out is recorded atomically with its child graph. Moving +that root back to Triage does not create another graph; ordinary prerequisite +links do not prevent a task's first decomposition. The completion marker survives +event retention until the task is deleted. This is not semantic deduplication of +independently created manual graphs, nor a repair for previously pruned history. + +When a new task omits its tenant, creation inherits the first nonempty tenant +among its parents, in supplied order. An explicit tenant (including the worker's +active tenant passed by tools) wins. Boards remain the hard isolation boundary. + **Manual** — `kanban.auto_decompose: false`. Triage tasks stay in triage until you act. Click the **⚗ Decompose** button on a card, run `hermes kanban decompose ` (or `--all`), or use `/kanban decompose ` from a chat. This matches the pre-decomposer behavior of the board, useful when you want full control over what runs when. **Important boundary:** Manual mode disables only the built-in Triage decomposer. It does not prevent a profile from calling `kanban_create`, and it does not disable creator-session wake-ups. With `kanban.auto_subscribe_on_create: true`, a task's terminal event resumes the originating agent with a synthetic status turn so it can inspect the handoff and decide whether genuinely new follow-up work is needed. Set `auto_subscribe_on_create: false` when task completion should remain passive. For provenance, built-in decomposer children use `created_by=auto-decomposer`; tasks created by a resumed profile carry that profile name instead.