跳转至

03 个人邮箱与大群分发

状态:可实施 更新日期:2026-09-01 上游契约:docs/PLAN.md §6.5、§6.7、§7.3、§7.8、§7.9、§9、§10、§18.1、§18.3 实施决策:ADR-0001(MailboxStore 选型)、ADR-0004(复合序号)、ADR-0007(一期 Redis)、ADR-0012(Fanout 无 Kafka 事务)

0. 文档边界

本文定义:MailboxStore 三种实现的落地形态、lane 子任务调度、dispatch 分块与 进度提交、成员 Bitmap 与槽位分配、漂移接管算法、REBUILD 执行器。

本文不得重新定义:mailbox_seq 格式与 lane 推导(§6.5)、事件组上限与 event_id/event_ordinal 生成(§6.7)、UserMailboxEntry 字段集(§7.3)、 水位推进规则(§9.2)。


1. MailboxStore 抽象

四原语见 §18.1.2,此处给出 Rust trait 形态。三种实现语义必须完全一致, 由 §18.1.2 的影子读比对验证。

#[async_trait]
pub trait MailboxStore: Send + Sync {
    /// 同一 (shard, lane, mailbox_seq) 的整个事件组要么全可见要么全不可见。
    /// 必须与 DispatchProgress 的分块进度按 §10.4.1 的顺序约束提交。
    async fn append_batch(&self, shard: ShardId, lane: LaneId,
                          seq: MailboxSeq, entries: &[Entry],
                          chunk: ChunkRef) -> Result<Durable>;

    /// after_seq 开区间下界,up_to_seq 闭区间上界;
    /// 切分只能在 mailbox_seq 边界,max_items / max_bytes 是软上限。
    async fn range_scan(&self, user: UserKey, after_seq: MailboxSeq,
                        up_to_seq: MailboxSeq, max_items: u32, max_bytes: u32)
                        -> Result<(Vec<Entry>, MailboxSeq /* covered_through */)>;

    /// 返回后 seq 之前的条目对读路径不可见;物理回收按实现分派。
    async fn truncate_before(&self, user: UserKey, seq: MailboxSeq) -> Result<()>;

    async fn watermark(&self, shard: ShardId, lane: LaneId)
                       -> Result<(MailboxSeq /* W */, MailboxSeq /* ttl_frontier */)>;
}

1.1 一期实现:Redis(ADR-0007)

键设计见 §18.1.2b。本节给出可执行细节。

部署边界:当前生产只支持 Redis 7.0.0+ standalone,要求 INFO server 中唯一、严格三段式 redis_version >= 7.0.0,并启用 AOF、appendfsync=everysec|always 与 maxmemory-policy noeviction。Redis Cluster 未实现,检测到 Cluster 端点必须拒绝启动; QIM_REDIS_ALLOW_UNSAFE_DEV=1 也只跳过持久化/淘汰策略,不绕过版本与 standalone 校验; 下文的日桶与 Lua 只描述 standalone 行为,不能被解读为 Cluster 兼容性。

mailbox:{tenant}:{user}:e{epoch}:{day}  zset   score = 精确 mailbox_seq,member = 版本化编码 Entry
mbmeta:{tenant}:{user}:e{epoch}         hash   retention_format_version=4、显式 trim、活动日缓存、expired_gap_boundary
mbdays:{tenant}:{user}:e{epoch}         zset   member = day,score = 该日桶固定绝对 expiry
mbseqbucket:{tenant}:{user}:e{epoch}    zset   member = 完整 seq hex,score = 事件组日桶固定绝对 expiry
wm:{shard}:{lane}                       hash   W、W_floor
dp:{shard}:{dispatch_id}                hash   canonical identity/manifest、chunk 进度、max observed offset、last_seen

day_bucket 取自 Entry.created_at(而非写入时刻)——created_at 来自 GroupDispatch.committed_at(§6.7 的确定性要求),因此重放落在同一个 key, 是同值覆盖而不是产生第二份。

AppendBatch 的 Lua 脚本按单用户原子写入,保证事件组不半可见,但绝不同时 标记 dispatch 进度:

