跳到主要内容

数据截至 (上游 commit c76af90d88f4)

Hub 的状态与同步:缓存、版本号、消息账本

30 秒导读: 前几章讲了三方拓扑怎么连(01)、本地/远程怎么接管(02)、手机上按"允许"怎么回到 CLI(03)。这一章回答另一个问题:hub 自己脑子里那份"世界状态"是怎么维护的,又是怎么推给所有客户端的。 三件事:状态放哪(store + cache),并发写谁赢(版本号),消息怎么记账(seq / invoked_at / scheduled_at)。

不覆盖:SSE 连接怎么建、订阅怎么过滤、心跳怎么送到浏览器,那是传输层的事,见 01 §7——本章只讲事件是怎么产生、怎么补上 namespace、怎么交给 SSEManager 的。runner 与机器侧怎么凭空开会话见 05


1. 先说清楚问题:谁是真相源

HAPI 有三方:本机跑着的 CLI、中间的 hub、手上的 Web/PWA。三方都想知道同一件事——这个会话现在活着吗、模型是什么、消息到哪一条了

如果让每一方各存一份自己认为对的状态,很快就会打架:CLI 断线重连后以为自己还是 5 分钟前的样子,手机上刚改过模型,hub 又在后台把超时会话标成了 inactive。

HAPI 的选择很干脆:hub 是唯一真相源。 其他两方都只是"提交意图 + 订阅结果"。

于是 hub 内部必须回答三个工程问题:

问题HAPI 的答案本章小节
状态放哪,重启后还在吗SQLite 落盘 + 内存热缓存的两层§2
两方同时改一份数据,谁赢版本号乐观并发,三态 ack§4
消息发出去没人接怎么办消息是账本不是事件,可重放§6

2. 顶层全景:四层部件

先看结构。从下往上读:磁盘 → 内存 → 门面 → 出口。

