跳转至

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 待补)