"""Intelligent-eval scheduler runtime(触发式执行的平台节奏)。 从 web/app.py 抽出(ADR-0009 / ADR-0011):平台掌控节奏——每 60s 扫描 executing 评估,watchdog 兜底、任务入队、决策日志补录,然后按需触发 无状态 OpenClaw agent(planner/worker)执行。web 层只负责 start/stop 本 runtime,不再知道扫描循环的细节。 对外 interface 只有三个入口:``scheduler_runtime.start()`` / ``.stop()`` / ``scan_once()``(单次扫描节拍,即测试面)。 """ import asyncio import logging import subprocess import time from typing import Callable, Optional from agenteval.storage.db import get_session _logger = logging.getLogger("agenteval") SCAN_INTERVAL_SECONDS = 60 # 触发指令必须带"立即执行"语义:`openclaw agent` 无 cron state,若只发 # skill 名,agent 会按 skill 的"跨节拍"设计决策后等下一拍而死锁。 WORKER_MSG = ( "执行 agenteval-intelligent-worker skill,立即完成当前任务,不要等待下一节拍:" "仅执行当前时间对应时段(time_distribution 中当前 offset 所在时段)内欠账的会话" "——按评估 started_at 与当前时间精确判断当前时段,只创建该时段计划内的会话," "绝不创建未来时段的会话,未来时段到期后平台会再次触发你;" "仅当所有时段计划内的会话都已达到终态(completed/failed/expired)时," "才调用 agenteval-intelligent-analyst skill 完成分析并生成报告;" "若还有未到期的未来时段,保持等待、本次不做任何操作。" ) PLANNER_MSG = ( "执行 agenteval-intelligent-planner skill,立即完成当前任务,不要等待下一节拍:" "为 planning 状态的智能评估读取输入、产出粗计划并提交平台审批。" ) async def _trigger_openclaw_agent( role: str, session_prefix: str, message: str, on_failure: Callable[[str], None], ) -> None: """触发原语:docker exec 唤醒一个无状态 headless OpenClaw agent。 方案③(避免外部 IM channel):`--deliver` 默认 false,不经 cron delivery channel。ADR-0011 孤儿 agent 双管之一:容器内命令包 `timeout 600`——agent 进程最多跑 10 分钟即被强杀。 必须用独立 session:`--agent main` 复用 main 持久 session,多次触发 累积上下文缓存(~12 万 token)后 LLM 不再执行 skill 的 API 步骤, 直接幻觉输出而不取任务。 失败(子进程缺失/超时/非零退出)调用 ``on_failure(error)`` 落账。 """ cmd = [ "docker", "exec", "openclaw-eval", # ADR-0011:容器内 timeout 强杀,agent 最长存活 10min "timeout", "600", "openclaw", "agent", "--agent", "main", "--session-id", f"{session_prefix}-{int(time.time())}", "-m", message, "--json", ] try: proc = await asyncio.to_thread( subprocess.run, cmd, capture_output=True, text=True, timeout=660, # 容器内 timeout 600 + 启动/回收余量 ) _logger.info("%s 触发完成 exit=%s", role, proc.returncode) if proc.returncode != 0: _logger.warning("%s 触发 stderr: %s", role, proc.stderr[-300:]) on_failure(f"exit={proc.returncode} stderr={proc.stderr[-200:]}") except Exception as exc: _logger.warning("%s 触发失败(忽略): %s", role, exc) on_failure(str(exc)) def _record_worker_trigger_failure(error: str) -> None: """ADR-0011:worker 触发失败落账到受影响评估(有待认领任务的评估)。""" from agenteval.intelligent_eval.lifecycle import ( eval_ids_with_pending_worker_tasks, record_trigger_failures, ) session = get_session() try: eval_ids = eval_ids_with_pending_worker_tasks(session) if eval_ids: record_trigger_failures(session, channel="worker", eval_ids=eval_ids, error=error) except Exception as exc: _logger.warning("触发失败落账失败(忽略): %s", exc) finally: session.close() def _record_planner_trigger_failure(error: str) -> None: """ADR-0011:planner 触发失败落账到所有 planning 评估。""" from sqlmodel import select from agenteval.intelligent_eval.lifecycle import record_trigger_failures from agenteval.intelligent_eval.models import IntelligentEvalStatus from agenteval.storage.db import IntelligentEvalDB session = get_session() try: eval_ids = [ row.id for row in session.exec( select(IntelligentEvalDB).where(IntelligentEvalDB.status == IntelligentEvalStatus.PLANNING.value) ).all() ] if eval_ids: record_trigger_failures(session, channel="planner", eval_ids=eval_ids, error=error) except Exception as exc: _logger.warning("触发失败落账失败(忽略): %s", exc) finally: session.close() async def trigger_worker(nudge_eval_ids: Optional[list[str]] = None) -> bool: """触发 OpenClaw headless agent 执行 worker skill。 平台只当"触发闹钟"——有冷却期外的候选评估时,用 docker exec 唤醒 agent,worker skill 由 agent 自主执行(决策/建会话/真实对话/上报)。 per-eval 触发冷却 10min(worker_trigger 决策日志记录)——agent 还在跑 的评估不重复触发,避免 60s 扫描节奏堆叠并发 agent。 nudge_eval_ids:analyst 兜底催促的评估(所有会话终态、无任务队列条目), 与候选任务评估合并去重后一同落账触发。 Returns: True 若确实触发了 agent(存在冷却期外的服务对象)。 """ from agenteval.intelligent_eval.lifecycle import ( record_worker_triggers, worker_trigger_candidates, ) session = get_session() try: trigger_ids = sorted(set(worker_trigger_candidates(session)) | set(nudge_eval_ids or [])) if not trigger_ids: return False # 触发前落账(含后续失败也计入冷却),防止 agent 在跑期间重复触发 record_worker_triggers(session, trigger_ids) finally: session.close() await _trigger_openclaw_agent("Worker", "agenteval-worker", WORKER_MSG, _record_worker_trigger_failure) return True async def trigger_planner() -> bool: """触发 OpenClaw headless agent 执行 planner skill。 方案③只自动化了 executing(worker)→ completed(analyst),planning 阶段由本函数补上:对 planning 状态评估触发 planner skill,planner 自会 拉取 planning 评估列表、读取四件套、产出粗计划并 PUT /plan 提交。 处理完评估离开 planning 后不再触发。 Returns: True 若确实触发了 agent(存在 planning 评估)。 """ # ADR-0011:触发前先落账 planner_trigger(含后续失败也计入双闸), # 0 个 planning 评估时不触发。 from agenteval.intelligent_eval.lifecycle import record_planner_triggers session = get_session() try: planning_count = record_planner_triggers(session) finally: session.close() if planning_count == 0: return False await _trigger_openclaw_agent("Planner", "agenteval-planner", PLANNER_MSG, _record_planner_trigger_failure) return True def _fire_and_forget(coro, name: str) -> None: """ADR-0011:触发改为 fire-and-forget——agent 最长跑 10min,await 会把 60s 扫描节奏拖到 10min+。派生 asyncio task,异常在完成回调中记录 (触发失败已由触发函数自身落账 trigger_failed)。""" def _done(task: asyncio.Task) -> None: if task.cancelled(): return exc = task.exception() if exc: _logger.warning("%s 触发任务异常(忽略): %s", name, exc) asyncio.create_task(coro).add_done_callback(_done) def scan_once() -> None: """单次扫描节拍:watchdog 兜底 → 入队 → 任务回收 → 决策日志补录 → analyst 催促 → fire-and-forget 触发 worker/planner。失败不阻断 (下次节拍继续)。""" try: session = get_session() try: from agenteval.intelligent_eval.decision_logs import supplement_decision_logs from agenteval.intelligent_eval.lifecycle import ( enforce_executing_ceiling, enforce_planning_gates, enforce_trigger_failure_gates, expire_stale_running_sessions, ) from agenteval.intelligent_eval.task_queue import ( requeue_stale_assigned_tasks, scan_and_enqueue_tasks, settle_tasks_for_finished_evals, ) r = requeue_stale_assigned_tasks(session) if r: _logger.info("卡死恢复:%d 个 assigned 任务重新入队", r) x = expire_stale_running_sessions(session) if x: _logger.info("会话过期:%d 个 running 会话 60 分钟无新轮次,置为 expired", x) g = enforce_planning_gates(session) if g: _logger.info("planning 双闸:%d 个评估超限判失败", g) c = enforce_executing_ceiling(session) if c: _logger.info("executing 兜底:%d 个评估超窗判失败", c) f = enforce_trigger_failure_gates(session) if f: _logger.info("触发失败判死:%d 个评估连续触发失败超限判失败", f) n = scan_and_enqueue_tasks(session) if n: _logger.info("智能评估扫描:入队 %d 个 Worker 任务", n) # 清理:评估已结束(非 executing)的 pending/assigned 任务回收为 completed settled = settle_tasks_for_finished_evals(session) if settled: _logger.info("任务回收:%d 个已结束评估的任务标记完成", settled) # 审计兜底:agent 未上报决策日志时,平台按评估状态补录 added = supplement_decision_logs(session) if added: _logger.info("决策日志兜底:补录 %d 条", added) finally: session.close() except Exception as exc: _logger.warning("智能评估扫描失败(忽略): %s", exc) # analyst 兜底(ADR-0011):末会话终态 10min 后平台催促 analyst 汇总报告, # 催促落账后把评估交给 worker 触发(worker 在全终态时会转 analyst 路径) nudge_ids: list[str] = [] try: from agenteval.intelligent_eval.lifecycle import ( evals_needing_analyst_nudge, record_analyst_nudge, ) session = get_session() try: nudge_ids = evals_needing_analyst_nudge(session) for eval_id in nudge_ids: record_analyst_nudge(session, eval_id) finally: session.close() if nudge_ids: _logger.info("analyst 兜底:催促 %d 个评估汇总报告", len(nudge_ids)) except Exception as exc: _logger.warning("analyst 兜底失败(忽略): %s", exc) # ADR-0011:触发 fire-and-forget,不阻塞 60s 扫描节奏(agent 最长跑 10min) _fire_and_forget(trigger_worker(nudge_eval_ids=nudge_ids), "Worker") _fire_and_forget(trigger_planner(), "Planner") class IntelligentEvalScheduler: """扫描循环的 runtime 单例:持有任务句柄,web 层只做 start/stop。""" def __init__(self) -> None: self._task: Optional[asyncio.Task] = None async def _loop(self) -> None: while True: scan_once() await asyncio.sleep(SCAN_INTERVAL_SECONDS) def start(self) -> None: if self._task is None or self._task.done(): self._task = asyncio.create_task(self._loop()) async def stop(self) -> None: if self._task is not None: self._task.cancel() try: await self._task except asyncio.CancelledError: pass self._task = None scheduler_runtime = IntelligentEvalScheduler()