atlas-iac/services/hermes/scripts/cli_lane_dispatch.py
2026-08-23 10:37:01 -03:00

205 lines
7.8 KiB
Python

#!/usr/bin/env python3
"""Concurrent board recovery, claim, and dispatch loop for CLI workers."""
from __future__ import annotations
import concurrent.futures
import os
import sys
import time
from collections.abc import Callable
from typing import Any
from cli_lane_board import _external, _record_board_access_error, _task_value
from cli_lane_capabilities import (
initialize_kanban_capabilities,
kanban_capabilities,
readiness_issue,
refresh_kanban_capabilities,
)
from cli_lane_config import (
BOARD_CORRUPTION_ERRORS,
DEFAULT_CLAIM_TTL,
EXTERNAL_PREFIX,
RESULT_SCHEMA,
RESULT_SCHEMA_PATH,
)
from cli_lane_execution import execute_claim
from cli_lane_files import atomic_json
from cli_lane_metrics import start_metrics_server
from cli_lane_recovery import _has_pending_finalization, recover_pending_finalizations
from cli_lane_retention import maybe_gc_lane_artifacts
def _board_slug(board: Any) -> str:
if isinstance(board, dict):
return str(board.get("slug") or board.get("id") or "")
return str(getattr(board, "slug", None) or getattr(board, "id", None) or board)
def _connect_healthy_board(kanban_db: Any, board: str) -> Any | None:
"""Open one board without letting localized storage faults stop other lanes."""
try:
return kanban_db.connect(board=board)
except Exception as error:
_record_board_access_error(board, error)
return None
def recover_orphans() -> None:
"""Return external running tasks to ready after a runner/pod restart."""
from hermes_cli import kanban_db
recover_pending_finalizations()
if not kanban_capabilities(kanban_db).exact_run_reclaim:
return
try:
boards = kanban_db.list_boards(include_archived=False)
except Exception as error:
_record_board_access_error("board-registry", error)
return
for raw_board in boards:
board = _board_slug(raw_board)
if not board:
continue
with kanban_db.scoped_current_board(board):
conn = _connect_healthy_board(kanban_db, board)
if conn is None:
continue
try:
for task in kanban_db.list_tasks(conn):
if (
_external(task)
and str(_task_value(task, "status", "")) == "running"
):
task_id = str(_task_value(task, "id"))
run_id = _task_value(task, "current_run_id", None)
if not isinstance(run_id, int):
continue
if _has_pending_finalization(board, task_id, run_id):
continue
kanban_db.reclaim_task(
conn,
task_id,
reason="direct CLI lane restarted; provider session will resume",
expected_run_id=run_id,
)
BOARD_CORRUPTION_ERRORS.pop(board, None)
except Exception as error:
_record_board_access_error(board, error)
finally:
conn.close()
def claim_ready(
active: set[tuple[str, str]],
limit: int,
eligible: Callable[[str, Any], bool] | None = None,
) -> list[tuple[str, str]]:
"""Atomically claim external ready tasks across all non-archived boards."""
from hermes_cli import kanban_db
claimed: list[tuple[str, str]] = []
if limit <= 0:
return claimed
for raw_board in kanban_db.list_boards(include_archived=False):
board = _board_slug(raw_board)
if not board:
continue
with kanban_db.scoped_current_board(board):
conn = _connect_healthy_board(kanban_db, board)
if conn is None:
continue
try:
kanban_db.recompute_ready(conn)
tasks = kanban_db.list_tasks(conn)
BOARD_CORRUPTION_ERRORS.pop(board, None)
for task in tasks:
task_id = str(_task_value(task, "id", ""))
assignee = str(_task_value(task, "assignee", "") or "")
if (
task_id
and not assignee
and str(_task_value(task, "status", "")) == "ready"
and kanban_db.assign_task(conn, task_id, "cli-auto")
):
task = kanban_db.get_task(conn, task_id)
assignee = "cli-auto"
if (
not task_id
or (board, task_id) in active
or not assignee.startswith(EXTERNAL_PREFIX)
or str(_task_value(task, "status", "")) != "ready"
or (eligible is not None and not eligible(board, task))
):
continue
try:
result = kanban_db.claim_task(
conn,
task_id,
ttl_seconds=DEFAULT_CLAIM_TTL,
claimer="direct-cli-lane",
)
except Exception:
continue
if result is not None:
claimed.append((board, task_id))
if len(claimed) >= limit:
return claimed
except Exception as error:
_record_board_access_error(board, error)
finally:
conn.close()
return claimed
def main() -> int:
"""Continuously bridge external Kanban lanes to provider CLIs.
The loop is exited only by process signals or unrecoverable exceptions;
every per-board and per-claim fault is bounded inside one pass.
"""
from hermes_cli import kanban_db
RESULT_SCHEMA_PATH.parent.mkdir(parents=True, exist_ok=True)
atomic_json(RESULT_SCHEMA_PATH, RESULT_SCHEMA, 0o644)
start_metrics_server(health_check=lambda: readiness_issue() is None)
capabilities = initialize_kanban_capabilities(kanban_db)
recover_orphans()
workers = max(1, min(int(os.environ.get("HERMES_CLI_LANE_CONCURRENCY", "4")), 8))
futures: dict[concurrent.futures.Future[None], tuple[str, str]] = {}
with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as pool:
while True:
previous_capabilities = capabilities
capabilities = refresh_kanban_capabilities(kanban_db)
if capabilities.ready and not previous_capabilities.ready:
recover_orphans()
for future in list(futures):
if future.done():
try:
future.result()
except Exception as error:
print(
f"worker future failed: {error}",
file=sys.stderr,
flush=True,
)
del futures[future]
recover_pending_finalizations()
maybe_gc_lane_artifacts()
active = set(futures.values())
try:
newly_claimed = (
claim_ready(active, workers - len(futures))
if capabilities.ready
else []
)
BOARD_CORRUPTION_ERRORS.pop("board-registry", None)
except Exception as error:
_record_board_access_error("board-registry", error)
newly_claimed = []
for board, task_id in newly_claimed:
future = pool.submit(execute_claim, board, task_id)
futures[future] = (board, task_id)
time.sleep(5)