# 智能评估 Cron 池实现规范 **版本**: v1.1 **日期**: 2026-08-11 **状态**: draft **决策依据**: ADR-0007 ## Problem Statement 当前智能评估的 OpenClaw 集成采用"多 Session"模式:每个评估对应多个无状态的 OpenClaw session,通过 API 传递上下文。这导致: 1. 上下文割裂:planner 的决策逻辑,evaluator 不知道 2. 无法追踪:哪个 OpenClaw session 做了什么,平台无视图 3. 重复读取:每个 session 都要重新读取评估详情 同时,如果改为"一个评估一个 cron",会导致 cron 爆炸(100 个评估 = 100 个 cron)。 ## Solution 采用 **Cron 池模式**:OpenClaw 维护一个 cron 池(5-20 个),平台维护任务队列,cron 从队列取任务执行,完成后归还。 ### 核心概念 - **Cron 池(Cron Pool)**:OpenClaw 侧的工作单元池,每个 cron 可以处理任意评估 - **任务队列(Task Queue)**:平台侧的待处理评估队列,按优先级排序 - **工作单元(Worker)**:一个 cron + 其 state,表示一个可用的执行单元 - **任务(Task)**:一个需要处理的评估,包含 eval_id、优先级、原因 --- ## 数据模型 ### 1. 平台侧:任务队列 ```python class EvalTaskQueueDB(SQLModel, table=True): """智能评估任务队列""" __tablename__ = "intelligent_eval_task_queue" id: str = Field(primary_key=True) # 任务 ID eval_id: str = Field(foreign_key="intelligent_evals.id", index=True) # 关联的评估 # 任务状态 status: str = Field(index=True) # pending / assigned / completed / failed priority: int = Field(index=True) # 优先级(越小越优先) reason: str # 需要处理的原因(如 "slot_due", "all_sessions_completed") # 分配信息 assigned_cron_id: str | None = None # 分配的 cron ID assigned_at: datetime | None = None # 分配时间 # 完成信息 completed_at: datetime | None = None error: str | None = None # 时间戳 created_at: datetime updated_at: datetime # 索引 __table_args__ = ( Index("idx_status_priority", "status", "priority"), Index("idx_eval_status", "eval_id", "status"), ) ``` ### 2. 平台侧:Cron 池状态 ```python class OpenClawCronPoolDB(SQLModel, table=True): """OpenClaw Cron 池状态(平台侧记录)""" __tablename__ = "openclaw_cron_pool" id: str = Field(primary_key=True) # 平台侧 ID openclaw_cron_id: str = Field(unique=True, index=True) # OpenClaw 侧的 cron ID # 状态 status: str = Field(index=True) # idle / busy / stuck current_eval_id: str | None = Field(foreign_key="intelligent_evals.id") # 当前处理的评估 # 心跳 last_active_at: datetime # 上次活跃时间 last_task_at: datetime | None = None # 上次取任务时间 # 统计 total_tasks_completed: int = 0 total_tasks_failed: int = 0 # 时间戳 created_at: datetime updated_at: datetime # 索引 __table_args__ = ( Index("idx_status_last_active", "status", "last_active_at"), ) ``` ### 3. OpenClaw 侧:Cron State ```json { "status": "idle | busy", "eval_id": "uuid | null", "started_at": "ISO8601 | null", "last_decision_at": "ISO8601", "completed_sessions": 0, "decisions_history": [ { "timestamp": "ISO8601", "decision": "execute_session | wait | start_analysis", "reason": "..." } ] } ``` --- ## API 契约 ### 1. 平台侧 API(供 OpenClaw 调用) #### 获取下一个任务 ```http GET /api/intelligent-evals/tasks/next X-API-Key: Response 200: { "task": { "id": "task_uuid", "eval_id": "eval_uuid", "priority": 1, "reason": "slot_due", "eval": { "id": "eval_uuid", "name": "...", "status": "executing", "plan": {...}, "started_at": "..." } } } Response 200 (无任务): { "task": null } ``` #### 标记任务完成 ```http POST /api/intelligent-evals/tasks/{task_id}/complete X-API-Key: Content-Type: application/json { "cron_id": "openclaw_cron_id", "success": true, "error": null } Response 200: { "success": true } ``` #### 上报 Cron 心跳 ```http POST /api/openclaw/crons/{cron_id}/heartbeat X-API-Key: Content-Type: application/json { "status": "busy", "current_eval_id": "eval_uuid" } Response 200: { "success": true } ``` ### 2. 平台侧管理 API(供前端/管理员调用) #### 查看 Cron 池状态 ```http GET /api/openclaw/cron-pool X-API-Key: Response 200: { "pool": { "total": 10, "idle": 3, "busy": 7, "stuck": 0, "min_size": 5, "max_size": 20 }, "crons": [ { "id": "platform_id", "openclaw_cron_id": "openclaw_id", "status": "busy", "current_eval_id": "eval_uuid", "last_active_at": "...", "total_tasks_completed": 15 } ] } ``` #### 手动扩容/缩容 ```http POST /api/openclaw/cron-pool/scale X-API-Key: Content-Type: application/json { "target_size": 15 } Response 200: { "success": true, "current_size": 15 } ``` --- ## 状态机 ### 1. 任务状态机 ``` pending → assigned → completed → failed ``` - **pending**:任务已创建,等待分配 - **assigned**:任务已分配给某个 cron - **completed**:任务已完成 - **failed**:任务失败(cron 卡死、评估取消等) ### 2. Cron 状态机 ``` idle → busy → idle → stuck → (平台介入) → idle ``` - **idle**:空闲,可以取任务 - **busy**:正在处理任务 - **stuck**:卡死(10 分钟未活跃) ### 3. 评估状态机(不变) ``` draft → planning → pending_approval → executing → completed → cancelled → failed ``` --- ## 核心流程 ### 1. 任务入队流程 ```python # 平台侧:定时扫描(每分钟) async def scan_and_enqueue_tasks(): """扫描所有 executing 评估,生成任务""" evals = get_executing_evals() for eval in evals: # 判断是否需要立即处理 if needs_attention(eval): # 检查是否已有 pending 任务(去重) existing = get_pending_task(eval.id) if existing: continue # 创建任务 task = EvalTaskQueueDB( eval_id=eval.id, status="pending", priority=calculate_priority(eval), reason=get_attention_reason(eval), created_at=utc_now() ) session.add(task) session.commit() def needs_attention(eval) -> bool: """判断评估是否需要立即处理""" # 1. 检查是否有时段到期 current_offset = utc_now() - eval.started_at for slot in eval.plan["time_distribution"]: if is_slot_due(slot, current_offset): return True # 2. 检查是否所有会话完成(需要开始分析) sessions = get_sessions(eval.id) if all(s.status == "completed" for s in sessions): return True return False def calculate_priority(eval) -> int: """计算优先级(越小越优先)""" priority = 100 # 时段到期的优先 if has_due_slot(eval): priority -= 50 # 欠账多的优先 deficit = calculate_session_deficit(eval) priority -= deficit * 10 # 等待时间长的优先 wait_minutes = (utc_now() - eval.started_at).total_seconds() / 60 priority -= min(wait_minutes / 10, 20) return max(priority, 1) ``` ### 2. Cron 取任务流程 ```python # OpenClaw 侧:Worker Skill 每分钟执行 async def worker_tick(): """Cron 每分钟唤醒""" # 1. 读取自己的 state state = get_cron_state() if state["status"] == "idle": # 2. 调平台 API 取任务 task = await fetch_next_task() if task is None: # 无任务,本节拍结束 return # 3. 更新 state 为 busy update_cron_state({ "status": "busy", "eval_id": task["eval_id"], "started_at": utc_now() }) # 4. 处理任务 await process_task(task) elif state["status"] == "busy": # 5. 继续处理当前任务 eval_id = state["eval_id"] await continue_processing(eval_id) async def process_task(task): """处理任务""" eval_id = task["eval_id"] # 读取评估详情 eval = await fetch_eval(eval_id) # 决策逻辑 decision = make_decision(eval) if decision == "execute_session": # 调用 evaluator skill await execute_session(eval) elif decision == "start_analysis": # 调用 analyst skill await start_analysis(eval) elif decision == "wait": # 等待,更新 state update_cron_state({"last_decision_at": utc_now()}) # 检查是否完成 if is_eval_completed(eval): # 标记任务完成 await complete_task(task["id"], success=True) # 归还 cron update_cron_state({"status": "idle", "eval_id": None}) ``` ### 3. 池管理流程 ```python # 平台侧:池管理器(每分钟执行) async def manage_pool(): """管理 cron 池""" pool = get_pool_status() # 1. 扩容 if pool["busy"] / pool["total"] > 0.8 and pool["total"] < MAX_POOL_SIZE: await scale_up(1) # 2. 缩容 if pool["idle"] > MIN_POOL_SIZE * 2 and pool["total"] > MIN_POOL_SIZE: await scale_down(1) # 3. 检测卡死的 cron stuck_crons = get_stuck_crons() # 10 分钟未活跃 for cron in stuck_crons: await handle_stuck_cron(cron) async def scale_up(count: int): """扩容""" for _ in range(count): # 调 OpenClaw CLI 创建 cron cron_id = await openclaw_client.create_cron( name=f"intelligent-eval-worker-{uuid()}", schedule="* * * * *", skill="agenteval-intelligent-worker", state={"status": "idle"} ) # 记录到平台 DB pool_db = OpenClawCronPoolDB( openclaw_cron_id=cron_id, status="idle", last_active_at=utc_now() ) session.add(pool_db) session.commit() async def handle_stuck_cron(cron): """处理卡死的 cron""" # 1. 标记 cron 为 stuck cron.status = "stuck" # 2. 如果有正在处理的任务,标记为 failed if cron.current_eval_id: task = get_assigned_task(cron.current_eval_id) if task: task.status = "failed" task.error = "Cron stuck" # 3. 重新入队(让其他 cron 处理) new_task = EvalTaskQueueDB( eval_id=cron.current_eval_id, status="pending", priority=1, # 高优先级 reason="cron_stuck_retry" ) session.add(new_task) # 4. 删除卡死的 cron await openclaw_client.delete_cron(cron.openclaw_cron_id) session.delete(cron) # 5. 创建新 cron 补充 await scale_up(1) session.commit() ``` --- ## 故障恢复策略 ### 1. OpenClaw 重启 **场景**:OpenClaw Gateway 重启 **影响**: - Cron state 持久化在 SQLite,不丢失 - 正在执行的 session 可能中断 **恢复**: 1. OpenClaw 重启后,cron 自动恢复 2. Cron 下次唤醒时,检查 state: - 如果 busy → 继续处理(从 state 恢复上下文) - 如果 idle → 正常取任务 ### 2. 平台重启 **场景**:AgentEvalTool 平台重启 **影响**: - 任务队列持久化在 DB,不丢失 - 正在处理的任务状态可能不一致 **恢复**: 1. 平台重启后,扫描所有 assigned 任务 2. 检查对应的 cron 是否还活跃: - 如果 cron 活跃 → 任务继续 - 如果 cron 不活跃 → 任务重新入队 ### 3. Cron 卡死 **场景**:Cron 10 分钟未活跃(未上报心跳) **检测**: ```python def detect_stuck_crons(): """检测卡死的 cron""" threshold = utc_now() - timedelta(minutes=10) stuck = session.query(OpenClawCronPoolDB).filter( OpenClawCronPoolDB.status == "busy", OpenClawCronPoolDB.last_active_at < threshold ).all() return stuck ``` **处理**: 1. 标记 cron 为 stuck 2. 任务重新入队 3. 删除卡死的 cron 4. 创建新 cron 补充 ### 4. 状态不一致 **场景**:平台 DB 与 OpenClaw state 不一致 **检测**: ```python async def reconcile_state(): """定期对账(每 5 分钟)""" # 1. 检查平台 DB 中 busy 的 cron,OpenClaw 侧是否存在 busy_crons = get_busy_crons() for cron in busy_crons: openclaw_state = await openclaw_client.get_cron_state(cron.openclaw_cron_id) if openclaw_state is None: # OpenClaw 侧不存在,标记为 stuck await handle_stuck_cron(cron) # 2. 检查 OpenClaw 侧 busy 的 cron,平台 DB 是否记录 openclaw_crons = await openclaw_client.list_crons() for oc_cron in openclaw_crons: platform_cron = get_cron_by_openclaw_id(oc_cron.id) if platform_cron is None: # 平台 DB 不存在,补充记录 await sync_cron_to_db(oc_cron) ``` --- ## 监控和告警 ### 1. 关键指标 ```python # 池使用率 pool_utilization = busy_crons / total_crons # 任务积压 pending_tasks_count = count_pending_tasks() # Cron 卡死率 stuck_rate = stuck_crons / total_crons # 平均任务处理时间 avg_task_duration = avg(completed_at - assigned_at) # 评估完成率 eval_completion_rate = completed_evals / total_evals ``` ### 2. 告警规则 | 指标 | 阈值 | 告警级别 | |------|------|---------| | 池使用率 | > 90% 持续 10 分钟 | Warning | | 任务积压 | > 50 个 | Warning | | Cron 卡死率 | > 10% | Critical | | 平均任务处理时间 | > 30 分钟 | Warning | | 评估完成率 | < 80% | Info | ### 3. 监控面板 前端新增「Cron 池监控」页面: - 池状态(total/idle/busy/stuck) - 任务队列(pending/assigned/completed) - Cron 列表(ID、状态、当前任务、最后活跃时间) - 任务历史(完成时间、成功率) --- ## 性能优化 ### 1. 数据库索引 ```sql -- 任务队列索引 CREATE INDEX idx_status_priority ON intelligent_eval_task_queue(status, priority); CREATE INDEX idx_eval_status ON intelligent_eval_task_queue(eval_id, status); -- Cron 池索引 CREATE INDEX idx_status_last_active ON openclaw_cron_pool(status, last_active_at); ``` ### 2. 缓存 ```python # 缓存池状态(1 分钟) @cache(ttl=60) def get_pool_status(): return calculate_pool_status() # 缓存任务队列(30 秒) @cache(ttl=30) def get_pending_tasks(): return fetch_pending_tasks() ``` ### 3. 批量操作 ```python # 批量标记任务完成 async def batch_complete_tasks(task_ids: list[str]): session.query(EvalTaskQueueDB).filter( EvalTaskQueueDB.id.in_(task_ids) ).update({"status": "completed", "completed_at": utc_now()}) session.commit() ``` --- ## 测试策略 ### 1. 单元测试 - 任务优先级计算 - 池扩容/缩容逻辑 - 卡死检测逻辑 - 状态对账逻辑 ### 2. 集成测试 - 任务入队 → cron 取任务 → 处理 → 完成 - Cron 卡死 → 任务重新入队 - 池满 → 新任务排队 - OpenClaw 重启 → 状态恢复 ### 3. 端到端测试 - 创建评估 → 分配 cron → 执行 → 完成 - 并发 20 个评估 → 池扩容 → 全部完成 - Cron 卡死 → 自动恢复 --- ## 迁移策略 ### 从当前方案迁移 1. **阶段 1**:部署 Cron 池,但保持旧的多 Session 模式 2. **阶段 2**:新评估使用 Cron 池,旧评估继续用旧模式 3. **阶段 3**:旧评估完成后,下线旧模式 ### 数据迁移 无需数据迁移(新功能,不影响现有数据) --- ## Out of Scope - 事件驱动唤醒(webhook):后续优化 - 优先级队列 UI:后续优化 - 跨评估对比:v1.2 功能 - 多对象绑定:v1.2 功能 --- ## Open Questions 1. **OpenClaw Gateway API 的具体端点?** - 需要验证或通过 CLI 封装 2. **Cron state 的 16KB 限制是否够用?** - 粗计划 JSON 大约 2-5KB - 决策历史可能需要截断 3. **池大小的最佳配置?** - 最小 5 个,最大 20 个,是否需要调整? 4. **任务优先级的权重?** - 时段到期 -50,欠账多 -10/个,等待时间 -1/10分钟 - 是否需要调整? --- ## References - ADR-0007: 智能评估 OpenClaw 集成采用 Cron 池模式 - CONTEXT.md: 智能评估词汇 - ADR-0003: 评估活动分期 - [OpenClaw Cron Jobs Documentation](https://docs.openclaw.ai/automation/cron-jobs)