Durable SSE Recovery with after_seq

解释实时 SSE 仅作为低延迟 transport,而 durable Runtime Event Log 作为恢复来源;客户端用 run identity、monotonic seq 与 after_seq 在断线后追赶缺失事件,并让 live/replay 共用同一 projection。

#type / synthesis #status / growing #tech / architecture #tech / dev / frontend #tech / dev / backend #resource / bodysense #resource / http #resource / typescript

[!info] related notes

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_seq connects 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 仍然能收敛

自测题

  1. 为什么 SSE 不应该是唯一 truth source?
  2. Append-before-deliver 解决了哪个 crash window?
  3. 为什么 (run_id, seq) 比 timestamp 更适合 recovery cursor?
  4. Cursor 应表示“收到过”还是“成功接受/投影过”?
  5. Empty recovery page 为什么不能当 terminal?
  6. 为什么一次 recovery page fetch 不一定完成 catch-up?
  7. Exactly-once 为什么通常不现实,什么模式更实用?
  8. Event Log、Checkpoint、Read Model 分别回答什么问题?
  9. Recovery timeout 为什么不能直接把 run 标成 failed?
  10. 为什么 AbortController 不等于 durable cancel?
  11. Terminal state 为什么必须唯一且不可逆?
  12. Live/recovery overlap 的最终正确性应该用什么 property 验证?
创建于 2026/8/22 更新于 2026/8/23