fix: retain Kanban decomposition identity and inherit parent tenants

This commit is contained in:
Teknium
2026-09-07 03:05:54 -07:00
parent 5dfea72f75
commit 258fa9741c
6 changed files with 185 additions and 26 deletions
+86
View File
@@ -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())
+4 -4
View File
@@ -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/<id>`` 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)
+29 -21
View File
@@ -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
+1 -1
View File
@@ -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,
)
@@ -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
@@ -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 <id>` (or `--all`), or use `/kanban decompose <id>` 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.