refactor(kanban): CLI micro-helpers, action dispatch tables, shared triage helpers; active_sessions dedupe

hermes_cli/kanban.py 3,565 -> 2,912; kanban_diagnostics 1,216 -> 996;
kanban_decompose 468 -> 393; kanban_transfer 478 -> 443; kanban_specify
264 -> 229; kanban_swarm 390 -> 378; active_sessions 871 -> 775. `hermes kanban
[sub] --help` byte-identical for all 55 parsers.

- kanban.py: _err / _print_json / _json_out / _fmt_counts / _bulk_apply /
  _obj_dict field tuples replace repeated print/JSON/exit-code blocks; action
  and board subcommand routing via dict dispatch; shared run-state and
  triage-sweep argparse blocks; argparse declarations re-packed (AST-identical).
- specify/decompose: one _run_triage_sweep driver, shared _extract_json_blob /
  _truncate / _profile_author / _title_body / _resolve_profile_from_cfg.
- diagnostics: rule helpers (_first_field / _latest_event_ts / _log_hint_action
  / _error_snippet), _rows_by_task fleet fetch; unreferenced DIAGNOSTIC_KINDS dropped.
- swarm: graph nodes share one create_task kwarg set.
- active_sessions: one _flock per platform, _pid_alive via _pid_liveness,
  shared _read_live_entries / _without_lease / _clean_metadata, table-driven
  strict registry validation.
- Docstrings/comments hand-compacted (AST-identical), invariants kept.
This commit is contained in:
Teknium
2026-09-02 12:09:01 -07:00
parent 841c975d43
commit ff660354f3
7 changed files with 1224 additions and 2350 deletions

View File

