错误、幂等与 SSE
处理 problem details、run 失败、事件续传与客户端重试
Problem details
非 2xx 响应使用 application/problem+json:
{
"type": "about:blank",
"title": "Quota Exceeded",
"status": 429,
"detail": "tenant queued-run quota exceeded",
"code": "quota_exceeded",
"request_id": "01J...",
"retryable": true
}| 字段 | 用途 |
|---|---|
status | HTTP 类别 |
code | 稳定机器码,业务分支的首选字段 |
request_id | 与服务端日志/trace 关联 |
retryable | 当前操作是否可使用退避重试 |
detail | 人类诊断信息,不要解析 |
常见稳定码包括 idempotency_conflict、quota_exceeded、join_timeout 以及由 schema、认证、并发和资源状态映射的 problem code。错误码新增时客户端应安全地归入未知错误,而不是失败解析。
Steering 相关错误码
code | HTTP | 触发条件 |
|---|---|---|
run_terminal | 409 | 对处于终止态(succeeded/failed/cancelled/timed_out/dead_letter)的 run 调用 /steer;不产生任何 durable 写入 |
run_superseded | 409 | 对已经被 resume 取代的旧 run_id 调用 /steer 或 /cancel;应改用新 run 的 id(见 metadata.resumed_from_run_id)。对 /cancel 而言,即使旧 run 的 status 仍是 paused 也会触发此错误——检查的是类型化的 superseded_by_run_id 列,而非 status,也不是可被用户写入的 metadata,因此历史 run 绝不会被误改为 cancelled,而其被 resume 出来的新 run 却仍在继续执行;该检查也无法通过客户端提交的 metadata key 伪造(详见运维文档"superseded_by_run_id 血缘标记"一节) |
payload_too_large | 413 | payload + metadata 序列化后超过大小上限(32KB,或服务端事件大小上限,取两者较小值);不产生任何 durable 写入 |
run_resume_conflict | 409 | 并发 resume 调用同一个已暂停的 run;只有一个调用会成功创建新 run,其余调用都会收到该错误 |
run_finalizing | 409 | run 的图已经执行结束、但最终状态尚未提交(worker 正在持久化刷新并收尾的短暂窗口)时调用 /steer;此时已经没有任何安全点可以消费该事件,因此直接拒绝,而不是接受一条永远不会被消费的记录。paused 结局不受影响——对已暂停的 run 发起 /steer 仍然正常 |
Run 业务错误
创建 run 成功返回 202 pending。之后的节点/schema/provider 错误会使资源变为 failed、timed_out 或 dead_letter,并写入:
{
"status": "failed",
"error": {
"code": "invalid_update",
"message": "...",
"retryable": false
}
}查询和 join 请求本身仍可返回 200。客户端必须同时检查 HTTP 状态与 run.status。
SSE 格式与续传
GET /v1/runs/{run_id}/stream
Accept: text/event-stream
Last-Event-ID: 17id: 18
event: node_completed
data: {"id":"...","run_id":"...","sequence":18,"kind":"node_completed","data":{},"created_at":"..."}- 事件在发送前写入 PostgreSQL;
sequence在 run 内从 1 单调递增;- 重连时发送最后已处理的
Last-Event-ID,服务端从下一条继续; - 客户端按
(run_id, sequence)去重,处理重复交付; - 以
: heartbeat开头的行是保活注释,应忽略; - run 到达 terminal 或
paused后,服务端关闭流。
Steering 生命周期事件
除了图执行过程中产生的普通事件(node_completed 等),运行中输入(steering)在其生命周期
中会产生三类专门的 SSE 事件:
kind | 何时产生 | 关键字段 |
|---|---|---|
run.steer.accepted | /steer 调用被 durably 接受时(202) | steering_event_id、source_event_id(若来自 resume 迁移)、sequence、kind |
run.steer.consumed | 图代码调用 drain_steering() 消费该事件,且这次消费被 commit_steering_consumptions_if_owned() durably 提交后 | steering_event_id、source_event_id、sequence、kind、queue_latency_seconds、node、namespace、task_id |
run.steer.superseded | 某条 steering 事件在从未被消费的情况下,其行 durably 转为 superseded 状态——原因可能是它所属的 paused run 被 resume(其仍 pending 的 steering 被迁移到新 run),也可能是它在该 run 到达终结 finalization 边界时仍处于 pending/delivered(见下文) | steering_event_id、source_event_id、sequence、kind、reason(resume_transfer 或 unconsumed_at_final_boundary)、superseded_by_run_id(resume_transfer 时有值,unconsumed_at_final_boundary 时为 null)、replacement_steering_event_id(resume_transfer 时有值,unconsumed_at_final_boundary 时为 null) |
run.steer.superseded 的 unconsumed_at_final_boundary 原因是一个刻意保持诚实、笼统的
标签——它并不声称该事件在与什么竞争。它涵盖了 worker 到达真正的终结 finalization 边界
时仍处于 pending/delivered 的每一条 steering 行,而这背后的成因不止一种:某次 /steer
调用确实可能恰好在 worker 关闭新准入之前完成 durable 提交(这是真实的竞争);但同样常见的
情况是,图从始至终就没有调用过 drain_steering() 去消费一条其实一直可见的行——这是普通的
应用层未消费,根本不是竞争。Runtime 目前无法从现有数据中区分这两种情况,因此它也不假装能
区分:两者都得到相同的 unconsumed_at_final_boundary 原因。与其仅因为存在一条 steering
事件就强制该 run 进入无限重试循环——这会惩罚那些本就合法地从不调用 drain_steering() 的
图——Runtime 会把这条行 durably 转为 superseded,并以之前已经算出的结果正常终结该 run。
这种情况下 run.steer.superseded 事件绝不会被静默丢弃:它总是与 run 的终结状态在同一个
durable 事务中一起提交,并且在重试下是幂等的(对已经 superseded 的行重复尝试转换不会
产生重复事件)。
字段说明:
steering_event_id:被消费的 steering 事件在当前 run 下的 id(如果该事件是从旧 run 迁移过来的,这是迁移后的新 id,不是客户端最初收到的 id);source_event_id:客户端最初/steer调用收到的根 id。同一根 steering 即使经过多次 pause/resume 迁移,该字段在每一跳都保持一致,客户端可用它把最终run.steer.consumed/run.steer.superseded与自己最初发出的请求关联起来。这一行为按具体 生命周期事件而定,并非统一规则:在run.steer.accepted与run.steer.consumed上,仅当 事件来自迁移时才出现该字段——直接提交(从未被迁移)的事件在这两类事件上没有此字段;run.steer.superseded则不同:它总是把source_event_id设为event.source_event_id or event.id,因此即便是一条直接提交、从未迁移过的事件,其 superseded 记录里也会带有一个根 id——就是它自己的 id;sequence:该 steering 事件在其所属 run 的 steering 序列里的序号(与 run 事件流的sequence是两个独立计数器);kind:透传自/steer请求体的kind;queue_latency_seconds:从该 steering 最初 durably accepted(包括跨 resume 累积的 等待时间)到被消费的秒数;node:消费该事件的节点名;namespace:该节点所在的 subgraph 命名空间路径(数组);task_id:区分同一 superstep 内并行/Send任务的具体 task。
durable source 是 PostgreSQL,Redis 只是可选加速器。 run.steer.accepted、
run.steer.consumed 和 run.steer.superseded 都是先 durably 写入 PostgreSQL 再对外可见的;Redis pub/sub 只用于
提示 worker "可能有新 steering 了",从而降低轮询延迟。即使 Redis 通知丢失(连接断开、
Redis 重启等),worker 仍会在下一个安全点检查/轮询 PostgreSQL 并最终消费——不会因为
Redis 通知丢失而永久遗漏一条 steering。同理,commit_steering_consumptions_if_owned() 是幂等的:
worker 因瞬时故障重试同一批 consumption 时,只有真正发生状态转换的事件才会产生新的
run.steer.consumed,不会为已经消费过的事件重复生成。
重试策略
- 只对
retryable=true、网络故障、429 和临时 5xx 使用有界指数退避; - 遵守
Retry-After; - 创建 run 重试复用相同
Idempotency-Key与完全相同请求; - SSE 重连复用最后 sequence,不创建新 run;
- 不自动 redrive 确定性 schema/tool 权限错误。
请求 key 相同但 body 不同会返回 409 idempotency_conflict。不要在冲突后随机换 key 并静默重试,否则可能掩盖调用方 bug 并产生重复业务操作。