Skip to content

feat: 分布式 agent 调度中心——任务队列 + 跨主机日志实时回流 #19

Description

@Camille1024

背景

当前 dashboard 只能监控本机运行的 agent,任务也只能在本地手动触发。随着形式化工作推进,需要支持:

  • 多台工作机各自跑 agent(absorber、prover 等)
  • 部署服务器统一调度任务、汇聚日志、展示进度
  • Dashboard 直接触发 agent 任务,无需登录工作机

目标架构

部署服务器 (kip.opensii.ai)              工作机 A / 工作机 B
┌─────────────────────────────┐          ┌──────────────────────┐
│  Dashboard (UI)             │          │  agent-runner daemon │
│  ├── 任务队列 /api/jobs      │◄─ poll ──│  每隔几秒轮询新任务   │
│  ├── 接收日志 /api/ingest    │◄─ push ──│  run.sh 运行中实时上报│
│  └── WebSocket 推送前端      │          │  仅需 HTTP 出口权限   │
└─────────────────────────────┘          └──────────────────────┘

工作机只需要能访问部署服务器的 HTTP 端口,不需要部署机 SSH 进工作机


分三步实现

Step 1:跨主机日志实时上报(/api/ingest

对应 #18,可独立先做。

run.sh 里的日志解析脚本每 emit 一条 JSONL event,同时 POST /api/ingest。服务器端:

  • 写入本地文件(保持现有目录结构)
  • 通过现有 WebSocket 广播给前端

配置方式:工作机设置环境变量即可,不影响本地模式。

KIP_INGEST_URL=https://kip.opensii.ai/dashboard/api/ingest
KIP_INGEST_TOKEN=<token>

Step 2:任务队列

.kip/state.db 新增 jobs 表:

CREATE TABLE jobs (
  id         TEXT PRIMARY KEY,
  kind       TEXT,   -- 'absorb' | 'prove' | ...
  payload    TEXT,   -- JSON: {hint_file, node_id, ...}
  status     TEXT,   -- 'queued' | 'claimed' | 'running' | 'done' | 'error'
  worker     TEXT,   -- 认领的工作机标识
  created_at TEXT,
  claimed_at TEXT,
  done_at    TEXT
);

新增 API:

方法 路径 说明
GET /api/jobs?status=queued 工作机轮询待领任务
POST /api/jobs/:id/claim 原子认领(防双抢)
POST /api/jobs 创建新任务(dashboard 或 CLI 触发)
GET /api/jobs/:id 查询任务状态

Step 3:工作机 agent-runner daemon

工作机上运行一个轻量 daemon(agents/runner/daemon.sh,约 50 行 bash):

while true; do
  job=$(curl -s "$SERVER/api/jobs?status=queued&kind=absorb" | jq '.[0]')
  if [ -n "$job" ] && [ "$job" != "null" ]; then
    id=$(echo $job | jq -r '.id')
    # 原子认领
    curl -s -X POST "$SERVER/api/jobs/$id/claim" \
      -H "Authorization: Bearer $TOKEN" \
      -d "{\"worker\": \"$(hostname)\"}"
    # 运行 agent,日志实时回流
    KIP_INGEST_URL="$SERVER/api/ingest" \
    KIP_JOB_ID="$id" \
      ./agents/blueprint-absorber/run.sh \
        "$(echo $job | jq -r '.payload.hint_file')"
  fi
  sleep 5
done

Step 4:Dashboard 触发入口

节点 drawer 中新增触发按钮(对应 PR #14 中 "out of scope" 的条目):

  • absorberdrafted 阶段节点显示 "Run absorber" → POST /api/jobs {kind:'absorb', payload:{hint_file}}
  • prover(未来):aligned 阶段节点显示 "Run prover" → POST /api/jobs {kind:'prove', payload:{node_id}}

任务创建后工作机 5 秒内认领,日志实时回流,drawer 内可看到 agent 运行状态。


实施顺序与工作量

步骤 依赖 估计工作量
Step 1:/api/ingest + run.sh hook ~1 天
Step 2:jobs 表 + API Step 1 ~1 天
Step 3:worker daemon Step 2 ~半天
Step 4:Dashboard 触发按钮 Step 2 ~半天

每步独立可用,不需要全部完成才能产生价值。建议从 Step 1 开始。

安全考虑

  • /api/ingest/api/jobs/claim 均需 Bearer token 校验
  • payload 中的文件路径在服务器侧做 containment check(防路径穿越)
  • worker 标识用 hostname,日志中可见来源机器

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions