From 41a26c8ba9724fd6b6083e1fbabc5706033e8433 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 17:04:54 -0700 Subject: [PATCH] refactor(kanban-db): compact long docstrings (rules kept) across origin/dispatch/connect --- hermes_cli/kanban_db.py | 124 +++++++++---------------------- hermes_cli/kanban_db_connect.py | 21 ++---- hermes_cli/kanban_db_dispatch.py | 99 ++++++++---------------- 3 files changed, 71 insertions(+), 173 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index f9fc26a563..d561d1295b 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -541,20 +541,13 @@ def _normalize_board_slug(slug: Optional[str]) -> Optional[str]: def kanban_home() -> Path: - """Return the shared Hermes root that anchors the kanban board. + """Shared Hermes root anchoring the board: ``HERMES_KANBAN_HOME`` if set, + else ``get_default_hermes_root()`` (```` for a + ``/profiles/`` HERMES_HOME, HERMES_HOME itself otherwise). - Resolution order: - - 1. ``HERMES_KANBAN_HOME`` env var when set and non-empty (explicit - override for tests and unusual deployments). - 2. ``get_default_hermes_root()``, which already returns ```` - when ``HERMES_HOME`` is ``/profiles/``, and returns - ``HERMES_HOME`` directly for Docker / custom deployments. - - The kanban board is shared across profiles **by design** (see the - module docstring). Resolving the kanban paths through the active - profile's ``HERMES_HOME`` would silently fork the board per profile, - which breaks the dispatcher / worker handoff. + The board is shared across profiles BY DESIGN; resolving through the + active profile's HERMES_HOME would silently fork it per profile and break + the dispatcher / worker handoff. """ override = os.environ.get("HERMES_KANBAN_HOME", "").strip() if override: @@ -728,23 +721,13 @@ def workspaces_root(board: Optional[str] = None) -> Path: def attachments_root(board: Optional[str] = None) -> Path: - """Return the directory under which task file attachments are stored. + """Per-board attachments root (``HERMES_KANBAN_ATTACHMENTS_ROOT`` wins). - Mirrors :func:`worker_logs_dir` / :func:`workspaces_root`: anchored - per-board so attachments don't leak between projects. Each task gets - its own ``/.../attachments//`` subdirectory. - - ``HERMES_KANBAN_ATTACHMENTS_ROOT`` pins the path directly (highest - precedence) for tests and unusual deployments. - - ``default`` uses ``/kanban/attachments/``; other boards use - ``/kanban/boards//attachments/``. - - Workers (which run with full file-tool access) read attached files - by the absolute path surfaced in :func:`build_worker_context`. On the - local terminal backend — the default for kanban — that path resolves - directly. Remote backends (Docker/Modal) need this directory mounted; - see the kanban docs. + ``default`` -> ``/kanban/attachments/``, other boards -> + ``/kanban/boards//attachments/``; each task gets its own + ``/`` subdir. Workers read attachments by the absolute path + surfaced in :func:`build_worker_context`, so remote terminal backends + (Docker/Modal) need this directory mounted. """ return _board_path("HERMES_KANBAN_ATTACHMENTS_ROOT", board, ("kanban", "attachments"), "attachments") @@ -3986,30 +3969,15 @@ def block_task( ) -> bool: """Transition ``running``/``ready`` → ``blocked`` (or route elsewhere). - ``kind`` (one of :data:`VALID_BLOCK_KINDS`, or ``None`` for a legacy - un-typed block) drives routing instead of every block landing in one - undifferentiated ``blocked`` bucket: - - * ``dependency`` — the task is only waiting on another task. It does NOT - sit in ``blocked`` (where a cron would keep "unblocking" it); it goes to - ``todo`` so the existing parent-gating / ``recompute_ready`` machinery - promotes it automatically once its parents finish. No human, no cron, no - retry storm. This is Dale's "Type 2 — dependency blocked". - - * ``needs_input`` / ``capability`` / ``None`` — "truly blocked" (Dale's - "Type 1"). Lands in ``blocked`` for a human. BUT: each time such a task - is re-blocked for the SAME kind after having been unblocked, the - unblock-loop counter (``block_recurrences``) increments. When it reaches - :data:`BLOCK_RECURRENCE_LIMIT`, the task is routed to ``triage`` instead - of ``blocked`` — breaking the cron-unblock ↔ worker-re-block loop and - forcing a human-in-the-loop triage decision. - - * ``transient`` — treated like a generic block for routing, but a worker - can use it to signal "this might clear on its own"; it still participates - in the loop breaker so a forever-flaky task eventually escalates. - - Returns True on any successful transition (to ``blocked``, ``todo``, or - ``triage``), False when the task wasn't in a blockable state. + ``kind`` (:data:`VALID_BLOCK_KINDS` or ``None`` = legacy un-typed) routes: + ``dependency`` -> ``todo`` (parent gating / ``recompute_ready`` promotes + it; never parked where a cron would keep "unblocking" it); everything + else -> ``blocked`` for a human, EXCEPT that a re-block for the SAME kind + after an unblock bumps ``block_recurrences`` and at + :data:`BLOCK_RECURRENCE_LIMIT` routes to ``triage`` instead, breaking the + cron-unblock ↔ worker-re-block loop. ``transient`` signals "may clear on + its own" but still counts toward the loop breaker so a forever-flaky task + escalates. Returns True on any transition, False when not blockable. """ if kind is not None and kind not in VALID_BLOCK_KINDS: raise ValueError( @@ -4826,24 +4794,12 @@ def decompose_triage_task( its assignee (typically the orchestrator profile) wakes back up to judge completion or spawn more work. - ``children`` is a list of dicts, each shaped like:: - - { - "title": "...", - "body": "...", # optional - "assignee": "profile-name", # optional, None -> default fallback - "parents": [0, 2], # indices into this same children list - } - - Returns the list of created child task ids (in input order) on - success. Returns ``None`` when: - - The root task does not exist - - The root task is not in ``triage`` - - A cycle would result (caller built a bad graph) - - Validation of titles/assignees happens inside the same write_txn as - the inserts so a malformed entry aborts the whole decomposition - cleanly (no orphan children). + ``children``: dicts of ``title`` (required), ``body``, ``assignee`` (None + -> default fallback) and ``parents`` (indices into this same list). + Returns the created child ids in input order, or ``None`` when the root + is missing / not in ``triage`` / the graph has a cycle. Title/assignee + validation runs inside the same write_txn as the inserts so a malformed + entry aborts the whole decomposition (no orphan children). """ if not children: return None @@ -5158,25 +5114,13 @@ def schedule_task( def build_worker_context(conn: sqlite3.Connection, task_id: str) -> str: """Return the full text a worker should read to understand its task. - Order: - 1. Task title (mandatory). - 2. Task body (optional opening post, capped at 8 KB). - 3. Prior attempts on THIS task (most recent ``_CTX_MAX_PRIOR_ATTEMPTS`` - shown; older attempts collapsed into a one-line summary). - Each attempt's ``summary`` / ``error`` / ``metadata`` capped at - ``_CTX_MAX_FIELD_BYTES`` each. - 4. Structured handoff results of every done parent task. Prefers - ``run.summary`` / ``run.metadata`` when the parent was executed - via a run; falls back to ``task.result`` for older data. Same - per-field cap. - 5. Cross-task role history for the assignee (most recent 5 - completed runs on other tasks). - 6. Comment thread (most recent ``_CTX_MAX_COMMENTS`` shown, older - collapsed). - - All caps exist so worker prompts stay bounded even on pathological - boards (retry-heavy tasks, comment storms). The per-field char cap - prevents a single 1 MB summary from dominating context. + Sections in order: header, body, attachments, prior attempts on this task, + done-parent handoffs (``run.summary``/``metadata``, falling back to + ``task.result`` for pre-runs data), the assignee's recent completed runs + on other tasks, comment thread. Every list is tail-capped (``_CTX_MAX_*``) + with the omitted head summarised, and every field is char-capped, so the + prompt stays bounded on pathological boards (retry storms, comment + storms, a single 1 MB summary). """ task = get_task(conn, task_id) if not task: diff --git a/hermes_cli/kanban_db_connect.py b/hermes_cli/kanban_db_connect.py index d182c2122b..29187e0177 100644 --- a/hermes_cli/kanban_db_connect.py +++ b/hermes_cli/kanban_db_connect.py @@ -837,21 +837,12 @@ def connect( ) -> sqlite3.Connection: """Open (and initialize if needed) the kanban DB. - WAL mode is enabled on every connection; it's a no-op after the first - time but keeps the code robust if the DB file is ever re-created. - - The first connection to a given path auto-runs :func:`init_db` so - fresh installs and test harnesses that construct `connect()` - directly don't have to remember a separate init step. Subsequent - connections skip the schema check via a module-level path cache. - - Path resolution: - - * ``db_path`` explicit → used as-is (legacy callers, tests). - * ``board`` explicit → resolves to that board's DB. - * Neither → :func:`kanban_db_path` resolves via - ``HERMES_KANBAN_DB`` env → ``HERMES_KANBAN_BOARD`` env → - ``/kanban/current`` → ``default``. + WAL is (re)enabled on every connection so a re-created file stays robust. + The first connection per path auto-runs :func:`init_db` (fresh installs + and tests need no separate init step); later ones skip the schema check + via ``_INITIALIZED_PATHS``. Path: explicit ``db_path`` as-is, else + ``board``, else :func:`kanban_db_path` (``HERMES_KANBAN_DB`` -> + ``HERMES_KANBAN_BOARD`` -> ``/kanban/current`` -> ``default``). """ path = db_path if db_path is not None else _kb.kanban_db_path(board=board) path.parent.mkdir(parents=True, exist_ok=True) diff --git a/hermes_cli/kanban_db_dispatch.py b/hermes_cli/kanban_db_dispatch.py index 7bae60e9ed..e1bd1597d2 100644 --- a/hermes_cli/kanban_db_dispatch.py +++ b/hermes_cli/kanban_db_dispatch.py @@ -625,27 +625,15 @@ def detect_stale_running( stale_timeout_seconds: int = 0, signal_fn=None, ) -> list[str]: - """Reclaim ``running`` tasks that show no progress (heartbeat) within the - staleness window. + """Reclaim ``running`` tasks with no heartbeat progress; returns their ids. - A task is considered stale when BOTH of these hold: - - 1. It has been running for longer than ``stale_timeout_seconds`` - (measured from the active run's ``started_at``, falling back to - ``tasks.started_at`` on older runs). - 2. Its ``last_heartbeat_at`` is older than - ``_STALE_HEARTBEAT_GAP_SECONDS`` (or NULL — never sent a heartbeat). - - On reclaim the task is restored to its source phase, the run is closed with - ``outcome='stale'``, and the host-local worker (if still running) is - terminated. - - Only considers ``status='running'`` tasks. Blocked tasks are never - candidates. Returns the list of reclaimed task IDs. - - ``stale_timeout_seconds=0`` disables the check entirely (returns ``[]`` - immediately). ``signal_fn`` is a test hook; defaults to ``os.kill`` - on POSIX. + Stale = running longer than ``stale_timeout_seconds`` (from the active + run's ``started_at``, else ``tasks.started_at``) AND ``last_heartbeat_at`` + older than ``_STALE_HEARTBEAT_GAP_SECONDS`` or NULL. The task returns to + its source phase, the run closes ``outcome='stale'`` and a live host-local + worker is terminated. Blocked tasks are never candidates; + ``stale_timeout_seconds=0`` disables the check. ``signal_fn`` is a test + hook (default ``os.kill`` on POSIX). """ if stale_timeout_seconds <= 0: return [] @@ -748,25 +736,17 @@ def detect_stale_running( def reconcile_orphaned_running( conn: sqlite3.Connection, ) -> list[str]: - """Reconcile ``running`` cards whose claim bookkeeping is broken. - - Tracked-state vs. reality divergence: a task can sit in - ``status='running'`` with ``claim_lock IS NULL`` or ``claim_expires IS - NULL`` (crash mid-claim, manual SQL, DB restore). None of the other - recovery paths ever touch such a card — ``release_stale_claims`` - requires a non-NULL ``claim_expires``, ``detect_crashed_workers`` - requires a host-local claim_lock + worker_pid, and - ``detect_stale_running`` is disabled by default — so the card shows - Running forever (a zombie). - - This pass finds those orphans, requeues them to ``ready`` with an - explanatory comment, closes any leaked run, and appends a - ``reconciled`` event. If the orphan row still records a live PID on - this host, requeueing is deferred to a later tick so we never spawn a - duplicate beside a possibly-alive worker. - - Returns the list of reconciled task ids. Safe to call every tick. + """Requeue ``running`` cards with broken claim bookkeeping; returns their ids. + A task can sit ``running`` with ``claim_lock``/``claim_expires`` NULL + (crash mid-claim, manual SQL, DB restore) and no other recovery path + touches it — ``release_stale_claims`` needs ``claim_expires``, + ``detect_crashed_workers`` needs a host-local lock + pid, + ``detect_stale_running`` is off by default — so it is a zombie forever. + Orphans go back to ``ready`` with an explanatory comment, a leaked run is + closed and a ``reconciled`` event appended; a row still recording a live + host-local PID is deferred to a later tick so no duplicate is spawned + beside a possibly-alive worker. Safe to call every tick. """ now = int(time.time()) reconciled: list[str] = [] @@ -1519,30 +1499,19 @@ def _has_spawnable(conn: sqlite3.Connection, status: str) -> bool: def has_spawnable_ready(conn: sqlite3.Connection) -> bool: - """Return True iff there is at least one ready+assigned+unclaimed task - whose assignee maps to a real Hermes profile. + """True iff a ready+assigned+unclaimed task maps to a real Hermes profile. - Used by the gateway- and CLI-embedded dispatchers' health telemetry to - decide whether ``0 spawned`` is a "stuck" condition (real spawnable - work waiting) or a "correctly idle" condition (only control-plane - lanes like ``orion-cc`` / ``orion-research`` waiting on terminals - that pull tasks via ``claim_task`` directly). - - Falls back to "any ready+assigned" if ``profile_exists`` is not - importable (e.g. partial install) — preserves the old behavior so - the warning still fires in degraded environments. + Health telemetry uses it to tell "stuck" (``0 spawned`` with spawnable + work) from "correctly idle" (only control-plane lanes waiting on terminals + that pull via ``claim_task``). Falls back to "any assigned" when + ``profile_exists`` is unimportable (partial install) so the warning still + fires when degraded. """ return _has_spawnable(conn, "ready") def has_spawnable_review(conn: sqlite3.Connection) -> bool: - """Return True iff there is at least one review+assigned+unclaimed task - whose assignee maps to a real Hermes profile. - - Mirror of :func:`has_spawnable_ready` for the review column — - used by the health telemetry to decide whether the dispatcher - should have spawned a review agent. - """ + """:func:`has_spawnable_ready` for the review column.""" return _has_spawnable(conn, "review") @@ -1766,18 +1735,12 @@ def dispatch_once( ) -> DispatchResult: """Run one dispatcher tick under the board's single-writer lock. - Thin wrapper around :func:`_dispatch_once_locked`. It acquires a - non-blocking, board-scoped dispatch lock (issue #35240) so that two - dispatchers pointed at the same ``kanban.db`` — e.g. the service- - managed gateway and a shell-spawned orphan that escaped the service - cgroup — can never run a reclaim/spawn/write tick concurrently and - race on WAL frames. The losing dispatcher returns an empty - ``DispatchResult`` with ``skipped_locked=True`` and does no DB writes; - the holder is already making progress on the same board. - - The lock is keyed off the board's resolved DB path, so unrelated - boards tick in parallel. See :func:`_dispatch_tick_lock` for the - cross-process / cross-platform mechanics. + Wraps :func:`_dispatch_once_locked` in the non-blocking board-scoped + :func:`_dispatch_tick_lock` so two dispatchers on one ``kanban.db`` (the + service-managed gateway plus an orphan that escaped its cgroup) never + race a write tick on WAL frames. The loser returns an empty + ``DispatchResult`` with ``skipped_locked=True`` and writes nothing; the + lock is keyed on the resolved DB path so unrelated boards tick in parallel. """ try: db_path = _kb.kanban_db_path(board=board)