Skip to content

feat(queue): guided 队列策略(运行中消息注入)+ SSE 阻塞读(流式实时化) - #1009

Closed
zgpnuaa wants to merge 2 commits into
xerrors:mainfrom
zgpnuaa:feat/guided-queue-and-sse
Closed

feat(queue): guided 队列策略(运行中消息注入)+ SSE 阻塞读(流式实时化)#1009
zgpnuaa wants to merge 2 commits into
xerrors:mainfrom
zgpnuaa:feat/guided-queue-and-sse

Conversation

@zgpnuaa

@zgpnuaa zgpnuaa commented Sep 10, 2026

Copy link
Copy Markdown
Contributor

现象

智能体长任务运行中,两个体验痛点:

  1. 不能中途修正:现有队列策略 enqueue(排队等 Run 结束)和 steer(结束当前 Run 让位)都覆盖不了「运行中补充一句修正,让它即时调整方向」——比如调研跑了 3 分钟补一句「只看 2024 年后文献」,期望它接着按新约束继续,而不是作废重来。
  2. 流式顿挫:事件流靠「非阻塞读 + sleep 轮询」,token 一阵一阵吐,即使活跃期也至少一个轮询间隔才取一次事件。

对标

Claude Code 的两条核心交互:① 长任务运行中直接输入补充消息,任务不中断、模型下一轮按新约束继续;② 流式「来一个 token 显示一个」,几乎无感知延迟。这两条是「智能体好不好用」的分水岭。

机理

  1. guided:上游已预留坑位但没实现——NOT_IMPLEMENTED_QUEUE_POLICIES = ("guided", "bridge")SteerMiddleware 挂了 before_model 钩子但只处理了 steer,缺「运行中注入消息」路径。
  2. SSE:轮询的固有缺陷是「事件到了也要等下一个轮询周期才被取到」,下界延迟就是轮询间隔(上游活跃 0.1s、空闲扩到 4s)。

改进方法

guided 队列

  • guidedNOT_IMPLEMENTED 移入 SUPPORTED_QUEUE_POLICIES,新增状态 REQUEST_STATUS_INJECTED = "injected"
  • SteerMiddleware.abefore_modeltake_pending_guided_messages(run_id) 把等待消息作为 HumanMessage 追加进 state 的 messages,模型下一轮即看到,Run 不终止
  • 读取-标记同一事务、靠 worker lease 单写者;从 raw_message 恢复原始 id 按 id 去重;注入失败静默退化。

SSE 阻塞读

  • 新增 blocking_read_run_stream_events,用 Redis Streams XREAD BLOCK:有事件毫秒级唤醒,无事件挂起最多 1s;保留 after_seq 游标续读语义,空闲零 CPU 空转。

效果

运行中发补充消息,模型下一轮按新约束修正,前序成果不丢失;流式延迟从轮询间隔降到毫秒级,token 连续吐出。附 10 个 guided/steer 单测 + 4 个阻塞读单测 + 4 个 SSE 行为单测。

@xerrors

xerrors commented Sep 10, 2026

Copy link
Copy Markdown
Owner

Codex Review:

预 review 基于 a73c0cc34e8a

发现 2 个会丢失用户输入的问题。

  1. [P1] guided 与 steer 同时等待时,已消费的 guided 更新被丢弃。 steer.py:16–20 先调用 _collect_guided_update(其中已把 Request 提交为 injected),发现 steer 后却只返回 {"jump_to":"end"},没有把 guided 的 messages 返回给 graph。注释所说“随 checkpoint 留给下一个 Run”不会发生;它既退出队列,也没有进入 checkpoint。已执行该 head 的真实 abefore_model 最小复现确认返回值不含 messages;新增 test_before_model_prefers_steer_jump_when_both_pending 反而把这个丢失行为写成了预期。请在决定让位前避免消费 guided,或保证 messages 更新随 jump 一起持久化,并验证下一 Run 真正收到消息。

  2. [P1] injected 提交早于 checkpoint,worker 崩溃会永久漏注入。 agent_request_queue_service.py:465–467 在将 HumanMessage 返回 middleware 之前就提交 injected 终态;此时 graph 尚未应用/保存更新。若 worker 在 commit 后、checkpoint 写入前退出,恢复时 list_pending_guided 只读取 queued,该消息不会再注入或派发。worker lease 只能解决执行 ownership,不能让这两个持久化点原子提交。请建立可恢复的消费确认/幂等重放,并用真实 PostgreSQL + checkpoint 的故障注入验证“提交后崩溃,恢复仍恰好收到一次消息”。

验证:该 head 独立快照执行 PYTHONPATH=package:server python -m pytest test/unit/agents/test_steer_middleware.py test/unit/services/test_agent_request_queue_service.py test/unit/services/test_agent_run_service.py test/unit/services/test_run_queue_service.py -q --disable-warnings --tb=short,131 passed;另执行了 guided+steer 返回值复现。第二项为持久化顺序静态确认,尚未做真实崩溃 E2E。新增 queue 测试使用 SQLite/fake,不证明 PostgreSQL/checkpoint 恢复语义;SSE 也未在真实 Redis/浏览器验证。请补 tracked decision,并考虑将 guided 与独立的 SSE 性能变更拆开审查。

