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 前后端状态系统。
[!info] related notes
- 所属 MOC: bodysense-moc、agent-runtime-moc
- 相关概念:
- 易混淆概念:
- 相关资源:
BodySense Consultation Streaming Architecture
范围:L3 不是重新学 React / SSE 基础
L3 也采用 Knowledge Delta。
下面这些已经有成熟基础笔记,不值得在 BodySense 课程里重新复制:
- SSE;
- Web Streams 与增量文本解码;
- NDJSON、SSE 与流式协议边界;
- TypeScript 静态类型与运行时校验;
- Zod 运行时 Schema 校验;
- React useReducer;
- React Context 与状态管理边界;
- TanStack Query 与 SSE 流式数据集成;
- AbortController 与异步取消。
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.ts和listRunEvents()当前仍有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-cancellation 与 cancellation-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 推荐阅读顺序
- bodysense-stream-event-trust-boundary
- bodysense-active-turn-state-machine
- durable-sse-recovery-with-after-seq
- agent-interrupt-resume-lifecycle
遇到基础断点再查:
- web-streams-and-incremental-text-decoding
- typescript-static-types-and-runtime-validation
- react-use-reducer
- react-context-and-state-management
- tanstack-query-with-sse-streaming
- abortcontroller-and-async-cancellation
建议源码阅读走廊
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 基础:
- 为什么
StreamEventTypeScript union 不能证明网络 JSON 合法? - 为什么 current streaming turn 要有独立
ActiveTurnState? - Reducer 为什么必须 pure,并把 effects 放到外面?
- 为什么 SSE 断线不等于 Agent run 结束?
after_seq恢复如何避免重复/遗漏?- 为什么空 durable page 不等于 completed?
- 为什么 live events 与 replay events 应走同一 reducer?
- TanStack Query 与 ActiveTurn 分别拥有哪类 state?
ask_user为什么是 interrupt/resume,而不是普通 user message?- 为什么
AbortController.abort()不等于 durable run cancellation?
如果这十个问题能稳定回答,L3 的核心知识增量已经掌握。