@@ -107,11 +107,8 @@ def active_session_limit_message(
max_sessions: int,
entries: Optional[list[dict[str, Any]]] = None,
) -> str:
# Name the holders: the slots are shared across CLI, desktop/TUI and the
# messaging gateway, so the surface that gets rejected is usually NOT the
# one squatting on them (idle desktop chats starving a Discord bot, say).
# Without this the message is unactionable and the only way to find out is
# reading runtime/active_sessions.json by hand.
# Name the holders: slots are shared across CLI, desktop/TUI and gateway,
# so the rejected surface is usually NOT the one squatting on them.
held = summarize_holders(entries or [])
detail = f" Held by: {held}." if held else ""
return (
@@ -124,42 +121,25 @@ def _registry_home(registry_home: str | Path | None = None) -> Path:
return Path(registry_home) if registry_home is not None else Path(get_hermes_home())
# WHY A REFUSAL IS REFUSED, in a form a caller can branch on.
#
# The two refusals mean opposite things to an automated client. Capacity is
# "the machine is busy, come back later". Ownership is "this specific session
# has a live owner, and writing to it would interleave with theirs".
#
# Callers used to have only the human-readable message, so anything that needed
# to DECIDE had to match prose -- which silently changes meaning whenever the
# wording is improved. The reason is the contract; the message is for people.
# Machine-readable refusal reasons: the reason is the contract, the message is
# for people. Capacity = "busy, come back later"; ownership = "this session
# has a live owner and writing would interleave with theirs".
SESSION_NOT_OWNED = "SESSION_NOT_OWNED"
MAX_CONCURRENT_SESSIONS = "MAX_CONCURRENT_SESSIONS"
# Ownership could not be PROVEN either way: the registry was unreadable or
# corrupt. Distinct from SESSION_NOT_OWNED on purpose -- "someone else owns
# this" and "I cannot tell who owns this" call for different operator action,
# and collapsing the second into a silent go-ahead is exactly the fail-open
# hole that let two writers share one session (#94595 review, blocker 2).
# Ownership could not be PROVEN (registry unreadable/corrupt). Deliberately
# distinct from SESSION_NOT_OWNED: collapsing "can't tell" into a silent
# go-ahead is the fail-open hole that let two writers share one session.
SESSION_COORDINATION_UNAVAILABLE = "SESSION_COORDINATION_UNAVAILABLE"
# Advertised through the gateway so a client can tell a build that enforces
# per-session exclusivity from one that does not.
#
# A module constant rather than a config flag, deliberately: it is true because
# try_acquire_active_session below performs the check atomically, so it cannot
# be turned on by an operator who has not got the enforcement, and cannot drift
# out of step with it without this file changing.
# Advertised through the gateway. A module constant, not a config flag: it is
# true because try_acquire_active_session checks atomically, so it cannot drift
# out of step with the enforcement without this file changing.
PER_SESSION_EXCLUSIVE_SUBMIT = True
class ActiveSessionRefusal(str):
"""A refusal message that also carries a machine-readable ``reason``.
A ``str`` subclass so every existing caller keeps working untouched -- they
format it, hand it back as a JSON-RPC message, or just test it for None --
while a caller that must act on WHICH refusal happened reads ``.reason``
instead of matching the wording.
"""
"""A refusal message (``str`` subclass, so existing callers are untouched)
that also carries a machine-readable ``reason``."""
reason: str
@@ -170,13 +150,10 @@ class ActiveSessionRefusal(str):
def _is_same_writer(entry: dict[str, Any], metadata: Optional[dict[str, Any]]) -> bool:
"""True when an existing lease belongs to the very caller now re-acquiring it.
Both halves are required. A pid alone would let two live sessions in one
process steal each other's lease -- which is a real hazard, since each holds
its own snapshot of the transcript. A live session id alone would let another
process with a coincidentally equal id do the same.
"""
"""True when an existing lease belongs to the very caller re-acquiring it.
Identity is (pid, live_session_id): pid alone would let two live sessions in
one process steal each other's lease; the live id alone would let another
process with an equal id do the same."""
try:
if int(entry.get("pid") or -1) != os.getpid():
return False
@@ -223,6 +200,19 @@ def _lease_paths(
return _state_path(home), _lock_path(home)
def _flock(fh, *, lock: bool) -> None:
"""Exclusive whole-file lock/unlock on ``fh`` (fcntl on POSIX, msvcrt on Windows)."""
if os.name == "nt":
import msvcrt
fh.seek(0)
msvcrt.locking(fh.fileno(), msvcrt.LK_LOCK if lock else msvcrt.LK_UNLCK, 1)
else:
import fcntl
fcntl.flock(fh.fileno(), fcntl.LOCK_EX if lock else fcntl.LOCK_UN)
class _FileLock:
def __init__(self, path: Path):
self.path = path
@@ -231,45 +221,21 @@ class _FileLock:
def __enter__(self):
self.path.parent.mkdir(parents=True, exist_ok=True)
self._fh = open(self.path, "a+b")
if os.name == "nt":
try:
import msvcrt
self._fh.seek(0)
msvcrt.locking(self._fh.fileno(), msvcrt.LK_LOCK, 1)
except Exception as exc:
self._fh.close()
self._fh = None
raise RuntimeError("active session file lock unavailable") from exc
else:
try:
import fcntl
fcntl.flock(self._fh.fileno(), fcntl.LOCK_EX)
except Exception as exc:
self._fh.close()
self._fh = None
raise RuntimeError("active session file lock unavailable") from exc
try:
_flock(self._fh, lock=True)
except Exception as exc:
self._fh.close()
self._fh = None
raise RuntimeError("active session file lock unavailable") from exc
return self
def __exit__(self, exc_type, exc, tb):
if self._fh is None:
return
if os.name == "nt":
try:
import msvcrt
self._fh.seek(0)
msvcrt.locking(self._fh.fileno(), msvcrt.LK_UNLCK, 1)
except Exception:
pass
else:
try:
import fcntl
fcntl.flock(self._fh.fileno(), fcntl.LOCK_UN)
except Exception:
pass
try:
_flock(self._fh, lock=False)
except Exception:
pass
try:
self._fh.close()
finally:
@@ -306,58 +272,51 @@ def _read_entries(path: Path, *, strict: bool = False) -> list[dict[str, Any]]:
seen_leases: set[str] = set()
for entry in valid:
lease_id = entry.get("lease_id")
session_id = entry.get("session_id")
pid = entry.get("pid")
if not isinstance(lease_id, str) or not lease_id.strip():
raise ActiveSessionRegistryError(
f"active session registry contains an invalid lease id: {path}"
)
if lease_id in seen_leases:
raise ActiveSessionRegistryError(
f"active session registry contains a duplicate lease id: {path}"
)
seen_leases.add(lease_id)
if not isinstance(session_id, str) or not session_id.strip():
raise ActiveSessionRegistryError(
f"active session registry contains an invalid session id: {path}"
)
if isinstance(pid, bool) or not isinstance(pid, (int, str)):
pid_int = 0
else:
try:
pid_int = int(pid)
except (TypeError, ValueError):
pid_int = 0
if pid_int <= 0:
raise ActiveSessionRegistryError(
f"active session registry contains an invalid pid: {path}"
)
surface = entry.get("surface")
if surface is not None and not isinstance(surface, str):
raise ActiveSessionRegistryError(
f"active session registry contains an invalid surface: {path}"
)
tracked = entry.get("track_liveness")
if tracked is not None and not isinstance(tracked, bool):
raise ActiveSessionRegistryError(
f"active session registry contains an invalid liveness marker: {path}"
)
metadata = entry.get("metadata")
if metadata is not None and not isinstance(metadata, dict):
raise ActiveSessionRegistryError(
f"active session registry contains invalid metadata: {path}"
)
process_start = entry.get("process_start_time")
parsed_process_start = _optional_float(process_start)
if process_start not in (None, "") and (
parsed_process_start is None or not math.isfinite(parsed_process_start)
# (problem-if-True predicate, message fragment) — checked lazily, in
# this order, so an unhashable lease id is reported before the dup check.
for bad, what in (
(lambda: not _nonblank_str(lease_id), "an invalid lease id"),
(lambda: lease_id in seen_leases, "a duplicate lease id"),
(lambda: not _nonblank_str(entry.get("session_id")), "an invalid session id"),
(lambda: _registry_pid(entry.get("pid")) <= 0, "an invalid pid"),
(lambda: not _optional_isinstance(entry.get("surface"), str), "an invalid surface"),
(lambda: not _optional_isinstance(entry.get("track_liveness"), bool), "an invalid liveness marker"),
(lambda: not _optional_isinstance(entry.get("metadata"), dict), "invalid metadata"),
(lambda: not _valid_process_start(entry.get("process_start_time")), "an invalid process start time"),
):
raise ActiveSessionRegistryError(
f"active session registry contains an invalid process start time: {path}"
)
if bad():
raise ActiveSessionRegistryError(
f"active session registry contains {what}: {path}"
)
seen_leases.add(lease_id)
return valid
def _nonblank_str(v: Any) -> bool:
return isinstance(v, str) and bool(v.strip())
def _optional_isinstance(v: Any, typ) -> bool:
return v is None or isinstance(v, typ)
def _registry_pid(pid: Any) -> int:
"""Registry pid as int; 0 for bools, non-int/str, or unparseable values."""
if isinstance(pid, bool) or not isinstance(pid, (int, str)):
return 0
try:
return int(pid)
except (TypeError, ValueError):
return 0
def _valid_process_start(v: Any) -> bool:
if v in (None, ""):
return True
parsed = _optional_float(v)
return parsed is not None and math.isfinite(parsed)
def _write_entries(path: Path, entries: list[dict[str, Any]]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_name(f"{path.name}.{os.getpid()}.{uuid.uuid4().hex}.tmp")
@@ -392,20 +351,24 @@ def _optional_float(value: Any) -> Optional[float]:
return None
def _pid_liveness(pid: Any, process_start_time: Any = None) -> Optional[bool]:
"""Return True/False for live/dead, or None when liveness is unknowable."""
def _pid_liveness(pid: Any, process_start_time: Any = None, *, lenient: bool = False) -> Optional[bool]:
"""Return True/False for live/dead, or None when liveness is unknowable.
``lenient`` never returns None: an unparseable pid or a failed existence
probe counts as dead, an unreadable current start time as alive.
"""
try:
pid_int = int(pid)
except (TypeError, ValueError):
return None
return False if lenient else None
if pid_int <= 0:
return None
return False if lenient else None
try:
from gateway.status import _pid_exists
exists = bool(_pid_exists(pid_int))
except Exception:
return None
return False if lenient else None
if not exists:
return False
expected_start = _optional_float(process_start_time)
@@ -413,32 +376,12 @@ def _pid_liveness(pid: Any, process_start_time: Any = None) -> Optional[bool]:
return True
current_start = _process_start_time(pid_int)
if current_start is None:
return None
return True if lenient else None
return abs(current_start - expected_start) < 0.001
def _pid_alive(pid: Any, process_start_time: Any = None) -> bool:
try:
pid_int = int(pid)
except (TypeError, ValueError):
return False
if pid_int <= 0:
return False
try:
from gateway.status import _pid_exists
exists = bool(_pid_exists(pid_int))
except Exception:
return False
if not exists:
return False
expected_start = _optional_float(process_start_time)
if expected_start is None:
return True
current_start = _process_start_time(pid_int)
if current_start is None:
return True
return abs(current_start - expected_start) < 0.001
return bool(_pid_liveness(pid, process_start_time, lenient=True))
def _prune_dead(
@@ -468,12 +411,9 @@ class ActiveSessionLease:
surface: str
enabled: bool = True
released: bool = False
# Registry paths pinned at acquisition time. A lease acquired under the
# root ``HERMES_HOME`` must release against the same registry even when
# ``release()`` runs inside a profile home override (native multiplex
# routes turns under ``_profile_runtime_scope``), otherwise the root
# entry survives until process exit and the session cap fills with
# phantom leases (#85431).
# Pinned at acquisition: a lease acquired under the root HERMES_HOME must
# release against the same registry even when release() runs inside a
# profile-home override, or phantom leases fill the session cap.
state_path: Optional[Path] = None
lock_path: Optional[Path] = None
track_liveness: bool = False
@@ -484,6 +424,33 @@ class ActiveSessionLease:
release_active_session(self)
def _clean_metadata(metadata: dict[str, Any]) -> dict[str, Any]:
return {str(k): v for k, v in metadata.items() if isinstance(k, str)}
def _without_lease(entries: list[dict[str, Any]], lease_id: str) -> list[dict[str, Any]]:
return [e for e in entries if str(e.get("lease_id") or "") != lease_id]
def _read_live_entries(
state_path: Path, *, track_liveness: bool, warn: str,
) -> Optional[tuple[list[dict[str, Any]], list[dict[str, Any]]]]:
"""``(raw, pruned)`` from the registry, or None when it is unreadable.
Liveness-tracked callers re-raise instead: they must not proceed on an
unprovable registry. Untracked callers get ``warn`` logged and decide
how to degrade themselves.
"""
try:
raw_entries = _read_entries(state_path, strict=True)
return raw_entries, _prune_dead(raw_entries, strict=track_liveness)
except ActiveSessionRegistryError:
if track_liveness:
raise
logger.warning(warn)
return None
def _lease_entry(
*,
lease_id: str,
@@ -505,9 +472,7 @@ def _lease_entry(
if track_liveness:
entry["track_liveness"] = True
if metadata:
entry["metadata"] = {
str(k): v for k, v in metadata.items() if isinstance(k, str)
}
entry["metadata"] = _clean_metadata(metadata)
return entry
@@ -522,27 +487,20 @@ def try_acquire_active_session(
) -> tuple[Optional[ActiveSessionLease], Optional[str]]:
"""Acquire an active-session slot.
Per-session exclusivity is CORRECTNESS and is enforced unconditionally:
at most one live owner may run a given stored session, whether or not an
operator configured ``max_concurrent_sessions`` (#94595). The concurrency
cap remains a resource POLICY and applies only when configured. Liveness
tracking keeps richer desktop lifecycle semantics; ``registry_home`` lets
profile-scoped backends share the owning profile's registry even when
launched from another home.
Per-session exclusivity is CORRECTNESS, enforced unconditionally: at most
one live owner per stored session. ``max_concurrent_sessions`` is resource
POLICY and applies only when configured. ``registry_home`` lets
profile-scoped backends share the owning profile's registry.
Returns ``(lease, None)`` on success and ``(None, refusal)`` otherwise,
where ``refusal`` is an :class:`ActiveSessionRefusal` carrying a
machine-readable ``reason``. Ownership uncertainty fails CLOSED: when the
registry cannot be read, the caller gets ``SESSION_COORDINATION_UNAVAILABLE``
rather than a silent go-ahead that could reopen the double-writer state.
Returns ``(lease, None)`` or ``(None, ActiveSessionRefusal)``. Ownership
uncertainty fails CLOSED with ``SESSION_COORDINATION_UNAVAILABLE``.
"""
max_sessions = resolve_max_concurrent_sessions(config)
lease_id = uuid.uuid4().hex
key = str(session_id or "")
# A session with no stored id yet cannot collide with another writer, and
# the strict registry schema (rightly) refuses entries with empty session
# ids. Nothing to fence, nothing to record: hand back a no-op lease.
# No stored id yet => nothing to fence or record (and the strict schema
# refuses empty session ids): hand back a no-op lease.
if not key and not track_liveness:
return ActiveSessionLease(
lease_id=lease_id,
@@ -560,22 +518,23 @@ def try_acquire_active_session(
)
state_path, lock_path = _lease_paths(registry_home=registry_home)
lease = ActiveSessionLease(
lease_id=lease_id,
session_id=key,
surface=str(surface),
state_path=state_path,
lock_path=lock_path,
track_liveness=track_liveness,
)
with _FileLock(lock_path):
try:
raw_entries = _read_entries(state_path, strict=True)
entries = _prune_dead(raw_entries, strict=track_liveness)
except ActiveSessionRegistryError:
if track_liveness:
raise
# A capacity cap could afford to degrade open -- worst case, more
# sessions than the operator wanted. Exclusivity cannot: "could not
# prove ownership" must never be collapsed into "no owner exists",
# because that silently reopens the exact double-writer state this
# fence guarantees against. Refuse, and say which file to fix.
logger.warning(
"Active-session registry is unavailable; refusing the session "
"rather than risking a concurrent writer"
)
# A capacity cap could degrade open; exclusivity cannot: "could not
# prove ownership" must never become "no owner exists".
loaded = _read_live_entries(
state_path, track_liveness=track_liveness,
warn="Active-session registry is unavailable; refusing the session "
"rather than risking a concurrent writer",
)
if loaded is None:
return None, ActiveSessionRefusal(
(
"Hermes could not read the active-session registry at "
@@ -584,47 +543,27 @@ def try_acquire_active_session(
),
SESSION_COORDINATION_UNAVAILABLE,
)
raw_entries, entries = loaded
pruned = len(raw_entries) - len(entries)
if pruned:
logger.info("Pruned %d stale active session lease(s)", pruned)
# Correctness first, and under the same lock that just pruned the dead
# owners -- so an owner that died is never mistaken for one that is
# running, and a live one is never overlooked.
#
# An empty key is exempt: a session with no stored id yet cannot collide
# with another, and treating "" as an identity would make every unsaved
# draft exclude every other one.
# Correctness first, under the same lock that just pruned dead owners.
# An empty key is exempt: treating "" as an identity would make every
# unsaved draft exclude every other one.
if key:
for index, existing in enumerate(entries):
if str(existing.get("session_id") or "") != key:
continue
# THE SAME WRITER IS NOT A SECOND WRITER.
#
# A live session that lost its lease reference -- its record was
# rebuilt in place, so the object holding the lease is unreachable
# while the session itself is still the one being driven -- would
# otherwise be fenced out of its own session by its own leak, and
# permanently: pruning only removes entries whose PROCESS is dead,
# and this process is very much alive.
#
# Identity here is (pid, live session id). Two processes never
# match, because their pids differ. Two live sessions inside one
# process never match, because their live ids differ. Only the
# exact same writer re-acquiring its own session matches, and that
# is re-entrancy rather than a concurrent writer.
# The same writer is not a second writer: a live session that
# leaked its lease reference would otherwise be fenced out of
# its own session permanently (pruning only removes entries
# whose PROCESS is dead). Re-entrancy, not concurrency.
if _is_same_writer(existing, metadata):
entries[index] = entry
_write_entries(state_path, entries)
return ActiveSessionLease(
lease_id=lease_id,
session_id=key,
surface=str(surface),
state_path=state_path,
lock_path=lock_path,
track_liveness=track_liveness,
), None
return lease, None
_write_entries(state_path, entries)
logger.info(
@@ -656,14 +595,7 @@ def try_acquire_active_session(
entries.append(entry)
_write_entries(state_path, entries)
return ActiveSessionLease(
lease_id=lease_id,
session_id=key,
surface=str(surface),
state_path=state_path,
lock_path=lock_path,
track_liveness=track_liveness,
), None
return lease, None
def release_active_session(lease: ActiveSessionLease) -> None:
@@ -673,25 +605,16 @@ def release_active_session(lease: ActiveSessionLease) -> None:
with _FileLock(lock_path):
if lease.released:
return
try:
raw_entries = _read_entries(state_path, strict=True)
entries = _prune_dead(raw_entries, strict=lease.track_liveness)
except ActiveSessionRegistryError:
if lease.track_liveness:
raise
logger.warning(
"Active-session registry is unavailable; preserving it while "
"releasing an untracked lease"
)
lease.released = True
return
kept = [
entry
for entry in entries
if str(entry.get("lease_id") or "") != lease.lease_id
]
if len(kept) != len(entries):
_write_entries(state_path, kept)
loaded = _read_live_entries(
state_path, track_liveness=lease.track_liveness,
warn="Active-session registry is unavailable; preserving it while "
"releasing an untracked lease",
)
if loaded is not None:
entries = loaded[1]
kept = _without_lease(entries, lease.lease_id)
if len(kept) != len(entries):
_write_entries(state_path, kept)
lease.released = True
@@ -717,17 +640,14 @@ def transfer_active_session(
# thread acquired the file lock. Never resurrect a durably removed lease.
if lease.released:
return False
try:
raw_entries = _read_entries(state_path, strict=True)
entries = _prune_dead(raw_entries, strict=lease.track_liveness)
except ActiveSessionRegistryError:
if lease.track_liveness:
raise
logger.warning(
"Active-session registry is unavailable; refusing to overwrite "
"it during lease transfer"
)
loaded = _read_live_entries(
state_path, track_liveness=lease.track_liveness,
warn="Active-session registry is unavailable; refusing to overwrite "
"it during lease transfer",
)
if loaded is None:
return False
entries = loaded[1]
updated = False
for entry in entries:
if str(entry.get("lease_id") or "") != lease.lease_id:
@@ -735,9 +655,7 @@ def transfer_active_session(
entry["session_id"] = new_session_id
entry["updated_at"] = time.time()
if metadata:
entry["metadata"] = {
str(k): v for k, v in metadata.items() if isinstance(k, str)
}
entry["metadata"] = _clean_metadata(metadata)
updated = True
break
if not updated and lease.track_liveness:
@@ -760,12 +678,10 @@ def transfer_active_session(
def release_orphaned_leases(live_lease_ids: set[str]) -> int:
"""Drop this process's registry entries that no live session owns.
``_prune_dead`` only reclaims leases whose owning process died. A server
that runs for days (``hermes dashboard`` / ``serve``) never trips that
check, so a lease whose session skipped teardown is held until restart.
The owning process is the only authority on which of its own leases are
real, so it drops the rest itself — exact, with no heartbeat write on the
turn path and no staleness threshold to tune.
``_prune_dead`` only reclaims leases of dead processes; a days-long server
never trips it, so a lease whose session skipped teardown is held until
restart. The owning process is the only authority on its own leases —
exact, no heartbeat on the turn path, no staleness threshold.
"""
pid = os.getpid()
state_path = _state_path()
@@ -774,14 +690,13 @@ def release_orphaned_leases(live_lease_ids: set[str]) -> int:
if not state_path.exists():
return 0
with _FileLock(_lock_path()):
try:
raw_entries = _read_entries(state_path, strict=True)
entries = _prune_dead(raw_entries)
except ActiveSessionRegistryError:
logger.warning(
"Active-session registry is unavailable; skipping orphaned-lease sweep"
)
loaded = _read_live_entries(
state_path, track_liveness=False,
warn="Active-session registry is unavailable; skipping orphaned-lease sweep",
)
if loaded is None:
return 0
entries = loaded[1]
kept = [
entry
for entry in entries
@@ -813,12 +728,9 @@ def active_session_liveness_guard(
*,
registry_home: str | Path | None = None,
) -> Iterator[bool]:
"""Hold the registry lock while reporting whether ``session_id`` is leased.
Keeping the lock across the caller's lifecycle mutation prevents a new
backend from acquiring a lease and reopening the row between the liveness
check and the corresponding ``end_session`` write.
"""
"""Hold the registry lock while reporting whether ``session_id`` is leased,
so no new backend can acquire a lease between the check and the caller's
``end_session`` write."""
target = str(session_id or "")
state_path, lock_path = _lease_paths(registry_home=registry_home)
with _FileLock(lock_path):
@@ -834,12 +746,9 @@ def release_active_session_liveness_guard(
lease: ActiveSessionLease,
session_id: str,
) -> Iterator[bool]:
"""Remove ``lease`` and hold its registry lock through a lifecycle write.
This makes automatic cleanup one atomic ownership decision: the local
runtime disappears, sibling liveness is checked, and the caller may end the
durable row before any new backend can acquire/reopen it.
"""
"""Remove ``lease`` and hold its registry lock through a lifecycle write,
making cleanup one atomic ownership decision (release, check siblings,
end the durable row) before any new backend can reopen it."""
if not lease.enabled or lease.released:
with active_session_liveness_guard(
session_id, registry_home=_registry_home_for_lease(lease)
@@ -850,13 +759,8 @@ def release_active_session_liveness_guard(
target = str(session_id or "")
state_path, lock_path = _lease_paths(lease)
with _FileLock(lock_path):
raw_entries = _read_entries(state_path, strict=True)
entries = _prune_dead(raw_entries, strict=True)
kept = [
entry
for entry in entries
if str(entry.get("lease_id") or "") != lease.lease_id
]
entries = _prune_dead(_read_entries(state_path, strict=True), strict=True)
kept = _without_lease(entries, lease.lease_id)
if len(kept) != len(entries):
_write_entries(state_path, kept)
lease.released = True

File diff suppressed because it is too large Load Diff

View File

@@ -11,40 +11,26 @@ 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 the shape of ``hermes_cli/kanban_specify.py``: lazy aux
client import inside the function, lenient response parse, never
raises on expected failure modes.
* The system prompt sees the *configured* profile roster — names plus
descriptions plus the default fallback. Profiles without a
description are still listed (with a note) so the decomposer can
match on name as a fallback, but the user has an obvious incentive
to describe them.
* ``fanout=false`` collapses to the same effect as ``kanban specify``:
we tighten the body and flip ``triage -> todo`` as a single task,
no children created. This makes ``decompose`` a strict superset of
``specify`` from the user's perspective.
* If the LLM picks an assignee that doesn't exist as a profile, we
rewrite it to the configured ``default_assignee`` (or the default
profile if unset). A child task NEVER ends up with ``assignee=None``.
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``.
"""
from __future__ import annotations
import json
import logging
import os
import re
from dataclasses import dataclass
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 _profile_author as _specify_author
logger = logging.getLogger(__name__)
@@ -136,37 +122,9 @@ class DecomposeOutcome:
new_title: Optional[str] = None
def _truncate(text: str, limit: int) -> str:
if len(text) <= limit:
return text
return text[: limit - 1] + "…"
def _extract_json_blob(raw: str) -> Optional[dict]:
if not raw:
return None
stripped = _FENCE_RE.sub("", raw.strip())
first = stripped.find("{")
last = stripped.rfind("}")
if first == -1 or last == -1 or last <= first:
return None
candidate = stripped[first : last + 1]
try:
val = json.loads(candidate)
except (ValueError, json.JSONDecodeError):
return None
if not isinstance(val, dict):
return None
return val
def _profile_author() -> str:
"""Mirror of ``hermes_cli.kanban._profile_author``."""
return (
os.environ.get("HERMES_PROFILE")
or os.environ.get("USER")
or "decomposer"
)
return _specify_author("decomposer")
def _load_config() -> dict:
@@ -177,31 +135,15 @@ def _load_config() -> dict:
return {}
def _resolve_orchestrator_profile(cfg: dict) -> str:
"""Resolve which profile owns the root/orchestration task after fan-out.
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.
Falls back to the active default profile when ``kanban.orchestrator_profile``
is unset, so a task is never stranded for lack of an orchestrator.
``orchestrator_profile`` owns the root after fan-out; ``default_assignee``
catches children the decomposer can't route.
"""
kanban_cfg = cfg.get("kanban", {}) if isinstance(cfg, dict) else {}
explicit = (kanban_cfg.get("orchestrator_profile") or "").strip()
if explicit:
try:
if profiles_mod.profile_exists(explicit):
return explicit
except Exception:
pass
# Fall back to the active default profile.
try:
return profiles_mod.get_active_profile_name() or "default"
except Exception:
return "default"
def _resolve_default_assignee(cfg: dict) -> str:
"""Resolve which profile catches child tasks the orchestrator can't route."""
kanban_cfg = cfg.get("kanban", {}) if isinstance(cfg, dict) else {}
explicit = (kanban_cfg.get("default_assignee") or "").strip()
explicit = (kanban_cfg.get(key) or "").strip()
if explicit:
try:
if profiles_mod.profile_exists(explicit):
@@ -215,12 +157,8 @@ def _resolve_default_assignee(cfg: dict) -> str:
def _build_roster() -> tuple[list[dict], set[str]]:
"""Return (roster_for_prompt, valid_assignee_names).
Each roster entry is ``{name, description, has_description}``. The
valid-set is used after the LLM responds to rewrite invalid
assignees to the default fallback.
"""
"""``(roster_for_prompt, valid_assignee_names)``; entries are
``{name, description, has_description}``."""
roster: list[dict] = []
valid: set[str] = set()
try:
@@ -255,11 +193,8 @@ def _normalize_assignee_choice(
default_assignee: str,
valid_names: set[str],
) -> str:
"""Return a valid assignee, falling back to ``default_assignee``.
Fan-out children and the single-task fallback should share the same
routing guarantee: promoted work must not be left unassigned.
"""
"""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()
@@ -274,13 +209,9 @@ def decompose_task(
author: Optional[str] = None,
timeout: Optional[int] = None,
) -> DecomposeOutcome:
"""Decompose a triage task into a graph of child tasks.
Returns an outcome describing what happened. Never raises for
expected failure modes (task not in triage, no aux client
configured, API error, malformed response, decomposer returned
fanout=true with empty task list) — those surface via ``ok=False``.
"""
"""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)
if task is None:
@@ -291,8 +222,8 @@ def decompose_task(
)
cfg = _load_config()
orchestrator = _resolve_orchestrator_profile(cfg)
default_assignee = _resolve_default_assignee(cfg)
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()
@@ -312,10 +243,8 @@ def decompose_task(
)
try:
# Route through call_llm so auxiliary.kanban_decomposer.* config
# (provider/model/base_url, extra_body, reasoning_effort, retries)
# all apply — the previous direct client.chat.completions.create()
# path dropped auxiliary.<task>.extra_body entirely (#35566).
# call_llm applies all auxiliary.kanban_decomposer.* config
# (provider/model/base_url, extra_body, reasoning_effort, retries).
resp = call_llm(
task="kanban_decomposer",
messages=[
@@ -337,7 +266,7 @@ def decompose_task(
except Exception:
raw = ""
parsed = _extract_json_blob(raw)
parsed = _extract_json_blob(raw, _FENCE_RE)
if parsed is None:
return DecomposeOutcome(task_id, False, "LLM returned malformed JSON")
@@ -346,10 +275,7 @@ def decompose_task(
if not fanout:
# Fall back to single-task spec promotion (same effect as specify).
new_title = parsed.get("title")
new_body = parsed.get("body")
title_val = new_title.strip() if isinstance(new_title, str) and new_title.strip() else None
body_val = new_body if isinstance(new_body, str) and new_body.strip() else None
title_val, body_val = _title_body(parsed)
assignee_val = None
if not task.assignee:
assignee_val = _normalize_assignee_choice(
@@ -385,8 +311,7 @@ def decompose_task(
task_id, False, "decomposer returned fanout=true with empty tasks list",
)
# Rewrite invalid assignees to the default fallback. Never leave a
# task with assignee=None — the user explicitly does not want that.
# Unknown assignees route to the default; never assignee=None.
children: list[dict] = []
for idx, entry in enumerate(raw_tasks):
if not isinstance(entry, dict):

View File

@@ -1,30 +1,13 @@
"""Kanban diagnostics — structured, actionable distress signals for tasks.
A ``Diagnostic`` is a machine-readable description of something that's wrong
with a kanban task: a hallucinated card id, a spawn crash-loop, a task
stuck blocked for too long, etc. Each one carries:
A ``Diagnostic`` carries a **kind** (canonical code the UI/tests match on), a
**severity**, title/detail text, and **actions** the dashboard renders as
buttons and the CLI as hints. Rules are stateless and read-only over
(task, events, runs, optional graph); callers compute on demand.
* A **kind** (canonical code; UI/tests match on this).
* A **severity** (``warning`` / ``error`` / ``critical``).
* A **title** (one-line human description) and **detail** (longer text).
* A list of **suggested actions** — structured entries the dashboard
turns into buttons and the CLI turns into hints.
Rules run over (task, recent events, recent runs, optional graph context) and
emit diagnostics. They are stateless and read-only — no DB writes. Callers compute
diagnostics on demand (on ``/board`` load, ``/tasks/:id`` fetch, or
``hermes kanban diagnostics``).
Design goals:
* Fixable-on-the-operator's-side signals only (missing config, phantom
ids, crash loop). Not "the provider returned 502 once" — that's a
transient runtime blip, not a diagnostic.
* Recoverable: every diagnostic comes with at least one suggested
recovery action the operator can actually take from the UI.
* Auto-clearing: when the underlying failure mode resolves (a clean
``completed`` event arrives, a spawn succeeds, the task gets
unblocked), the diagnostic stops firing. The audit event trail stays.
Design goals: operator-fixable signals only (not a one-off provider 502);
every diagnostic has at least one recovery action; diagnostics auto-clear
when the failure mode resolves (the audit event trail stays).
"""
from __future__ import annotations
@@ -35,9 +18,7 @@ import json
import time
# Severity rungs, ordered least → most urgent. The UI colors them
# amber (warning), orange (error), red (critical). Sorted outputs put
# critical first so operators see the worst fires at the top.
# Least → most urgent; sorted outputs put critical first.
SEVERITY_ORDER = ("warning", "error", "critical")
@@ -52,24 +33,10 @@ def severity_at_or_above(severity: Optional[str], threshold: Optional[str]) -> b
@dataclass
class DiagnosticAction:
"""A single recovery action attached to a diagnostic.
The ``kind`` determines how both the UI and CLI render it:
* ``reclaim`` / ``reassign`` — POST to the matching /tasks/:id/*
endpoint; dashboard wires into the existing recovery popover.
* ``unblock`` — PATCH status back to ``ready`` (for stuck-blocked
diagnostics).
* ``cli_hint`` — print/copy a shell command (e.g.
``hermes -p <profile> auth``). No HTTP side effect.
* ``open_docs`` — deep-link to the docs URL named in ``payload.url``.
* ``comment`` — nudge the operator to add a comment (for
stuck-blocked tasks that need human input).
``suggested=True`` marks the action as the recommended first step;
the UI highlights it. Multiple actions can be suggested if they're
equally valid.
"""
"""A recovery action. ``kind`` drives rendering: ``reclaim``/``reassign``
POST to /tasks/:id/*; ``unblock`` PATCHes status to ready; ``cli_hint``
shows ``payload.command``; ``open_docs`` links ``payload.url``; ``comment``
nudges the operator. ``suggested=True`` = recommended first step."""
kind: str
label: str
@@ -122,20 +89,10 @@ class Diagnostic:
# ---------------------------------------------------------------------------
def _task_field(task, name, default=None):
"""Read a field from a task regardless of representation.
Callers pass sqlite3.Row (dict-like with [] but no attribute
access), kanban_db.Task dataclasses (attribute access), or plain
dicts (both). This normalises them so rule functions don't have
to branch on type each time.
"""
"""Read a field from a sqlite3.Row, a kanban_db.Task dataclass, or a dict."""
if task is None:
return default
# sqlite Row + plain dicts both support mapping access; Row also
# supports .keys().
try:
# Row raises IndexError if the key isn't a column in the query;
# dicts return default via .get. Handle both.
if hasattr(task, "keys") and name in task.keys():
return task[name]
except Exception:
@@ -169,20 +126,42 @@ def _event_ts(ev) -> int:
return int(t or 0)
def _first_field(task, primary: str, legacy: str, default=None):
"""``task[primary]`` unless it is None, else ``task[legacy]`` (old DB rows)."""
v = _task_field(task, primary, None)
return v if v is not None else _task_field(task, legacy, default)
def _latest_event_ts(events: Iterable[Any], kinds: set[str]) -> int:
"""Max ``created_at`` over events whose kind is in ``kinds`` (0 if none)."""
latest = 0
for ev in events:
if _event_kind(ev) in kinds:
latest = max(latest, _event_ts(ev))
return latest
def _log_hint_action(task_id: str) -> DiagnosticAction:
return DiagnosticAction(
kind="cli_hint",
label=f"Check logs: hermes kanban log {task_id}",
payload={"command": f"hermes kanban log {task_id}"},
suggested=True,
)
def _error_snippet(last_err) -> str:
"""First 500 chars of the error (with ellipsis), or "" when absent."""
err_text = (last_err or "").strip() if last_err else ""
return err_text[:500] + ("…" if len(err_text) > 500 else "") if err_text else ""
def _active_hallucination_events(
events: Iterable[Any],
kind: str,
) -> list[Any]:
"""Return events of ``kind`` that have no ``completed``/``edited``
event *strictly after* them. Walks chronologically: each clean
event resets the accumulator; each matching event gets appended.
Events must be sorted by id (i.e. arrival order); callers pass the
task's full event list which the DB already returns in that order.
"""
# Events arrive sorted by id asc (chronological). Walk once, track
# which hallucination events are still "active" (no clean event
# supersedes them).
"""Events of ``kind`` with no ``completed``/``edited`` event strictly after
them. Requires id-sorted (arrival-order) input, which the DB provides."""
active: list[Any] = []
for ev in events:
k = _event_kind(ev)
@@ -191,9 +170,9 @@ def _active_hallucination_events(
elif k == kind:
active.append(ev)
return active
# Standard always-available actions. Every diagnostic can offer these as
# fallbacks regardless of kind — they're the two baseline recovery
# primitives the kernel supports.
# Baseline recovery primitives every diagnostic can fall back on.
def _generic_recovery_actions(task: Any, *, running: bool) -> list[DiagnosticAction]:
out: list[DiagnosticAction] = []
if running:
@@ -214,22 +193,16 @@ def _generic_recovery_actions(task: Any, *, running: bool) -> list[DiagnosticAct
# Rule implementations
# ---------------------------------------------------------------------------
# Each rule takes (task, events, runs, now_ts, config) and returns
# zero or more Diagnostic instances. ``events`` / ``runs`` are lists of
# kanban_db.Event / kanban_db.Run (or plain dicts matching the same
# shape — for test convenience).
# Each rule: (task, events, runs, now_ts, config) -> list[Diagnostic].
# ``events``/``runs`` are kanban_db rows/dataclasses or same-shaped dicts.
RuleFn = Callable[[Any, list[Any], list[Any], int, dict], list[Diagnostic]]
def _aux_slot_explicit(slot: Any) -> bool:
"""Return True if the auxiliary slot has user-supplied non-default fields.
Defaults from ``DEFAULT_CONFIG`` use ``provider: "auto"`` with empty
model/base_url/api_key — that path falls through to the main model. An
"explicit" config is one where the user actively set a provider (not
"auto"), or supplied a model / base_url / api_key.
"""
"""True if the aux slot was user-configured: provider other than "auto",
or any of model/base_url/api_key set (the default falls through to the
main model)."""
if not isinstance(slot, dict):
return False
provider = str(slot.get("provider") or "").strip().lower()
@@ -242,12 +215,9 @@ def _aux_slot_explicit(slot: Any) -> bool:
def _main_model_visible(raw_config: Any) -> bool:
"""Best-effort check that a main model is configured.
Diagnostics runs in the dashboard process which may not share the CLI's
runtime state, so we read the raw config dict. If we cannot prove the
main model is set, we err on the side of NOT firing the diagnostic.
"""
"""Best-effort "a main model is configured" from the raw config dict (the
dashboard process may not share CLI runtime state). Unprovable => False,
which errs toward NOT firing the diagnostic."""
if not isinstance(raw_config, dict):
return False
model_cfg = raw_config.get("model")
@@ -264,17 +234,9 @@ def _main_model_visible(raw_config: Any) -> bool:
def triage_aux_status(config: Optional[dict]) -> Optional[dict]:
"""Inspect raw config and report whether triage paths look configured.
Returns ``None`` when config context is unavailable (suppress diagnostic
to avoid noisy false positives in tests / low-level callers). Otherwise
returns a dict with:
- ``auto_decompose``: bool — whether the dispatcher auto-runs decompose
- ``decomposer_explicit``: bool — user-supplied decomposer slot
- ``specifier_explicit``: bool — user-supplied specifier slot
- ``main_model_visible``: bool — main model can serve as auto fallback
"""
"""Report whether the triage aux paths look configured: ``{auto_decompose,
decomposer_explicit, specifier_explicit, main_model_visible}``. ``None``
when no config context is present (keeps low-level callers/tests silent)."""
if not isinstance(config, dict):
return None
@@ -285,9 +247,7 @@ def triage_aux_status(config: Optional[dict]) -> Optional[dict]:
aux = config.get("auxiliary")
kanban_cfg = config.get("kanban") if isinstance(config.get("kanban"), dict) else {}
# Have we been handed any config context at all? When neither auxiliary
# nor kanban nor model keys are present, the caller is a low-level test
# passing {} — stay silent.
# No auxiliary/kanban/model keys at all => a low-level caller passing {}.
if (
not isinstance(aux, dict)
and not kanban_cfg
@@ -323,14 +283,8 @@ def _positive_int(value: Any, default: int) -> int:
def _rule_hallucinated_cards(task, events, runs, now, cfg) -> list[Diagnostic]:
"""Blocked-hallucination gate fires: a worker called kanban_complete
with created_cards that didn't exist or weren't created by the
completing profile. Task stayed in its prior state; the operator
needs to decide how to proceed.
Auto-clears when a successful completion (or edit) follows the
blocked event.
"""
"""A worker's kanban_complete named created_cards that don't exist / weren't
its own; the completion was blocked. Clears on a later completion/edit."""
hits = _active_hallucination_events(events, "completion_blocked_hallucination")
if not hits:
return []
@@ -343,13 +297,11 @@ def _rule_hallucinated_cards(task, events, runs, now, cfg) -> list[Diagnostic]:
if pid not in phantom_ids:
phantom_ids.append(pid)
running = _task_field(task, "status") == "running"
actions: list[DiagnosticAction] = []
actions.append(DiagnosticAction(
actions = [DiagnosticAction(
kind="comment",
label="Add a comment explaining what to do",
suggested=False,
))
actions.extend(_generic_recovery_actions(task, running=running))
)] + _generic_recovery_actions(task, running=running)
return [Diagnostic(
kind="hallucinated_cards",
severity="error",
@@ -370,22 +322,12 @@ def _rule_hallucinated_cards(task, events, runs, now, cfg) -> list[Diagnostic]:
def _rule_triage_aux_unavailable(task, events, runs, now, cfg) -> list[Diagnostic]:
"""A triage task cannot leave triage without an auxiliary helper.
With the auto-decompose dispatcher (kanban.auto_decompose, default True),
triage tasks fan out via ``auxiliary.kanban_decomposer`` and fall back to
``auxiliary.triage_specifier`` when the decomposer returns ``fanout=false``.
With auto-decompose off, the user must run ``hermes kanban specify``,
which only needs ``auxiliary.triage_specifier``.
The default slot is ``provider: auto`` → auto-falls back to the main model,
so this rule only fires when:
- the relevant slot is explicitly set to something broken, OR
- the auto fallback has no main model to fall back to.
Config context is required; pass {} from tests to keep the rule silent.
"""
"""A triage task can't leave triage without a usable aux model. With
auto-decompose on the primary slot is ``auxiliary.kanban_decomposer``
(specifier as fallback); off, it is ``auxiliary.triage_specifier``. The
default ``provider: auto`` falls back to the main model, so this fires only
when the slot isn't explicit AND no main model is visible. Requires config
context ({} keeps it silent)."""
if _task_field(task, "status") != "triage":
return []
@@ -421,9 +363,6 @@ def _rule_triage_aux_unavailable(task, events, runs, now, cfg) -> list[Diagnosti
"`hermes kanban specify`, which uses auxiliary.triage_specifier."
)
# The primary slot is usable when either: it was explicitly configured by
# the user, OR the default `provider: auto` can fall back to the main
# model. If both fail, we have a real configuration gap.
if primary_explicit or main_visible:
return []
@@ -482,12 +421,8 @@ def _rule_triage_aux_unavailable(task, events, runs, now, cfg) -> list[Diagnosti
def _rule_prose_phantom_refs(task, events, runs, now, cfg) -> list[Diagnostic]:
"""Advisory prose-scan: the completion summary mentions ``t_<hex>``
ids that don't resolve. Non-blocking; surfaced as a warning only.
Auto-clears when a fresh clean completion arrives AFTER the
suspected event.
"""
"""Advisory: the completion summary mentions ``t_<hex>`` ids that don't
resolve. Warning only; clears on a later clean completion."""
hits = _active_hallucination_events(events, "suspected_hallucinated_references")
if not hits:
return []
@@ -516,32 +451,14 @@ def _rule_prose_phantom_refs(task, events, runs, now, cfg) -> list[Diagnostic]:
def _rule_repeated_failures(task, events, runs, now, cfg) -> list[Diagnostic]:
"""Task's unified ``consecutive_failures`` counter is climbing —
something about this task+profile combo is broken and each retry
fails the same way. Triggers regardless of the specific failure
mode (spawn error, timeout, crash) because operationally they
all look the same: the kernel keeps retrying and the operator
needs to intervene.
"""``consecutive_failures`` >= cfg["failure_threshold"] (legacy key
``spawn_failure_threshold``), regardless of failure mode — the kernel keeps
retrying and the operator must intervene. Runtime callers derive the
threshold from ``kanban.failure_limit`` so it doesn't lag the breaker.
Threshold: cfg["failure_threshold"]. Runtime callers should derive
this from ``kanban.failure_limit`` unless the user explicitly set a
diagnostics threshold, so the signal does not lag behind the
dispatcher's circuit breaker.
Accepts the legacy ``spawn_failure_threshold`` config key for
back-compat.
Terminal statuses are exempt: a done/archived card has nothing left
to retry, so a lingering failure streak is history, not a signal.
(``complete_task`` resets the counter, but a manual done — e.g. a
dashboard drag — ends no run and used to leave the flag stuck.)
A fresh attempt in flight (``running``) is also exempt: retrying a
task should clear the stale failure banner until this attempt also
resolves. Otherwise a card that's actively trying again still shows
"failed Nx", which reads as a current failure. It re-fires if the new
run fails too (status leaves ``running`` with a recorded outcome).
"""
Exempt: done/archived (a manual done ends no run, so the streak is history)
and running (a retry in flight must not read as a current failure; re-fires
if it fails too)."""
if _task_field(task, "status") in ("done", "archived", "running"):
return []
threshold = _positive_int(cfg.get(
@@ -549,26 +466,13 @@ def _rule_repeated_failures(task, events, runs, now, cfg) -> list[Diagnostic]:
cfg.get("spawn_failure_threshold", 3),
), 3)
failure_limit = _positive_int(cfg.get("failure_limit"), threshold)
# Read the new unified counter name, with a fallback to the legacy
# column name so this rule keeps working against old DB rows the
# caller somehow materialised without running the migration.
failures = (
_task_field(task, "consecutive_failures", None)
if _task_field(task, "consecutive_failures", None) is not None
else _task_field(task, "spawn_failures", 0)
)
failures = _first_field(task, "consecutive_failures", "spawn_failures", 0)
if failures is None or failures < threshold:
return []
last_err = (
_task_field(task, "last_failure_error", None)
if _task_field(task, "last_failure_error", None) is not None
else _task_field(task, "last_spawn_error", None)
)
last_err = _first_field(task, "last_failure_error", "last_spawn_error")
assignee = _task_field(task, "assignee")
# Classify the most recent failure by peeking at run outcomes so
# the title + suggested action can be specific without a separate
# per-outcome rule.
# Most recent failure outcome makes the title/action specific.
ordered_runs = sorted(runs, key=lambda r: _task_field(r, "id", 0))
most_recent_outcome = None
for r in reversed(ordered_runs):
@@ -596,19 +500,13 @@ def _rule_repeated_failures(task, events, runs, now, cfg) -> list[Diagnostic]:
# to diagnose; reclaim/reassign are the recovery levers.
task_id = _task_field(task, "id")
if task_id:
actions.append(DiagnosticAction(
kind="cli_hint",
label=f"Check logs: hermes kanban log {task_id}",
payload={"command": f"hermes kanban log {task_id}"},
suggested=True,
))
actions.append(_log_hint_action(task_id))
actions.extend(_generic_recovery_actions(
task, running=_task_field(task, "status") == "running",
))
severity = "critical" if failures >= threshold * 2 else "error"
err_text = (last_err or "").strip() if last_err else ""
err_snippet = err_text[:500] + ("…" if len(err_text) > 500 else "") if err_text else ""
err_snippet = _error_snippet(last_err)
outcome_label = {
"spawn_failed": "spawn",
"timed_out": "timeout",
@@ -651,29 +549,14 @@ def _rule_repeated_failures(task, events, runs, now, cfg) -> list[Diagnostic]:
def _rule_repeated_crashes(task, events, runs, now, cfg) -> list[Diagnostic]:
"""The worker spawns fine but keeps crashing mid-run. Check the last
N runs' outcomes; N consecutive ``crashed`` without a successful
``completed`` means something about the task + profile combo is
broken (OOM, missing dependency, tool it needs is down).
"""Trailing run outcomes show >= cfg["crash_threshold"] (default 2)
consecutive ``crashed`` with no ``completed``/``reclaimed`` between. Fires
earlier than ``repeated_failures`` for a crash-specific heads-up and
suppresses itself when the unified rule is about to fire.
Threshold: cfg["crash_threshold"] (default 2).
Narrower than ``repeated_failures`` — fires earlier (2 crashes vs 3
total failures) so the operator gets a crash-specific heads-up
before the unified rule kicks in. Suppresses itself when the
unified rule is also about to fire, to avoid double-flagging.
Terminal statuses are exempt for the same reason as
``repeated_failures`` — with one extra wrinkle: this rule reads run
history, and a manual done (dashboard drag) appends no ``completed``
run to break the crash streak, so the flag was permanent (#kanban
desktop dogfood). Done means done.
``running`` is exempt too: a fresh attempt is in flight, and its
in-flight run (no outcome yet) doesn't break the trailing crash scan,
so a retried card kept showing "crashed Nx" over an active run. The
banner re-fires if the new attempt also crashes.
"""
Exempt: done/archived (a manual done appends no completed run, so the
streak would be permanent) and running (an in-flight run has no outcome
and wouldn't break the scan)."""
if _task_field(task, "status") in ("done", "archived", "running"):
return []
failure_threshold = int(cfg.get(
@@ -702,29 +585,19 @@ def _rule_repeated_crashes(task, events, runs, now, cfg) -> list[Diagnostic]:
# A success (or manual reclaim) breaks the streak.
break
else:
# Other outcomes (timed_out, blocked, spawn_failed, gave_up)
# aren't crash signals — don't count them, but they also
# don't break the crash streak.
# Other outcomes neither count as crashes nor break the streak.
continue
if consecutive < threshold:
return []
task_id = _task_field(task, "id")
actions: list[DiagnosticAction] = []
if task_id:
actions.append(DiagnosticAction(
kind="cli_hint",
label=f"Check logs: hermes kanban log {task_id}",
payload={"command": f"hermes kanban log {task_id}"},
suggested=True,
))
actions.append(_log_hint_action(task_id))
running = _task_field(task, "status") == "running"
actions.extend(_generic_recovery_actions(task, running=running))
severity = "critical" if consecutive >= threshold * 2 else "error"
# Put the actual error up-front so operators see WHAT broke without
# having to open the logs. Truncate defensively — these can be huge
# (full tracebacks).
err_text = (last_err or "").strip() if last_err else ""
err_snippet = err_text[:500] + ("…" if len(err_text) > 500 else "") if err_text else ""
# Error up-front so operators see WHAT broke without opening the logs.
err_snippet = _error_snippet(last_err)
if err_snippet:
title = f"Agent crashed {consecutive}x: {err_snippet.splitlines()[0][:160]}"
detail = (
@@ -751,14 +624,9 @@ def _rule_repeated_crashes(task, events, runs, now, cfg) -> list[Diagnostic]:
def _rule_review_dependency_deadlock(task, events, runs, now, cfg) -> list[Diagnostic]:
"""Detect a legacy review handoff that starves downstream children.
Older workers were instructed to sticky-block an implementation with a
``review-required:`` reason. A separately modelled reviewer child cannot
promote until that parent is terminal, so the lane has no autonomous next
step. This compatibility diagnostic is graph-aware but deliberately leaves
both the dependency graph and the user's sticky block unchanged.
"""
"""Legacy review handoff starving children: the implementation is
sticky-blocked with a ``review-required:`` reason while todo children wait
for it to be terminal. Graph-aware; deliberately mutates nothing."""
if _task_field(task, "status") != "blocked":
return []
@@ -829,21 +697,13 @@ def _rule_review_dependency_deadlock(task, events, runs, now, cfg) -> list[Diagn
def _rule_stuck_in_blocked(task, events, runs, now, cfg) -> list[Diagnostic]:
"""Task has been in ``blocked`` status for too long without a comment.
Threshold: cfg["blocked_stale_hours"] (default 24).
Surfaced as a warning so humans know there's a pending unblock.
"""
"""Blocked for >= cfg["blocked_stale_hours"] (default 24) with no comment
or unblock since the last ``blocked`` event."""
hours = float(cfg.get("blocked_stale_hours", 24))
status = _task_field(task, "status")
if status != "blocked":
return []
# Find the most recent ``blocked`` event.
last_blocked_ts = 0
for ev in events:
if _event_kind(ev) == "blocked":
t = _event_ts(ev)
last_blocked_ts = max(last_blocked_ts, t)
last_blocked_ts = _latest_event_ts(events, {"blocked"})
if last_blocked_ts == 0:
return []
age_hours = (now - last_blocked_ts) / 3600.0
@@ -879,29 +739,17 @@ def _rule_stuck_in_blocked(task, events, runs, now, cfg) -> list[Diagnostic]:
def _rule_block_unblock_cycling(task, events, runs, now, cfg) -> list[Diagnostic]:
"""Task has cycled through blocked → unblocked many times — the
``unblock`` is not fixing the underlying problem and the worker
keeps re-blocking for substantially the same reason.
``_rule_stuck_in_blocked`` resets its timer on any ``commented`` /
``unblocked`` event, so a task that cycles every few minutes is
invisible to it regardless of how many times it cycles (#29747
gap 1). This rule complements that one by counting block→unblock
cycles in a sliding window.
Threshold: cfg["block_cycle_threshold"] (default 3) cycles within
cfg["block_cycle_window_seconds"] (default 24h).
"""
""">= cfg["block_cycle_threshold"] (default 3) blocked-after-unblocked
cycles within cfg["block_cycle_window_seconds"] (default 24h). Complements
``_rule_stuck_in_blocked``, whose timer any unblock resets, so fast cyclers
are invisible to it."""
threshold = _positive_int(cfg.get("block_cycle_threshold"), 3)
window_seconds = float(cfg.get("block_cycle_window_seconds", 24 * 3600))
cycle_cutoff = now - window_seconds
# Walk events chronologically (arrival order — callers pre-sort by
# id, which is the canonical chronological order; ``created_at``
# alone is insufficient because multiple events can share the same
# second). Count "blocked after unblocked" transitions: every time
# a blocked event follows at least one unblocked event since the
# last cycle was counted, that's a new cycle.
# Walk in id (arrival) order — created_at alone can't order events that
# share a second. A blocked event after >= 1 unblocked since the last
# counted cycle is a new cycle.
cycles = 0
seen_unblock_since_last_cycle = False
initial_blocked_ts = 0
@@ -956,66 +804,31 @@ def _rule_block_unblock_cycling(task, events, runs, now, cfg) -> list[Diagnostic
def _rule_stranded_in_ready(task, events, runs, now, cfg) -> list[Diagnostic]:
"""Task has been in ``ready`` status for too long without any worker
claiming it.
Threshold: cfg["stranded_threshold_seconds"] (default 1800 = 30 min).
Catches every "task waiting for a worker that never comes" case
without caring WHY:
* Operator typo'd the assignee — no profile or external worker matches.
* Profile was deleted, leaving its tasks stranded.
* External worker pool (Codex CLI, Claude Code lane, custom daemon)
is down, hung, or wasn't started.
* Dispatcher is misconfigured (wrong board, wrong HERMES_HOME).
Pre-rule, all of these silently rotted in ``skipped_nonspawnable`` —
the dispatcher correctly skipped them (good — no respawn loop) but
nobody surfaced the fact that operator-actionable work was
accumulating. The rule fires when a ready task's promoted-to-ready
timestamp is older than the threshold AND the assignee is non-empty
(truly unassigned tasks have their own ``skipped_unassigned`` signal
on the dispatcher and a different operator response).
The signal is age-based on purpose: it's identity-agnostic, so it
works for Hermes profiles, registered lanes, external workers, and
typos uniformly. No registry to curate, no per-board allowlist.
"""
"""Assigned, unclaimed, ``ready`` for >= cfg["stranded_threshold_seconds"]
(default 30 min). Deliberately age-based and identity-agnostic so it
catches typo'd assignees, deleted profiles, and down external worker
pools alike without a registry to curate. Unassigned tasks are excluded —
the dispatcher's ``skipped_unassigned`` already covers them."""
threshold_seconds = float(
cfg.get("stranded_threshold_seconds", 30 * 60)
)
status = _task_field(task, "status")
if status != "ready":
return []
# Skip tasks with a live claim — they're being worked on, even if
# the worker hasn't reported progress yet (run-level liveness
# extends the claim TTL; we don't want to second-guess that here).
# A live claim means it's being worked on even without progress yet.
if _task_field(task, "claim_lock"):
return []
assignee = _task_field(task, "assignee") or ""
if not assignee.strip():
# Unassigned tasks: the dispatcher's ``skipped_unassigned`` is
# already the right signal. A separate diagnostic here would
# double-flag the same condition.
return []
# Find the most recent event that put this task into ready.
# ``created`` covers tasks born ready; ``promoted`` covers parent-
# done auto-promotion; ``reclaimed`` covers TTL/crash recovery;
# ``unblocked`` covers human-driven resumes.
READY_TRANSITION_KINDS = {
"created", "promoted", "reclaimed", "unblocked",
}
last_ready_ts = 0
for ev in events:
if _event_kind(ev) in READY_TRANSITION_KINDS:
t = _event_ts(ev)
last_ready_ts = max(last_ready_ts, t)
# Most recent event that put the task into ready.
last_ready_ts = _latest_event_ts(
events, {"created", "promoted", "reclaimed", "unblocked"},
)
# Fallback: if no qualifying event exists (very old task or events
# truncated), fall back to ``created_at`` on the task row. Better
# to occasionally over-flag an ancient task than miss a stranded one.
# No qualifying event (old task / truncated events): fall back to
# created_at — over-flagging an ancient task beats missing a stranded one.
if last_ready_ts == 0:
last_ready_ts = int(_task_field(task, "created_at", default=0) or 0)
if last_ready_ts == 0:
@@ -1031,9 +844,7 @@ def _rule_stranded_in_ready(task, events, runs, now, cfg) -> list[Diagnostic]:
else:
age_str = f"{int(age_seconds / 60)}m"
# Severity escalates with age. Below 2x threshold = warning;
# 2x – 6x = error; beyond 6x = critical (something is clearly
# broken, not just slow).
# Escalate with age: <2x threshold warning, 2x-6x error, >6x critical.
if age_seconds >= threshold_seconds * 6:
severity = "critical"
elif age_seconds >= threshold_seconds * 2:
@@ -1078,8 +889,7 @@ def _rule_stranded_in_ready(task, events, runs, now, cfg) -> list[Diagnostic]:
)]
# Registry — order matters: rules higher on the list render first when
# severity ties. Add new rules here.
# Order matters: earlier rules render first on severity ties.
_RULES: list[RuleFn] = [
_rule_hallucinated_cards,
_rule_triage_aux_unavailable,
@@ -1093,21 +903,6 @@ _RULES: list[RuleFn] = [
]
# Known kinds (for the UI's filter / legend / i18n keys). Update when
# rules are added.
DIAGNOSTIC_KINDS = (
"hallucinated_cards",
"triage_aux_unavailable",
"prose_phantom_refs",
"repeated_failures",
"repeated_crashes",
"review_dependency_deadlock",
"stuck_in_blocked",
"block_unblock_cycling",
"stranded_in_ready",
)
DEFAULT_CONFIG = {
# Match the dispatcher default (kanban.failure_limit) so repeated-failure
# diagnostics do not lag behind the default auto-block threshold.
@@ -1116,21 +911,16 @@ DEFAULT_CONFIG = {
"spawn_failure_threshold": 2,
"crash_threshold": 2,
"blocked_stale_hours": 24,
# Stranded-task threshold. 30 min by default — below that, the
# signal is dominated by tasks that are about to be claimed on the
# next dispatcher tick (default 60s) and would just be noise.
# Below 30 min the signal is dominated by tasks about to be claimed on
# the next dispatcher tick.
"stranded_threshold_seconds": 30 * 60,
}
def config_from_kanban_config(kanban_cfg: Optional[dict]) -> dict:
"""Build diagnostics config from the runtime ``kanban`` config section.
``kanban.diagnostics.failure_threshold`` remains an explicit override.
Otherwise, derive the repeated-failure threshold from
``kanban.failure_limit`` so CLI/dashboard diagnostics match the
dispatcher's actual circuit-breaker threshold.
"""
"""Diagnostics config from the ``kanban`` section. ``kanban.diagnostics.
failure_threshold`` is an explicit override; otherwise the threshold is
``kanban.failure_limit`` so diagnostics match the dispatcher's breaker."""
kanban_cfg = kanban_cfg or {}
diag_cfg = dict(kanban_cfg.get("diagnostics") or {})
diag_cfg.setdefault(
@@ -1146,13 +936,9 @@ def config_from_kanban_config(kanban_cfg: Optional[dict]) -> dict:
def config_from_runtime_config(raw_config: Optional[dict]) -> dict:
"""Build diagnostics config from the full Hermes runtime config.
Carries through ``kanban``, ``auxiliary``, and ``model`` keys so triage-
aware rules can inspect the active aux-helper and main-model state.
Folds the ``kanban`` block through ``config_from_kanban_config`` so the
repeated-failure threshold derivation still applies.
"""
"""Diagnostics config from the full runtime config: folds ``kanban`` through
``config_from_kanban_config`` and carries ``kanban``/``auxiliary``/``model``
through for the triage-aware rules."""
raw_config = raw_config or {}
if not isinstance(raw_config, dict):
return {}
@@ -1177,12 +963,8 @@ def compute_task_diagnostics(
config: Optional[dict] = None,
graph: Optional[dict] = None,
) -> list[Diagnostic]:
"""Run every rule against a single task's state and return a
severity-sorted list of active diagnostics.
Sorting: critical first, then error, then warning; ties broken by
most-recent ``last_seen_at``.
"""
"""Run every rule for one task; critical first, then error, warning; ties
broken by most-recent ``last_seen_at``."""
now_ts = int(now if now is not None else time.time())
config = config or {}
cfg = {**DEFAULT_CONFIG, **config}
@@ -1202,9 +984,7 @@ def compute_task_diagnostics(
try:
out.extend(rule(task, events, runs, now_ts, cfg))
except Exception:
# A broken rule must never crash the dashboard. Rule bugs
# get caught in tests; in production we'd rather drop the
# diagnostic than 500 a whole /board request.
# A broken rule must never 500 a whole /board request.
continue
severity_idx = {s: i for i, s in enumerate(SEVERITY_ORDER)}
out.sort(

View File

@@ -1,32 +1,13 @@
"""Kanban triage specifier — flesh out a one-liner into a real spec.
Used by ``hermes kanban specify [task_id | --all]``. Takes a task that
lives in the Triage column (a rough idea, typically only a title), calls
the auxiliary LLM to produce:
``hermes kanban specify [task_id | --all]`` asks the auxiliary LLM for a
tightened title + concrete body for a Triage task, then flips it
``triage -> todo`` via ``kanban_db.specify_triage_task``.
* A tightened title (optional — only replaces if the model proposes a
materially different one)
* A concrete body: goal, proposed approach, acceptance criteria
and then flips the task ``triage -> todo`` via
``kanban_db.specify_triage_task``. The dispatcher promotes it to
``ready`` on its next tick (or immediately if there are no open parents).
Design notes
------------
* This module intentionally mirrors ``hermes_cli/goals.py`` — same aux
client pattern, same "empty config => skip, don't crash" tolerance.
Keeps the surface area tiny and the failure modes predictable.
* The prompt is a short system + user pair. We ask for JSON with
``{title, body}``; if parsing fails, we fall back to treating the
whole response as the body and leave the title untouched. No
retry loop — one shot, keep cost bounded.
* Structured output / JSON mode is not requested explicitly so the
specifier works on providers that don't implement it. The parse
is lenient (tolerates markdown code fences around the JSON).
Mirrors ``hermes_cli/goals.py``: same aux-client pattern, same "empty config
=> skip, don't crash" tolerance. One shot, no retry loop. JSON mode is not
requested (works on providers without it); the parse is lenient and falls
back to "whole reply is the body" so a malformed reply never strands a task.
"""
from __future__ import annotations
@@ -108,13 +89,12 @@ def _truncate(text: str, limit: int) -> str:
_FENCE_RE = re.compile(r"^\s*```(?:json)?\s*|\s*```\s*$", re.IGNORECASE)
def _extract_json_blob(raw: str) -> Optional[dict]:
"""Lenient JSON extraction — tolerates fenced code blocks and
leading/trailing whitespace. Returns None if nothing parses."""
def _extract_json_blob(raw: str, fence_re: re.Pattern = _FENCE_RE) -> Optional[dict]:
"""Lenient JSON object extraction: strip code fences, take the first ``{``
to the last ``}``. None if nothing parses to a dict."""
if not raw:
return None
stripped = _FENCE_RE.sub("", raw.strip())
# Greedy: find the first `{` and last `}` and try that slice.
stripped = fence_re.sub("", raw.strip())
first = stripped.find("{")
last = stripped.rfind("}")
if first == -1 or last == -1 or last <= first:
@@ -124,19 +104,24 @@ def _extract_json_blob(raw: str) -> Optional[dict]:
val = json.loads(candidate)
except (ValueError, json.JSONDecodeError):
return None
if not isinstance(val, dict):
return None
return val
return val if isinstance(val, dict) else None
def _profile_author() -> str:
def _nonblank(v) -> Optional[str]:
return v if isinstance(v, str) and v.strip() else None
def _title_body(parsed: dict) -> tuple[Optional[str], Optional[str]]:
"""``(title, body)`` from an LLM reply: title stripped, body verbatim,
either None when missing/blank."""
title = _nonblank(parsed.get("title"))
return (title.strip() if title else None), _nonblank(parsed.get("body"))
def _profile_author(default: str = "specifier") -> str:
"""Mirror of ``hermes_cli.kanban._profile_author``. Kept local to
avoid a circular import when kanban.py imports this module."""
return (
os.environ.get("HERMES_PROFILE")
or os.environ.get("USER")
or "specifier"
)
return os.environ.get("HERMES_PROFILE") or os.environ.get("USER") or default
def specify_task(
@@ -145,13 +130,9 @@ def specify_task(
author: Optional[str] = None,
timeout: Optional[int] = None,
) -> SpecifyOutcome:
"""Specify a single triage task and promote it to ``todo``.
Returns an outcome describing what happened. Never raises for expected
failure modes (task not in triage, no aux client configured, API
error, malformed response) — those surface via ``ok=False`` so the
``--all`` sweep can continue past individual failures.
"""
"""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)
if task is None:
@@ -174,9 +155,8 @@ def specify_task(
)
try:
# Route through call_llm so auxiliary.triage_specifier.* config
# (provider/model/base_url, extra_body, reasoning_effort, retries)
# all apply — the direct-create path dropped extra_body (#35566).
# call_llm applies all auxiliary.triage_specifier.* config
# (provider/model/base_url, extra_body, reasoning_effort, retries).
resp = call_llm(
task="triage_specifier",
messages=[
@@ -206,9 +186,7 @@ def specify_task(
new_title: Optional[str]
new_body: Optional[str]
if parsed is None:
# Fall back: treat the whole reply as the body, leave title as-is.
# Worst case the user edits afterward — still better than stranding
# the task in triage on a malformed LLM reply.
# Whole reply becomes the body; the user can edit afterward.
stripped_raw = raw.strip()
if not stripped_raw:
return SpecifyOutcome(
@@ -217,16 +195,7 @@ def specify_task(
new_title = None
new_body = stripped_raw
else:
title_val = parsed.get("title")
body_val = parsed.get("body")
new_title = (
title_val.strip()
if isinstance(title_val, str) and title_val.strip()
else None
)
new_body = (
body_val if isinstance(body_val, str) and body_val.strip() else None
)
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"
@@ -241,8 +210,7 @@ def specify_task(
author=author or _profile_author(),
)
if not ok:
# Race: someone else promoted / archived the task between our
# read above and the write. Report, don't crash.
# Race: promoted/archived between our read and the write.
return SpecifyOutcome(
task_id, False, "task moved out of triage before promotion"
)
@@ -250,10 +218,7 @@ def specify_task(
def list_triage_ids(*, tenant: Optional[str] = None) -> list[str]:
"""Return task ids currently in the triage column.
``tenant`` narrows the sweep; ``None`` returns every triage task.
"""
"""Task ids in the triage column; ``tenant`` narrows the sweep."""
with kb.connect_closing() as conn:
tasks = kb.list_tasks(
conn,

View File

@@ -19,6 +19,7 @@ from __future__ import annotations
from dataclasses import dataclass, field
import json
import sqlite3
import time
from typing import Any, Iterable, Optional
from hermes_cli import kanban_db as kb
@@ -83,16 +84,12 @@ def _activate_root_inline(
) -> bool:
"""Inline blocked→done CAS flip + event insert for the swarm root.
Runs INSIDE create_swarm's outer write_txn, so it must not call
``kb.complete_task`` — that helper opens its own transaction and fires
post-commit side effects (workspace cleanup, failure-counter clear,
``recompute_ready``) that would execute while the outer transaction can
still roll back. Instead we do the minimal durable writes here and let
the caller run ``recompute_ready`` after the outer commit.
Runs INSIDE create_swarm's write_txn, so it must not call
``kb.complete_task`` (own transaction + post-commit side effects that
would run while the outer txn can still roll back). The caller runs
``recompute_ready`` after the outer commit.
"""
import time as _time
now = int(_time.time())
now = int(time.time())
cur = conn.execute(
"""
UPDATE tasks
@@ -179,9 +176,8 @@ def create_swarm(
raise RuntimeError("could not activate the completed swarm topology")
activated = True
if activated:
# Outside the outer transaction: promote the root's children now
# that its 'done' flip is durable (recompute_ready opens its own
# txn and must never run under an open write_txn).
# After commit: recompute_ready opens its own txn and must never run
# under an open write_txn.
kb.recompute_ready(conn)
root = kb.get_task(conn, created.root_id)
run = kb.latest_run(conn, created.root_id)
@@ -249,9 +245,8 @@ def _create_swarm_uncommitted(
workspace_path=workspace_path,
)
# If idempotency returned an existing non-archived root, do not duplicate the
# swarm graph. Recover the topology from the root's latest blackboard, if it
# was created by this helper previously.
# Idempotency may return an existing root: recover its topology from the
# blackboard instead of duplicating the graph.
existing = latest_blackboard(conn, root).get("topology")
if isinstance(existing, dict):
worker_ids = [str(x) for x in existing.get("worker_ids", []) if x]
@@ -266,61 +261,54 @@ def _create_swarm_uncommitted(
)
context_suffix = _swarm_context(root, goal)
worker_ids: list[str] = []
for spec in worker_specs:
worker_id = kb.create_task(
common = dict(
created_by=created_by, tenant=tenant,
workspace_kind=workspace_kind, workspace_path=workspace_path,
)
worker_ids = [
kb.create_task(
conn,
title=spec.title,
body=(spec.body or "") + context_suffix,
assignee=spec.profile,
created_by=created_by,
parents=[root],
tenant=tenant,
priority=spec.priority or priority,
workspace_kind=workspace_kind,
workspace_path=workspace_path,
skills=spec.skills or None,
max_runtime_seconds=spec.max_runtime_seconds,
**common,
)
worker_ids.append(worker_id)
for spec in worker_specs
]
verifier_body = (
"Review every worker handoff and blackboard update. Gate the swarm: "
"complete only with metadata {\"gate\": \"pass\"} when evidence is "
"sufficient; otherwise block with exact missing work."
+ context_suffix
)
verifier = kb.create_task(
conn,
title=verifier_title,
body=verifier_body,
body=(
"Review every worker handoff and blackboard update. Gate the swarm: "
"complete only with metadata {\"gate\": \"pass\"} when evidence is "
"sufficient; otherwise block with exact missing work."
+ context_suffix
),
assignee=verifier_assignee,
created_by=created_by,
parents=worker_ids,
tenant=tenant,
priority=priority,
workspace_kind=workspace_kind,
workspace_path=workspace_path,
skills=["requesting-code-review"],
**common,
)
synthesizer_body = (
"Synthesize the verified worker outputs into the final deliverable. "
"Do not start until the verifier has passed the gate."
+ context_suffix
)
synthesizer = kb.create_task(
conn,
title=synthesizer_title,
body=synthesizer_body,
body=(
"Synthesize the verified worker outputs into the final deliverable. "
"Do not start until the verifier has passed the gate."
+ context_suffix
),
assignee=synthesizer_assignee,
created_by=created_by,
parents=[verifier],
tenant=tenant,
priority=priority,
workspace_kind=workspace_kind,
workspace_path=workspace_path,
skills=["humanizer"],
**common,
)
created = SwarmCreated(root, worker_ids, verifier, synthesizer)

View File

@@ -14,30 +14,21 @@ source board's slug)::
attachments/<task>/… attachment blobs (unless --no-attachments)
logs/<task>.log worker logs (only with --include-logs)
Two things make this more than a ``tar czf`` of the board directory.
Two things make this more than ``tar czf`` of the board directory:
**The database is live.** Kanban runs in WAL mode and a dispatcher may be
mid-write, so copying ``kanban.db`` off the filesystem yields a torn
snapshot that is missing whatever still sits in the ``-wal`` file. Export
goes through SQLite's online-backup API instead, which produces a
consistent single-file image of a database that is being written to.
* **The database is live** (WAL mode, dispatcher may be mid-write), so the
export uses SQLite's online-backup API for a consistent image instead of a
file copy that would miss the ``-wal`` sidecar.
* **Rows carry machine-local state** — claims, PIDs, heartbeats, absolute
paths, gateway chat subscriptions, session ids. Shipping them verbatim
would import claims owned by a stranger's process or push events into a
stranger's Telegram thread. Scrubbed on export and re-scrubbed on import
(an archive is untrusted); see :func:`_scrub_local_state` and
:func:`_relocate_imported_rows`.
**Rows carry machine-local state.** Claims, PIDs, heartbeats, absolute
workspace and attachment paths, gateway chat subscriptions, and session
ids are all meaningful only on the machine that wrote them. Shipping them
verbatim is how an imported board arrives holding claims owned by a
process on somebody else's laptop, or starts pushing task events into a
stranger's Telegram thread. Everything machine-local is scrubbed on the
export side (so the archive itself never carries it) and defensively
re-scrubbed on import; see :func:`_scrub_local_state` and
:func:`_relocate_imported_rows`.
Imports always land as a **new** board — the slug auto-suffixes on
collision — so an import can never mutate a board that is already there.
That also means an imported board is never ``default``, which is what
lets the import side ignore the default board's split on-disk layout
(``<root>/kanban.db`` beside ``<root>/kanban/attachments/``) and put
everything inside one ``boards/<slug>/`` directory.
Imports always land as a **new** board (slug auto-suffixes on collision), so
an import never mutates an existing board and is never ``default`` — which
lets the importer ignore the default board's split on-disk layout.
"""
from __future__ import annotations
@@ -74,13 +65,8 @@ _DISPATCHABLE_STATUSES = ("ready", "running", "todo", "scheduled")
# ---------------------------------------------------------------------------
def _snapshot_db(source: Path, target: Path) -> None:
"""Write a consistent copy of ``source`` to ``target``.
Uses SQLite's online-backup API rather than a file copy: in WAL mode
a just-committed page can still live in the ``-wal`` sidecar, so
copying only ``kanban.db`` loses recent writes and can produce a
torn image if the dispatcher commits mid-copy.
"""
"""Consistent copy of ``source`` via the online-backup API (a file copy
would miss pages still in the ``-wal`` sidecar and could tear)."""
src = sqlite3.connect(str(source))
try:
dst = sqlite3.connect(str(target))
@@ -93,14 +79,9 @@ def _snapshot_db(source: Path, target: Path) -> None:
def _scrub_local_state(conn: sqlite3.Connection) -> None:
"""Strip machine-local runtime state. Caller owns the transaction.
Runs on the export side so the archive itself never carries another
machine's claims, PIDs, or — the one that actually matters for a
board shared with someone else — the gateway chat ids subscribed to
its task events. Repeated on import because an archive is untrusted
input.
"""
"""Strip machine-local runtime state (claims, PIDs, and above all the
gateway chat ids subscribed to task events). Caller owns the transaction.
Run on export and again on import (an archive is untrusted input)."""
conn.execute("DELETE FROM kanban_notify_subs")
conn.execute(
"""
@@ -157,12 +138,9 @@ def export_board(
include_attachments: bool = True,
include_logs: bool = False,
) -> dict[str, Any]:
"""Export ``board`` to a ``tar.gz`` archive. Returns a summary dict.
``output_path`` may be given with or without the ``.tar.gz`` suffix.
Workspaces are never included: they are git worktrees and scratch
trees that are large, machine-local, and rebuilt on demand.
"""
"""Export ``board`` to a ``tar.gz`` (suffix optional on ``output_path``);
returns a summary dict. Workspaces are never included — large,
machine-local, rebuilt on demand."""
slug = kb._normalize_board_slug(board) or kb.get_current_board()
if not kb.board_exists(slug):
raise ValueError(f"board {slug!r} does not exist")
@@ -241,12 +219,8 @@ def export_board(
# ---------------------------------------------------------------------------
def _available_slug(preferred: str) -> str:
"""Return ``preferred``, or the first free ``<preferred>-N`` variant.
``default`` always reports as existing, so an archive exported from a
default board naturally lands as ``default-2`` instead of colliding
with the importer's own default board.
"""
"""``preferred`` or the first free ``<preferred>-N``. ``default`` always
exists, so a default-board export lands as ``default-2``."""
if not kb.board_exists(preferred):
return preferred
# Leave headroom for the suffix inside the 64-char slug limit.
@@ -295,22 +269,17 @@ def _read_board_metadata(path: Path) -> dict[str, Any]:
def _relocate_imported_rows(
conn: sqlite3.Connection, slug: str
) -> tuple[dict[str, int], list[str]]:
"""Re-anchor an imported board's rows to this machine.
"""Re-anchor an imported board's rows to this machine; returns
``(stats, warnings)``.
Returns ``(stats, warnings)``. Three things move:
* Attachment rows are repointed at this board's attachments tree.
Rows whose blob did not travel (an export made with
``--no-attachments``) are dropped, because a row pointing at a file
that does not exist breaks download in every UI that lists it.
* Workspace paths are cleared. ``scratch`` tasks regenerate one under
this board on the next claim, so they are simply reset. ``dir`` and
``worktree`` tasks cannot be resolved without a path that means
something here, so any that are still dispatchable are parked in
``triage`` — otherwise the dispatcher claims them, fails to build a
workspace, and burns them straight into the failure breaker.
* Runtime state is scrubbed again. Export already did this, but an
archive is an untrusted input and the cost is one UPDATE.
* Attachment rows are repointed at this board's tree; rows whose blob
did not travel (``--no-attachments``) are dropped, since a dangling row
breaks download in every UI.
* Workspace paths are cleared. ``scratch`` regenerates on next claim;
dispatchable ``dir``/``worktree`` tasks are parked in ``triage``,
otherwise the dispatcher claims them, fails to build a workspace, and
burns them into the failure breaker.
* Runtime state is scrubbed again (untrusted input, one UPDATE).
"""
warnings: list[str] = []
now = int(time.time())
@@ -389,12 +358,8 @@ def import_board(
*,
activate: bool = False,
) -> dict[str, Any]:
"""Import a board archive as a new board. Returns a summary dict.
``slug`` overrides the name from the archive. Either way the final
slug auto-suffixes if it is taken, so an import never merges into or
overwrites an existing board.
"""
"""Import an archive as a NEW board (``slug`` overrides the archive's;
either way it auto-suffixes if taken). Returns a summary dict."""
archive = Path(archive_path).expanduser()
if not archive.exists():
raise FileNotFoundError(f"archive not found: {archive}")