04 会话列表¶
状态:一期 durable dispatch 投影与
UserConversationState会话集合兜底(ADR-0021)已实现; 旧 mailbox-tail 压缩器章节仅作阶段二设计参考 更新日期:2026-09-26 上游契约:docs/PLAN.md§6.9.3、§7.5~§7.7、§12、§25.3 实施决策:ADR-0007(一期 Redis MailboxStore)
0. 文档边界¶
本文定义:投影压缩器实现、内存快照层、分页游标编码、未读重算算法、角标聚合。
本文不得重新定义:SESSION_DELTA 的绝对值语义与版本源、未读权威定义式、
UserSessionProjection / UserBadgeState 字段集、§12.9 控制事件表的行为。
1. 三层模型的实现归属¶
ConversationHead MessageStore(ScyllaDB) docs/02 §2.4
UserConversationState ScyllaDB 会话集合的持久权威
UserSessionProjection qsession + Redis checkpoint 可重建物化视图
UserBadgeState 目标态角标投影器(当前未落地) 可由 MessageStore + read state 重建
上表是 docs/PLAN.md 的目标权威模型。一期实现(ADR-0021,2026-09-26):
qsession 从 Redis convs:{user} 投影取得会话集合,再通过
MessageStore::conversation_summaries 读取会话头与定义式未读;群会话按成员可见区间裁剪——
base = max(read_conversation_seq, joined_at),已退群/被移除者的会话头与未读只看 left_at 之前
(经 MessageStore::visible_conversation_summary)。UserConversationState 一期只实现会话集合子集:
ScyllaDB 表 user_conversation_state 在会话首次进入某用户集合时登记一次(投影 Lua 同步写入 Redis
待登记队列,后台任务批量落表);checkpoint 越过 dispatch 保留窗口或投影整体丢失时递增投影纪元,
每个用户首次拉列表时从该表合并回投影。已知缺口:越窗期间首次出现且之后无消息的会话无法恢复。
UserBadgeState 尚未落地。
除 §2.2 明确描述当前 qsession 的部分外,§2.1 与 §3~§9 均是目标算法、协议和验收契约;
其中 mailbox-tail、projection_mailbox_seq、UserBadgeState、内存快照分页尚未按本文形态落地。
权威分层(§9.6.2,最容易被实现者搞反):
会话集合的权威 = UserConversationState
新设备即使没有任何投影,也能从这里拿到全量会话
投影 = 排序与未读的物化视图,可缺失、可整表重建
因此目标态 PULL_SESSION_LIST 的数据来源是“投影为主 + UserConversationState 兜底”,
而不是"只读投影"。投影缺失的会话必须以 ConversationHead join 出来的值填充。
2. 投影压缩器(目标态与当前实现对照)¶
2.1 阶段二目标驱动方式(§12.3.1)¶
输入 = MailboxNode 本地邮箱写入流的 tail,按 mailbox_seq 递增
**明令禁止按用户主键全表扫描** —— 那等于把"每次登录再计算"换成常驻全量扫描
loop {
// 消费上界必须是已物化水位,不能超前
let hi = min_over_lanes(watermark); // §12.3.1
let batch = mailbox_tail.read_until(hi);
for e in batch {
agg.entry((e.user_id, e.conversation_id)).apply(e);
dirty_users.insert(e.user_id); // RoaringBitmap
}
if elapsed >= projection_compaction_window || batch.len() >= 512 {
flush(&mut agg, &mut dirty_users); // per-user WriteBatch
}
}
消费上界必须是 min_over_lanes(W):若压缩器越过水位消费,一个重试晚完成的
低 seq 条目会被"mailbox_seq > projection_mailbox_seq 才应用"的规则直接丢弃,
投影永久少算一条。
dirty_users 丢失是安全的:它只是"哪些用户需要 flush"的优化,
丢失后退化为按当前聚合表全量 flush,不影响正确性。因此不进检查点。
2.2 当前一期实现:qsession 独立消费 dispatch,会话头读时现取¶
一期没有邮箱尾部压缩器。qsession 独立静态消费 dispatch 日志,以
cluster_id + topic_id + partition_count 绑定 Redis checkpoint;一条 record 的全部
recipients 幂等写入成功后才 CAS 推进该 partition 的下一 offset,启动追到初始高水位后
才监听。mailbox 发出的 Projection 只作低延迟补充,不是唯一事实源;断线、Redis
部分成功或进程崩溃都由同一 dispatch 重放收敛。
UserSessionProjection.projection_mailbox_seq(§7.6 的单调版本源)仍未按用户落地,
当前的可恢复顺序边界是 qsession 的 per-partition dispatch checkpoint;因此上面那条
"mailbox_seq > projection_mailbox_seq 才应用"规则不直接用于一期实现。
这本身不致命——致命的是当时投影行里还存着会话头
(latest_conversation_seq / last_activity_id / last_message_id)。会话头随
每条消息变化,于是"谁最后写谁赢"就要求投影事件按 conversation_seq 有序到达。
而它从来不是有序的:writer 对每个请求各自 spawn、writer→mailbox 走非键控
连接池、mailbox 对每条 dispatch 再 spawn(上限 DISPATCH_INFLIGHT_LIMIT)。
乱序落盘让会话头倒退,用户看到会话列表里的最后一条消息是倒数第二条
(2026-08-19 实测;低速率下更容易复现,因为消费者被逐条唤醒,乱序的两条
分属不同批次,同批合并根本碰不到)。
一期的解法是会话头哪儿都不存,两个键各归其位:
convs:{user} zset 成员 = hex(cid),分值 = 最后一条消息的 created_at(毫秒)
写用 ZADD ... GT —— 分值只增不减是 Redis 原生语义
hist:{cid} zset score = conversation_seq,member = 消息(writer 维护)
末条即会话头:member 里有 message_id,score 就是最新 seq
- 写路径零额外成本:
ZADD ... GT替代的就是原来那一次写,命令数不变;GT让活跃度天然单调,乱序事件写不小已有的分值,新会话照常加入。 顺序要求不是被防住了,是不存在了。 - 读路径反而更快:
ZREVRANGE convs:{user} 0 N直接给出按活跃度排序的 首页,重度用户(5000 会话,附录 B.5)不再"全量载入 + 内存排序 + 截断"。 本页每个会话再取两条命令(ZREVRANGE hist 0 0 WITHSCORES取会话头、ZCOUNT算定义式未读),一次流水线往返,规模封顶在session_list_page_limit。 - 页内仍按 §6.9.3 精确重排一次:zset 分值是毫秒,同毫秒内不区分,
而契约要求
last_activity_id DESC + conversation_id ASC稳定断连。
为什么不给会话头单独建一个键:试过让 writer 的分配脚本顺带写
head:{cid}——INCR cseq:{cid} 是会话的串行化点,同一次原子执行里写头,
单调性由构造保证、零比较。设计上很干净,代价却不可接受:单线程 Redis
已近饱和,每消息多一次写就把排队推上去,SEND_ACK P50 6.3 → 10.0 ms、
端到端 P99 239 → 443 ms(同机 A/B,只增删那一行)。这与 §12 反复出现的
"热路径禁加 Redis 侧保险"是同一条教训:近饱和时 ρ 的微增就是排队的剧增。
没有残留窗口:会话头不落库,也就不存在"落库顺序"问题;活跃度由 GT
原生保证单调;两者都与 SessionProjection 进程的存活无关,压测中途重启
session 服务,会话列表的最新消息仍然正确(已实测)。
3. 未读¶
3.1 权威定义式(§12.5.1)¶
unread(u,c) = |{ m : base(u,c) < m.conversation_seq <= head.latest_conversation_seq
∧ m.sender_id ≠ u
∧ m.counts_unread
∧ ¬recalled ∧ ¬deleted }|
base(u,c) = max(read_conversation_seq,
hidden_before_conversation_seq,
joined_at_conversation_seq,
deleted_before_conversation_seq)
unread_count 是该式的缓存;read_conversation_seq 是权威输入;冲突以定义式为准。
下界是开区间(base < seq),与 §7.5 的 visible() 一致(§12.5.1 的开下界约定)。
3.2 三态判定(unread_base_seq)¶
match (new_base, proj.unread_base_seq) {
// (a) base 前进且窗口内条目仍在手 -> 增量扣减
(nb, ob) if nb > ob && have_entries_in(ob..nb) => {
let (dec_unread, dec_mention) = count_in(ob..nb);
proj.unread_count = proj.unread_count.saturating_sub(dec_unread);
proj.mention_count = proj.mention_count.saturating_sub(dec_mention);
recompute_mention_first_if_needed();
proj.unread_base_seq = nb;
}
// (b) base 前进到 latest -> 清零
(nb, _) if nb >= head.latest_conversation_seq => {
proj.unread_count = 0; proj.mention_count = 0;
proj.mention_first_conversation_seq = None;
proj.unread_base_seq = nb;
}
// (c) 其余 -> 有界重算
_ => recompute_bounded(),
}
(a) 的可判定前置:have_entries_in(ob..nb) 的判据是
"该区间完整落在当前压缩窗口持有的条目集合内",即 ob >= 本窗口起始 mailbox_seq。
不满足就走 (c),不允许猜。
(a) 必须同步处理 mention:只扣 unread_count 不扣 mention_count 会让
提及数永久偏高。
3.3 有界重算¶
latest - base <= unread_precise_limit(200)
-> 向 MessageStore 发一次单分区范围读精确重算(docs/02 §4.3 的分组读)
> 200
-> unread_count = 200, unread_exact = false,与 UI 的 99+ 截断闭环
3.4 reconcile 触发点¶
新设备 / 登录、检查点回滚重放、MAILBOX_DIRTY、
read_conversation_seq 被跨设备前推到非最新位置(此时必须**重算**而不是清零)
后台对"有未读且 24h 无变更"的会话按 unread_audit_sample_rate 抽样对账
指标 unread_drift_ratio 入 §24
3.5 撤回与删除的扣减¶
条件(缺一不可,字段来自 §7.3 的 target_*):
target_sender_id != 本人
target_flags.counts_unread == true
target_conversation_seq > base(u,c)
目标消息此前未被扣减过(由 projection_mailbox_seq 水位保证幂等)
clamp: unread_count = max(0, unread_count - 1)
缺 target_sender_id 就会撤回自己发的消息时错误扣减本人未读——
这正是 §7.3 引入三个 target_* 字段的原因。
4. 角标聚合(§7.7)¶
total_unread = Σ unread_count over { membership_state=ACTIVE ∧ ¬archived
∧ (¬muted ∨ 租户策略 include_muted) }
total_mention = Σ mention_count over { membership_state=ACTIVE }
(静音会话仍计入,与"静音只影响通知"一致)
muted_unread = Σ unread_count over { muted 会话 }
增量维护而非每次全量求和:flush 时对每个变更会话计算
delta = new_unread - old_unread,累加到 UserBadgeState。
每次 reconcile(§3.4)后强制全量重算一次,纠正累积漂移。
下发:在线走 BADGE_UPDATE,离线走 APNs/FCM 的 aps.badge(docs/09)。
AUTH_OK 携带登录快照(含 muted_unread)。
5. SESSION_DELTA¶
5.1 绝对值语义(§12.6)¶
版本源 = projection_mailbox_seq(不新增计数器)
客户端规则:projection_mailbox_seq 更大才应用,绝对值直接覆盖,否则丢弃
合并:只允许"同 conversation_id 后帧整体覆盖前帧"的丢弃式合并
**禁止任何累加型语义合并**
窗口 session_delta_merge_window(100~200ms)
缺口自愈:客户端发现 driving_mailbox_seq 与本地不连续时
**不做补偿计算**,直接 PULL_MAILBOX 重算
driving_event_id + driving_mailbox_seq 仅用于缺口定位与 §24 追踪,
不参与版本比较(§12.6.2)。
5.2 为什么不能用增量¶
推送通道是至少一次且可丢弃的(§11.3 超软水位即停推)。增量语义在
MAILBOX_DIRTY 重放与连接替换丢帧两条正常路径下必然双加或少算,且永久漂移。
REACTION_UPDATE(§13.6.2)出于完全相同的理由也是绝对值帧。
6. 内存快照层与分页¶
6.1 快照¶
排序**不发生在存储层**:last_activity_id 只是值列(§7.6)
SessionProjection 服务在内存中按 §6.9.3 的排序键建快照:
置顶区:pin_rank ASC, last_activity_id DESC, conversation_id ASC
普通区:last_activity_id DESC, conversation_id ASC
冷用户按需一次分区读全量载入(<= max_conversations_per_user 5000 行)
snapshot_revision:服务端载入时分配的单调版本号,首次取值 1,TTL snapshot_ttl(5 min)
这避免了把 last_activity_id 作为聚簇键带来的墓碑风暴——它每条消息都变。
6.2 keyset 分页游标¶
page_cursor = base64url( protobuf{
snapshot_revision : u64
zone : PINNED | NORMAL
pin_rank : u32 // zone=PINNED 时有效
last_activity_id : bytes16
conversation_id : uuid
} )
下一页条件(NORMAL 区):
(last_activity_id, conversation_id) < (cursor.last_activity_id, cursor.conversation_id)
—— 严格字典序,保证无重复无遗漏
服务端必须校验 snapshot_revision:与当前快照不符时返回新 revision 的第一页,
并置 projection_complete=false 提示客户端重新同步。
6.3 分页中途的结构性变更(§12.8.2)¶
置顶 / 取消置顶 / 归档 / 退群 —— 改变的是会话的**分区归属或排序键**,
不是内容。这类变更必须**强制失效 snapshot_revision**,
因为 keyset 分页的前提是排序键在遍历期间稳定。
内容变更(新消息、未读变化)不失效 revision:
客户端通过 SESSION_DELTA 按 conversation_id 覆盖合并即可补回
这也是 SESSION_DELTA 必须是绝对值帧的又一个理由
7. 会话复活与删除(§12.10)¶
清空聊天记录 写 hidden_before_conversation_seq = head.latest_conversation_seq
会话**保留在列表中**,内容与未读按 base 清零
删除会话 写 deleted_before_conversation_seq = head.latest_conversation_seq
会话**移出列表**;新消息到达时复活(只展示 deleted_before 之后的)
目标态重建规则(§12.2;当前 `UserConversationState` 兜底实现缺口见 §1):
deleted_before_conversation_seq >= head.latest_conversation_seq -> 不产出投影行
hidden_before 只影响可见内容与未读 base,不影响列表成员资格
目标设计让两个水位写不同字段,是为了让投影整表删除后重建能得到确定结果—— 写同一字段则"清空过的会话"与"删除过的会话"持久状态不可区分。
8. 容量(§25.3)¶
每压缩窗口投影持久写入量 ≈ 窗口内活跃用户数 × 该用户窗口内发生变化的不同会话数
自检断言的前提:压缩比 = 窗口内同一 (user, conversation) 的平均事件数。
长尾场景该值趋近 1,此时写入量与"消息数 × 收件人数"恒等 ——
**这是模型固有下界,不是实现错误**(§25.3 已说明)。
§2.3 成功标准第 5 条:单成员持久投影写次数
<= ceil(消息时间跨度 / projection_compaction_window) + 1
(+1 为消息流结束后的收尾 flush)
9. 验收¶
J-1【发布阻断】投影写压缩比
1000 人群连续 1000 条消息(20 msg/s,跨度 50 s,窗口 5 s):
每用户每会话持久投影写次数 <= ceil(50/5)+1 == 11,压缩比 1000/11 >= 90
J-2【发布阻断】SESSION_DELTA 绝对值幂等
同一帧重复投递 100 次:客户端 unread_count 恒等于帧内绝对值(不是 100 倍)
乱序投递旧版本帧:应用次数 == 0
代码级断言:SESSION_DELTA 处理函数不出现 "+=" 型未读更新
J-3 压缩器不越水位
构造一个重试晚完成的低 seq 条目:
压缩器消费上界恒 <= min_over_lanes(W);该条目最终被计入,投影少算次数 == 0
J-4 撤回不误扣自己的未读
撤回本人发送的消息:本人 unread_count 变化 == 0
撤回他人 counts_unread 消息且 seq > base:unread_count 减 1,且不为负
J-5 分页一致性
5000 会话用户遍历全部页:无重复无遗漏
遍历中途置顶一个会话:snapshot_revision 失效,客户端重新拉第一页
遍历中途来新消息:revision 不失效,通过 SESSION_DELTA 覆盖合并补回
J-6 投影整表重建
删除全部 UserSessionProjection 后重建:
排序、未读、提及逐字段与删除前一致(deleted/hidden 两个水位的语义可区分)
J-7 角标不撕裂
在 flush 中途注入崩溃:UserBadgeState.total_unread 与
Σ UserSessionProjection.unread_count 的差异 == 0(同批提交验证)
10. 待办¶
[ ] unread_audit_sample_rate 的默认值与对账任务实现
[ ] 内存快照的驱逐策略(按 snapshot_ttl + LRU,需与 max_conversations_per_user 对齐)
[ ] MAILBOXED->APPLIED 到达率对账任务(§2.4 引用了它但 §19.1 未定义,见 PLAN 待补)