BodySense Consultation Streaming Architecture

L3 的差分学习总览:从 POST SSE 的不可信网络字节开始,经 StreamEvent trust boundary、ActiveTurn reducer、durable Runtime Event Log、after_seq recovery、TanStack Query reconciliation 与 interrupt/resume,形成可恢复的 Consultation 前后端状态系统。

#type / synthesis #status / growing #tech / ai #tech / architecture #tech / dev / frontend #resource / bodysense #resource / agent #resource / react #resource / typescript

[!info] related notes

BodySense Consultation Streaming Architecture

范围:L3 不是重新学 React / SSE 基础

L3 也采用 Knowledge Delta

下面这些已经有成熟基础笔记,不值得在 BodySense 课程里重新复制:

L3 真正要学的是:

这些基础技术在一个 production-shaped Agent conversation runtime 里怎样组合成“live stream 可交互、断线可恢复、刷新可回放、状态不串台”的系统。

当前端到端数据流

当前真实 Consultation 路径可以压缩成:

User sends message

React / assistant-ui adapter
        ↓ POST
Go ConsultationRuntime

Python Consultation Agent runtime
        ↓ semantic events
Go normalizes + persists public StreamEvent

Live SSE ---------------------------┐
        ↓                            │
Frontend parser                     │
        ↓                            │
StreamEvent                         │
        ↓                            │
ActiveTurnReducer                   │
        ↓                            │
ActiveTurnState                     │
        ↓                            │
Streaming UI                        │
                                     │ network drop
Durable Runtime Event Log <----------┘
        ↓ GET after_seq
Replay same StreamEvents

same reducer / handlers

recover live UI state

这张图是 L3 的中心。

L3 的四个核心知识增量

1. Network Bytes → Trusted StreamEvent

共享 TypeScript contract 定义了 StreamEvent union,但 TypeScript 类型在 runtime 消失。

因此:

JSON.parse(...)
as StreamEvent

只让编译器闭嘴,并没有证明网络数据真的合法。

正确心智模型:

Network bytes
→ JSON unknown
→ runtime validation
→ trusted StreamEvent
→ reducer/components

[!warning] 当前实现状态 useSSEProcessor.tslistRunEvents() 当前仍有 JSON.parse(...) as StreamEvent 的 cast 路径,代码注释也明确说明 runtime validation 应位于 network boundary。这是 L3 要理解并后续补齐的真实边界缺口,不能把它误记为“已经完成验证”。

详见 bodysense-stream-event-trust-boundary

2. ActiveTurnState 是“当前流式 turn”的唯一 UI 真值

以前流式 AI UI 很容易有多份互相竞争的状态:

assistant-ui streaming text
local tool call state
local citation state
TanStack Query persisted messages
interaction card state

BodySense 的目标架构把当前未完成 turn 聚合成:

ActiveTurnState
├─ runId / conversationId
├─ text
├─ toolCallsById
├─ citationsByKey
├─ knowledgeGapsByKey
├─ redFlag
├─ pendingInteraction
├─ extractedInfoByBodyPart
├─ status
├─ finalParts
└─ lastSeqByType

所有 live events 经过纯 reducer:

State + StreamEvent
→ New State + Effects

详见 bodysense-active-turn-state-machine

3. SSE 是 Transport,不是 Durable Truth

Live SSE 连接可能:

  • Wi-Fi 中断;
  • tab 切换;
  • proxy 断流;
  • 浏览器网络错误。

但一个 Agent run 不应该因为 HTTP reader 消失就等同于业务执行消失。

Go 把关键 public events 写入 durable runtime_events

conversation_id
run_id
seq
channel
type
ids
payload
replayable

前端在 live stream 断开且没有看到 terminal event 时:

GET /runs/{run_id}/events?after_seq=N

继续追赶缺失 events。

详见 durable-sse-recovery-with-after-seq

4. ask_user 是 Interrupt,不是一条普通聊天消息

ask_user 的正确生命周期不是:

Assistant: 你疼吗?
User: 疼

简单文本拼接。

它是 Agent runtime 的控制流:

running
→ interaction.required
→ interrupted / waiting_user
→ answer recorded
→ resume input
→ run.resumed
→ continue same logical task
→ completed

用户回答不是一个新的独立 chat task,而是恢复被中断 run 的输入

详见 agent-interrupt-resume-lifecycle

三类前端 State:不要混成一个大 Store

L3 需要非常明确地区分:

A. Durable Server State

例如:

conversation
messages
persisted run events
BodyState
Diagnosis

适合 TanStack Query / 后端 read model。

B. Active Runtime Projection

例如当前正在生成的:

text delta
tool running
pending interaction
current citation
stream status

适合 ActiveTurnState

C. Pure UI State

例如:

panel collapsed
selected tab
input draft
hover state

属于 local React state。

不要把:

network durable truth
+
current stream projection
+
UI presentation

塞进同一全局 store。

TanStack Query 与 ActiveTurn 怎么协作

