数据截至 (上游 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 迁移、装配各子 store | hub/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 |
EventPublisher | 把 SyncEvent 发给进程内订阅者 + SSE | hub/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,initSchema按1→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-added 或 session-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 / getDeliverableMessagesAfter | MessageService |
| 会话生命周期 | archiveSession / renameSession / deleteSession / resumeSession / reopenSession / spawnSession | Cache + RpcGateway |
| 反向 RPC 门面 | approvePermission / readSessionFile / runRipgrep / listSkills … | RpcGateway(见 03) |
| 便签(scratchlist) | listScratchlistEntries / createScratchlistEntry … | store.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_version、agent_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::31if (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),它的逻辑很有意思:
success和version-mismatch都要把返回的值和版本写回本地(:38-50)——冲突时 hub 捎回来的就是最新真值,白拿不用;- 然后
version-mismatch才throw(: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 为什么是版本号,不是锁
三条理由,都能在代码里找到落点:
- CLI 会断线重连。 心跳走的是
socket.volatile.emit(cli/src/api/apiSession.ts:1197),断线期间直接丢弃。如果 hub 给 CLI 发过锁,CLI 掉线后这把锁谁来解?版本号没有"持有者",掉线不欠债。 - hub 自己也会改同一行。
renameSession(hub/src/sync/sessionCache.ts:827)、markSessionArchivedFromHub(:617)、clearSessionArchiveMetadata(:716)都是 hub 侧直接写 metadata。它们用的是同一套三态 + 重试,最多试 5 次(METADATA_RETRY_ATTEMPTS,:14),5 次还冲突就抛错让 HTTP 层返回 409/5xx。 - 冲突本来就罕见。 真正高频的是心跳,而心跳走的是不带版本号的旁路(见 §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_FIELDS | path、host | 会话身份,写一次就不该被后续局部更新抹掉 |
ROUTING_FIELDS | flavor、machineId | 路由必需,丢了就找不到该发给谁 |
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)。