refactor(kanban-db): compact long docstrings (rules kept) across origin/dispatch/connect
This commit is contained in:
+34
-90
@@ -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()`` (``<root>`` for a
|
||||
``<root>/profiles/<name>`` 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 ``<root>``
|
||||
when ``HERMES_HOME`` is ``<root>/profiles/<name>``, 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 ``<root>/.../attachments/<task_id>/`` subdirectory.
|
||||
|
||||
``HERMES_KANBAN_ATTACHMENTS_ROOT`` pins the path directly (highest
|
||||
precedence) for tests and unusual deployments.
|
||||
|
||||
``default`` uses ``<root>/kanban/attachments/``; other boards use
|
||||
``<root>/kanban/boards/<slug>/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`` -> ``<root>/kanban/attachments/``, other boards ->
|
||||
``<root>/kanban/boards/<slug>/attachments/``; each task gets its own
|
||||
``<task_id>/`` 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:
|
||||
|
||||
@@ -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 →
|
||||
``<root>/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`` -> ``<root>/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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user