当前代码的思路是:

During stream
→ ActiveTurn owns current turn projection
→ selected domain effects may patch query cache

At stream completion
→ invalidate thread / conversations / workspace
→ refetch durable server truth

在 patch server cache 前还会:

cancelQueries(queryKey)

避免:

old in-flight GET response

覆盖刚从 stream 得到的新 projection。

这不是“TanStack Query 和 SSE 只能二选一”,而是:

TanStack Query
= durable server-state cache

ActiveTurn
= ephemeral/live execution projection

详见 tanstack-query-with-sse-streaming

Live 与 Replay 必须尽量共用 Reducer

一个非常重要的设计目标:

Live SSE Event

Reducer

Durable Replayed Event

same Reducer

如果 live path 与 replay path 使用两套 UI 解释逻辑,迟早会出现:

在线时显示 A
刷新后显示 B

当前 recoverDurableRunEvents() 就复用同一组 handlers / event dispatch。

这是 Event Log 架构最大的收益之一:

UI state is a projection of events, not a pile of ad-hoc callbacks.

seq 与幂等

每个 durable runtime event 有序号:

seq = 1, 2, 3, ...

客户端记录 maxSeq,恢复时请求:

after_seq = maxSeq

服务器查询:

seq > after_seq
ORDER BY seq ASC

于是可以做到:

already seen events
→ skip

missing events
→ replay

当前 ActiveTurn reducer 还使用 lastSeqByType 做前端幂等保护,因为现有后端存在不同 seq space 混合的历史现实;这属于项目级 compatibility detail,而不是理想协议必须永远如此。

为什么空 Event Page 不等于 Run 已结束

Durable recovery 中:

GET after_seq=3
→ []

不能立刻判:

run completed

因为可能只是:

Agent 还在运行,下一个 event 尚未 commit。

所以 recovery 必须等待显式 terminal event:

stream.done
or
stream.error

这是“状态缺失”与“终止事实”的典型区别。

Transport Lifetime 与 Agent Run Lifetime 分开

Go 当前通过 detached/durable execution context,使 Agent execution 不完全绑定原 HTTP connection cancel。

这体现:

HTTP/SSE connection lifetime

Agent run lifetime

如果网络断开,run 仍可以继续,后续由 durable log 恢复客户端。

这是 production Agent 和普通“请求断了任务就没了”的 API 的明显差异。

Cancellation:课程目标与当前实现要区分

L3 roadmap 中包含:

AbortController / cancellation

知识库已有 abortcontroller-and-async-cancellationcancellation-propagation

但当前 Consultation web feature 中,没有发现显式 AbortController 或用户级 cancel-run path。

所以这里应记成:

Concept: should understand
Current BodySense implementation: not yet a complete explicit cancellation path

不要把“课程要求理解 cancellation”误写成“项目已经完成 end-to-end cancellation”。

后续真正补齐时需要区分:

Cancel transport/read

Cancel durable Agent run

前者只停止前端 reader;后者需要明确的 server-side command/state transition。

L3 推荐阅读顺序

  1. bodysense-stream-event-trust-boundary
  2. bodysense-active-turn-state-machine
  3. durable-sse-recovery-with-after-seq
  4. agent-interrupt-resume-lifecycle

遇到基础断点再查:

建议源码阅读走廊

Trust Boundary

packages/contracts/src/stream-events.ts
apps/web/src/features/consultation/hooks/useSSEProcessor.ts
apps/web/src/features/consultation/services/consultationService.ts

Active Turn Projection

apps/web/src/features/consultation/runtime/activeTurnReducer.ts
apps/web/src/features/consultation/context/ActiveTurnContext.tsx
apps/web/src/features/consultation/hooks/useAssistantChatRuntime.ts

Durable Recovery

apps/api/internal/model/runtime_event.go
apps/api/internal/service/runtime_event_service.go
apps/web/src/features/consultation/runtime/durableRunRecovery.ts

Interrupt / Resume

state.interaction.required
state.interaction.answered
run.interrupted
run.resumed
resumeInteractionStream()

L3 完成标准

能解释以下问题即可,不需要重复背 React 基础:

  1. 为什么 StreamEvent TypeScript union 不能证明网络 JSON 合法?
  2. 为什么 current streaming turn 要有独立 ActiveTurnState
  3. Reducer 为什么必须 pure,并把 effects 放到外面?
  4. 为什么 SSE 断线不等于 Agent run 结束?
  5. after_seq 恢复如何避免重复/遗漏?
  6. 为什么空 durable page 不等于 completed?
  7. 为什么 live events 与 replay events 应走同一 reducer?
  8. TanStack Query 与 ActiveTurn 分别拥有哪类 state?
  9. ask_user 为什么是 interrupt/resume,而不是普通 user message?
  10. 为什么 AbortController.abort() 不等于 durable run cancellation?

如果这十个问题能稳定回答,L3 的核心知识增量已经掌握。

创建于 2026/8/22 更新于 2026/8/22