-- KEYS: 该用户涉及的 mailbox 日桶 + mbmeta + mbdays + mbseqbucket
-- ARGV: 编码后的 (key_index, 精确 log_offset, member) 三元组序列
-- 1) 逐条 ZADD(同 score 同 member 为幂等覆盖)
-- 2) 固定每桶绝对过期并记录活动日 proof;同一 seq 的日桶身份不可改变
-- 3) 两阶段验证/提升已过期日到 expired_gap_boundary,验证失败前不修改任何 proof
-- 4) 只接受 v4;V1~V3、身份索引缺失或跨日重放冲突均拒绝写入

worker 只有在某个 (lane, chunk) 的全部用户 append 成功后,才单独调用 mark_chunk_completed;该方法按 begin_dispatch 持久化的 canonical manifest 校验。 所有必需 chunk 完成后再以 current W 的 CAS 推进到本次 observed offset。顺序不可颠倒: 先条目、后 chunk、最后水位;任一失败都不得把水位越过未物化数据。

DispatchProgress 默认不 GC(replay safety window = u64::MAX);只有部署方提供 经审计的日志重放窗口,且所有 lane 的 W/W_floor 越过最大 observed offset、距最后一次 重复记录超过该窗口时才可删除。旧 mbmeta 没有桶最大序号/绝对过期证明,必须迁移或 从分发日志重放,不能把历史 TTL 数据静默解释成空邮箱。

日桶 TTL 到期与显式 TruncateBefore 是两条不同路径。自然到期只把已完整验证的 活动日最大序号合并进 expired_gap_boundary,显式裁剪才推进 user_trim_seq;两者都会 令跨过边界的旧游标返回 CursorExpired。提升脚本先完整验证全部到期日 proof,再统一 写 gap/HDEL/ZREM,任何损坏前零写入。RangeScan 在读取前后各验证一次,避免桶在 校验与 ZRANGE 之间到期后被误报为空且推进 covered。活动日 proof 最多 31 天,事件组 身份索引使用绝对 EXPIREAT;V4 以前的布局不能安全混读或混写。

为什么当前必须拒绝 Redis Cluster:这不是给 mailbox key 补局部 hash tag 就能完成的 问题。GroupMembership 的成员版本/快照会跨槽;MessageStore 的全局 {commit} 键会把 提交热路径压到单个热槽;现有 ConnectionManager 也没有 Cluster 路由、MOVED/ASK 重试 与拓扑刷新。因而跨槽 Lua 不能伪装成原子提交,按 slot 分组更不能形成完整正确性证明。 Cluster 支持必须另行完成状态机、迁移和端到端验证后才可评估。

RangeScan:

1. 从 mbmeta 的真实 min/max day 得到仍在保留窗口内的 day_bucket 集合
2. 每个 key 一次 ZRANGEBYSCORE ({after} {up_to} LIMIT 0 N),pipeline 并发
3. 归并排序后按 §9.3.3 规则 2 在 mailbox_seq 边界切分:
   若最后一个 mailbox_seq 的事件组被截断,回退到上一个完整组的 seq
4. covered_through_seq 按 §9.3.3 规则 2 取值(扫到 up_to_seq 则取 up_to_seq)

TruncateBefore:对仍存活的 day key 执行 ZREMRANGEBYSCORE 细粒度删除, 随后 HSET mbmeta user_trim_seq。自动 TTL 到期不调用该路径,也不改标量 trim。

持久性缺口:standalone Redis 的 AOF everysec 有 1 秒窗口(always 更强),唯一兜底是 Redpanda 分发日志重放(§18.1.2b);noeviction 是提交与邮箱键不被逐出而破坏该证明的前提。 实现约束:Redis 与分发日志的保留期必须联合校验, log_retention_days 覆盖不到 Redis 的崩溃窗口时,ADR-0007 的兜底论证失效。

1.2 阶段一实现:ScyllaDB

DDL 见 §18.1.3。与 Redis 实现的两处关键差异:

原子性   跨分区无法原子批量(entry 分区键是 (tenant,user),
         DispatchProgress 是 (tenant,shard,dispatch_id))
         -> 必须走 §10.4.1 的等价方案:严格顺序(先条目后进度)+ 确定性幂等重放

TruncateBefore  纯逻辑操作:只推进 user_trim_seq,**禁止下发任何 CQL DELETE**
                物理回收完全由 default_time_to_live + TWCS 整文件过期承担(§18.1.2)

