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.
706 lines
17 KiB
Markdown
706 lines
17 KiB
Markdown
# 智能评估 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 的 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)
|