【免费下载链接】nasiko
Developer Control Plane for your AI Agents
导读:本文以 server/src/agents/BUILD_WORKER.md 为骨架,深入剖析 nasiko(Developer Control Plane for your AI Agents)中负责 Agent 构建的核心组件 ——
build_worker。它用一张 Postgres 表 +FOR UPDATE SKIP LOCKED实现了一个零外部中间件(无 Redis、无 Celery)的进程内异步任务队列,支撑 OCI 镜像推送、git clone + buildkit 构建、版本回滚等耗时 5~30 分钟的操作。读完本文,你将掌握该队列的表结构设计、唤醒机制、两阶段 panic 隔离执行、尝试计数与三层恢复策略、多副本安全模型,以及如何向队列新增一种任务类型。
为什么需要任务队列
Agent 构建操作——OCI 镜像推送、git clone+ buildkit 调用、版本回滚——往往需要5 到 30 分钟。如果 HTTP handler 同步执行这类操作,请求连接会被长时间占用,这在生产环境中非常脆弱:负载均衡器有超时限制、客户端可能随时断开。
nasiko 的解法是引入一个Postgres 支撑的进程内任务队列:HTTP handler 只负责把任务写入build_jobs表并立即返回(202 Accepted),真正的构建由后台的 build worker 异步执行;前端通过轮询agent_builds状态列来获取进度。
核心设计要点:
build_jobs表即队列:没有独立的队列服务,队列就存在于数据库行中,天然具备持久化与事务能力。FOR UPDATE SKIP LOCKED即分布式锁:多个 server replica 并发抢任务时,只有第一个能锁定某一行,其余自动跳过。- 零外部依赖:不引入 Redis、Celery 或任何外部 broker(build_worker.rs 的依赖只有
sqlx::PgPool、tokio::sync::mpsc与uuid)。
从源码注释(build_worker.rs)可以看到,worker 主循环run在服务启动时被 spawn 一次(每个 replica 一个),接收新任务通知的同时,以 5 秒为周期兜底轮询、以 10 分钟为周期清扫卡死任务。
数据库 schema:一张表就是一个队列
build_jobs表定义在 migrations/0001_schema.sql,核心 DDL 如下:
CREATE TABLE build_jobs ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), agent_id UUID REFERENCES agents(id) ON DELETE CASCADE, owner_id UUID NOT NULL, payload JSONB NOT NULL, -- BuildJobPayload (tagged enum) status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending','in_progress','done','failed')), attempt INTEGER NOT NULL DEFAULT 0, error_msg TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), picked_at TIMESTAMPTZ, -- set when claimed; NULL = pending or done completed_at TIMESTAMPTZ -- set when done or failed ); -- Partial index covers only the statuses the worker queries. CREATE INDEX build_jobs_work_queue ON build_jobs(status, created_at) WHERE status IN ('pending', 'in_progress');几个值得注意的设计决策:
status用 TEXT + CHECK 约束,而不是 Postgres 枚举类型——0003_mcp.sql 的注释解释了这一惯例:避免枚举迁移的麻烦,同时保留约束力。- 部分索引(partial index):索引只覆盖
pending和in_progress两种 worker 实际查询的状态,done/failed的历史行不会拖累索引大小,符合"队列短、历史长"的负载特征。 picked_at语义:任务被认领时写入,为NULL表示从未被认领(pending)或已完成(done)。它是判断任务是否"卡死"的唯一时间依据。ON DELETE CASCADE:Agent 在构建中途被删除时,对应任务行随之消失。执行器随后对该任务的UPDATE影响 0 行(no-op),代码中对此做了明确容忍——build_worker.rs 的注释写着"the UPDATE above is a no-op (0 rows affected) and that's fine"。
一个队列服务两类目标:Agent 与 MCP Connector
build_jobs并不只服务 Agent 构建。在 0003_mcp.sql 中,表被扩展用于 MCP Server 构建——MCP 构建没有agents行,因此增加了connector_id列,并靠一个 CHECK 约束保证每个任务恰好引用一个真实目标:
ALTER TABLE build_jobs ADD COLUMN connector_id UUID REFERENCES mcp_connectors(id) ON DELETE CASCADE, ADD CONSTRAINT chk_build_jobs_one_target CHECK ( (agent_id IS NOT NULL AND connector_id IS NULL) OR (agent_id IS NULL AND connector_id IS NOT NULL) ); CREATE INDEX idx_build_jobs_connector ON build_jobs(connector_id) WHERE connector_id IS NOT NULL;对应地,worker 侧的结构体BuildJob(build_worker.rs)中agent_id与connector_id均为Option<Uuid>,由该 CHECK 保证二者不会同时为Some。后续的恢复逻辑(recover_stuck_jobs、reset_panicked_job)也据此分支处理两种目标。
任务负载变体:BuildJobPayload标签化枚举
任务的详细参数并不拆成多张表,而是整体序列化进payload这一列 JSONB。BuildJobPayload是一个#[serde(tag = "kind")]枚举,定义在 upload.rs,kind标签决定了 worker 的派发路径(dispatch path),所有字段都内嵌在 payload 中——worker 执行时不需要回查数据库。
文档中的变体表如下:
| Variant | Trigger | Executor |
|---|---|---|
Upload | New agent uploaded via zip | execute_upload_and_deployinupload.rs |
Update | Agent source updated (new version) | execute_agent_updateinupdate.rs |
Rollback | Agent rolled back to a previous version | execute_agent_rollbackinupdate.rs |
StandaloneBuild | Build from a GitHub URL (no zip) | execute_buildinbuild/routes.rs |
对照 upload.rs 的实际定义,当前仓库中枚举实际包含7 个变体(文档表格之外还有 3 个):
| Variant | 触发场景 | 执行器 | 关键负载字段 |
|---|---|---|---|
Upload | 通过 zip 上传新 Agent(POST /api/agents/upload-and-deploy) | execute_upload_and_deploy(upload.rs) | zip_path、image_tag、ports、env、writable、writable_path |
Update | Agent 源码更新(PUT /api/agents/{id}/update) | execute_agent_update(update.rs) | new_version、prev_version、changelog、zip_path: Option<String>(None表示 GitHub 重新部署) |
Rollback | 回滚到先前版本(POST /api/agents/{id}/rollback) | execute_agent_rollback(update.rs) | target_version、target_image_tag、reason |
StandaloneBuild | 从 GitHub URL 构建镜像但不部署(POST /api/build/builds) | execute_build(build/routes.rs) | github_url、source_key、version_tag |
Clone | 旧的 GitHub clone-and-deploy(POST /api/github/clone) | execute_clone_and_deploy | 源码注释明确:legacy 变体,实际已不再入队,仅为兼容在途任务保留 |
GithubClone | GitHub clone-and-deploy(当前版本) | execute_github_clone_and_deploy | repo_full_name、branch、version_override、prior_version/prior_image/prior_status |
McpServerUpload | MCP Server 上传/从 GitHub 构建(POST /api/mcp/connectors/upload等) | execute_mcp_server_build(mcp/build.rs) | connector_id、source: McpBuildSourcePayload(Zip/Github)、加密的env |
每个变体都携带一个build_id: Uuid——即agent_builds(或mcp_connector_builds)表的主键。执行结束后,worker回读*_builds表中的status来判断成败,而不是依赖执行器的返回值;执行器把结果写进*_builds行,worker 再据此终结build_jobs行(build_worker.rs)。这一"结果以数据库为准"的设计让任务状态与构建状态天然一致。
值得注意的负载字段设计细节:
Update::owner_id与Update::agent_owner_id是两个不同的字段:前者是触发更新的调用者(superuser 或 ACL 授权者可能更新他人的 Agent),仅用于状态/审计追踪;后者才是 Agent 的真正属主,必须注入 LLM router JWT 并持久化到agent_deployments.owner_id,保证 Agent 始终解析真实属主的 secrets 与 llm_config(upload.rs)。- 带
#[serde(default)]的字段(如writable_path)保证在字段新增之前入队的旧任务行仍能正常反序列化,体现向后兼容设计。 McpServerUpload的env在落库前按属主密钥加密,worker 执行时通过decrypt_build_secrets解密,明文绝不以静止状态存在于 payload 中(build_worker.rs),与 Agent LLM secrets "执行时注入、不持久化"(RUN-5)的既有原则一脉相承。
端到端生命周期
文档给出了完整的生命周期示意,结合源码可以还原每一个环节:
HTTP handler │ INSERT INTO build_jobs (payload, status='pending') │ build_tx.send(()) ← wake notification (capacity-64 mpsc) │ return 202 Accepted ▼ build worker (background) ├─ select! { notify | 5s poll | 10min recovery_tick } │ └─ drain loop claim_next_job() ← SELECT … FOR UPDATE SKIP LOCKED + UPDATE in_progress │ Ok(None) → break (queue empty) │ Ok(Some) → continue ├─ attempt cap check ← fail immediately if attempt ≥ MAX_ATTEMPTS (3) └─ tokio::task::spawn(execute_claimed_job()) │ Ok(()) → check for next job │ is_panic() → reset_panicked_job() immediately └─ cancelled → break (server shutting down)入队端:写库 + 唤醒 + 立即返回
以 Agent 更新为例(update.rs),HTTP handler 的入队动作是原子化的三步:
sqlx::query("INSERT INTO build_jobs (agent_id, owner_id, payload) VALUES ($1, $2, $3)") .bind(agent_id) .bind(owner_id) .bind(serde_json::to_value(&payload).expect("serialize update payload")) .execute(&state.db) .await?; let _ = state.build_tx.send(()).await; // 唤醒 worker ( StatusCode::ACCEPTED, Json(UpdateAgentResponse { build_id, agent_id, new_version, previous_version: current_version, status: "queued", }), ).into_response()客户端立刻收到202 Accepted和status: "queued",随后即可通过agent_builds轮询状态。当前仓库中所有入队点(build_tx.send)位于:
- upload.rs(Upload)
- update.rs(Update)与 update.rs(Rollback)
- build/routes.rs(StandaloneBuild)
- mcp/build.rs(McpServerUpload)
(原文档给出的行号为upload.rs:365、update.rs:287/697、build/routes.rs:211,与当前仓库代码存在偏差,应以实际代码位置为准。)
消费端:select! 三路唤醒 + drain 循环
worker 主循环(build_worker.rs)用tokio::select!同时等待三种信号,任一信号触发后都进入 drain 循环——连续认领任务直到队列为空,而不是只处理一个就回去睡觉,避免突发批量任务时每个任务都引入 5 秒延迟。drain 循环内部分两个阶段执行(详见下文"两阶段执行")。
唤醒机制:三个信号兜底
| # | 信号 | 机制 | 作用 |
|---|---|---|---|
| 1 | mpsc 通知 | build_tx.send(()),channel 容量 64,发送为 fire-and-forget(let _ = build_tx.send(()).await) | HTTP handler 入队后立即唤醒 worker;channel 满说明 worker 已在 drain,丢弃唤醒也无妨 |
| 2 | 5 秒兜底轮询 | tokio::time::sleep(Duration::from_secs(5)) | 捕获丢失的通知(如发送瞬间 channel 恰好满员),把尾部延迟控制在 5 秒内,无需心跳列 |
| 3 | 10 分钟恢复 tick | tokio::time::interval_at,首 tick 在启动后 10 分钟(启动时已内联执行过恢复),MissedTickBehavior::Skip | 周期性调用recover_stuck_jobs,让崩溃 replica 遗留的任务在多 replica 集群中也能被回收,即使没有任何 replica 重启 |
channel 与 spawn 的装配点在 state.rs:mpsc::channel(64)创建(build_tx, build_rx)对,随后tokio::spawn(crate::agents::build_worker::run(worker_state, build_rx))。
两阶段执行:claim 与 execute 的 panic 隔离
构建执行器可能 panic——比如编解码器里的一个unwrap、原生代码中的坏指针。若不做隔离,一个 panic 的构建会击穿整个 worker 协程,把所有排队任务永久卡死。
修复方案是把claim(认领)与execute(执行)分离:
- Phase 1(claim):在 worker 外层任务中直接运行,只做最小化的数据库读写,不存在 panic 风险;
- Phase 2(execute):被包进
tokio::task::spawn,panic 只杀死这个子任务,不会波及 worker 循环。
核心代码(build_worker.rs):
// Phase 1 — claim (outer task, no spawn) let job = claim_next_job(&state).await?; let job_id = job.id; let old_attempt = job.attempt; // pre-increment; DB now holds old_attempt + 1 // Phase 2 — execute (spawned task; panic here does not kill the worker) match tokio::task::spawn(async move { execute_claimed_job(state_clone, job).await }).await { Ok(()) => { /* done */ } Err(ref e) if e.is_panic() => reset_panicked_job(&state.db, job_id, old_attempt).await, Err(_) => break, // task cancelled (server shutdown) }关键细节:job_id和old_attempt在spawn 之前就被捕获到外层作用域,因此 panic 分支可以立即基于它们执行恢复,无需任何共享状态——这正是两阶段分离的价值:恢复动作不需要等待 60 分钟的卡死阈值,而是毫秒级即时响应。
claim_next_job本身(build_worker.rs)是一个事务:
let mut tx = state.db.begin().await?; let job = sqlx::query_as::<_, BuildJob>( "SELECT id, agent_id, connector_id, payload, attempt FROM build_jobs WHERE status = 'pending' ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 1", ).fetch_optional(&mut *tx).await?; let Some(job) = job else { tx.rollback().await?; return Ok(None); // 队列为空 }; sqlx::query( "UPDATE build_jobs SET status = 'in_progress', picked_at = now(), attempt = attempt + 1 WHERE id = $1", ).bind(job.id).execute(&mut *tx).await?; tx.commit().await?; Ok(Some(job))认领与标记in_progress+ 自增attempt在同一事务内完成,保证"认领即占用"的原子性。按created_at排序实现 FIFO 公平调度。
尝试计数:attempt的完整语义
build_jobs.attempt记录一个任务被认领的次数,它的生命周期有三处关键节点:
- 递增:
claim_next_job在把任务标记为in_progress的同一事务中执行attempt = attempt + 1。返回的BuildJob.attempt是递增前的值;此时数据库里已是attempt + 1。 - 上限检查:
MAX_ATTEMPTS = 3(build_worker.rs)。认领后立即检查——若old_attempt >= 3,任务被标记failed且不执行(build_worker.rs)。 - panic 重置:
reset_panicked_job接收old_attempt;若old_attempt + 1 >= 3则永久失败,否则重置为pending等待立即重试(build_worker.rs)。
一个容易忽略的细节:当任务因超过上限被永久标记failed后,它会掉出recover_stuck_jobs的清扫范围(该查询只扫描仍为in_progress的行),因此 worker 必须同步把 Agent/MCP Connector 驱动到终态(调用fail_agent_terminal/fail_mcp_connector_terminal),否则目标会永远停留在building/deploying的非终态,等待中的 SSE/轮询将无限期挂起。这正是 build_worker.rs 注释中 RUN-4 修复的回归点。
恢复策略:即时 + 周期,两层互补
两套机制协同处理卡死任务:
即时恢复 ——reset_panicked_job
由 drain 循环的is_panic()分支触发,在 panic 发生后毫秒级动作:
old_attempt >= MAX_ATTEMPTS → mark 'failed' otherwise → reset to 'pending' (picked_at = NULL)永久失败路径同样会调用fail_agent_terminal终结 Agent 状态,避免agent_builds行停留在building。
周期恢复 ——recover_stuck_jobs
在启动时和每 10 分钟运行一次,处理两类场景:崩溃 replica 遗留的in_progress任务,以及从未 panic 但也没完成的任务(网络分区、OOM、构建中途 SIGKILL)。阈值STUCK_JOB_MINS = 60(build_worker.rs),是最大构建超时(Docker 与 K8s 均默认 30 分钟)的 2 倍,给大镜像留足余量、避免误伤。对应 SQL:
-- Permanently fail exhausted jobs UPDATE build_jobs SET status = 'failed', error_msg = 'max attempts exceeded', completed_at = now() WHERE status = 'in_progress' AND picked_at < now() - make_interval(mins => 60) AND attempt >= 3; -- Reset remaining stuck jobs for retry UPDATE build_jobs SET status = 'pending', picked_at = NULL WHERE status = 'in_progress' AND picked_at < now() - make_interval(mins => 60) AND attempt < 3;实现要点(build_worker.rs):
- 第一条查询带
RETURNING agent_id, connector_id,这样被永久失败的构建能立即驱动目标进入终态; - 两条语句都用
make_interval(mins => $2)绑定STUCK_JOB_MINS常量,阈值只存在于一处 Rust 常量中,避免在两条 SQL 字符串里重复硬编码; - 第二条查询用
rows_affected()记录被重置的任务数,仅在有实际影响时输出 warn 日志。
恢复矩阵
| 场景 | 恢复途径 | 延迟 |
|---|---|---|
| 执行器 panic | reset_panicked_job(is_panic()分支) | 即时 |
| 服务崩溃于构建中途 | 下次启动时recover_stuck_jobs(内联执行) | 下次重启 |
| 多 replica:peer 崩溃且无重启 | 每 10 分钟recover_stuck_jobs | ≤ 70 分钟 |
| 通知丢失(channel 满) | 5 秒兜底轮询 | ≤ 5 秒 |
多副本安全:SKIP LOCKED 与幂等恢复
FOR UPDATE SKIP LOCKED保证同一时刻只有一个 replica 能认领某个任务。两个 replica 同时调用claim_next_job时,后到者会跳过已被锁定的行,转而认领另一个 pending 任务,或返回Ok(None)。配合事务(先 SELECT 锁定、再 UPDATE 标记),认领过程无竞态。
recover_stuck_jobs的两条 UPDATE 也可以安全地并发执行:多个 replica 会对同一批in_progress且超时的行执行相同的UPDATE … WHERE,Postgres 按行 last-write-wins——第二个 replica 对第一个已重置的行执行更新是 no-op,天然幂等。
另外一个隐性保障:worker 循环在认领阶段对数据库错误(Err(e))的处理是记录 error 日志并 break,等待下一轮信号再次尝试,不会因单次 DB 抖动而崩溃。
启动与关闭语义
启动:run()在 state.rs 被tokio::spawn一次(每个 replica 一个实例)。recover_stuck_jobs在进入第一个select!之前内联执行——确保上一个 replica 遗留的任务在第一个 HTTP 请求到达前就已被重新排队,最大化恢复速度。
关闭:服务收到 SIGTERM 后,Tokio 会 dropmpsc::Sender(build_tx)。worker 的notify.recv()分支返回None,当前 drain 循环迭代结束后干净退出(build_worker.rs)。正在执行的tokio::task::spawn任务会运行至完成(Tokio 默认关闭行为),因此 SIGTERM 前一刻开始的构建仍会完成,不会半途丢失。
关键常量速查
| 常量 | 值 | 用途 |
|---|---|---|
MAX_ATTEMPTS | 3 | 永久失败前的最大认领次数 |
STUCK_JOB_MINS | 60 | 判定任务卡死的时间阈值(2× 最大构建超时) |
| Channel 容量 | 64 | mpsc::channel(64)——worker 已在 drain 时丢弃唤醒信号 |
| 兜底轮询 | 5 s | tokio::time::sleep(5s)——捕获丢失的通知 |
| 恢复间隔 | 10 min | interval_at+MissedTickBehavior::Skip(首 tick 在启动后 10 分钟,启动时已内联恢复) |
这些常量集中在 build_worker.rs 文件顶部,配以注释说明设计理由(如STUCK_JOB_MINS必须超过 runtime 的build_timeout,2× 余量避免大镜像误判)。
如何新增一种任务类型
worker 循环本身是任务类型无关的——它只负责认领、派发、恢复,不关心具体业务。新增类型只需 4 步(BUILD_WORKER.md 原文 + 源码印证):
- 在 upload.rs 的
BuildJobPayload中添加一个变体,#[serde(tag = "kind")]会自动写入新的kind标签;记得给新增字段加#[serde(default)]以兼容旧行。 - 编写一个 async 执行器函数,接收
AppState(或其拆分字段);从现有执行器(如execute_agent_update、execute_github_clone_and_deploy)可以看到,执行器内部直接读写state.db、state.runtime、state.config等,构建结果写回*_builds表。 - 在 build_worker.rs 的
execute_claimed_job中新增一个 match 分支调用该执行器;分支开头可解构负载字段,build_id用于最终回读构建状态。 - 在触发构建的 HTTP handler 中:
INSERT INTO build_jobs(payload 序列化为 JSONB),随后state.build_tx.send(()).await唤醒 worker,返回202 Accepted。
不需要修改 worker 循环本身——认领、attempt 上限、panic 隔离、卡死恢复对任何新变体自动生效。单元测试层面,可参照 build_worker.rs 中已有的两个测试:mcp_server_upload_payload_round_trips_through_json验证 payload 经 JSONB 序列化往返后build_id与label不变;decrypt_build_secrets_recovers_encrypted_values_and_drops_bad_ones验证加密 env 的解密与坏密文静默丢弃行为。
小结
nasiko 的 build worker 是一个"以数据库为骨、以 Tokio 为肌"的精巧设计:用一张build_jobs表承载队列状态机,用FOR UPDATE SKIP LOCKED化解多副本竞争,用 mpsc + 轮询 + 周期清扫三路信号覆盖从即时唤醒到崩溃恢复的全场景,用 claim/execute 两阶段分离实现 panic 隔离,再用attempt计数与 60 分钟阈值把"卡死"收敛为可预测的终态。它证明了在 Agent 编排这类需要持久化、可靠性与多副本一致性的场景下,Postgres 本身就可以是足够好的任务队列——不需要为队列引入额外的中间件。
【免费下载链接】nasiko
Developer Control Plane for your AI Agents
相关推荐
PDFMathTranslate异步翻译实现:基于Celery的任务队列设计
PDFMathTranslate异步翻译实现:基于Celery的任务队列设计 在处理大型PDF学术论文翻译时,同步翻译模式常导致页面卡顿和超时问题。本文将详解P
AI 应用人工智能NLPOCRPostGraphile 后台任务(Background Tasks)实战指南:基于 Graphile Worker 的异步任务队列方案
PostGraphile 后台任务(Background Tasks)实战指南:基于 Graphile Worker 的异步任务队列方案 本文是 PostGra
后端API网关MoneyPrinter 架构解析:基于 Postgres 队列与 API/Worker 分离的视频生成流水线设计
MoneyPrinter 架构解析:基于 Postgres 队列与 API/Worker 分离的视频生成流水线设计 导读 本文以 MoneyPrinter 仓库
后端人工智能大模型本地部署媒体生成音视频
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考