1.3 阶段二实现:自研 LSM

rust-rocksdb(ADR-0006)。条目与 DispatchProgress 的 chunk 位在同一个 WriteBatch 内 fsync,天然原子,不需要 1.2 的等价方案。

代价是要自建复制、检查点、接管——ADR-0001 的四条切换判据成立前不做。


2. 成员 Bitmap 与槽位

2.1 槽位分配(§7.9)

slot_id 在 (tenant_id, group_id) 内单调递增分配,**永不复用**
退群只置 released_at 并从当前版本 Bitmap 清位
槽位空洞由 RoaringBitmap 的稀疏压缩吸收,不需要紧凑化

重编号条件:max(slot_id) > member_count × slot_compaction_ratio(4)
  必须在一次**停写窗口**内完成,并强制递增 membership_version
  重编号期间所有旧版本 Bitmap 一律作废

为什么永不复用:slot 复用会让旧版本 Bitmap 中的某一位在新版本里指向另一个人, 用旧 membership_version 展开时直接跨用户错投。这是 §7.9 的核心约束。

2.2 在线 Bitmap 必须与成员 Bitmap 同源

两者共用同一套 (tenant, group) -> slot 映射,否则求交结果无意义。
在线 Bitmap 由 PresenceDirectory 的订阅流维护(§5.4),
语义是"该用户至少一个设备在线"的**快速过滤器**;
per-device 明细一律取自 presence 缓存,不从 Bitmap 推断(§16.1 的 PD-1 不变量)。

2.3 版本存储与缓存

GroupMembershipVersion 按 (tenant, group, version) 分片存 RoaringBitmap
  小群(< 1 万成员):直接存 ScyllaDB blob 列
  大群:blob 超行阈值时转存 S3,ScyllaDB 只留指针

序列化必须用 RoaringFormatSpec 的 portable 格式(§7.9 契约性约束)

节点侧按 (group, version, shard) 做不可变 LRU 缓存 —— 不可变意味着无需失效逻辑
保留期:membership_version_retention >= max(dispatch_progress_retention,
        log_retention_days) + 安全余量(§10.4.3 的单向依赖链)

3. 分发执行器

3.1 主循环

consume(dispatch) ->
  1. dispatch_id 去重(dispatch_progress_retention 窗口内)
  2. fencing_epoch 过滤:按 (tenant, conversation) 维护已见最大值的单调过滤器,
     小于该值整条丢弃并计入 stale_epoch_dispatch_dropped_total(§19.2.1 校验点 2)
  3. 加载 membership_version 的本分片 Bitmap(LRU)
  4. 成员边界过滤(§10.1 第 9 步)
  5. 按 lane_id 分组 -> 至多 lane_count(64) 个子任务,只为非空 lane 生成
  6. 每子任务按 mailbox_writebatch_max_entries(2000) 分块
  7. 逐块 append_batch,达到持久性契约后推进 W[lane]

步骤 2 的过滤器为什么在消费侧:broker producer 身份不能表达 (tenant, conversation, fencing_epoch) 的业务所有权,且 ADR-0012 已取消 Fanout Kafka 事务与 transactional_id 前提。§19.2.1 的两个强制校验点是 ConversationHead 条件更新与 GroupDispatch.fencing_epoch 消费侧过滤;Writer 自我 fencing 只缩小故障窗口,不能替代任一点。

实现状态(2026-09-01):目标 GroupDispatch 要求携带 fencing_epoch,但当前 Rust 结构及 MailboxNode 尚未实现该字段与过滤,属于发布阻断。不得以删除本步骤或恢复 ProducerFenced 文字来掩盖缺口。

3.2 lane 子任务调度

目标/现状边界:本节与 §3.3 描述目标 lane 隔离调度。当前实现可按 lane/chunk 拆分写入,但 finish_dispatch 等整条 record 的全部 chunk durable 后才把 64 lane 统一推进到同一 observed_log_offset;空 lane 不会提前越过。故 M-1 当前未通过, 不能把 packed 64-lane 存储格式等同于 lane 独立推进语义。

子任务之间**完全独立**:不共享状态、不互相等待、可并发执行
每个 lane 内部串行(同 lane 的 mailbox_seq 必须按序物化,否则 W[lane] 无法连续推进)

