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)