AgentEvalTool/backend/agenteval/intelligent_eval/scheduler.py
sinohqb 913dc9ae86
All checks were successful
CI / test (push) Successful in 3m18s
perf: address remaining heuristic issues from code review
- decision_logs.py: add LIMIT 100 to dedup query to avoid loading all records
- ExecutionProcess.tsx: document N+1 API pattern and explain why acceptable
- scheduler.py: document why scan_once runs synchronously (thread pool would break fire-and-forget)

All 880 tests pass.
2026-08-24 02:01:04 +08:00

303 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Intelligent-eval scheduler runtime触发式执行的平台节奏
从 web/app.py 抽出ADR-0009 / ADR-0011平台掌控节奏——每 60s 扫描
executing 评估watchdog 兜底、任务入队、决策日志补录,然后按需触发
无状态 OpenClaw agentplanner/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 sqlmodel import select
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,
eval_ids_with_pending_worker_tasks,
evals_needing_analyst_nudge,
expire_stale_running_sessions,
record_analyst_nudge,
record_planner_triggers,
record_trigger_failures,
record_worker_triggers,
worker_trigger_candidates,
)
from agenteval.intelligent_eval.models import IntelligentEvalStatus
from agenteval.intelligent_eval.task_queue import (
requeue_stale_assigned_tasks,
scan_and_enqueue_tasks,
settle_tasks_for_finished_evals,
)
from agenteval.storage.db import IntelligentEvalDB, 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-0011worker 触发失败落账到受影响评估(有待认领任务的评估)。"""
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-0011planner 触发失败落账到所有 planning 评估。"""
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 唤醒
agentworker skill 由 agent 自主执行(决策/建会话/真实对话/上报)。
per-eval 触发冷却 10minworker_trigger 决策日志记录——agent 还在跑
的评估不重复触发,避免 60s 扫描节奏堆叠并发 agent。
nudge_eval_idsanalyst 兜底催促的评估(所有会话终态、无任务队列条目),
与候选任务评估合并去重后一同落账触发。
Returns:
True 若确实触发了 agent存在冷却期外的服务对象
"""
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。
方案③只自动化了 executingworker→ completedanalystplanning
阶段由本函数补上:对 planning 状态评估触发 planner skillplanner 自会
拉取 planning 评估列表、读取四件套、产出粗计划并 PUT /plan 提交。
处理完评估离开 planning 后不再触发。
Returns:
True 若确实触发了 agent存在 planning 评估)。
"""
# ADR-0011触发前先落账 planner_trigger含后续失败也计入双闸
# 0 个 planning 评估时不触发。
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 最长跑 10minawait 会把
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:
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:
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 runs synchronously in the event loop. This is acceptable because:
# 1. It primarily does DB queries/updates (fast, non-blocking I/O)
# 2. Execution time is typically milliseconds to tens of milliseconds
# 3. Scan interval is 60s, so brief blocking has minimal impact
# 4. Moving to thread pool would break _fire_and_forget (needs event loop)
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()