#!/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 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_recovery import _has_pending_finalization, recover_pending_finalizations from cli_lane_retention import maybe_gc_lane_artifacts OWNED_WORKSPACES_ONLY = os.environ.get( "HERMES_CLI_LANE_OWNED_WORKSPACES_ONLY", "" ).strip().lower() in {"1", "true", "yes", "on"} def _owns_local_workspace(task: Any) -> bool: """Keep legacy/dirty worktree tasks on their existing single-host owner.""" return bool(str(_task_value(task, "workspace_path", "") or "").strip()) 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" and (not OWNED_WORKSPACES_ONLY or _owns_local_workspace(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, ) -> 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 (OWNED_WORKSPACES_ONLY and not _owns_local_workspace(task)) 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.""" from hermes_cli import kanban_db RESULT_SCHEMA_PATH.parent.mkdir(parents=True, exist_ok=True) atomic_json(RESULT_SCHEMA_PATH, RESULT_SCHEMA, 0o644) 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: 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) return 0