@xerrors

xerrors commented Sep 10, 2026

Copy link
Copy Markdown
Owner

前者的逻辑已经在 steer 实现了,后者是性能和体验的取舍,目前(上游活跃 0.1s、空闲扩到 4s)的体验已经很好了。如果觉得可以优化,可以将退避上限调整为0.5s。

使用 Redis Streams XREAD BLOCK 会一直占用 SSE 连接,高用户场景下 SSE 资源会被占满。

@xerrors xerrors closed this Sep 10, 2026
@zgpnuaa

zgpnuaa commented Sep 10, 2026

Copy link
Copy Markdown
Contributor Author

感谢 review,两个丢失路径都成立,已修复(最新 head ffb73b6a)。

[P1] 让位即丢失 → 先判让位、不消费 guided

确认:abefore_model 先调 _collect_guided_update(其中已把请求提交为 injected),命中 steer 后只返回 {"jump_to": "end"}——injected 意味着不再参与派发,消息既不进 checkpoint 也不回队列。原注释「随 checkpoint 留给下一个 Run」不成立,test_before_model_prefers_steer_jump_when_both_pending 确实把这个丢失行为写成了预期。

修复:先判定 steer;命中即返回让位,完全不触碰 guided,请求保持 queued,由下一个 Run 正常注入或派发。对应测试已改为断言「让位时不调用收割函数」(taken == [])。

[P1] injected 早于 checkpoint → 以 checkpoint 为已生效判据做幂等重放

确认:injectedHumanMessage 交给 middleware 之前就提交,graph 尚未应用/保存;worker 在这两步之间退出后,只读 queuedlist_pending_guided 永远捞不到它。

修复:把 injected 语义收窄为「已提交注入,待 checkpoint 确认」,收割时用当前 graph state 的消息 id 对账:

请求状态 消息 id 在 state 中 行为
injected 已生效,跳过(不重复注入)
injected 否(提交后崩溃) 重新注入(幂等重放)
queued 正常注入
无可用对账 id 不跳过,宁可重放也不丢消息

判据用 Message.extra_metadata.raw_message.id(入库时生成的稳定 LangChain 消息 id,注入时被原样恢复)。仓储 list_pending_guided 增加 statuses 参数(默认仍 ("queued",)),仅收割路径传 ("queued", "injected")
中间件 abefore_model 相应把 _state_message_ids(state) 传给收割函数。

证据

  • 单元(test_steer_middleware.py):让位不消费、state id 透传、四种重放判定(已生效跳过 / 未落盘重放 / queued 正常 / 无 id 不跳过)。
  • 真实 PostgreSQL 集成测试test_agent_request_queue_concurrency.py::test_guided_injected_request_remains_replayable_until_checkpoint_confirms):落一条 queued guided 请求 → 提交 injected(模拟提交后崩溃)→ 默认查询返回空(复现丢失路径),带 statuses=("queued","injected") 的恢复查询返回该请求。
  • tracked decision:docs/develop-guides/decisions/implemented/2026-09-10-guided-queue-delivery-guarantees.md。全量 unit 1939 passed / 52 skipped。

未验证范围:没有做「提交 injected → 杀 worker 进程 → 重启续跑」的进程级故障注入 E2E;上面用真实 PostgreSQL 覆盖了持久化与恢复判定这一步。若你希望以进程级故障注入作为合并前提,我可以按 test/e2e 的 deterministic replay 拓扑补一个。

关于拆分:guided 与 SSE 阻塞读在 run_queue_service.py / agent_run_service.py 里确有耦合(阻塞读依赖事件流的游标语义),拆开会引入同文件的双 PR 冲突,所以暂留在一个 PR;如果你更希望分两次审查,我可以按 run_queue_service 与 SSE 消费端为界拆成两个分支。SSE 侧此前已在浏览器端验证过流式与终态收尾。

@zgpnuaa

zgpnuaa commented Sep 10, 2026

Copy link
Copy Markdown
Contributor Author

接受关闭,两个判断都成立,感谢说明。

  1. guided 与 steer 重叠:guided 的运行中消息注入确实是 steer 已覆盖的语义,重做一遍是冗余。幂等重放那部分(injected 待 checkpoint 确认、按消息 id 对账防丢)是我对「提交后崩溃」窗口的加固,但既然 steer 是上游权威方案,这部分不需要独立存在。

  2. SSE 阻塞读的资源占用:XREAD BLOCK 长期占用 SSE 连接、高并发下占满连接池的担忧成立,这点我没有充分权衡连接数上限,接受「保持轮询」的取舍。当前的 0.1s 活跃 / 4s 空闲退避确实已经够流畅。

关于「退避上限调整为 0.5s」:如果后续观察到空闲→活跃切换的迟滞,我会按这个方向单独提一个很小的轮询参数 PR,不再动阻塞读。

这个 PR 就不再推了,谢谢 review。

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants