AgentEvalTool/.scratch/intelligent-eval-cron-pool/spec.md
sinohqb 1aa453ef0a feat(intelligent-eval): add cron pool data model and task queue API
Implement Ticket 01 of intelligent eval cron pool architecture (ADR-0007):

- Add 4 new tables: task_queue, cron_pool, config_snapshots, decision_logs
- Implement task enqueueing logic with priority calculation
- Implement task assignment and completion APIs
- Add unit tests (9) and integration tests (7)
- Update CONTEXT.md with new vocabulary
- Add ADR-0007 documenting cron pool architecture decision

All 760 tests passing.
2026-08-12 02:13:21 +08:00

706 lines
17 KiB
Markdown
Raw Permalink 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.

# 智能评估 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: <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: <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: <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: <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: <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 的 cronOpenClaw 侧是否存在
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)