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

17 KiB
Raw Permalink Blame History

智能评估 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 PoolOpenClaw 侧的工作单元池,每个 cron 可以处理任意评估
  • 任务队列Task Queue:平台侧的待处理评估队列,按优先级排序
  • 工作单元Worker:一个 cron + 其 state表示一个可用的执行单元
  • 任务Task:一个需要处理的评估,包含 eval_id、优先级、原因

数据模型

1. 平台侧:任务队列

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 池状态

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

{
  "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 调用)

获取下一个任务

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
}

标记任务完成

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 心跳

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 池状态

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
    }
  ]
}

手动扩容/缩容

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. 任务入队流程

# 平台侧:定时扫描(每分钟)
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 取任务流程

# 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. 池管理流程

# 平台侧:池管理器(每分钟执行)
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 分钟未活跃(未上报心跳)

检测

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 不一致

检测

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. 关键指标

# 池使用率
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. 数据库索引

-- 任务队列索引
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. 缓存

# 缓存池状态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. 批量操作

# 批量标记任务完成
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