服务端:从 issue 到 run 的判定、任务状态机与路由
30 秒导读: 用户在看板上给一个 issue 换指派人、把它拖出 backlog,或者一条定时规则到点了——服务端要在这一瞬间回答两个问题:这次写操作该不该启动一个 agent run?该给谁跑? 本章讲清楚这个判定怎么做(单一谓词
WillEnqueueRun)、任务被创建后在数据库里怎么走完一生(queued → dispatched → running → terminal状态机,以及多个 daemon 抢同一个任务时怎么不打架),以及一整套「为什么没跑」的稳定错误码怎么在不泄密的前提下解释结果。
本章属于 Multica 讲解系列。上游把活派到服务端之后,真正执行发生在别处: 本地守护进程怎么认领并执行、agent 运行时怎么把十几种 CLI 抽象成一种执行。想先看全局,回到 index。
1. 这是什么(零基础也能懂)
一句话定义
这是 Multica 服务端的「派单中枢」:它决定每一次会改动 issue 的写操作要不要变成一个 agent 的活儿,把活儿写进一张任务队列表,然后管这张表里每条任务从排队到终结的全过程。
它解决谁的什么问题
把 agent 当成团队里的「人」来用,就会撞上一堆调度问题:
- 我把 issue 从 A agent 改派给 B agent,B 要不要立刻开跑?改派本身算不算一次「开工信号」?
- 一个 agent 同时最多能跑几个任务?超了怎么办?
- 两台 daemon(两台开发者的笔记本)同时在线,同一个任务会不会被两边同时抢走、跑两遍?
- 一条评论触发 agent 跑,agent 跑完又发评论——会不会自己触发自己,无限循环?
- 一个「私有 agent」别人看不见,那当别人试图触发它时,服务端拒绝的理由要怎么写才既能解释又不暴露「这个 agent 存在」?
这一章就是这些问题的服务端答案。
它管的东西:一张队列表
所有「agent 要干的活」都落在一张表 agent_task_queue 里。一行就是一个任务(task),核心字段:
| 字段 | 含义 |
|---|---|
status | 生命周期状态(queued / dispatched / running / completed / failed / cancelled …) |
agent_id | 谁来跑 |
runtime_id | 跑在哪个 runtime(哪台 daemon)上 |
issue_id / chat_session_id / autopilot_run_id | 这活儿的来源(三选一,或都空 = quick-create) |
priority | 优先级,认领时高优先先出队 |
originator_user_id / accountable_user_id | 谁的授权、谁担责(见 §8 attribution) |
一句话直觉
把它想成餐厅后厨的挂单夹:前台(HTTP 写操作)判断「这单要不要下给厨房」,判断为是就把小票夹上挂单夹(INSERT 一行 queued 任务);厨师(daemon)来取单时,一把把小票从夹子上「原子地」揭下来(UPDATE ... status='dispatched'),同一张小票不可能被两个厨师同时揭走。
2. 顶层全景(它大概怎么转)
先看一次「人给 issue 改状态 → agent 真的跑起来」的完整链路。怎么读这张图:从上到下是时间顺序;左半是「判定+入队」,右半是「认领+执行」,中间隔着一张数据库表。
┌──────────────────────── 服务端(本章) ────────────────────────┐
HTTP │ │
写操作 │ ①判定「要不要跑/给谁跑」 ②把活写进队列表 │
────► │ WillEnqueueRun(单一谓词) ──► CreateAgentTask │
(改派/ │ · 私有 agent 门禁 INSERT status='queued' │
改状态)│ · self-loop 抑制 │ │
│ · pending 去重 ▼ │
│ ┌───────────────────────┐ │
│ │ agent_task_queue │ │
│ Autopilot / cron / webhook │ (一行 = 一个 task) │ │
│ ────► DispatchAutopilot ──►│ queued│dispatched│... │ │
│ (自动建 issue/直发) └───────────────────────┘ │
└───────────────────────────────────────┬──────────────────────┘
│ daemon 轮询 / 被唤醒
▼
┌──────────────────── daemon 侧(见 02 章) ─────────────────────┐
│ ③认领:ClaimTask / ClaimTasksForRuntimes │
│ UPDATE ... status='dispatched' ← 原子 CAS,抗并发抢占 │
│ ④上报:StartTask→running→CompleteTask/FailTask(终态) │
└───────────────────────────────────────────────────────────────┘
部件一句话职责
| 部件 | 干什么 | 在哪 |
|---|---|---|
WillEnqueueRun | 单一谓词:这次 issue 写操作要不要跑、给谁跑 | internal/service/issue_trigger.go:89 |
IssueTriggerProbe | 把「私有 agent 门禁」「self-loop 判断」这些请求级检查注入谓词 | internal/service/issue_trigger.go:36 |
CreateAgentTask | 真正 INSERT 一行 queued 任务 | enqueueIssueTask, internal/service/task.go:973 |
ClaimTask | 单 agent 认领:容量判定 + 原子出队 | internal/service/task.go:2009 |
ClaimTasksForRuntimes | 一次给一整台机器的多个 runtime 批量认领 | internal/service/task.go:2296 |
ReasonCode | 跨层稳定错误码枚举,解释「为什么没跑」且不泄密 | internal/dispatch/reason.go |
DispatchAutopilot | 定时/webhook 触发:准入 → 建 run → 建 issue 或直发任务 | internal/service/autopilot.go:114 |
applyAttributionFallback | fail-closed:解析不出担责的人就拒绝入队 | internal/service/task.go:573 |
主线走一遍(高层)
- 用户 PATCH 一个 issue(改派或改状态)。Handler 先在 HTTP 边界把 UUID 解析、把私有 agent 门禁走一遍(§9)。
- Handler 调
WillEnqueueRun(§3):传入「这个 issue 写完之后长什么样」,谓词回答(给谁跑, 要不要跑)。 - 要跑 →
dispatchIssueRun→EnqueueTaskForIssueWithHandoff→CreateAgentTask落一行queued(途中做 attribution,§8)。 - 广播
task:queued,唤醒对应 runtime 的 daemon。 - daemon 调
ClaimTask/批量认领(§4):容量够就用一条UPDATE把任务原子地从queued改成dispatched。 - daemon 准备好后
StartTask(→running),跑完CompleteTask/FailTask(→终态)。
3. 单一真理:WillEnqueueRun 谓词
它要解决的小问题
「改派 / 改状态 / 新建」三种 issue 写操作,过去各自有一段「要不要开跑」的判断代码,四个入口(单条更新、批量、创建、preview 预览)慢慢长歪了:有的漏了 squad 分支、有的漏了 self-loop、口径不一致。结果:preview 告诉你「会启动 2 个 run」,真写的时候只跑了 1 个。这就是 MUL-3375。
思路
把判断收敛成一个纯谓词,让所有写路径和 preview 端点逐字共用同一个函数。谓词只回答两件事:会不会跑(bool)、给谁跑(IssueRunTrigger)。凡是判断,一个地方改,所有入口一起动。
依据:internal/service/issue_trigger.go:64-89(函数注释明确说它替代了「四个入口飘移」的旧实现)。
输入怎么建模
谓词不直接吃 HTTP 请求,而是吃一个「写完之后的 issue 快照 + 哪几个 字段被动过」:
// internal/service/issue_trigger.go:44 —— 一次 issue 写操作的抽象
type IssueTriggerInput struct {
Issue db.Issue // 写「之后」的形态
PrevStatus string // 写「之前」的状态
IsCreate bool // 全新 issue(无旧任务可取消、无 self-loop)
AssigneeChanged bool // 这次写动了指派人吗
StatusChanged bool // 这次写动了状态吗
}
判定分两级:先看「哪种开工信号」,再看「目标能不能跑」
第一级——这次写算不算一次开工信号(source),只有两种能开工:
| source | 触发条件 | 关键规则 |
|---|---|---|
assign | 创建 or 改派(`IsCreate | |
status | 已指派的 issue 从 backlog 提升到活动状态 | 且新状态不是 done/cancelled;要过 self-loop 抑制 |
两者都不是 → 直接返回「不跑」。依据:WillEnqueueRun 的 switch,internal/service/issue_trigger.go:99-115。
第二级——目标(agent 或 squad)当下能不能跑。这里三道闸门顺序卡:
目标是 agent / squad?
│
▼
①目标存在且可跑? agent: RuntimeID 有效 且 未 archived
│ squad : leader 过 AgentReadiness(runtime online)
▼
②私有 agent 门禁? canAccess(agent) —— 见下「门禁为什么在这留个钩子」
│
▼
③(仅 status 源)已有 pending 任务? hasPendingRun → 有就不跑(去重)
│
▼
给谁跑:agent 自己 / squad 的 leader
依据:agent 分支 internal/service/issue_trigger.go:117-134;squad 分支 :136-163。
私有 agent 门禁:为什么在谓词里留个钩子却传「全放行」
IssueTriggerProbe.CanAccessAgent 是私有 agent 的访问闸(internal/service/issue_trigger.go:37)。微妙点在于写路径和 preview 传的东西不一样:
- 写路径:门禁其实已经在 HTTP 边界执行过了(改派时
validateAssigneePair,squad 时canEnqueueSquadLeader),所以写路径给谓词传一个「全放行」的探针(CanAccessAgent: nil→ 被当作allowAllAgents),避免把同一道闸重复跑、或把它下沉进 service 层。依据:internal/handler/issue_trigger.go:24-31的issueTriggerWriteProbe。 - preview 端点:preview 不经过写边界的那道闸,所以它必须给谓词传真的门禁,否则 preview 会把「一个私有 agent 现在能跑」这个事实泄露给一个根本看不见它的成员。依据:
internal/handler/issue_trigger.go:37-47的issueTriggerPreviewProbe,里面调canInvokeAgent。
一句话:门禁只判一次,但 preview 因为绕过了写边界,得自己补上这一判。
self-loop 抑制:agent 别把自己再触发一遍
只有 status 源会问这个问题(create/assign 从不问)。判断逻辑:如果发起这次「提升出 backlog」写操作的,正是那个已经在这个 issue 上跑着的 agent 自己,就抑制,不再开新 run。
- 谓词侧:
if probe.IsSelfLoop != nil && probe.IsSelfLoop() { return 不跑 },internal/service/issue_trigger.go:109-111。 - 实现侧:探针的
IsSelfLoop指向 handler 的isAgentRunningOnIssue——它读请求头X-Task-ID,查这个 task 的issue_id是否就是当前 issue。非 agent actor 直接返回 false。依据:internal/handler/issue.go:3121-3141。
pending-task 去重:和数据库唯一索引对齐
status 源在开跑前还要问「这个 agent 是不是已经有一个 pending 任务在这个 issue 上了」。这不是随手加的检查,而是镜像数据库那条 partial 唯一索引:
-- migrations/037_fix_pending_task_unique_index.up.sql
CREATE UNIQUE INDEX idx_one_pending_task_per_issue_agent
ON agent_task_queue (issue_id, agent_id)
WHERE status IN ('queued', 'dispatched');
也就是说:同一个 (issue, agent) 最多只能有一条「未开跑」的任务。谓词里 hasPendingRun 提前判这个,好让 preview 不承诺一个「唯一索引反正会合并掉」的假 run。为什么只有 status 源判、assign 源不判?
status源(backlog→活动)可能撞上一个「issue 还在 backlog 时被 @提及 塞进来的 pending 任务」,得判,免得 preview 空承诺。assign源不判:创建面对全新 issue 没有旧任务;改派已经不再取消已有任务(#4963/MUL-4113),真撞上就靠唯一索引让INSERT静默 no-op,该 agent 照样有一个 pending run。
依据:注释与实现 internal/service/issue_trigger.go:76-88 与 :126、:155;hasPendingRun 本体 :171(注意它失败时 fail-closed 到「有 pending」,让 preview 宁可少承诺)。
一处诚实说明:
WillEnqueueRun是「单条更新 / 批量更新」写路径 + preview 的单一真理(internal/handler/issue.go:2962、:3460两处逐字共用同一谓词与探针)。issue「创建」的写侧入队走的是另一条并行镜像service.IssueService.maybeEnqueueOnAssign(internal/service/issue.go:481,内部用shouldEnqueueAgentTask/shouldEnqueueSquadLeaderOnAssign);WillEnqueueRun的IsCreate分支主要服务 preview 端点的创建预览,与那条镜像保持口径一致。这两条创建路径的口径对齐没有像更新路径那样收敛进同一个函数。
4. task 生命周期与认领竞态
状态机:一条 task 的一生
怎么读:实线是正常流转,虚线是异常/恢复回边。每个转移都由一条带 WHERE status IN (...) 的 SQL 做,状态不对就 0 行、悄悄不动。
(INSERT)
│
▼
┌───────────────► queued ◄─────────────┐
│ requeue │ │ reclaim
│ (finalize 失败) │ ClaimAgentTask │ (响应丢失,
│ ▼ │ 超恢复窗口)
│ dispatched ───────────────►┘
│ │ │
│ (等本地目录锁) │ │ StartAgentTask
│ ▼ ▼
│ waiting_local_directory ──► running
│ │ │ │
│ complete │ │ │ fail
│ ▼ │ ▼
└──────────────────── completed │ failed ──┐
│ │ 可重试原因
cancelled ◄┘ │ (超时/离线…)
▼
建 retry 子任务(同事务)
关键点:状态转移全靠 SQL 的条件 UPDATE,而不是「先读后写」。这是抗并发的根 。
| 转移 | SQL | 守卫条件 | 位置 |
|---|---|---|---|
| 认领 | ClaimAgentTask | status='queued' + 无同 (issue,agent) 活跃任务 | pkg/db/queries/agent.sql:493 |
| 开跑 | StartAgentTask | status IN ('dispatched','waiting_local_directory') | agent.sql:635 |
| 完成 | CompleteAgentTask | status='running' | agent.sql:667 |
| 失败 | FailAgentTask | status IN ('dispatched','running','waiting_local_directory') | agent.sql:739 |
认领的两把锁:容量 + 每-(issue,agent) 串行
ClaimTask 在一个事务里做两件卡:
第一把——容量判定(MaxConcurrentTasks)。 先 FOR UPDATE 锁住 agent 行,再数它当前活跃任务数,满了就直接放弃:
// internal/service/task.go:2029-2040 —— 容量闸
running, err := qtx.CountRunningTasks(ctx, agentID) // 数 dispatched+running+waiting
if running >= int64(agent.MaxConcurrentTasks) {
outcome = "no_capacity"
return nil // 不认领,静默返回
}
CountRunningTasks 数的是 dispatched|running|waiting_local_directory 三态(agent.sql:921)。GetAgentForClaimUpdate 的 FOR UPDATE(agent.sql:925)让并发认领在容量判断上串行——这正是竞态测试 TestClaimTaskConcurrentCapacityRespected 保证的:两个 worker 同抢、max=1,最终恰好 1 个被认领(internal/service/task_claim_race_test.go:41,断言 :91)。
第二把——每-(issue,agent) 串行,写在 SQL 里。 ClaimAgentTask 的 WHERE 用 NOT EXISTS 排除「同 agent 已有 dispatched/running 任务撞在同一个 issue(或同一个 chat_session,或同为 quick-create)上」的情况,再 FOR UPDATE SKIP LOCKED + LIMIT 1 原子出队:
-- pkg/db/queries/agent.sql:507-530(节选骨架)
WHERE id = (
SELECT atq.id FROM agent_task_queue atq
WHERE atq.agent_id = $1 AND atq.status = 'queued'
AND NOT EXISTS ( ...同 issue/chat/quick-create 已有活跃任务... )
ORDER BY atq.priority DESC, atq.created_at ASC
LIMIT 1 FOR UPDATE SKIP LOCKED)
SKIP LOCKED 是并发认领不打架的关键:两台 daemon 同时扫,各自跳过对方已锁的行,永不重复认领同一行。
批量认领:一台机器一次拉多个 runtime
ClaimTasksForRuntimes(internal/service/task.go:2296,MUL-4257)是机器级批量版:一个 daemon 托管多个 runtime,与其「每个 runtime 一次 HTTP」,不如一次请求把整台机器的活都认了。它保持单 runtime 的每一步语义,只是「集合化」:
1. 批量提升到点的 deferred 任务(一条 UPDATE) PromoteDueDeferredTasksForRuntimes
2. 批量回收「响应丢失」的 stale-dispatched(一条 UPDATE) ReclaimStaleDispatchedTasksForRuntimes
—— 必须在空缓存检查之前,因为丢响应会让任务离开 queued 态
3. 空缓存短路 + 采样失效版本号(EmptyClaim)
4. 一条 SELECT 拉出非空 runtime 集的候选
5. 把仍为空的 runtime 标记 empty(下次空轮询跳过 Postgres)
6. 按「不同 agent」逐个走 ClaimTask(复用 §4 的串行+容量+派发副作用)
批量认领的部分成功(partial success):宁可少返回,不可重复认领
批量的难点在于步骤 2/6 各自在独立事务里已经把一些任务派发到服务端了。如果第 4 步或第 6 步中途报错,直接返回 500 会让 daemon 走 HTTP 回退、把同一批空槽再认领一遍——正是这个 PR 要消灭的双重认领。所以规则是:
只要已经认到过任务(
len(claimed)>0),就把已认的部分交回去,报错的候选留在队列等下次轮询。
依据:候选查询失败时 internal/service/task.go:2389-2393;单个 ClaimTask 失败时 :2433-2437。竞态测试 TestClaimTasksForRuntimes_PartialSuccessOnSecondAgentClaimFailure / ...OnCandidateQueryFailureAfterReclaim(internal/service/task_batch_claim_partial_test.go:63、:95)钉住这个行为。
批量里还有一处 runtime 归属守卫:ClaimAgentTask 只按 agent 选行,可能选中一个属于「另一台 daemon 的 runtime」上的更高优任务;批量里显式判 runtimeInSet,不属于本机的跳过(task.go:2449),那个误派发的任务由属主 daemon 下次轮询的 reclaim 路径回收。
claim/complete 竞态与恢复
三种「响应丢失 / 慢处理」的恢复,全靠带 CAS 条件的 SQL,不靠时钟猜:
| 场景 | 机制 | 守卫 |
|---|---|---|
| 认领响应到 daemon 前,payload 组装失败 | RequeueAgentTaskAfterClaimFailure 立刻放回 queued | CAS 带 dispatched_at,老处理器回滚不了新 reclaim(agent.sql:558) |
认领成功但响应never到 daemon(daemon 没 StartTask 确认) | ReclaimStaleDispatchedTaskForRuntime 超恢复窗口后重投 | status='dispatched' 且 started_at IS NULL 且过了恢复窗口(agent.sql:576) |
| 认领慢/准备久 | prepare-lease 短租约续期,ExtendAgentTaskPrepareLease | 只在 dispatched/waiting 且未 start 时续(agent.sql:622) |
恢复窗口 claimResponseRecoveryWindow = 90s 是刻意大于 daemon 两段超时(claim 30s + start 30s)之和(internal/service/task.go:150-153)。CompleteAgentTask 用 status='running' 这 个 CAS 保证「只有真在跑的任务能被完成」——一个被中途 cancel 的任务,complete 会 0 行、悄悄失效(agent.sql:667-671)。
5. admission / dispatch reason codes:既解释又不泄密
它要解决的小问题
当一次触发没有变成 run,得给调用方一个理由。但理由不能用「人类可读的失败字符串」去反向猜(脆、会漂),更不能泄露一个私有 agent 是否存在——否则枚举这些理由就能探出别人藏起来的 agent。
思路:一套跨层稳定枚举,在「决定不跑的那个分支」当场定,原样带到响应
internal/dispatch 是个叶子包(无内部依赖),所以做决定的 service 层和序列化到线上的 handler 层共用同一枚举、永不漂移。一枚 ReasonCode 在阻断/跳过的那个 if 分支就定死,一路原样传出去,绝不从失败文本里逆向工程。依据:internal/dispatch/reason.go:1-12 包注释。
枚举一览
| ReasonCode | 含义 | 关键脱敏语义 |
|---|---|---|
queued / coalesced / deferred | 成功路径三态 | — |
invocation_not_allowed | 发起者无权触发该目标 | 故意不区分「目标私有」和「目标不存在」 |
target_unavailable | 目标跑不了(archived/删了/leader 解析不出/无指派) | — |
runtime_offline | 目标获准但 runtime 未绑定/不在线 | — |
attribution_blocked | fail-closed 工作区解析不出担责的人,拒跑(见 §8) | — |
already_active | 已有活跃/pending run,这次触发没合并 | — |
self_trigger_suppressed | 目标被 self-trigger 守卫刻意不再触发,且无活跃 run 兜底 | 不是权限拒绝,但也不是成功:什么都没新跑 |
internal_error | 意外服务端错误 | — |
依据:internal/dispatch/reason.go:17-48。其中 invocation_not_allowed 的注释白纸黑字写「Deliberately generic — it does not distinguish 'target is private' from 'target does not exist'」(reason.go:24-26)——这就是脱敏的核心。
它怎么被「当场定、原样传」
以 autopilot 手动触发为例,reason code 在准入分支定出,一路返回,handler 直接塞进响应,不碰失败文本:
// internal/handler/autopilot.go:2027-2034 —— 原样带出,不逆向
resp := runToResponse(*run)
if reasonCode != "" {
c := string(reasonCode)
resp.ReasonCode = &c // JSON: "reason_code"
}
而错误到 reason code 的映射用类型判断,绝不子串匹配:dispatchFailReasonCode 用 errors.Is(err, ErrAttributionFailClosed) 判出 attribution_blocked,其余归 internal_error(internal/service/autopilot.go:555-560)。
6. Squad 路由:把活派给 leader,leader 再委派
直觉:把 squad 当成「对外一个人」
Squad(小队)对外是一个可被指派的实体,但真正干活的是它的 leader agent。所以「把 issue 指派给 squad」在服务端等价于「把 issue 派给 leader」,leader 拿到后再在自己的运行里把子活委派出去(委派发生在 agent 运行时,不在本章)。
路由怎么走
WillEnqueueRun 的 squad 分支已经把「给谁跑」解析成 leader:它 GetSquadInWorkspace → 取 LeaderID → 过 AgentReadiness → 过门禁 → 返回 AssigneeType:"squad", AgentID: squad.LeaderID(internal/service/issue_trigger.go:136-163)。真正入队的副作用则统一走 handler 的 enqueueSquadLeaderTask:
issue 指派给 squad
│ dispatchIssueRun(trigger.AssigneeType=="squad")
▼
enqueueSquadLeaderTask (handler/squad.go:1033)
│ ① 取 squad → leader
│ ② canEnqueueSquadLeader(leader, actor, originator) ← 私有 leader 门禁
│ ③ HasPendingTaskForIssueAndAgent(issue, leader) ← 去重
▼
EnqueueTaskForSquadLeaderWithHandoff → 打上 is_leader_task=true / squad_id
依据:internal/handler/squad.go:1033-1087。任务上打 is_leader_task=true 是为了让下游 self-trigger 守卫能区分「agent 是以 leader 身份发的评论(跳过)」还是「以 worker 身份发的(不跳过)」——这对「既是 leader 又是自己 squad 的 worker」的 agent 是必需的(internal/service/task.go:1072-1085)。
门禁:A2A 由「链顶的人」判,而不是直接的 agent
私有 leader 的触发闸 canEnqueueSquadLeader 转手 canInvokeAgent(internal/handler/agent_access.go:394-399)。canInvokeAgent 的精髓:判权用「有效发起人」——member actor 就是他自己,agent/system actor 则用链顶的人类 originator,绝不信直接的 agent principal(internal/handler/agent_access.go:48-54)。这挡住了「用户 U 触发 A、A 再 @ B」 这种 agent 私自搭桥绕过 owner 白名单的路子(:32-35)。
其中一个刻意放的口子:public_to workspace 目标,允许 workspace 内部的 agent/system 自动化在没有人类 originator 时也能触发——但只放 workspace 这一种 target,member/team target 仍然 fail-closed,一个无归属的 agent/system 触发永远匹配不上某人的「指定人」授权(internal/handler/agent_access.go:72-107,MUL-3963)。
为什么 leader 的 originator 要精确解析
enqueueSquadLeaderTask 里门禁判的那个 originator,必须和入队路径最终写到 leader 任务行上的那个人一致,否则会漂:一个 agent 创建、正确继承了 originator 的 issue(MUL-4305),若门禁这里用空 originator 去判,就会被误拒。所以 member 作者就是自己;agent/system 触发则用 OriginatorForIssueTask 按 issue 的 origin 链解析,和 EnqueueTaskForSquadLeader* 用的完全一样(internal/handler/squad.go:1042-1057)。