调度器:
  每分片一个 lane_count 长度的任务队列数组
  worker 池按 lane 取任务,同一 lane 同时只有一个 worker
  大群 dispatch 的 64 个子任务可被 64 个 worker 并行消费

超时:单子任务超过 mailbox_subtask_timeout(30s) 只告警并降级,
      **不得无限期拖住 W[j]** —— 转入慢速重试队列,W[j] 停在该 seq 之前
停滞:lane_stall_alert(30s) 告警;lane_stall_failover(120s) 触发租约漂移接管

3.3 水位推进

W[j] 可推进到 S,当且仅当对所有 mailbox_seq ∈ (W_old[j], S],
该 seq 在 lane j 上派生的子任务集合为空,或全部达到持久性契约(§9.2.2)

实现:每 lane 维护一个 (seq -> 未完成块数) 的有序小 map
      块完成时递减,归零则该 seq 可推进
      从 W_old+1 开始连续扫描已归零的 seq,推进到第一个未归零处停止

推进后必须同批写 W_floor[j](§10.4.2)——它是接管时的单调性下界。

"无子任务"与"子任务已完成"在 lane 内等价——这是 lane 能消除队头阻塞的关键: 一个只发单聊的用户所在 lane 若不含某大群成员,该 lane 的水位直接越过大群的 mailbox_seq 继续推进。


4. 接管与恢复

4.1 漂移接管(ADR-0007 的默认形态)

T0  ShardRegistry 检测租约超时
T1  等待租约自然过期(不是心跳超时就抢,防脑裂,§19.2)
T2  递增 shard_epoch,授予新节点
T3  新节点重算 W[j]:扫 DispatchProgress 得到每 lane 的连续完成上界
T4  校验 重算 W[j] >= W_floor[j],不满足 -> 拒绝服务该 lane + P1(不可自愈)
T5  起始位点 := min_j(W[j]) 对应 log_offset + 1     ← 不取消费者组已提交位点
T6  重放:dispatch_id 去重 -> DispatchProgress 跳过已完成块 -> 未完成块重做
T7  逐 lane 追平后开放服务;未追平返回 SHARD_MOVED
T8  旧游标按 EpochBoundary 换发 CURSOR_REBASED

T5 是正确性关键:用消费者组已提交位点会在"先提交位点后物化"的崩溃下 永久跳过某个 dispatch,该 lane 的 W[j] 就此停滞(由 lane_watermark_stall_ms 告警暴露,但需人工回退 offset)。按水位推导则自愈。

一期(Redis)与阶段一(ScyllaDB)下权威数据在节点外,T6 的重放窗口只覆盖 min_j(W[j]) 到当前,通常是秒级;阶段二才需从 S3 检查点重建。

4.2 REBUILD 执行器(§9.6)

CURSOR_EXPIRED 后的客户端侧流程,服务端只提供数据源:

0. 客户端重新 AUTH(携带空游标),服务端按新设备路径下发 AUTH_OK
1. PULL_SESSION_LIST 取会话集合(权威来源是 UserConversationState,
   不是 UserSessionProjection —— 后者可缺失、可重建,§9.6.2)
2. 每会话 PULL_HISTORY(anchor = 本地该会话最大连续 conversation_seq)
3. 按 message_id 去重合并到本地
4. 未读改由 latest_conversation_seq 与 read_conversation_seq 计算
5. 窗口外提及数不保证精确,unread_exact = false
6. 游标从 sync_to_seq 起算(**不是 trim_watermark**,否则会重复拉一遍窗口内数据)

撤回与编辑通过每会话最近一页 MessageRecord 的当前状态自然收敛, 不需要重放邮箱事件流。


5. 事件组与批次切分

事件组 = 同一 (user_id, mailbox_seq) 下的全部条目
硬上限:<= 8 条且 <= 8 KiB,且**每种 event_type 至多 1 条**(§6.7)
        当前 4 种类型下实际至多 4 条,8 为未来类型预留

event_ordinal = event_type 的固定优先级值(1:1 映射)
        MESSAGE=0, MENTION=1, MEMBERSHIP=2, CONTROL=3
event_id = blake3(tenant_id, dispatch_id, user_id, event_type)[0:8]

