Durable SSE Recovery with after_seq
解释实时 SSE 仅作为低延迟 transport,而 durable Runtime Event Log 作为恢复来源;客户端用 run identity、monotonic seq 与 after_seq 在断线后追赶缺失事件,并让 live/replay 共用同一 projection。
[!info] related notes
- 所属 MOC: bodysense-moc、frontend-engineering-moc
- 相关概念: bodysense-active-turn-state-machine、bodysense-stream-event-trust-boundary
- 易混淆概念: transport connection lifetime vs durable run lifetime
- 相关资源: bodysense-consultation-streaming-architecture、agent-interrupt-resume-lifecycle
Durable SSE Recovery with after_seq
一句话定义
Durable SSE Recovery 是把 SSE 当作低延迟实时运输层,而不是唯一事实来源;关键 StreamEvent 同时写入 durable Runtime Event Log,客户端断线后用 after_seq 从最后已消费位置继续追赶,并将 replay event 送入与 live event 相同的 reducer。
Agent Run
↓
Durable Event Log
├────────→ Live SSE
│ ↓
│ Browser
│ X network drop
│
└────────→ GET events?after_seq=N
↓
Browser catches up
为什么 Production Streaming 的目标不是“永远不断线”
任何跨网络的长连接都应该假设会断:
- Wi-Fi 切换;
- 手机网络切换;
- proxy timeout;
- browser tab suspend;
- deployment/restart;
- client crash;
- NAT / load balancer reset。
因此真正健壮的目标不是:
connection must never break
而是:
connection may break
but execution identity and durable history survive
and client can converge again
这和数据库系统里“磁盘/进程可能失败,所以需要 WAL/transaction log”的思路很像:不要把传输可靠性误当成业务持久性。
SSE 为什么不能承担 Durable Truth
SSE 很适合:
- 低延迟;
- 单向服务器推送;
- 流式文本;
- 事件连续展示。
但连接天然脆弱:
- Wi-Fi 切换;
- 移动网络抖动;
- proxy idle timeout;
- 浏览器 reader error;
- tab suspend;
- 前端 crash/reload。
如果架构是:
Agent Event
→ only SSE
→ Browser
那么 browser 没收到就永远丢失。
因此 production-shaped runtime 需要:
Execution event
→ durable append
→ live delivery
即使 live transport 失败,历史仍在。
Append-before-deliver:Durability Ordering 为什么重要
一个关键设计问题:
Event 先写 durable log,还是先推 SSE?
Deliver-first 风险
1. SSE 给浏览器 seq=8
2. process crash
3. durable append 没发生
用户短暂看到过 seq=8,但刷新后服务器历史里不存在它。
这叫“幽灵事件”。
Append-first 更容易恢复
1. persist seq=8
2. commit
3. deliver seq=8 to SSE
即使第 3 步前网络断了:
recovery after_seq=7
→ still gets seq=8
代价是 live latency 多一点持久化开销。
因此对需要恢复/审计的重要 runtime event,一般更倾向:
durable before externally observable
或至少确保 delivery 与 durability 有明确一致性协议。
不是每个 Token 都必须成为高成本 Durable Domain Event
这里要避免另一个极端:
LLM 每一个字符 delta
→ 独立昂贵 DB transaction
可以根据需求选择:
- 每个 semantic event 持久化;
- delta batching;
- periodic checkpoints;
- final message persistence + selected control events durable。
关键问题是:
哪些事件一旦丢失会破坏恢复、authority、interaction 或用户可见历史?
例如:
- interaction.required;
- red flag;
- tool lifecycle;
- terminal state;
通常比每一个字符 delta 更需要可靠 identity 和持久化。
Runtime Event Identity
当前 Go Runtime Event 具备类似:
conversation_id
run_id
seq
channel
type
ids
payload
replayable
created_at
其中最关键的三个恢复坐标是:
conversation_id
run_id
seq
因为“恢复一个会话”还不够精确。
同一 conversation 可以有多个 runs:
Conversation C1
├─ Run R1
├─ Run R2
└─ Run R3
所以恢复必须知道:
当前缺的是哪个 run 的哪一段 event log?
为什么 Event Identity 不应该只靠时间戳
时间戳不能稳定表达:
- strict ordering;
- duplicate identity;
- same-millisecond events;
- retry/replay;
- cursor boundary。
所以:
(run_id, seq)
比:
created_at > T
更适合做 durable cursor。
Timestamp 可以用于观察和展示,但不是主要逻辑 cursor。
seq 是什么
理想上,一个 run 内:
seq = 1, 2, 3, 4, ...
它表达逻辑事件顺序,而不是 wall-clock 时间。
客户端消费:
1, 2, 3
然后连接断开。
本地记录:
maxSeq = 3
恢复请求:
GET /events?after_seq=3
服务端查询:
WHERE seq > 3
ORDER BY seq ASC
返回:
4, 5, 6
这就是一个典型 cursor-based catch-up protocol。
seq 的 Namespace 必须明确
最理想的定义之一:
seq is strictly monotonic within run_id
那么:
(R7, 12)
是明确事件位置。
如果系统历史上存在:
不同 event type 各自 seq
或者:
数据库 seq 与 Python internal seq 混合
就要明确 adapter/compatibility layer。
不要让 seq 这个字段名字看起来一样,就默认它们属于同一个 ordering domain。
Canonical Seq 为什么比 Per-Type Seq 更容易做 Recovery
Canonical run-local seq 可以回答:
E10 happened before E11
无论它们分别是:
- text;
- tool;
- interaction;
- red flag。
这对:
- after_seq;
- replay;
- projection;
- failure attribution;
都更简单。
Per-type seq 只能表达局部顺序,需要额外规则才能恢复全局因果顺序。
为什么叫 after_seq 而不是 from_seq
after_seq=N 的语义非常清楚:
N 已经处理过
只返回 N 之后的事件
如果用含糊的 from=N,容易产生:
N 是 inclusive 还是 exclusive?
恢复协议中 off-by-one 很容易制造:
- duplicate;
- lost event。
所以 cursor contract 要明确。
Cursor 是“客户端已经成功投影到哪里”,不只是“收到过哪里”
一个微妙区别:
socket received seq=8
不一定等于:
client successfully validated/reduced/applied seq=8
如果 parser 在 seq=8 失败,却把 maxSeq 先更新到 8,recovery 会从 8 以后开始,导致坏事件永久跳过。
更稳妥的 cursor 语义是:
last successfully accepted/projected seq
也就是先:
validate
→ reduce/apply accepted event
→ advance cursor
具体 effect 是否也必须成功后才能 advance,需要根据 effect durability/idempotency 设计决定。
Live Stream 如何记录最后位置
当前前端 SSE parser 会跟踪:
maxSeq
每解析一个 event:
maxSeq = max(maxSeq, event.seq)
正常看到:
stream.done
就不需要 recovery。
如果发生:
network/read failure
+
没有 stream.done
+
没有已知 business stream.error
才进入 durable catch-up。
这很重要,因为:
business error event
≠
transport error
两者的恢复策略不同。
为什么 Transport Error 不等于 Run Failure
假设:
Server Agent continues running
Browser Wi-Fi drops
浏览器看到的是:
reader.read() throws
它只能证明:
我失去了这条 transport。
不能证明:
Agent run failed
所以:
Transport Lifetime
≠
Run Lifetime
生产 Agent runtime 必须明确区分这两个生命周期。
Socket EOF 也不等于 Business Terminal
即使不是 error,而是:
reader.done = true
也可能是:
- proxy closed;
- server process restarted;
- intermediary timeout;
不一定是 run.completed。
所以终止最好由明确 durable event/state 表达:
run.completed
run.failed
run.cancelled
stream.done
而不是依赖 TCP/SSE EOF 解释业务语义。
Empty Page 为什么不代表完成
恢复时可能:
GET after_seq=5
→ events=[]
一个很自然但错误的实现是:
没有新 event
→ completed
但真实情况可能是:
Agent still reasoning
next event not committed yet
所以 durable recovery 需要 polling/catch-up:
empty page
→ wait
→ retry same cursor
直到看到显式终止事实:
stream.done
or
stream.error
[!important] Absence of new event ≠ terminal state
这和 Evidence 领域的:
absence of evidence ≠ evidence of absence
有相似的 epistemic 结构。
为什么必须等待 persisted Terminal Event
如果 termination 只由连接状态判断:
socket closed → done
那么 network drop 与正常结束无法区分。
正确模型:
stream.done
= explicit business/runtime terminal event
connection EOF/error
= transport observation
这是 Event Protocol 比“纯文本流结束”更强的地方。
Recovery 本身也是一个状态机
可以显式建模:
LIVE
↓ transport lost
RECOVERING(after_seq=N)
├─ page has events → apply → advance cursor → continue
├─ empty, run nonterminal → wait/poll
├─ terminal event → DONE/FAILED/CANCELLED
└─ timeout/budget exhausted → UNKNOWN/RECOVERY_FAILED
这比在一个 catch 里写几次 fetch 更容易测试和理解。
Catch-up / Live Race:恢复时新事件还在产生
假设:
Client consumed 1..5
network drops
恢复第一次查询:
after_seq=5
→ returns 6..10
但服务器同时已经产生:
11..13
客户端应用到 10 后必须继续:
after_seq=10
直到:
- 看到 terminal;
- 或追平且重新建立 live stream;
- 或 run status 明确完成。
一次 page fetch 不是自动“恢复完成”。
Page Size 与 Cursor Loop
如果 event log 很长:
limit=100
恢复逻辑应该:
fetch after_seq=N
apply ordered page
N = max accepted seq
repeat
而不是只拿第一页。
这本质上是 cursor pagination,只不过终点是动态增长的 event log。
Live 重连与 Polling Catch-up 怎么衔接
一种策略:
1. catch up durable events to current high-water mark
2. reopen live stream with cursor/resume token
另一种:
live stream immediately reconnects
+ durable catch-up in parallel
+ reducer dedupes by run/seq
第二种延迟低,但 race 更复杂。
无论哪种,都不要依赖“绝不会 overlap”。
因此前端要有:
at-least-once delivery assumption
+ idempotent projection
Exactly-once Delivery 为什么很难真正保证
经典场景:
server sends seq=8
client receives and applies
ACK/connection dies before server knows
重新连接时 server/client 很难同时确定:
seq=8 到底已经处理了吗?
所以分布式系统更实际的模式是:
at-least-once delivery
+
stable event identity
+
idempotent projection/effects
而不是宣称网络层提供 magical exactly-once。
Live 与 Replay 为什么要共用 Handler/Reducer
当前 recovery 把 durable events 再次送入同一组 handlers:
Live SSE
↓
handlers
↓
ActiveTurnReducer
Durable Page
↓
same handlers
↓
ActiveTurnReducer
好处是:
1. Projection Consistency
在线看到什么
≈
断线恢复后看到什么
2. 少两套业务逻辑
否则需要:
applyLiveToolCall()
applyReplayToolCall()
长期一定漂移。
3. Reducer tests 同时保护 live/replay
如果 event semantics 本身稳定,运输来源就不应改变 UI 语义。
Live 与 Replay 共用 Reducer,但不代表共用所有 Transport Metadata
例如:
SSE retry delay
HTTP polling interval
network error object
这些属于 transport state,不应强塞进 StreamEvent domain union。
保持:
transport-specific shell
→ normalized validated event
→ shared reducer
可以避免领域协议被 HTTP/SSE 实现细节污染。
幂等:Recovery 必须允许一定程度重复
现实里可能出现:
Live 已收到 seq=6
network error happened before maxSeq persisted locally
Recovery 又返回 seq=6
因此 projection 应尽量:
same identity event replay
→ state converges, no duplicate card
方式包括:
- cursor
after_seq; - event identity;
- keyed upsert;
- reducer seq guard;
- idempotent effects。
不能把“理论上不会重复”当成唯一安全假设。
Effect Idempotency 比 UI 去重更重要
重复 citation card 只是视觉问题。
更危险的是 replay 重复触发:
analytics
query cache mutation
notification
client-side action
所以 reducer/effect contract 应明确哪些 effect:
- pure projection only;
- idempotent upsert;
- exactly-once business command(通常不应该由 replayable public event 自动执行)。
Replayable event 不应该重新触发现实 side effect。
当前项目中的 lastSeqByType
当前 ActiveTurn reducer 还有:
lastSeqByType
用于处理项目历史上不同 seq spaces 混合的现实。
学习时应把它标为:
compatibility / implementation detail
而核心通用原则仍然是:
stable run identity
+
ordered event cursor
+
idempotent projection
Durable Event Log、Checkpoint 与最终 Read Model 的区别
三者很容易混淆。
Runtime Event Log
回答:
run 过程中发生了什么?
适合:
- recovery;
- audit;
- replay;
- projection。
LangGraph / Runtime Checkpoint
回答:
Agent 执行机当前内部应该从哪里继续?
它是 execution-resume state,不一定适合作为用户业务历史 API。
Durable Thread / Read Model
回答:
页面正常加载时,现在应该展示什么稳定结果?
它是面向查询的业务 projection。
所以:
Event Log ≠ Checkpoint ≠ Read Model
虽然三者都可以持久化。
为什么普通页面不应该每次重放所有 Delta
如果一个长会话有几十万 text delta,普通 GET 每次:
load all events
→ replay from zero
会越来越贵。
更合理:
Durable finalized messages/read model
→ normal page load
Runtime events
→ active run catch-up / audit / targeted replay
事件日志可能还需要 retention/compaction policy。
Event Retention 与 Replayability
如果 event log 不是永久保存,就要定义:
how long active-run recovery is guaranteed?
which semantic events are retained longer?
can old text deltas be compacted into final message?
例如:
raw text deltas retain 7 days
final message forever
control/safety events forever
具体策略由产品决定,但不能让 recovery contract 暗中依赖一个随时可能被清理的日志。
Recovery 与 Refresh 的区别
In-flight Recovery
当前 run 还可能继续:
live stream drops
→ after_seq catch-up
→ keep rendering current ActiveTurn
Full Page Refresh
应用重新启动后:
load durable thread/read model
+
if there is an active run
→ resume/catch up runtime events
这两种入口不同,但都依赖明确的 server truth。
Refresh 时 Cursor 从哪里来
如果 browser memory 已丢失,不能只依赖 local maxSeq。
可以从:
- durable read model last applied event seq;
- server active-run projection;
- persisted browser storage(若安全合适);
- current run event log + projection reconciliation。
关键是定义一个可靠 high-water mark。
如果没有,refresh recovery 可能只能从 run start 重放,代价更高但仍应正确。
Recovery Timeout 是资源边界,不是业务结论
当前 helper 有有限 timeout / poll interval。
如果 timeout:
recovery timed out
只能说明:
客户端在允许的恢复窗口内没有确认终止。
它不等于:
Agent definitely failed
后续 UI 可以提供:
- retry sync;
- reload thread;
- explicit run status query。
这是 transport/recovery budget 与 business state 的区别。
Recovery Timeout 后应该进入 unknown/reconnecting,而不是伪造 failed
如果没有 durable run.failed,最诚实的状态可能是:
run status unknown to client
UI 可以展示:
“连接已中断,正在重新同步”
或:
“暂时无法确认运行状态,点击重新同步”
而不是把 transport uncertainty 转成业务 failure。
Cancellation 不要和 Recovery 混
Network drop
→ recover
而:
User explicitly cancels run
→ durable cancellation command/state transition
不是同一件事。
单纯 AbortController.abort() 只保证浏览器停止等待 fetch;如果 server run 被设计为 detached durable execution,它可能仍继续。
所以 end-to-end cancel 需要:
Client cancel intent
→ server command
→ run state = cancelled
→ persisted terminal event
→ frontend projection
当前 BodySense Consultation 尚未看到完整显式用户 cancel-run path,因此这是概念目标而不是已完成能力。
为什么 Cancel Event 也要进入同一个 Durable Ordering
如果用户 cancel 与模型完成几乎同时发生:
seq=20 cancel accepted
seq=21 run.cancelled
或:
seq=20 run.completed
cancel arrives too late
必须有明确 ordering/authority。
否则客户端可能看到:
cancelled + completed
两个互相冲突 terminal states。
Canonical run event ordering 能帮助定义哪一个 transition 最终生效。
Terminal State 应当唯一且不可逆
典型 invariant:
completed
failed
cancelled
三者最多一个成为 durable terminal state。
terminal 后迟到的非修复事件:
must not reopen run
这是 event log 和 ActiveTurn reducer 都应该保护的 contract。
Run Isolation
after_seq 不能脱离 run identity:
错误:
GET /conversation/C1/events?after_seq=10
如果 C1 已开始 R2,seq=10 可能属于 R1。
更稳:
conversation_id + run_id + after_seq
或者 endpoint 天然 scoped 到 run。
同一个 conversation 的不同 run cursor 不能混用。
Security / Authorization 也属于 Recovery Contract
Durable event endpoint 不能因为“只是恢复”就放宽权限。
需要验证:
current user owns conversation/run
requested run belongs to conversation
no cross-user event leakage
尤其 event payload 可能包含:
- health information;
- tool inputs;
- safety state。
Recovery API 是正式数据读取边界,不是内部 debug endpoint。
测试清单
Cursor
after_seq=3
→ only events > 3
Cursor advancement
malformed seq=4:
must not silently advance to 4 then skip it forever
Empty page
空 page 后应继续 poll,不提前 done。
Pagination
超过 page size:
must continue with new cursor
Duplicate
重复 durable event 不产生重复 tool/citation/interaction UI 或 effect。
Terminal
只有 durable done/error/cancelled 等 terminal contract 结束 catch-up。
Network drop
live 收到 1..4 后断开,durable 5..8 能恢复到等价 final projection。
Catch-up race
recovery 6..8 与 reconnect live 7..10 overlap:
final projection still equals ordered unique 1..10
Run isolation
C1/R1 的 cursor 不得用于 C1/R2。
Terminal immutability
completed + late text.delta
→ no reopen
Append/deliver crash window
模拟:
persist success + delivery fail
→ recovery gets event
如果实现有反向 ordering,则必须有等价的可靠协议来证明不会丢失。
Authorization
用户不能读取其他 conversation/run 的 event log。
Property-style Recovery 测试
可以构造一个 canonical event list:
E1...E20
随机模拟:
- 在任意位置断线;
- 任意重复部分 live events;
- recovery 分页;
- live/recovery overlap。
最终要求:
projection(recovered delivery)
==
projection(canonical ordered event log)
这种 property test 很适合发现边缘 off-by-one / duplicate bugs。
一个完整 Recovery Race 例子
R7
seq 1..5 live applied
network drop
server persists 6,7,8,9
recovery after_seq=5 returns 6..8
client applies, cursor=8
live reconnect receives 8,9,10
8 duplicate ignored
9/10 applied
recovery after_seq=8 later also returns 9,10
duplicates ignored
seq=11 run.completed
terminal projection reached
最终 UI 应等价于:
fold(E1..E11)
而不是取决于“哪条 transport 先到”。
最终心智模型
SSE gives immediacy;Event Log gives durability;
after_seqconnects the two.
更完整:
Durable append
→ live at-least-once delivery
→ validated/idempotent projection
→ connection may fail
→ cursor catch-up
→ explicit durable terminal
→ reconcile read model
真正的 production streaming 不是“连接尽量别断”,而是:
连接可以断
事件可以重复
恢复可以和 live overlap
但业务执行和最终 projection 仍然能收敛
自测题
- 为什么 SSE 不应该是唯一 truth source?
- Append-before-deliver 解决了哪个 crash window?
- 为什么
(run_id, seq)比 timestamp 更适合 recovery cursor? - Cursor 应表示“收到过”还是“成功接受/投影过”?
- Empty recovery page 为什么不能当 terminal?
- 为什么一次 recovery page fetch 不一定完成 catch-up?
- Exactly-once 为什么通常不现实,什么模式更实用?
- Event Log、Checkpoint、Read Model 分别回答什么问题?
- Recovery timeout 为什么不能直接把 run 标成 failed?
- 为什么 AbortController 不等于 durable cancel?
- Terminal state 为什么必须唯一且不可逆?
- Live/recovery overlap 的最终正确性应该用什么 property 验证?