┌──────────────────────────────────────┐
CLI (Socket.IO)───▶│ socket/handlers/cli/* 事件入口层 │
Web (REST) ──────▶│ 校验 payload + 校验 namespace 归属 │
└───────────────────┬──────────────────┘
│ 只调门面
┌───────────────────▼──────────────────┐
│ SyncEngine —— 对外唯一门面 │
│ (2400+ 行,自己几乎不存状态) │
└──┬─────────┬──────────┬──────────┬───┘
│ │ │ │
┌──────────▼──┐ ┌────▼──────┐ ┌─▼────────┐ ┌▼────────────┐
│SessionCache │ │MachineCache│ │Message │ │RpcGateway │
│会话热态 │ │机器热态 │ │Service │ │(见第 03 章) │
└──────┬──────┘ └────┬───────┘ └─┬────────┘ └─────────────┘
│ │ │
└──────┬──────┴───────────┘
│ ↘ 直接 emit 给 /cli 房间
┌─────────▼─────────┐
│ Store (bun:sqlite)│ ← 落盘: sessions / machines /
│ 子 store 装配 │ messages / message_epochs / …
└─────────┬─────────┘
│ 状态变了就发事件
┌─────────▼─────────┐
│ EventPublisher │──▶ SSEManager ──▶ 所有 Web 客户端
└───────────────────┘

各部件一句话职责:

部件干什么文件
Store打开 SQLite、跑 schema 迁移、装配各子 storehub/src/store/index.ts:61
SessionCache会话的内存热态(active/thinking 只活在这里)hub/src/sync/sessionCache.ts:17
MachineCache机器的内存热态(含 CPU/内存 health 快照)hub/src/sync/machineCache.ts:42
MessageService消息账本:写入、分页、取消、定时释放hub/src/sync/messageService.ts:158
SyncEngine对外唯一门面 + 5 秒后台节拍hub/src/sync/syncEngine.ts:173
EventPublisherSyncEvent 发给进程内订阅者 + SSEhub/src/sync/eventPublisher.ts:6

2.1 Store:七个子 store,八张必备表

Store 构造函数里干的事很朴素:建目录(0o700)、开库、打开 WAL、跑迁移,然后 new 出七个子 store 挂成只读字段(hub/src/store/index.ts:85-128)。

子 store 和表不是一一对应的——MessageStore 一个人管两张表,所以子 store 是七个,而开工前必须存在的表是八张(REQUIRED_TABLES,hub/src/store/index.ts:46-59):

Store ── 7 个子 store / 8 张必备表
├─ sessions SessionStore 表 sessions 会话行(含两个版本号列)
├─ machines MachineStore 表 machines 机器行(同构:也是两个版本号)
├─ messages MessageStore 表 messages 消息账本
│ 表 message_epochs 分页游标的 epoch
├─ users UserStore 表 users
├─ push PushStore 表 push_subscriptions Web Push 订阅
├─ fcm FcmStore 表 fcm_devices 设备 token
└─ scratchlist ScratchlistStore 表 session_scratchlist 每会话便签

三个细节值得记:

  • schema 版本是硬编码的阶梯:SCHEMA_VERSION = 25,initSchema1→2→…→25 一步步跑迁移函数,版本对不上直接抛错而不是猜(hub/src/store/index.ts:45:319-392)。
  • 缺表就拒绝开工:assertRequiredTablesPresent 拿这八张表名去 sqlite_master 里查一遍,少哪张就把哪张名字写进错误消息,让人去做离线迁移——而不是等运行时某条 SQL 崩掉(hub/src/store/index.ts:1069-1083)。
  • 跨表的事务只有一处:recordMessagesConsumed 把"标记消息已消费"和"更新会话活跃时间"包进一个 db.transaction,并在事务内断言写成功了才返回(hub/src/store/index.ts:137-159)。这是消息 ack 路径的原子点。

2.2 两个 Cache:内存里的"当前样子"

Cache 不是为了性能,而是为了存磁盘上根本没有的字段

对照一下 StoredSession(磁盘行,hub/src/store/types.ts:1)和 Session(对外视图,hub/src/sync/sessionCache.ts:178-213):

字段磁盘有内存有说明
metadata / agentState + 各自 version落盘的核心状态
model / effort / serviceTier运行时配置,落盘
active / activeAt有(但只写 false,见 §5.3)活跃判定
thinking / thinkingAt没有纯瞬时态
backgroundTaskCount没有后台任务计数
permissionMode / collaborationMode只在 metadata 里留偏好当前实际模式

refreshSession 就是这两者的桥:从 store 读行 → 用 Zod schema 收窄 metadata/agentState/todos保留内存里已有的瞬时字段 → 覆盖回 Map → 发一条 session-addedsession-updated(hub/src/sync/sessionCache.ts:126-217)。

注意这一行的写法:

active: existing?.active ?? stored.active,

hub/src/sync/sessionCache.ts:186内存里的活跃状态优先于磁盘——因为磁盘那份是"上次判死时写下的",不是当前事实。

2.3 SyncEngine:一个几乎不存状态的门面

SyncEngine 有 2400 多行,但它自己只持有五个 Set/Map(hub/src/sync/syncEngine.ts:180-203),其余全是转发。构造函数一眼看完它的装配顺序(:163-181):

this.eventPublisher = new EventPublisher(sseManager, (e) => this.resolveNamespace(e))
this.sessionCache = new SessionCache(store, this.eventPublisher)
this.machineCache = new MachineCache(store, this.eventPublisher)
this.messageService = new MessageService(store, io, this.eventPublisher,)
this.rpcGateway = new RpcGateway(io, rpcRegistry)
this.reloadAll()
this.inactivityTimer = setInterval(() => this.expireInactive(), 5_000)

能力分组来看这 2400 行,就不用逐行走读了:

能力组代表方法实际干活的是谁
读视图getSessions / getSessionByNamespace / getMachines / resolveSessionAccess两个 Cache
CLI 事件入口handleSessionAlive / handleSessionReady / handleSessionEnd / handleMachineAlive / recordSessionActivity两个 Cache
消息sendMessage / getMessagesPage / cancelQueuedMessage / getDeliverableMessagesAfterMessageService
会话生命周期archiveSession / renameSession / deleteSession / resumeSession / reopenSession / spawnSessionCache + RpcGateway
反向 RPC 门面approvePermission / readSessionFile / runRipgrep / listSkillsRpcGateway(见 03)
便签(scratchlist)listScratchlistEntries / createScratchlistEntrystore.scratchlist
后台节拍expireInactive(5 秒一跳)自己

resolveNamespace 是门面上一个不起眼但关键的钩子:任何事件出门前,若自身没带 namespace,就按 sessionId/machineId 反查补上(hub/src/sync/syncEngine.ts:241-252)。SSE 层靠它做多租户隔离广播——01 §7.2 里"事件没有 namespace 就一律丢",丢的就是这里漏补的那些。

2.4 EventPublisher:出门的唯一闸口

它只有 34 行,做两件事:给事件补 namespace,然后同时喂给进程内订阅者和 SSEManager(hub/src/sync/eventPublisher.ts:20-33)。监听器抛错被 try/catch 吞掉,不让一个坏订阅者拖垮广播。


3. 一次写入的完整走向(先看主线,再抠细节)

以"CLI 报告 agent 换了模型"为例,走一遍:

CLI: socket.emit('update-metadata', {sid, expectedVersion, metadata})

▼ hub/src/socket/handlers/cli/sessionHandlers.ts:176
Zod 校验 payload → resolveSessionAccess(sid) 校验 namespace 归属

▼ store.sessions.updateSessionMetadata(...)
UPDATE sessions SET metadata=?, metadata_version=version+1
WHERE id=? AND namespace=? AND metadata_version=@expectedVersion

├─ changes==1 ──▶ ack {result:'success', version, metadata}
│ ├─ 广播 'update' 给同房间的其它 CLI socket
│ └─ onWebappEvent → SyncEngine.handleRealtimeEvent
│ └─ sessionCache.refreshSession → EventPublisher → SSE

└─ changes==0 ──▶ 回读当前行 → ack {result:'version-mismatch', version, metadata}
CLI 用新版本重试(backoff)

三个要点,接下来三节分别展开:版本号(§4)、心跳(§5)、消息账本(§6)。


4. 乐观并发:一条 UPDATE 决定谁赢

4.1 要解决的小问题

同一个会话的 metadata,可能同时被两个人改:

  • CLI 检测到 agent 分配了新的 claudeSessionId,要写进去;
  • 你在手机上点了"重命名",hub 侧的 renameSession 也在写同一个字段块。

谁后写谁赢没问题,怕的是丢更新:后写的人拿着 5 秒前读到的旧快照,把对方刚写的字段整个覆盖掉。

4.2 思路:把"我以为的版本"写进 WHERE

不加锁,也不做长事务。做法是给每个可并发字段配一个版本号列(metadata_versionagent_state_version、机器侧的 runner_state_version),更新时把期望版本写进 WHERE:

UPDATE sessions
SET metadata = @field_value, metadata_version = metadata_version + 1
WHERE id = @id AND namespace = @namespace AND metadata_version = @expectedVersion

SQLite 返回的 changes 就是裁判:等于 1 说明你手上的版本还是最新的,写入成功;等于 0 说明有人抢先了。

4.3 原理演示

# 示意,非源码:乐观并发的最小骨架
def update_versioned(row_id, new_value, expected_version):
changes = db.execute(
"UPDATE t SET v=?, ver=ver+1 WHERE id=? AND ver=?", # 版本写进 WHERE
new_value, row_id, expected_version
).changes
if changes == 1:
return ("success", expected_version + 1, new_value) # 我赢了
current = db.query_one("SELECT v, ver FROM t WHERE id=?", row_id)
if current is None:
return ("error",) # 行没了
return ("version-mismatch", current.ver, current.v) # 顺手把最新值捎回去

重点看最后一行:冲突时不是简单返回失败,而是把"当前真值 + 当前版本"一起还给调用方,让它不必再发一次读请求就能重试。

4.4 真实实现

updateVersionedField 是这个招式的唯一实现,会话和机器的四个可并发字段全走它(hub/src/store/versionedUpdates.ts:20-61):

  • WHERE … AND ${versionField} = @expectedVersion::31
  • if (result.changes === 1) return { result: 'success', version: expectedVersion + 1, … }::40
  • 失败分支回读当前行,返回 version-mismatch + 当前值::44-58
  • 整个函数被 try/catch 包住,异常一律降级成 { result: 'error' }::59

结果类型就是三态(hub/src/store/types.ts:101-104):

result含义调用方该做什么
success写入成功,version 是新版本更新本地版本号,收工
version-mismatch别人抢先了,version/value 是当前真值用返回值重算,重试
error行不存在或 SQL 抛错上报错误,不要重试

调用侧四个入口,形状完全一致:updateSessionMetadata(hub/src/store/sessions.ts:277)、updateSessionAgentState(:321)、updateMachineMetadata(hub/src/store/machines.ts:188)、updateMachineRunnerState(:150)。

4.5 三态怎么跨到 CLI

协议里把这三态原样搬到了 socket ack 上(shared/src/socket.ts:184-208):

export type UpdateMetadataAck =
| { result: 'error'; reason?: SocketErrorReason }
| { result: 'version-mismatch'; version: number; metadata: unknown | null }
| { result: 'success'; version: number; metadata: unknown | null }

UpdateStateAck 同构,只是字段叫 agentState;机器侧还有 MachineUpdateMetadataAck / MachineUpdateStateAck(:166-190)。hub 侧的 handler 就是把 store 的三态逐个映射成 ack(hub/src/socket/handlers/cli/sessionHandlers.ts:255-261)。

CLI 侧的消费统一走 applyVersionedAck(cli/src/api/versionedUpdate.ts:23-64),它的逻辑很有意思:

  • successversion-mismatch 都要把返回的值和版本写回本地(:38-50)——冲突时 hub 捎回来的就是最新真值,白拿不用;
  • 然后 version-mismatchthrow(:52-54)。

抛出去给谁接?给 backoff:

this.metadataLock.inLock(async () => {
await backoff(async () => {
const updated = handler(this.metadata ?? {}) // 基于刚写回的最新值重算
const answer = await this.socket.emitWithAck('update-metadata', {
sid: this.sessionId,
expectedVersion: this.metadataVersion, // 已被上一轮 ack 刷新
metadata: updated
})
applyVersionedAck(answer, {})
})
})

cli/src/api/apiSession.ts:1276-1319(updateAgentState:958 同构)。backoff 是无限重试 + 随机化指数退避(250ms→1s,cli/src/utils/time.ts:12-41)。外层还套了 AsyncLock,保证同一个 CLI 进程内不会自己跟自己抢。

于是完整闭环是:

CLI 本地值 v3 ──emit(expectedVersion=3)──▶ hub 库里已是 v5

◀──ack{version-mismatch, v5, 真值}──┘
applyVersionedAck 把 v5 和真值写回本地
throw → backoff 等 250ms~1s
CLI 用 handler 在 v5 上重算 ──emit(expectedVersion=5)──▶ 成功 v6

4.6 为什么是版本号,不是锁

三条理由,都能在代码里找到落点:

  1. CLI 会断线重连。 心跳走的是 socket.volatile.emit(cli/src/api/apiSession.ts:1197),断线期间直接丢弃。如果 hub 给 CLI 发过锁,CLI 掉线后这把锁谁来解?版本号没有"持有者",掉线不欠债。
  2. hub 自己也会改同一行。 renameSession(hub/src/sync/sessionCache.ts:827)、markSessionArchivedFromHub(:617)、clearSessionArchiveMetadata(:716)都是 hub 侧直接写 metadata。它们用的是同一套三态 + 重试,最多试 5 次(METADATA_RETRY_ATTEMPTS,:14),5 次还冲突就抛错让 HTTP 层返回 409/5xx。
  3. 冲突本来就罕见。 真正高频的是心跳,而心跳走的是不带版本号的旁路(见 §5),根本不参与这套并发控制。

4.7 一个容易漏的细节:metadata 是合并写不是覆盖写

updateSessionMetadata 并不是把 CLI 给的对象直接落盘。它先在事务里读出旧 metadata,跑一遍 mergeSessionMetadata,再把合并结果交给 updateVersionedField(hub/src/store/sessions.ts:289-320)。

合并规则是"新的没提到就沿用旧的",只针对三组字段(hub/src/store/sessions.ts:123-133):

字段组成员为什么要沿用
PARSE_IDENTITY_FIELDSpathhost会话身份,写一次就不该被后续局部更新抹掉
ROUTING_FIELDSflavormachineId路由必需,丢了就找不到该发给谁
SIMPLE_RESUME_TOKENS各 flavor 的 resume id丢了就再也 resume 不回原会话

外加一条特例:cursorSessionProtocol 必须和 cursorSessionId 配对沿用——若新写显式给了新 cursorSessionId,旧协议就不能带过来(preserveCursorProtocolPair,:100-120)。这条特例服务的正是 06 §4.2 讲的 cursor 双协议选路:协议标记和会话 id 一旦错配,旧会话就会被送去走 ACP 路径而载入失败。

正因为落盘值可能和 CLI 提交的值不同,广播给同房间其它 CLI 的必须是合并后的值,源码注释专门点了这件事(hub/src/socket/handlers/cli/sessionHandlers.ts:272-278)。


5. 活跃判定与心跳:怎么知道会话还活着

5.1 CLI 每 2 秒喊一声

AgentSessionBase 构造函数里立刻喊一次,然后挂一个 2 秒的 setInterval(cli/src/agent/sessionBase.ts:81-84):

this.client.keepAlive(this.thinking, this.mode, this.getKeepAliveRuntime());
this.keepAliveInterval = setInterval(() => {
this.client.keepAlive(this.thinking, this.mode, this.getKeepAliveRuntime());
}, 2000);

心跳不是空包,它顺带捎上"当前运行时配置":thinkingmode(local/remote)、permissionModemodelmodelReasoningEfforteffortserviceTiercollaborationMode(协议定义见 shared/src/socket.ts:260-272)。状态一变还会立刻补喊一次,不等下一个 tick(onThinkingChange :85onModeChange :94)。

发送用的是 socket.volatile.emit(cli/src/api/apiSession.ts:1197)——断线期间丢掉就丢掉,2 秒后还有下一发。

5.2 hub 侧收到之后

链路:sessionHandlers.ts:356 校验 → SyncEngine.handleSessionAlive(:395)→ SessionCache.handleSessionAlive(sessionCache.ts:344)。

SessionCache.handleSessionAlive 里做四件事,顺序值得记:

第一步,先给时间戳消毒。 clampAliveTime 把未来时间压回 now,把超过 10 分钟前的时间直接判为无效返回 null(hub/src/sync/aliveTime.ts:1-7)。无效就整包丢弃。

第二步,更新活跃与思考态。 activeAtMath.max(只前进不回退),thinking 叠加一个 15 秒宽限窗(sessionCache.ts:380-391,宽限窗见 §5.5)。

第三步,按需把运行时配置落盘。 每个字段都先过 isStaleRuntimeKeepAlive 再写(:253-291):

if (payload.model !== undefined && !this.isStaleRuntimeKeepAlive(session.id, 'model', t)) {}

这个守卫解决的是一个真实竞态:你在手机上刚把模型改了,hub 通过 applySessionConfig 立刻记下"model 在 T 时刻被改过"(markRuntimeConfigUpdated,:584-592);此时一条 T 之前发出的、还带着旧模型的心跳到货了。守卫比较时间戳,旧于最后一次改动的心跳字段一律忽略(:594-601),否则界面会闪回旧值。

第四步,决定要不要广播。 不是每 2 秒都发 SSE,而是三选一才发(:301-305):

触发条件为什么
从不活跃变活跃状态跃迁,必须让 UI 知道
thinking 变了,或任一运行时配置变了有实质变化
距上次广播超过 10 秒保底刷新 activeAt

否则每个会话每 2 秒推一次,多开几个会话 SSE 就成噪声了。

5.3 一个反直觉的事实:active=true 从不落盘

翻遍 hub/src 只有两处调用 setSessionActive,而且都传 false:handleSessionEnd(sessionCache.ts:613)和 expireInactive(:480)。建行时插入的也是 active 常量 0(hub/src/store/sessions.ts:245)。

所以:"活着"是纯内存事实,只有"死了"才落盘。 后果很清爽——hub 重启后所有会话都从"不活跃"起步,最多 2 秒后各 CLI 的心跳会把还活着的那些重新点亮,而崩溃时残留的"假活跃"不会被持久化下来骗人。

5.4 三个生命周期事件

事件谁发hub 侧入口效果
session-aliveCLI 每 2ssessionCache.ts:344刷新 activeAt / thinking / 运行时配置
session-readyCLI 完成 agent 初始化,可以收 promptsyncEngine.ts:530记入 sessionReadyIds,作为 Pi/Cursor 恢复流程的判据
session-endCLI 退出前,带 SessionEndReasonsyncEngine.ts:548落盘 inactive + 清 thinking + 发 session-ended

SessionEndReason 是五值枚举:completed | terminated | error | handoff | cleared(shared/src/schemas.ts:13)。倒数第二个 handoff 就是 02 §2.8 讲的反向交班——会话结束了,但不是失败,紧接着会有一个本地终端把它接过去;最后的 cleared 是 OpenCode /clear 触发的收尾(cli/src/opencode/runOpencode.ts:579),它会让 session-end 跳过队列清扫。

handleSessionEnd 有一个早退守卫:如果 !active && !thinking 就直接返回(sessionCache.ts:608-610),避免重复的 end 事件把 activeAt 反复往前推。

结束时 SyncEngine 那一层还多做了几件收尾(syncEngine.ts:548-598):发 session-ended、按条件重跑去重(§7.3)、清掉 Pi 恢复相关的三个 Set。

5.5 thinking 的 15 秒宽限窗

你在手机上发了条消息,从"POST 成功"到"CLI 真的开始跑"之间有个空窗。这段时间 CLI 的心跳还带着 thinking=false,界面上的转圈会闪一下。

解法:markMessageQueued 在消息入队时把 thinking 强制置真,并记一个 now + 15s 的宽限截止(QUEUED_MESSAGE_THINKING_GRACE_MS,sessionCache.ts:9:342-367)。宽限期内 CLI 报 thinking=false 也不信(:244)。

但有些消息是同步处理完的(比如被 CLI 拦截的斜杠命令),根本不会有 thinking=true 跟上来,转圈就会卡满 15 秒。所以协议里加了个显式开关:CLI 在 messages-consumed 里带 clearQueuedThinkingGrace: true,hub 立刻撤销宽限(cli/src/api/apiSession.ts:1225-1235sessionHandlers.ts:418-420sessionCache.ts:486)。

5.6 机器侧:同构,参数不同

MachineCacheSessionCache 是一个模子:handleMachineAlive 同样先 clampAliveTime、同样 10 秒广播节流、同样 Math.max 推进 activeAt(machineCache.ts:198-225)。

差别在两点:

  • 机器心跳可以捎带 health(load1m / cpuPercent / memoryPercent / cpuCount / uptimeSeconds)。health 只活在内存里,refreshMachine 从 store 重建时是 health: existing?.health ?? null(machineCache.ts:183);广播判据里还要额外比一次"显示值是否真的变了"(healthDisplayChanged,:24-40),避免小数抖动刷屏。
  • runnerState 走的是和会话 agentState 完全同构的版本号通道(machineHandlers.ts:97-140),供 runner 用(见 05)。

5.7 全部时间常数一览

名目位置
CLI 心跳间隔2 秒cli/src/agent/sessionBase.ts:82
hub 后台节拍5 秒hub/src/sync/syncEngine.ts:223
会话广播节流10 秒hub/src/sync/sessionCache.ts:450
会话判死阈值30 秒hub/src/sync/sessionCache.ts:628
机器判死阈值45 秒hub/src/sync/machineCache.ts:228
心跳时间戳容忍窗10 分钟hub/src/sync/aliveTime.ts:5
排队消息 thinking 宽限15 秒hub/src/sync/sessionCache.ts:9
取消消息等 CLI ack500 毫秒hub/src/sync/messageService.ts:622
metadata 重试上限5 次hub/src/sync/sessionCache.ts:14

那个 5 秒节拍一跳干三件事(syncEngine.ts:924-964):会话判死 → 机器判死 → 释放到点的定时消息。第三件其实和前两件无关,源码注释坦白说"只是搭个便车,省一个定时器"。


6. 消息账本:为什么消息不是"事件"而是"账本"

6.1 四个字段撑起全部语义

messages 表结构很小(hub/src/store/index.ts:440-459),但有四个字段扛着全部状态机:

含义谁写
seq会话内单调序号,MAX(seq)+1addMessage 插入时
local_id客户端幂等键,也是唯一的 ack 路径发送方给
invoked_atCLI 已消费的时刻;NULL = 还在排队CLI 的 messages-consumed ack
scheduled_at定时发送时刻;NULL = 立即消息sendMessage 传入

配套两条索引直接编码了业务查询:按位置分页的 idx_messages_session_position(COALESCE(invoked_at, created_at) DESC, seq DESC),以及只索引待发定时消息的部分索引 idx_messages_scheduled_pending(WHERE scheduled_at IS NOT NULL AND invoked_at IS NULL)。

addMessage 里三条不变量值得抄走(hub/src/store/messages.ts:103-178):

  • 没有 localId 就没有 ack 路径,所以这类消息插入时直接盖上 invoked_at = now,否则会永远卡在"排队中"(:77-80);
  • 顺着这条推论,定时消息必须有 localId,没有就直接抛(:56-58);
  • localId 时先查重,命中就原样返回旧行——天然幂等(:60-67)。

6.2 sendMessage:一次写入,两条出口

sendMessage(sessionId, {text, localId, scheduledAt, …})

├─ 拒绝: scheduledAt + attachments 组合(附件目录会话结束就被清)

▼ store.messages.addMessage → 拿到 seq

├─ isFutureScheduled? ── 是 ──▶ 不发 CLI,等 5 秒节拍来释放
│ │
│ 否 ──▶ io.of('/cli').to(`session:<id>`).emit('update', …)

└─ 无论如何 ──▶ publisher.emit({type:'message-received', …}) ──▶ 所有 Web 端

源码在 hub/src/sync/messageService.ts:846-950。三个细节:

  • 给 CLI 的那条不走 EventPublisher,是直接对 /cli 命名空间的 session:<id> 房间 emit(:637)。给 Web 的才走 publisher(:641)。两条出口物理隔离。
  • isFutureScheduledDate.now() 是在 addMessage 之后重新取的,注释解释了原因:插入前取的 now 会让"刚好卡在边界上"的 scheduledAt 被误判成未来(:616-619)。
  • 定时消息带附件被显式拒绝(:587-589),因为附件放在会话上传目录里,会话结束就被 cleanupUploadDir 清掉,等它到点再发只会指向已删文件。

SyncEngine.sendMessage 在此之上再补两下:markMessageQueued(点亮 thinking,§5.5)和 recordSessionActivity(syncEngine.ts:997-1021)。

6.3 定时消息:故意不停地重发

releaseMatureScheduledMessages 每 5 秒扫一次 scheduled_at <= now AND invoked_at IS NULL 的行,逐条 emit 给 CLI(messageService.ts:1051-1092)。

关键在它不写 invoked_at(:728 有醒目注释)。后果是:只要 CLI 没 ack,这条消息每 5 秒就会被重发一次。

这不是 bug,是设计。因为 invoked_at 的唯一写入者是 CLI 的 ack,所以 hub 崩溃重启、CLI 重连、消息在网络上丢了——统统自动补发,不需要额外的重试队列。

代价是"已到点"这个 SSE 提示会重复,于是用一个进程内 Set 去重:每个 localId 每个 hub 进程只发一次 scheduled-matured(scheduledMatureNotifiedLocalIds,:86:707-709)。

6.4 会话结束时的清扫:为什么故意跳过所有定时消息

CLI 退出时,还排在队里的消息永远等不到 ack 了,浮动条会一直挂着。所以 session-end 会触发一次清扫(sessionHandlers.ts:508-513messageService.ts:968-981)。

但它只扫立即消息(scheduled_at IS NULL)。所有定时行——不管到没到点——一律跳过。源码在两处写了同一个结论(hub/src/store/messages.ts:588-600hub/src/sync/messageService.ts:952-967):

若这里给一条"已到点的定时消息"盖上 invoked_at,下一轮 mature 扫描(过滤条件是 invoked_at IS NULL)就再也看不见它了;等 CLI 重新挂上来,用户那条定时 prompt 就被静默吞掉。

一句话:invoked_at 的写入者必须唯一(只有 CLI 的 ack),任何抄近路都会破坏可重放性。

6.5 取消一条排队消息:五个分支

取消是这一章最绕的路径,因为它要跟"CLI 可能已经把消息 shift 出队了"抢时间。cancelQueuedMessage 先只查不删(messageService.ts:476-681):

lookupQueuedMessage(sessionId, messageId)

├─ absent ────────────▶ 当作已取消
├─ invoked ───────────▶ 已被消费,回完整的 invoked 行让 UI 校正
└─ queued
├─ 无 localId ─────────────▶ 直接删(本就没有 ack 路径)
├─ scheduledAt > now ──────▶ 直接删(从没发给过 CLI,问也白问)
├─ 房间里 CLI 数为 0 ──────▶ 直接删 + 复查是否被抢先
└─ 有 CLI 在线
emit('update', {t:'cancel-queued-message'}) 等 500ms ack
├─ removed ─────────▶ 删行 + 发 message-cancelled
└─ not-found/timeout ▶ 盖 invoked_at + 发 messages-consumed
(当作"已发送",别让消息凭空消失)

几个刻意为之的地方:

  • 先问 CLI 再删行。 顺序反过来的话,CLI 那边可能已经把消息喂给 agent 了,DB 却把行删了——用户看到消息消失,agent 却真的收到了。
  • CLI 离线时立刻删,否则 CLI 重连后会通过 seq 回填把这条已取消的消息重新捞起来(:428-435)。删完还要复查一次:万一 CLI 恰好在这两步之间连上并 ack 了,DELETE 会因为 invoked_at IS NULL 守卫而空转,此时按"已消费"处理(:442-455)。
  • ack 回调先看 responses 再看 err(:549-558):重连重叠期房间里可能有两个 CLI socket,一个超时一个成功,只要有人确认删掉了就算成功。

6.6 读取侧:分页、epoch、CLI 回填

三条读路径,服务三种客户端:

方法给谁特点
getMessagesPageWeb 线程视图COALESCE(invoked_at, created_at), seq 双键游标分页;最新页额外把所有未消费的排队行带外塞进来(messageService.ts:314-321)
getDeliverableMessagesAfterCLI 重连回填seq > ? 拉,且排除未到点的定时行(hub/src/store/messages.ts:365-385)
getQueuedStateWeb 浮动条只查一批 localIdinvoked_at 状态

第二条的过滤是必须的:不然 CLI 一重连就把未来的定时消息当普通消息全吃了,绕过了 mature 扫描这唯一通道(hub/src/store/messages.ts:357-364 的注释)。

message_epochs 表是给分页游标兜底的。会话历史被合并(§7.2)时 seq 会整体重排,此时 bumpMessageEpoch 把 epoch +1(hub/src/store/messages.ts:517-524)。客户端带着旧 epoch 来翻页,getMessagesPage 发现对不上就返回一个 reset: true 的最新页,让客户端丢弃本地缓存重来(messageService.ts:280-292)。


7. 会话生命周期的高级操作

前面几节讲的是"数据怎么流"。这一节是几个改变会话身份的操作,全部挂在 SyncEngine 上,实现落在 SessionCache

7.1 归档、重命名、删除

操作入口实质
归档syncEngine.ts:1662 archiveSessionRPC 杀掉 CLI 进程 → 再走一遍 handleSessionEnd
重命名syncEngine.ts:1860sessionCache.ts:827只改 metadata 里的 name,touchUpdatedAt: false
删除syncEngine.ts:1872sessionCache.ts:1051删行(消息靠外键级联)+ 清四个内存 Map + 异步删附件文件

archiveSession 有个细节值得单独说:如果杀进程的 RPC 抛 RpcTargetMissingError(CLI 早就没了,但内存里的 active 还没来得及对账),不当失败处理,而是转由 hub 自己把 lifecycleState 写成 archived(markSessionArchivedFromHub,sessionCache.ts:778-825)。这个 hub 侧写入本身也是幂等的:已经是 archived 就直接返回,避免重置 lifecycleStateSinceRpcTargetMissingError 为什么值得单独定义一个类型,见 03 §4.4。

删除有个硬前置:if (session.active) throw new Error('Cannot delete active session')(sessionCache.ts:1057-1059)。

这三个写 metadata 的操作全部用 §4.6 那套"读快照 → 试写 → mismatch 就 refresh 重试,最多 5 次"的循环。

7.2 合并:把两个会话的历史并成一条

两个入口,共用一份实现(sessionCache.ts:1088-1102mergeSessionData :901):

  • mergeSessions(old, new):搬完历史后删掉旧会话行;
  • mergeSessionHistory(old, new):只搬历史,保留旧会话行(旧的还活着时用这个)。

mergeSessionData 的执行顺序是被踩过坑的(:901-1010):

1. mergeSessionMessages 消息改挂到新 session_id
2. transfer scratchlist ← 必须在 deleteSession 之前!
3. 搬移附件文件 路径里嵌了旧 session id,要 re-key
4. 合并 metadata 版本号写,失败最多重试 2 次
5. 补齐 model / effort 等空字段 新的为 null 才用旧的

第 2 步的位置有注释说明原因:session_scratchlist.session_id 上有 ON DELETE CASCADE,晚一步做便签就被级联删光了。

消息层的 mergeSessionMessages 自己也有两处巧思(hub/src/store/messages.ts:1041-1167):

  • 目标会话的 seq 整体加上源会话的 oldMaxSeq,腾出前半段给旧消息,保证合并后仍单调;
  • local_id 撞车时,把旧那份local_id 清空并 invoked_at = COALESCE(invoked_at, created_at)——清空 local_id 等于切断了它的 ack 路径,不顺手盖上 invoked_at 它就会永远挂在浮动条上。

7.3 去重:同一个 agent 线程只留一条

同一个 Claude/Codex 线程可能对应多条 HAPI 会话行(重连、恢复、runner 重启都可能新建行)。deduplicateByAgentSessionId 负责收敛(sessionCache.ts:1513-1600)。

判据是 extractAgentSessionId:按固定优先级从 metadata 里取第一个非空的 agent 会话 id(codexSessionIdclaudeSessionIdgeminiSessionIdopencodeSessionIdgrokSessionIdcursorSessionIdpiSessionId,:1221-1232)。

策略:

候选 = 同 namespace 下该 agent id 相同的所有会话

├─ 候选 ≤ 1 ────────────▶ 收工
├─ 活跃候选 > 1 ────────▶ 什么都不做(hub 不知道用户正在看哪一个)
└─ 排序: active 优先 → updatedAt 新的优先 → 触发本次去重的那个优先
保留第一个作为 target,其余:
├─ 还活着 ──▶ mergeSessionHistory(保留行,只搬历史)
└─ 已不活跃 ▶ mergeSessions(搬完删行)

并发保护是一对 Set:deduplicateInProgress / deduplicatePending(:21-22:1246-1254)。同一 agent id 的第二次触发不会并发跑,而是记一个"待跑"标记,等当前这轮结束后再跑一遍——因为一个会话可能在第一轮进行中变成不活跃,第二轮才有资格删它。

触发点有四处,全在 SyncEngine:心跳(:408)、session-ready(:422)、session-end(:450)、5 秒判死之后(:757,且按 activeAt 倒序处理,保证同批过期的重复会话里最新那条被留下)。

7.4 兜底:从历史消息里反查 agent 会话 id

想 resume 一个会话,必须有 agent 侧的会话 id。正常情况它在 metadata 里,但也可能因为崩溃或旧版本没写上。

于是 resolveAgentResumeId 有个 ?? 兜底(syncEngine.ts:2389-2412):

if (flavor === 'codex') {
return metadata.codexSessionId ?? this.recoverCodexSessionIdFromMessages(session.id, namespace)
}

return metadata.claudeSessionId ?? this.recoverClaudeSessionIdFromMessages(session.id, namespace)

两条恢复路径的做法不同:

flavor方法怎么找
ClauderecoverClaudeSessionIdFromMessages(:1926)最近 200 条消息从新往旧扫,递归钻进 content/datasession_id/sessionId(extractClaudeSessionId :1988)
CodexrecoverCodexSessionIdFromMessages(:1937)同样倒扫,但只认 role=agent + content.type='codex' + scope 角色为 parent 的事件里的 thread id(extractCodexParentThreadId :1961)

Codex 那条多两道闸:

  • 遇到 "Context was reset" 事件立刻返回 null(isCodexContextResetMessage :1951)——上下文已经重置,再往前找到的 thread id 已经无效了;
  • 同一条事件里若出现多个不一致的 thread id,宁可返回 null 也不猜(:1976-1985)。

两条路径都必须过 normalizeAgentSessionId 的 UUID 正则(:2014-2022),找到后用版本号写回 metadata 落盘(persistRecoveredAgentSessionId :2024),下次就不用再扫了。这个 id 最终被 05 §8 的 resolveLocalResumeTarget 拿去拼 resume 参数。


8. CLI 事件入口层:三个文件的分工

hub 收 CLI socket 事件的地方在 hub/src/socket/handlers/cli/registerCliHandlers(index.ts:56)是装配点,它先做两件事,再把 socket 分给三个注册函数。

先做的两件事:

  1. 构造两个 access resolver(index.ts:61-87)。它们的返回值分三态:ok / access-denied(行存在但不属于本 namespace)/ not-found。这是多租户隔离的地基——每一个事件处理器的第一步都是调它
  2. 按握手参数入房间:socket.join('session:<id>')socket.join('machine:<id>'),且必须先过 access 检查(index.ts:89-98)。前面那些 io.of('/cli').to('session:<id>') 的定向广播,靠的就是这里的房间。

然后三个文件各管一摊:

文件管什么事件特点
sessionHandlers.ts:103messageupdate-metadataupdate-statesession-alivesession-readymessages-consumedsession-end唯一带 Zod schema 校验 + ack 回调的一组
machineHandlers.ts:42machine-alivemachine-update-metadatamachine-update-state与 session 版同构,只是换了表和字段名
rpcHandlers.ts:13rpc-registerrpc-unregister只有 29 行:往 RpcRegistry 里登记/注销方法名(见 03)

message 事件(CLI 上报 agent 输出)的处理器最值得看一眼(sessionHandlers.ts:106-233):它在写库之后,还要从消息内容里顺手抽三样东西——TodoWrite 待办列表、team state 增量、后台任务启停增量——各自更新会话字段并触发一次 session-updated。也就是说,hub 会解析 agent 的输出内容来维护结构化状态,而不只是当二进制字节转发。


9. 巧妙之处(可以抄走的)

1. 冲突时把真值一起还回去。 大多数乐观并发实现只回一个"冲突了",调用方还得再读一次。HAPI 的 version-mismatch 直接带上当前版本和当前值(versionedUpdates.ts:44-58),CLI 侧 applyVersionedAck 在 throw 之前先把它写回本地(cli/src/api/versionedUpdate.ts:38-50),于是重试天然是"基于最新值重算",少一个 round-trip。

2. 只落盘"死",不落盘"活"。 active=true 从不写库(§5.3)。重启即失忆,2 秒后靠心跳自愈——比"重启后要对账一遍谁还活着"简单太多。

3. invoked_at 只有一个写入者。 定时消息的释放路径故意不写 invoked_at(messageService.ts:1086),会话结束的清扫故意跳过所有定时行(:662-667)。守住这条不变量,"没 ack 就重发"就成了免费的容错,不需要单独的重试队列。

4. 心跳走旁路,不参与并发控制。 2 秒一次的高频状态用 volatile.emit 发、hub 侧 10 秒节流一次广播、字段级用时间戳挡回滚(isStaleRuntimeKeepAlive);低频的 metadata/agentState 才走版本号 ack。高频路径和一致性路径分开,两边都不难受。

5. 取消消息先问后删。 500ms 的 ack 窗口买来的是"要么真取消、要么标记为已发送",绝不会出现"UI 显示已取消、agent 却真收到了"(messageService.ts:622-680)。

6. metadata 是合并写,且合并规则只覆盖三组关键字段。 身份、路由、resume token 这三类"丢了就回不来"的字段被显式列白名单沿用(hub/src/store/sessions.ts:52-66:122-132),其余字段仍是覆盖语义——比全量深合并更可预测。


10. 边界与局限

诚实地说几条:

  • 单进程假设。 版本号 + changes===1 在单个 hub 进程内是可靠的;跨进程共享同一个 SQLite 文件时,热缓存不会互相通知(machineCache.ts:98-102 的注释明确说重试"只在另一个进程写同一个库时才有意义")。想水平扩多个 hub 进程,这层缓存需要重做。
  • 去重是尽力而为。 合并失败被 catch {} 静默吞掉,注释直说"重复项留着,靠 web 侧兜底隐藏"(sessionCache.ts:1591-1593)。两个都活跃的重复会话则完全不处理。
  • 定时消息的取消有 5 秒窗口。 释放路径每 5 秒扫一次,若 CLI 恰好在你按取消前把它 shift 出队了,结果是"标记为已发送"而不是"取消"(messageService.ts:1044-1050 的注释把这个契约写死了)。
  • 心跳丢失 = 判死。 30 秒内没心跳就标 inactive,网络抖动超过 30 秒会看到会话"假死"再复活。volatile.emit 决定了心跳不排队重发。
  • 合并会改写 seqmessage_epochs 让客户端整页 reset 来兜底,代价是合并后所有正在翻页的客户端都要丢缓存重来。
  • 消息表没有分区/归档。 单会话消息全量堆在一张表里,导出时分两档拦截:按字节估算超 SESSION_EXPORT_MAX_BYTES 的直接拒(too-large),只超条数上限 SESSION_EXPORT_MESSAGE_LIMIT 的降级为 warningforce 可强行导出(messageService.ts:240-256)。

11. 代码地图

主题文件符号
SQLite 装配、必备表与迁移阶梯hub/src/store/index.tsStoreSCHEMA_VERSIONREQUIRED_TABLESinitSchemaassertRequiredTablesPresent
消息消费的跨表事务hub/src/store/index.tsrecordMessagesConsumed
乐观并发核心hub/src/store/versionedUpdates.tsupdateVersionedField
三态结果类型hub/src/store/types.tsVersionedUpdateResult
会话行读写 + metadata 合并hub/src/store/sessions.tsupdateSessionMetadatamergeSessionMetadatapreserveCursorProtocolPairsetSessionActive
机器行读写hub/src/store/machines.tsupdateMachineMetadataupdateMachineRunnerState
消息账本hub/src/store/messages.tsaddMessagegetDeliverableMessagesAftergetMatureScheduledMessagesgetImmediateQueuedLocalMessagesmergeSessionMessagesbumpMessageEpoch
会话热缓存hub/src/sync/sessionCache.tsSessionCacherefreshSessionhandleSessionAlivehandleSessionEndexpireInactiveisStaleRuntimeKeepAlive
会话高级操作hub/src/sync/sessionCache.tsrenameSessiondeleteSessionmergeSessionDatadeduplicateByAgentSessionIdmarkSessionArchivedFromHub
机器热缓存hub/src/sync/machineCache.tsMachineCachehandleMachineAliverefreshMachinehealthDisplayChanged
心跳时间戳消毒hub/src/sync/aliveTime.tsclampAliveTime
对外门面 + 5 秒节拍hub/src/sync/syncEngine.tsSyncEnginehandleRealtimeEventexpireInactiveresolveNamespace
agent 会话 id 反查hub/src/sync/syncEngine.tsresolveAgentResumeIdrecoverClaudeSessionIdFromMessagesextractCodexParentThreadIdpersistRecoveredAgentSessionId
消息服务hub/src/sync/messageService.tsMessageServicesendMessagecancelQueuedMessagereleaseMatureScheduledMessagessweepImmediateQueuedOnSessionEndgetMessagesPage
事件出口hub/src/sync/eventPublisher.tsEventPublisher
CLI 事件入口装配hub/src/socket/handlers/cli/index.tsregisterCliHandlersresolveSessionAccess
会话事件处理hub/src/socket/handlers/cli/sessionHandlers.tsregisterSessionHandlershandleUpdateMetadatahandleUpdateState
机器事件处理hub/src/socket/handlers/cli/machineHandlers.tsregisterMachineHandlers
RPC 注册hub/src/socket/handlers/cli/rpcHandlers.tsregisterRpcHandlers
协议: 三态 ack + 事件签名shared/src/socket.tsUpdateMetadataAckUpdateStateAckClientToServerEvents
协议: 事件与补丁 schemashared/src/schemas.tsSyncEventSchemaSessionPatchSchemaSessionEndReasonSchema
CLI 侧 ack 消费cli/src/api/versionedUpdate.tsapplyVersionedAck
CLI 侧带退避的写cli/src/api/apiSession.tsupdateMetadataupdateAgentStatekeepAlive
CLI 侧退避实现cli/src/utils/time.tscreateBackoffexponentialBackoffDelay
CLI 侧 2 秒心跳cli/src/agent/sessionBase.tsAgentSessionBasestopKeepAlive

下一章: Runner 守护进程:从手机上凭空开一个新会话 —— 本章讲的 machine-aliveRunnerState 版本号通道,到那边才会看到它们真正的用途。