174 lines
6.9 KiB
Python
174 lines
6.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 sys
|
|
import time
|
|
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
|
|
|
|
|
|
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) -> 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"
|
|
):
|
|
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
|