运行中输入(Steering)的运维
迁移上线顺序、PostgreSQL 作为真相来源、延迟、重试与回滚边界
迁移上线顺序
运行中输入在基础 schema 之上新增了四个 forward-only 的 Alembic revision:
| Revision | 新增内容 |
|---|---|
0002_steering | run_steering_events 表及其索引,以及读写它的 worker/server 代码路径 |
0003_steering_source_event | run_steering_events 上的 source_event_id 列 + 局部索引,用于跨多次 pause/resume 迁移关联同一条 steering |
0004_steering_closed | runs.steering_closed——finalization 准入闸门;详见下方"steering_closed finalization 闸门" |
0005_superseded_by_run_id | runs.superseded_by_run_id——被 resume 覆盖的类型化控制面血缘标记;把它从可被用户写入的 metadata 中移出,使其无法被伪造。详见下方"superseded_by_run_id 血缘标记" |
按标准 Alembic 路径应用:
alembic upgrade head四者都是纯新增(ADD COLUMN IF NOT EXISTS、CREATE INDEX IF NOT EXISTS、可为空/
NOT NULL DEFAULT FALSE 列),可以在部署新的 server/worker 二进制之前先对存活数据库
执行——早于 0004/0005 的旧二进制根本不会读写 steering_closed/superseded_by_run_id,
行为与这些列不存在时完全一致。零停机上线顺序:
- 在旧的 server/worker 二进制仍在提供服务时执行
alembic upgrade head(新增run_steering_events、source_event_id、steering_closed与superseded_by_run_id); - 再滚动发布新的 server 与 worker 二进制。
不要颠倒这个顺序:在列存在之前就尝试写入 source_event_id/steering_closed/
superseded_by_run_id 的二进制会插入/更新失败。因为本项目的迁移是 forward-only 的(downgrade() 会直接抛错),schema
层面没有自动回滚路径——见下方"升级与回滚边界"。
PostgresRepository.setup() 也会直接应用同一批 SQL 文件(测试和嵌入式/开发环境会用到)。
这条路径不能替代在真实部署中驱动 Alembic revision 链:生产上线应该且只应该执行
alembic upgrade head,并且必须走到 0005_superseded_by_run_id,而不是停在
0002_steering、0003_steering_source_event 或 0004_steering_closed。
steering_closed finalization 闸门
一旦 worker 的图执行到达真正的终态结果,该 run 可能还会短暂地继续读作
running/cancelling,同时 worker 把图已经消费掉的 steering durably flush 出去、
提交最终状态——此时已经没有更多图安全点可以消费任何新输入了。runs.steering_closed
是持有该 run 的 worker(与状态写入一样,基于自己的 lease_owner/attempt 加锁)在执行
结束的瞬间、final flush 开始之前设置的一个布尔值。
- 当
steering_closed为false时,POST /steer会像0004之前一样被正常接受; - 一旦当前 attempt 的
steering_closed变为true,一个全新的POST /steer(尚无匹配 Idempotency-Key 记录)会被拒绝,返回409 run_finalizing——这是一个稳定的、 与终态 run 的409不同的独立错误,方便客户端区分"这个 run 正在收尾"和"这个 run 已经彻底结束了"; - 一次重放的
POST /steer——即(tenant_id, run_id, Idempotency-Key)与某个已经 durably 接受过的事件完全相同——永远返回该已存在事件的当前状态,绝不会返回新的409,无论steering_closed/终态/superseded 状态如何。见下方"幂等重放在所有准入闸门 前都是安全的"; paused结果永远不会关闭闸门:暂停并不像终态那样"没有安全点剩下"——对一个 paused run 调用/steer是设计中刻意支持的行为(其 pending steering 会在/resume时迁移到 新 run 上);- 一次被放弃的投递尝试关闭的闸门不会泄漏到下一次尝试:
claim_run在每次新的 claim 上 都会把steering_closed重置为false,所以被重试的 run 会重新正常接受 steering。
仅仅关闭闸门本身并不足以防止一条永久 pending 的 steering 行:/steer 请求有可能在
close_steering() 之前赢得该 run 行的锁——即使此时图早已执行完毕。真正能捕获这种情况的
是 worker 在 flush 之后的检查(finalize_run_with_steering_disposition_if_owned(),见
Worker._execute)——如果 final flush 已经提交了它能提交的一切之后,仍然存在
pending/delivered 的行,该行会被 durably 转为 superseded(原因为
unconsumed_at_final_boundary——这是一个刻意保持笼统的标签,并不声称发生了竞争;它
同样覆盖更常见的情况:图从始至终就没有调用过 drain_steering()),并伴随一条
run.steer.superseded 生命周期事件,二者与 run 的终结状态写入在同一个
lease/attempt-fenced 事务中一起提交——不再是两次分别提交、可能在崩溃或 lease 丢失时
彼此脱节的写入。随后该 run 会以之前已经算出的结果正常终结——不会仅仅因为存在一条
steering 事件就走 retry/dead-letter 循环,因为那样会给一个本就合法地从不调用
drain_steering() 的图强加一次不必要的重跑。一个终态 run 自身的终结决策绝不能让一条
未被消费的 steering 事件被静默搁置、没有任何 durable、可观测的归宿——一次由图自身计算出的
普通 failed/dead_letter 结果,会像 succeeded/cancelled 一样封闭准入并 supersede
遗留行。唯一更窄的例外是:某次 run 的最终 steering flush 本身就无法持久提交
(_retry_or_dead_letter()),这条路径根本不会走到
finalize_run_with_steering_disposition_if_owned(),而是转而让这次投递尝试走普通的
retry/dead-letter 流程——只有在这种情况下,一个 dead_letter 的 run 才可能合法地仍然携带
pending 的 steering 行,作为刻意保留的 durable 历史记录,之后由 /redrive 像任何其他排队
中的 steering 一样重新拾起。但 /redrive 只接受 failed/dead_letter 状态的 run——一次
普通的、由客户端发起的 cancel 不在这条恢复路径之内,因此留在 cancelled run 上的
pending steering 目前没有任何自动的未来投递机会。完整策略见 steering 概念页"终结流程与
steering 准入边界"一节。
superseded_by_run_id 血缘标记
当 /resume 从一个 paused 的 run 创建出后继 run 时,旧 run 需要一种可靠的方式来表示
"我已经被取代了——请对新 run 执行 steer/cancel/resume"。早期实现曾把这个标记存放在
metadata.superseded_by_run_id 中——普通的、用户可写的 Run metadata。由于 /steer、
/cancel 以及并发 resume 的冲突闸门都把这个 key 当作权威状态来读取,调用方可以在创建 run
时伪造它(POST /v1/runs 携带 metadata: {"superseded_by_run_id": "..."}),从而在一个从未
真正被 resume 过的 run 上错误触发 409 run_superseded / 409 run_resume_conflict。
runs.superseded_by_run_id(由 0005_superseded_by_run_id 新增)是一个类型化、可为空的列,
取代该 key 成为唯一的真相来源:
- 它只在创建后继 run、迁移旧 run 待处理 steering 的同一个受 fence 保护的 resume 事务内
(
resume_run_with_pending_steering()及其 PostgreSQL 对应实现)被原子地写入一次—— 永远不会从RunCreate.metadata派生,也不能通过它写入。 /steer的准入检查、/cancel的 superseded 检查,以及 resume 冲突闸门,现在都直接读取run.superseded_by_run_id;它们都不再为此决策查询metadata。- 调用方在普通
RunCreate上提交metadata.superseded_by_run_id现在对控制流完全没有 影响——这个 key 是惰性的,最多作为普通用户 metadata 存储,永远不会被拷贝进类型化列。 - 升级兼容性:升级到当前版本之后,一个升级前就存在、恰好携带同名用户 metadata key 的历史
run,不会被重新解释为系统血缘状态——类型化列对每一条既有记录都从
NULL开始,只有未来 真正发生的 resume 才会填充它。
幂等重放在所有准入闸门前都是安全的
submit_steering()(内存版和 PostgreSQL 版仓库都一样)会先检查是否已经存在匹配
(tenant_id, run_id, Idempotency-Key) 的事件,然后才走终态/steering_closed/
superseded 这些准入闸门,而不是反过来。如果客户端的 POST /steer 已经 durably
提交,但 202 响应从未真正送达(连接中断、代理超时),使用相同 key 的标准幂等重试会
返回该事件当前的持久化状态——id 和 sequence 都不变——即使该 run 此后已经进入
finalizing、terminal 或 superseded 状态。只有一个真正全新的 key(这个 run 从未见过
的 key)才会受这些闸门约束;那才是一次真正的准入决策,仍然会得到相应的 409。
PostgreSQL 是真相来源,Redis 只是加速器
每一条 steering 事件是否存在及其状态,都存放在 PostgreSQL 的 run_steering_events 表里。
Redis pub/sub 只用于缩短"事件变为 durable"和"worker 注意到它"之间的延迟——从不会被用来
判断一个事件是否存在,它的丢失(Redis 重启、连接断开、NOTIFY 被漏掉)也从不会丢失一条
事件。运维可以把 Redis 视为对 steering 正确性完全可丢弃的组件:即使 Redis 宕机或被驱逐,
steering 仍能正常工作,只是 worker 的发现方式退化为在每个安全点周期性轮询 PostgreSQL,
而不是近乎即时的唤醒。
Worker heartbeat 对 steering 投递延迟的影响
一个 worker 只会在其持有租约的 run 到达下一个安全点时(节点开始、下一 superstep、重试、 resume 之后)去拉取该 run 新产生的 steering。因此投递延迟主要受两个因素支配:
- 该 run 多久到达一次安全点。 一个正在做单次长耗时 tool 调用的节点,在调用返回之前
看不到调用期间提交的 steering——直到该节点结束,或节点自己主动调用了
drain_steering(); - Redis 通知不可用时的 worker 轮询/heartbeat 间隔。 更大的间隔用更低的 PostgreSQL
轮询负载换取更高的 steering 延迟;可参照
docs/api.md与部署时lingxigraph worker的相关参数调优。
如何诊断长期 pending 的 steering
如果某条 steering 事件停留在 pending 状态的时间远超预期:
- 先确认该 run 确实已被 claim 并在推进:
GET /v1/runs/{run_id}——一个仍处于pending的 run 还没有 worker 在执行它,其上的 steering 在等待被 claim,而不是 bug; - 检查该 run 的图代码是否真的调用了
runtime.drain_steering(),以及调用频率—— steering 投递完全取决于图是否选择去查看,这是应用自选安全点模型的一部分(详见 Steering 概念页); - 直接查询
run_steering_events该行的status、created_at、source_event_id, 看它是否是从更早、曾经 paused 的 run 迁移过来的——最终run.steer.consumed事件上的queue_latency_seconds会包含所有累积的暂停等待时间,这是预期行为,不是 bug; - 检查 worker 日志/指标中是否有反复的
commit_steering_consumptions_if_owned()失败——该操作重试是安全的(幂等,并且基于 worker 自己的lease_owner/attempt加锁,见"steering_closedfinalization 闸门"),但持续失败说明存在值得告警的数据库 连通性问题; - 如果客户端在
/steer上收到409 run_finalizing,说明请求到达时该 run 的图已经 执行完毕——检查当前 attempt 的runs.steering_closed(true表示闸门已关闭),并对 一个全新的 run 重试,而不是同一个 run。这里用相同 Idempotency-Key 重放永远是安全的, 即使在闸门关闭之后也一样——见"幂等重放在所有准入闸门前都是安全的"; - 如果某条 steering 事件的
run.steer.superseded事件数据里reason是unconsumed_at_final_boundary(而不是resume_transfer),这正是终结边界修复按预期工作:该事件 已经 durably 被接受,但该投递尝试已经没有安全点可以消费它了,所以它被赋予了一个 durable、可观测的superseded归宿,而不是让 run 带着一条永久 pending 的行走向终态, 也不会被强加一次不必要的重跑。这本身不是 bug;但如果同一个 run 反复出现unconsumed_at_final_boundary,就值得在应用/图层面排查为什么 steering 总是在图到达终止节点之后 才到达。
正常的 run.steer.accepted → consumed 生命周期
POST /steerdurably 插入一条pending行,并(202 响应中)产生run.steer.accepted;- 持有该 run 租约的 worker 发现这条行(通过 Redis 提示或 PostgreSQL 轮询),把它摄入
该 run 的内存
SteeringChannel; - 图代码在某个安全点调用
drain_steering()并收到它; - worker durably 提交这次消费(
commit_steering_consumptions_if_owned(),基于 worker 自己的lease_owner/attempt加锁),在同一事务中把该行原子地转为consumed并产生run.steer.consumed。一个已经不再持有该 lease/attempt 的 worker(比如在 heartbeat 或 final flush 期间发生了 lease 被接管)什么都不会写入——从这一刻起该行 完全由新的持有者负责。
这条链路中没有任何一步是可选或可打乱顺序的;客户端应该只把第 4 步的
run.steer.consumed 事件当作"确实被处理"的证据。
为什么 worker/DB 瞬时故障不会丢失一次 consumption,以及如何观察重试
状态更新与 run.steer.consumed 事件的插入发生在同一个 PostgreSQL 事务内。如果事务已经
提交,但连接在 worker 看到确认之前断开,worker 的重试逻辑会重新发送同一批 consumption。
因为 commit_steering_consumptions_if_owned() 是幂等的——只为这次调用中真正发生状态
转换的行产生 lifecycle 事件,并且只有在调用方仍然持有该 run 当前 lease/attempt 时才会
写入任何东西——针对已经 consumed 的行重放同一批次是安全的空操作:状态保持
consumed,也不会产生第二条 run.steer.consumed。运维可以在 worker 日志里观察到这种
模式:同一批 steering id 反复出现 commit_steering_consumptions_if_owned 调用,但
run_events 里对应的 run.steer.consumed 数量并没有增长——只要重试最终会停止,这就是
预期行为,而不是卡住的重试循环。
升级与回滚边界
本项目的迁移是 forward-only 的:每个 revision 的 downgrade() 都会直接抛错,而不是
尝试撤销变更。这是刻意的设计:真正回滚 run_steering_events 表或其新增列,要么会丢失
durable 的 steering 历史,要么需要把正在进行的 run 与新二进制已经不再预期的 schema
重新对齐,代价过高。实践含义:
- 永远不要把"回滚 schema"当作回滚策略的一部分。如果确实需要回滚一次部署,应该回滚
二进制、保持 schema 停在当前(更靠前)的 revision——
0002_steering/0003_steering_source_event/0004_steering_closed/0005_superseded_by_run_id都是纯 新增,早于对应特性的旧二进制会直接忽略多出来的表/列,不会因此启动失败; - schema 兼容不等于语义兼容。 旧二进制在新 schema 上不会崩溃,但它并不理解
steering_closed或superseded_by_run_id——它不会执行 finalization 的 steering admission 门控,也不会执行当前版本的 resume-lineage cancel/steer 拒绝逻辑。应该把 新的 steering/lineage 控制语义视为"只有当所有控制面实例都升级到新二进制后才可用", 而不是 schema 迁移一落地就生效。如果一个二进制已经真正执行过这些语义(例如某次 resume 已经写入了superseded_by_run_id),把它回滚成旧版本并不会撤销这次转换—— 旧二进制只是从此不再强制执行它。如果必须在 steering 正在被使用时回滚二进制,应先 quiesce 或等待 in-flight run 结束,不要假设旧二进制能保留同样的保证; - 因为没有 schema 回滚路径,
0002_steering/0003_steering_source_event/0004_steering_closed/0005_superseded_by_run_id应该先在 staging 环境(或 CI 的test_postgres_integration.pyPostgres job)中验证过,再应用到生产环境,而不是依赖 "出问题了还能撤销"。