atlas-iac/services/hermes/scripts/cli_lane_dispatch.py

257 lines
9.9 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 sqlite3
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
import supervisor_state
OWNED_WORKSPACES_ONLY = os.environ.get("HERMES_CLI_LANE_OWNED_WORKSPACES_ONLY", "").lower() == "true"
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 _direct_lane_eligible(board: str, task: Any) -> bool:
"""Reserve signed PR continuations for the mediated execution pool."""
task_id = str(_task_value(task, "id", "") or "")
if not task_id or (OWNED_WORKSPACES_ONLY and not str(_task_value(task, "workspace_path", "") or "").strip()):
return False
try:
return supervisor_state.get_child(board, task_id) is None
except (OSError, ValueError) as error:
_record_board_access_error(board, error)
return False
def _listed_boards(kanban_db: Any) -> list[Any]:
"""Fall back to board-directory discovery when one metadata read is unreadable."""
try:
return list(kanban_db.list_boards(include_archived=False))
except (OSError, sqlite3.Error) as error:
_record_board_access_error("board-registry", error)
root = getattr(kanban_db, "boards_root", None)
metadata = getattr(kanban_db, "read_board_metadata", None)
if not callable(root) or not callable(metadata):
return []
try:
candidates = [path for path in root().iterdir() if not path.is_symlink() and path.is_dir()]
except (OSError, sqlite3.Error) as error:
_record_board_access_error("board-registry", error)
return []
boards: list[Any] = []
for path in sorted(candidates, key=lambda item: item.name):
try:
if not ((path / "kanban.db").is_file() or (path / "board.json").is_file()):
continue
board = metadata(path.name)
except (OSError, sqlite3.Error, ValueError) as error:
_record_board_access_error(path.name, error)
continue
if isinstance(board, dict) and not board.get("archived") and _board_slug(board) == path.name:
boards.append(board)
return boards
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
boards = _listed_boards(kanban_db)
if not boards:
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"
and str(_task_value(task, "claim_lock", "") or "") == "direct-cli-lane"
and _direct_lane_eligible(board, task)
):
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,
claimer: str = "direct-cli-lane",
priority: Callable[[str, Any], int] | 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 _listed_boards(kanban_db):
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 = list(kanban_db.list_tasks(conn))
if priority is not None:
tasks.sort(key=lambda task: priority(board, task))
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 (
not task_id
or (board, task_id) in active
or str(_task_value(task, "status", "")) != "ready"
or (eligible is not None and not eligible(board, task))
):
continue
if (
task_id
and not assignee
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 not assignee.startswith(EXTERNAL_PREFIX)
):
continue
try:
result = kanban_db.claim_task(
conn,
task_id,
ttl_seconds=DEFAULT_CLAIM_TTL,
claimer=claimer,
)
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), _direct_lane_eligible)
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)