LingxiGraph
运维与项目

运行中输入(Steering)的运维

迁移上线顺序、PostgreSQL 作为真相来源、延迟、重试与回滚边界

迁移上线顺序

运行中输入在基础 schema 之上新增了四个 forward-only 的 Alembic revision:

Revision新增内容
0002_steeringrun_steering_events 表及其索引,以及读写它的 worker/server 代码路径
0003_steering_source_eventrun_steering_events 上的 source_event_id 列 + 局部索引,用于跨多次 pause/resume 迁移关联同一条 steering
0004_steering_closedruns.steering_closed——finalization 准入闸门;详见下方"steering_closed finalization 闸门"
0005_superseded_by_run_idruns.superseded_by_run_id——被 resume 覆盖的类型化控制面血缘标记;把它从可被用户写入的 metadata 中移出,使其无法被伪造。详见下方"superseded_by_run_id 血缘标记"

按标准 Alembic 路径应用:

alembic upgrade head

四者都是纯新增(ADD COLUMN IF NOT EXISTSCREATE INDEX IF NOT EXISTS、可为空/ NOT NULL DEFAULT FALSE 列),可以在部署新的 server/worker 二进制之前先对存活数据库 执行——早于 0004/0005 的旧二进制根本不会读写 steering_closed/superseded_by_run_id, 行为与这些列不存在时完全一致。零停机上线顺序:

  1. 在旧的 server/worker 二进制仍在提供服务时执行 alembic upgrade head (新增 run_steering_eventssource_event_idsteering_closedsuperseded_by_run_id);
  2. 再滚动发布新的 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_steering0003_steering_source_event0004_steering_closed

steering_closed finalization 闸门

一旦 worker 的图执行到达真正的终态结果,该 run 可能还会短暂地继续读作 running/cancelling,同时 worker 把图已经消费掉的 steering durably flush 出去、 提交最终状态——此时已经没有更多图安全点可以消费任何新输入了。runs.steering_closed 是持有该 run 的 worker(与状态写入一样,基于自己的 lease_owner/attempt 加锁)在执行 结束的瞬间、final flush 开始之前设置的一个布尔值。

  • steering_closedfalse 时,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 状态的时间远超预期:

  1. 先确认该 run 确实已被 claim 并在推进:GET /v1/runs/{run_id}——一个仍处于 pending 的 run 还没有 worker 在执行它,其上的 steering 在等待被 claim,而不是 bug;
  2. 检查该 run 的图代码是否真的调用了 runtime.drain_steering(),以及调用频率—— steering 投递完全取决于图是否选择去查看,这是应用自选安全点模型的一部分(详见 Steering 概念页);
  3. 直接查询 run_steering_events 该行的 statuscreated_atsource_event_id, 看它是否是从更早、曾经 paused 的 run 迁移过来的——最终 run.steer.consumed 事件上的 queue_latency_seconds 会包含所有累积的暂停等待时间,这是预期行为,不是 bug;
  4. 检查 worker 日志/指标中是否有反复的 commit_steering_consumptions_if_owned() 失败——该操作重试是安全的(幂等,并且基于 worker 自己的 lease_owner/attempt 加锁,见"steering_closed finalization 闸门"),但持续失败说明存在值得告警的数据库 连通性问题;
  5. 如果客户端在 /steer 上收到 409 run_finalizing,说明请求到达时该 run 的图已经 执行完毕——检查当前 attempt 的 runs.steering_closedtrue 表示闸门已关闭),并对 一个全新的 run 重试,而不是同一个 run。这里用相同 Idempotency-Key 重放永远是安全的, 即使在闸门关闭之后也一样——见"幂等重放在所有准入闸门前都是安全的";
  6. 如果某条 steering 事件的 run.steer.superseded 事件数据里 reasonunconsumed_at_final_boundary(而不是 resume_transfer),这正是终结边界修复按预期工作:该事件 已经 durably 被接受,但该投递尝试已经没有安全点可以消费它了,所以它被赋予了一个 durable、可观测的 superseded 归宿,而不是让 run 带着一条永久 pending 的行走向终态, 也不会被强加一次不必要的重跑。这本身不是 bug;但如果同一个 run 反复出现 unconsumed_at_final_boundary,就值得在应用/图层面排查为什么 steering 总是在图到达终止节点之后 才到达。

正常的 run.steer.accepted → consumed 生命周期

  1. POST /steer durably 插入一条 pending 行,并(202 响应中)产生 run.steer.accepted
  2. 持有该 run 租约的 worker 发现这条行(通过 Redis 提示或 PostgreSQL 轮询),把它摄入 该 run 的内存 SteeringChannel
  3. 图代码在某个安全点调用 drain_steering() 并收到它;
  4. 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_closedsuperseded_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.py Postgres job)中验证过,再应用到生产环境,而不是依赖 "出问题了还能撤销"。

On this page