fix(kanban): text dispatch output and both "dispatcher stuck" warnings name the hold reason
`hermes kanban dispatch` (plain output), the standalone daemon's stuck warning and the gateway's embedded dispatcher stuck warning all reported a bare `Spawned: 0` / "0 workers spawned" while the respawn guard held every ready card — the reason existed only as a `respawn_guarded` task event visible via `hermes kanban tail`. Operators watching the gateway health warning for 73+ ticks (#111910) had nothing to act on. - `kanban_db_dispatch.describe_suppression()` renders the guard reasons per task plus rate_limited / skipped_locked / memory_pressure for one or more DispatchResults, so the CLI daemon and gateway warnings share one wording: `Last tick held back: active_pr=1, memory_pressure=elevated.` - plain `dispatch` output prints `Guarded (<reason>): <task id>` and the tick-level holds, mirroring the JSON fields. - kanban docs: how to see why a ready card is not spawning. Co-authored-by: Steven Saehrig <trac3r726@users.noreply.github.com> Part of #111910
This commit is contained in:
@@ -272,6 +272,7 @@ class GatewayKanbanWatchersMixin:
|
||||
# broken PATH, missing venv, or credential loss.
|
||||
bad_ticks = 0
|
||||
last_warn_at = 0
|
||||
results: Optional[list] = None
|
||||
dispatcher = _KanbanDispatcher(_kb, settings)
|
||||
|
||||
logger.info("kanban dispatcher: embedded in gateway (interval=%.1fs)", interval)
|
||||
@@ -304,12 +305,13 @@ class GatewayKanbanWatchersMixin:
|
||||
bad_ticks = bad_ticks + 1 if ready_pending and not any_spawned else 0
|
||||
now = int(time.time())
|
||||
if bad_ticks >= _HEALTH_WINDOW and now - last_warn_at >= 300:
|
||||
held = _kbd.describe_suppression(res for _slug, res in (results or []))
|
||||
logger.warning(
|
||||
"kanban dispatcher stuck: ready queue non-empty for "
|
||||
"%d consecutive ticks but 0 workers spawned. Check "
|
||||
"%d consecutive ticks but 0 workers spawned.%s Check "
|
||||
"profile health (venv, PATH, credentials) and "
|
||||
"`hermes kanban list --status ready`.",
|
||||
bad_ticks,
|
||||
bad_ticks, f" Last tick held back: {held}." if held else "",
|
||||
)
|
||||
last_warn_at = now
|
||||
except asyncio.CancelledError:
|
||||
|
||||
@@ -20,6 +20,7 @@ from dataclasses import field
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from typing import Callable
|
||||
from typing import Iterable
|
||||
from typing import Mapping
|
||||
from typing import Optional
|
||||
from typing import TYPE_CHECKING
|
||||
@@ -135,6 +136,35 @@ class DispatchResult:
|
||||
Reclaim/promotion bookkeeping still ran; deferred tasks stay queued."""
|
||||
|
||||
|
||||
def describe_suppression(results: Iterable[Optional["DispatchResult"]]) -> str:
|
||||
"""One line naming why the tick(s) held ready work back, or ``""``.
|
||||
|
||||
``active_pr=1, recent_success=2, rate_limited=1, skipped_locked=1,
|
||||
memory_pressure=critical`` — the respawn-guard reasons counted per task
|
||||
plus the tick-level holds. Feeds the "dispatcher stuck" warnings of the
|
||||
CLI daemon and the embedded gateway dispatcher, which otherwise report a
|
||||
bare zero-spawn count while ``hermes kanban tail`` is the only place the
|
||||
guard reason is written (#111910).
|
||||
"""
|
||||
counts: dict[str, int] = {}
|
||||
pressure: Optional[str] = None
|
||||
for res in results:
|
||||
if res is None:
|
||||
continue
|
||||
for _task_id, reason in res.respawn_guarded:
|
||||
counts[reason] = counts.get(reason, 0) + 1
|
||||
if res.rate_limited:
|
||||
counts["rate_limited"] = counts.get("rate_limited", 0) + len(res.rate_limited)
|
||||
if res.skipped_locked:
|
||||
counts["skipped_locked"] = counts.get("skipped_locked", 0) + 1
|
||||
if res.memory_pressure:
|
||||
pressure = res.memory_pressure
|
||||
parts = [f"{k}={v}" for k, v in sorted(counts.items())]
|
||||
if pressure:
|
||||
parts.append(f"memory_pressure={pressure}")
|
||||
return ", ".join(parts)
|
||||
|
||||
|
||||
# Bounded registry of recently-reaped worker exits, filled by the reap loop in
|
||||
# ``dispatch_once`` and read by ``detect_crashed_workers`` to classify a dead-pid
|
||||
# task. Entry: ``pid -> (raw_wait_status, reaped_at_epoch)``; raw status kept so
|
||||
|
||||
@@ -143,6 +143,14 @@ def _cmd_dispatch(args: argparse.Namespace) -> int:
|
||||
f"Skipped (non-spawnable assignee — terminal lane, OK): "
|
||||
f"{', '.join(res.skipped_nonspawnable)}"
|
||||
)
|
||||
for tid, reason in res.respawn_guarded:
|
||||
print(f"Guarded ({reason}): {tid}")
|
||||
if res.rate_limited:
|
||||
print(f"Rate-limited (released to ready, no failure counted): {', '.join(res.rate_limited)}")
|
||||
if res.skipped_locked:
|
||||
print("Skipped: another dispatcher holds this board's lock (no writes this tick)")
|
||||
if res.memory_pressure:
|
||||
print(f"Memory pressure {res.memory_pressure}: new workers restricted this tick")
|
||||
return 0
|
||||
|
||||
|
||||
@@ -212,10 +220,12 @@ def _cmd_daemon(args: argparse.Namespace) -> int:
|
||||
if health_state["bad_ticks"] >= HEALTH_WINDOW:
|
||||
now = int(time.time())
|
||||
if now - health_state["last_warn_at"] >= 300:
|
||||
held = kbd.describe_suppression([res])
|
||||
held = f" Last tick held back: {held}." if held else ""
|
||||
print(
|
||||
f"[{_fmt_ts(now)}] WARN dispatcher stuck: ready queue non-empty for "
|
||||
f"{health_state['bad_ticks']} consecutive ticks but 0 workers spawned "
|
||||
f"successfully. Check profile health (venv, PATH, credentials) and `hermes "
|
||||
f"successfully.{held} Check profile health (venv, PATH, credentials) and `hermes "
|
||||
f"kanban list --status ready` / `hermes kanban list --status blocked` for "
|
||||
f"recent spawn_failed tasks.",
|
||||
file=sys.stderr, flush=True,
|
||||
|
||||
@@ -26,3 +26,51 @@ def test_mixin_defines_kanban_methods():
|
||||
assert hasattr(GatewayKanbanWatchersMixin, m), f"mixin missing {m}"
|
||||
|
||||
|
||||
def test_gateway_dispatcher_stuck_warning_names_guard_reason(monkeypatch, caplog):
|
||||
"""The embedded dispatcher's "stuck" warning names the respawn-guard reason
|
||||
holding the ready queue (#111910) instead of a bare zero-spawn count."""
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
import gateway.kanban_watchers as kw
|
||||
from hermes_cli import kanban_db_dispatch as kbd
|
||||
|
||||
held = kbd.DispatchResult(respawn_guarded=[("t_held", "active_pr")])
|
||||
runner = object.__new__(kw.GatewayKanbanWatchersMixin)
|
||||
runner._running = True
|
||||
monkeypatch.setattr(runner, "_kanban_dispatcher_boot", lambda: (lambda: {}, object(), {}))
|
||||
|
||||
class _Dispatcher:
|
||||
def __init__(self, *a, **k):
|
||||
pass
|
||||
|
||||
def tick_once(self):
|
||||
return [("board", held)]
|
||||
|
||||
def ready_nonempty(self):
|
||||
return True
|
||||
|
||||
ticks = {"n": 0}
|
||||
|
||||
async def _direct(fn, *args):
|
||||
return fn(*args)
|
||||
|
||||
async def _sleep(_delay):
|
||||
ticks["n"] += 1
|
||||
if ticks["n"] > kw._HEALTH_WINDOW:
|
||||
runner._running = False
|
||||
|
||||
monkeypatch.setattr(kw, "_KanbanDispatcher", _Dispatcher)
|
||||
monkeypatch.setattr(kw, "_resolve_dispatcher_settings", lambda cfg, kb: type("S", (), {"interval": 1.0})())
|
||||
monkeypatch.setattr(kw, "_to_thread_process_service", _direct)
|
||||
monkeypatch.setattr(kw, "_kanban_dispatch_allowed", lambda: True)
|
||||
monkeypatch.setattr(kw, "_resolve_auto_decompose_settings", lambda load_config: (False, 0))
|
||||
monkeypatch.setattr(kbd, "reap_worker_zombies", lambda: [])
|
||||
monkeypatch.setattr(kw.asyncio, "sleep", _sleep)
|
||||
|
||||
with caplog.at_level(logging.WARNING, logger=kw.logger.name):
|
||||
asyncio.run(asyncio.wait_for(runner._kanban_dispatcher_watcher(), timeout=5.0))
|
||||
|
||||
stuck = [r.getMessage() for r in caplog.records if "dispatcher stuck" in r.getMessage()]
|
||||
assert stuck, [r.getMessage() for r in caplog.records]
|
||||
assert "Last tick held back: active_pr=1." in stuck[0]
|
||||
|
||||
@@ -506,6 +506,45 @@ def test_dispatch_json_exposes_suppression_reasons(
|
||||
assert payload["memory_pressure"] == "elevated"
|
||||
|
||||
|
||||
def test_dispatch_text_and_daemon_stuck_warning_name_guard_reason(
|
||||
monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str]
|
||||
) -> None:
|
||||
"""The plain `hermes kanban dispatch` output and the standalone daemon's
|
||||
"dispatcher stuck" warning both say WHY a ready card was held (#111910):
|
||||
a guarded card must not look like an idle tick with `Spawned: 0`."""
|
||||
res = kbd.DispatchResult(respawn_guarded=[("t_held", "active_pr")], memory_pressure="elevated")
|
||||
monkeypatch.setattr(kanban_ops.kbd, "dispatch_once", lambda *a, **k: res)
|
||||
|
||||
class _Conn:
|
||||
def __enter__(self):
|
||||
return None
|
||||
|
||||
def __exit__(self, *exc):
|
||||
return False
|
||||
|
||||
monkeypatch.setattr(kanban_ops.kbc, "connect_closing", lambda: _Conn())
|
||||
assert kanban_ops._cmd_dispatch(
|
||||
SimpleNamespace(dry_run=True, max=None, failure_limit=kbd.DEFAULT_FAILURE_LIMIT, json=False)
|
||||
) == 0
|
||||
out = capsys.readouterr().out
|
||||
assert "Guarded (active_pr): t_held" in out
|
||||
assert "Memory pressure elevated" in out
|
||||
|
||||
def _fake_daemon(*, interval, max_spawn, failure_limit, on_tick):
|
||||
for _ in range(6): # HEALTH_WINDOW consecutive bad ticks
|
||||
on_tick(res)
|
||||
|
||||
monkeypatch.setattr(kanban_ops.kbd, "run_daemon", _fake_daemon)
|
||||
monkeypatch.setattr(kanban_ops.kbd, "has_spawnable_ready", lambda conn: True)
|
||||
monkeypatch.setattr(kanban_ops.kb, "init_db", lambda *a, **k: None)
|
||||
assert kanban_ops._cmd_daemon(
|
||||
SimpleNamespace(force=True, interval=5, max=None, failure_limit=2, verbose=False, pidfile=None)
|
||||
) in (0, None)
|
||||
err = capsys.readouterr().err
|
||||
assert "dispatcher stuck" in err
|
||||
assert "Last tick held back: active_pr=1, memory_pressure=elevated." in err
|
||||
|
||||
|
||||
def test_review_dispatch_preserves_task_skills_and_adds_reviewer_skill(
|
||||
kanban_home: Path, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
|
||||
@@ -929,6 +929,8 @@ hermes kanban create "nightly backup audit" \
|
||||
|
||||
The dispatcher refuses to re-spawn a ready task when it hit a quota/auth/429 error on the previous run (`blocker_auth`), or completed a run successfully within the guard window (`recent_success`), or a recent task comment links to a GitHub PR (`active_pr`). This prevents repeat worker storms on the same bug or task while a human catches up. See the `respawn_guarded` row in the [event reference](#event-reference).
|
||||
|
||||
To see why a ready card is not spawning, run `hermes kanban dispatch --dry-run` — the output lists `Guarded (<reason>): <task id>` per held card (and `respawn_guarded`, `rate_limited`, `skipped_locked`, `memory_pressure` with `--json`). The gateway's and the standalone daemon's "dispatcher stuck" warning also names what the last tick held back, e.g. `Last tick held back: active_pr=1`.
|
||||
|
||||
### Drag-to-delete and bulk delete (dashboard)
|
||||
|
||||
The dashboard exposes a **trash drop zone** on the kanban page — drag any card into it to delete the task (cascades through `task_events`, child links, and subscriptions). A confirmation prompt protects against accidents. Bulk delete is also reachable via `DELETE /api/plugins/kanban/tasks` with a JSON body `{"ids": ["t_abc", "t_def", ...]}`.
|
||||
|
||||
Reference in New Issue
Block a user