Workflow LLM Stream Chunk SSE 投影方案 #2945
Replies: 2 comments
|
总体上我赞成这个方向:Workflow LLM stream progress 应继续进入统一的 committed observation / Projection Pipeline,由 role actor 保持 authority,并在 commit 前 coalescing;不应在 Host 增加 raw stream 旁路或进程内 session registry。 不过我建议把当前结论标记为 direction approved, changes required before implementation,至少需要补齐以下几点: 1. 不要把几种可靠性保证混在一起committed progress 只能说明 source fact 已进入 role actor event store,并允许既有 projection failure 机制处理已记录的失败。 建议明确列出 guarantee matrix,分别说明:
否则“committed / replay / dedupe / durable observability”容易被误解为同一层保证。 2.
|
|
我按最新 现状证据:
所以缺口不是 “chunk 到 SSE frame 的 mapper 不存在”,而是 producer 产出的 raw topology event 不会进入生产 workflow projection session relay。mapper 支持 raw payload 只说明“如果这个 payload 到了 mapper,它能转”,但生产链路的入口仍然是 committed observation。 我建议把实现边界收敛成:
因此我会把本 discussion 的落点表述为:mapper 已具备部分 frame 转换能力;仍需要把 workflow LLM stream progress 从 transient topology publication 升级为 actor-owned、typed、coalesced committed observation,并继续走现有 Workflow Execution Projection Session 链路。 |
Uh oh!
There was an error while loading. Please reload this page.
Related: #2915
背景
在本地 mainnet host 验证 workflow chat / service chat stream 时,前端没有观察到预期的 LLM chunk。进一步看下来,问题不在 SSE writer,也不在 mapper 是否能把 chunk 转成前端 frame,而在 producer 产出的事件契约和 projection session relay 的输入不一致。
当前
WorkflowRoleGAgent会把 LLM chunk 作为 transient topology publication 发给 parent:对应 proto:
EventEnvelopeToWorkflowRunEventMapper已经知道如何把WorkflowLlmStreamChunkEvent.delta_content映射成WorkflowRunEventEnvelope.TextMessageContent,也能把delta_reasoning_content映射成aevatar.llm.reasoningcustom event。但是生产 projection session relay 只消费 committed observation envelope,过滤的是
CommittedStateEventPublished。raw topologyWorkflowLlmStreamChunkEvent不会进入这条链路,所以 mapper 虽然支持该类型,live host 的 workflow SSE 仍然可能看不到 chunk。设计结论
Workflow LLM stream chunk 应继续走现有 Workflow Execution Projection Session 链路:
推荐让 workflow/role actor 产出 typed、run-scoped、coalesced 的 committed LLM stream observation,再由统一 projection session 链路转发。raw provider chunk 不应逐条写入 event store。
这个方向保持了统一 CQRS / AGUI Projection Pipeline,也避免新增 Host 旁路订阅、进程内
sessionId -> sinkregistry 或第二条 observability path。ProjectionSessionEventHub 的边界
ProjectionSessionEventHub<TEvent>应理解为按rootActorId + sessionId路由的 live fan-out transport。它不写 source actor event store,不写业务 actor state,也不写 document/read model。是否进入 event store,取决于 hub 之前的 source fact 是否已经由 actor
PersistDomainEventAsync提交;是否进入 projection state,取决于 projection scope 是否为了 attachment、watermark、failure replay 记录自己的 projection 运行态。这两者都不是 hub 自身的数据持久化。所以项目里可以存在不进入存储的 live-only 路径,例如 voice realtime frame。但 workflow LLM chunk 如果选择 live-only,就不能同时声称具备 committed projection、replay、dedupe 和 durable observability 语义;如果选择 committed observation,就应先 coalescing,再提交合并后的 progress。
路径对照
WorkflowChatRunInteractionService通过WorkflowExecutionProjectionPort绑定 live sink,SSE 写出WorkflowRunEventEnvelopeDefaultCommandInteractionService、AGUI session projector 和ProjectionSessionEventHub<AGUIEvent>RoleGAgent提交RoleChatSessionProgressedEvent(session_id, sequence, typed payload),NyxIdChatSessionEventProjector只消费 committed envelopeTextMessage*ConversationGAgent,通过平台 reply/update carrier 发回外部平台,不是浏览器 SSERecord*command 回投LlmSessionGAgent,再提交LlmStreamChunkObserved等 recorder event存储语义
需要区分三类 payload:
如果把 raw provider chunk 逐条写入 event store,系统会重复承担多层成本:
现有
LlmSessionGAgent的LlmStreamChunkObserved模型适合专门的 LLM session observation actor,但不应不加区分地复制到 workflow run state,用来承接每个 provider frame。workflow 更适合参考RoleGAgent的 progress 模型:typed progress event 可以提交,但 state 只保存单调 progress watermark 和 terminal session facts。Coalescing 成本边界
Coalescing 必须发生在
PersistDomainEventAsync之前。如果 raw provider chunk 已经逐条 committed,再由 projector 或 SSE writer 合并,这种后置 coalescing 只能减少前端 frame 数,不能减少 event store 写入、source actor state transition、committed-state publication、projection scope watermark 或 replay/activation 成本。
推荐链路:
不推荐链路:
Actor 内 coalescer 可以临时保存尚未 flush 的 chunk 信息,但它应是当前 execution flow 内的 transient buffer,不是跨节点事实源。推荐使用 handler 局部
StringBuilder/ coalescer;不推荐把 pending raw chunk buffer 放入 actorState,也不应在 Host/Application/projection service 中维护sessionId -> bufferdictionary。如果 actor 或进程在 flush 前崩溃,未 flush 的 raw delta 可能丢失。这个前提只有在 raw provider chunk 不被定义为 durable fact 时才成立。flush 间隔越短,写入次数越多但丢失窗口越小;flush 间隔越长,写入次数越少但前端粒度更粗。
建议实现
1. 定义 committed stream progress contract
新增或演进 workflow execution proto,不复用 raw topology 语义。事件应包含:
run_idstep_idsession_idrole_actor_idsequenceobserved_atoneof payload,例如 text delta、reasoning delta、usage、text ended,以及后续 tool lifecycle payload。这些字段是核心业务语义,不应放入 generic metadata bag。
2. 从 producer 提交 coalesced progress
在
WorkflowRoleGAgent或明确选择的 producer actor 内做 actor-owned coalescing,然后提交合并后的 progress。保留现有 terminalWorkflowLlmInvocationCompletedEvent和RoleChatSessionCompletedEvent路径,继续由它们作为最终 output、tool calls、receipts、usage 和 failure 的权威事实。Flush 边界建议:
3. 保持 producer state 小而稳定
Reducer 只记录 latest accepted stream sequence 和幂等所需的最小 active stream identity,不把每个 chunk 追加进 actor state,不在 query 时通过 replay chunk event 重建最终 workflow result。
4. 把 committed progress 映射成 workflow SSE frame
扩展
EventEnvelopeToWorkflowRunEventMapper,继续基于 committed observed envelope 映射:WorkflowRunEventEnvelope.TextMessageContentaevatar.llm.reasoningWorkflowRunEventEnvelope.UsageWorkflowRunEventEnvelope.TextMessageEndHost 不应新增特殊订阅 raw actor stream 的代码。
5. 验证 live SSE
建议补充测试覆盖:
aevatar.llm.reasoning发送。ProjectionSessionEventHub仍然不承担存储职责。待决策点
实现前需要明确 stream observation authority:
WorkflowRoleGAgent中提交 progress,并让 role actor 继续作为 producer authority。WorkflowRunGAgent中提交 run-scoped observation,把 workflow run actor 定义为 stream observation authority。默认推荐第一种,因为它遵循现有
RoleGAgentcommitted progress 模型,并避免把 role output 重新归属到 parent actor。只有当 workflow 产品语义明确要求所有 live stream facts 都由 workflow run actor 直接拥有时,第二种才适合采用。非目标
ProjectionSessionEventHub写 actor state、event store、document store 或任何 read model。All reactions