FanoutCoordinator 超出上限时必须拆成多条 dispatch(多个 mailbox_seq),禁止超发。
这条上限是"批次切分只能在 mailbox_seq 边界"(§9.3.3 规则 2)可行的前提。

切分实现:

// 从归并后的条目流中切出一个批次
// 不变量:返回的批次末尾必然是一个完整事件组
fn cut_batch(entries: &[Entry], max_items: u32, max_bytes: u32) -> usize {
    let mut n = 0; let mut bytes = 0;
    let mut last_complete = 0;          // 最后一个完整组的结束位置
    for (i, e) in entries.iter().enumerate() {
        if i > 0 && e.mailbox_seq != entries[i-1].mailbox_seq {
            last_complete = i;          // 组边界
            if n >= max_items || bytes >= max_bytes { return last_complete; }
        }
        n += 1; bytes += e.encoded_len();
    }
    entries.len()
}
// max_items / max_bytes 是软上限:单个事件组即使超限也必须整组返回,
// 只受 max_frame_bytes 硬约束(§9.3.3)

6. 大群成本控制(§10.5)

平台 fanout 预算
    platform_fanout_entries_per_sec = Σ(msg_rate_g × N_g) + 单聊事件/s
    mailbox_shard_count >= platform_fanout_entries_per_sec / per_node_entry_budget

租户配额(§8.2 L3)
    tenant_fanout_quota 令牌桶,按 GroupMembershipVersion.member_count 扣减
    超限返回 FANOUT_QUOTA_EXCEEDED —— **发送侧拒绝或排队,绝不静默丢邮箱引用**

降级档(默认关闭)
    mailbox_write_policy = mention_only,准入判据是 α 判据(§10.2.3),
    不是单一人数阈值;启用需 ADR-0003

一期(千人群、R_avg ≈ 22)下 mention_only 的 α 判据不成立,不实现。


7. 验收

M-1【发布阻断】lane 隔离
  同分片上构造:一个 10 万人群 dispatch 展开中,同时给纯单聊用户投递
  纯单聊用户的 W[lane(U)] 推进不等待大群 lane
  W[lane(U)] 的推进延迟 P99 <= lane_watermark_advance_p99(2s)
  当前统一推进 64 lane 的实现必须判失败,禁止仅检查水位数组长度

M-2【发布阻断】接管起始位点
  注入"先提交位点后物化"的崩溃,触发漂移接管:
  起始位点 == min_j(W[j]) 的 offset + 1,且 <= 崩溃点
  该 dispatch 被完整重放,收件人条目缺失数 == 0
  重放结果与崩溃前已写部分逐字节一致

M-3 W_floor 下界
  人为损坏 DispatchProgress 使重算 W[j] < W_floor[j]:
  该 lane 拒绝服务 + P1;向客户端下发低于 W_floor 的水位次数 == 0

M-4 事件组不可切分
  构造 4 条组,令 max_items=1:批次仍完整返回该组全部 4 条
  续拉后无重复无遗漏;主备/重放后组内次序不变

M-5 槽位不复用
  用户退群再加群:新 slot_id > 旧 slot_id;
  用旧 membership_version 展开不得命中新用户(跨用户错投次数 == 0)

M-6【发布阻断】Redis Cluster 启动拒绝
  对三节点 Redis Cluster 配置启动 writer / mailbox:
  进程必须非零退出、readiness 不得变绿,且不得写入任何提交、邮箱或进度状态。
  该负向门禁通过前,不能把 Redis Cluster 标为已验证或已支持。

M-7 REBUILD 游标起点
  CURSOR_EXPIRED 后完成 REBUILD:游标从 sync_to_seq 起算
  不出现"重复拉取 trim_watermark 到 sync_to_seq 区间"的行为

8. 待办

[x] Lua 脚本落地与 EVALSHA 缓存管理(`crates/qim-store/lua/` + `redis/scripts.rs`)
[ ] Redis Cluster 支持:当前明确不实现;必须启动拒绝。若未来重新立项,需独立完成
    状态机、迁移、Cluster 路由与端到端故障验证,局部 hash tag 不得作为完成标志。
[ ] per_node_entry_budget 实测回填(附录 B.7)
[ ] 槽位重编号的停写窗口操作手册(进 docs/06)