跳转至

Q-IM 消息持久化与完整客户端架构

日期:2026-08-23 状态:历史冻结稿;Fanout Kafka 事务部分已被 ADR-0012 取代 适用范围:main、store/redis、store/scylla,以及 Rust SDK、CLI、C FFI、UniFFI、Tauri GUI

当前契约以 docs/PLAN.md、ADR-0012 与 arch_20260901_current_log_driven_summary.md 为准。本文出现的 TransactionalFanout、transactional.id、send_offsets_to_transaction 等内容仅记录 2026-08-23 当时方案,不代表现行实现。

Design Intent

  • MessageCommitter:独占“幂等占位 → 固化坐标 → MessageRecord → Outbox → COMMITTED”的提交状态机,对调用方只暴露 commit 与稳定结果,隐藏 Redis/CQL/Redpanda 的恢复细节。
  • CommitLog + FanoutCoordinator:把可靠扩散从 writer 的进程内 TCP 队列移到 Redpanda;事务性地把一条 Outbox 转成确定性 MailboxShard dispatch,并提交输入位点。
  • MailboxMaterializer:直接以分发日志 partition/offset 构造稳定 mailbox_seq,重放只会同值覆盖;水位和 W_floor 决定接管起点与服务就绪状态。
  • MessageStore / MailboxStore:共享语义、后端各自实现;Redis 与 ScyllaDB 的 key、表、LWT、TTL、槽位和一致性级别不泄漏到 writer、fanout 或客户端。
  • ClientRuntime:连接、同步、游标、pending、历史与宿主磁盘路径形成一条恢复链;GUI、CLI、C FFI、UniFFI 仅适配宿主,不复制 SDK 状态机。

What Changed

  1. 原实现把 dedup 首次写成最终态,并在历史与可靠分发前返回 SEND_ACK;新设计只允许 CommittedMessage 生成 ACK,未完成状态可由请求或后台恢复器续做。
  2. 原实现把 writer 到 mailbox 的一次 TCP 写成功当成可靠分发;新设计以持久 Outbox 和事务分发日志作为唯一恢复源。
  3. 原实现每次 dispatch 重试重新分配邮箱序号;新设计令日志 offset 成为稳定邮箱 offset,消除重试重复和“只剩 pending 号、没有 payload”的永久空洞。
  4. 原 Redis 扫描按日期拼接、用完整 u64 作浮点 score;新设计把精确序号编码与跨桶有界归并封装在 Redis MailboxStore 内。
  5. 原客户端各宿主存在未 Connect、默认内存库、游标身份未持久化、REBUILD 未接线和批量事件丢失;新设计统一由 ClientRuntime 契约驱动,宿主只提供稳定路径和事件呈现。
  6. 原 CREATE_GROUP 可被任意认证用户用于修改已有群;新设计区分“客户端创建新群”和“受信管理面成员变更”,并在所有读写入口重新校验成员权限。

1. 问题陈述

当前主链路可以在客户端收到成功 ACK 后永久丢失历史或邮箱投递:writer 的 dedup、历史写入和 mailbox TCP dispatch 不是一个可恢复提交状态机,mailbox 也没有持久分发输入可供失败后重放。Redis 与 ScyllaDB 实现还分别存在序号精度、跨桶排序、LWT 初始化竞态、段租约空洞和生产一致性缺口。客户端虽已有 SDK 状态机与本地存储接口,但 GUI/UniFFI 未真正连接,默认路径使用内存库,游标、pending 类型和 REBUILD 链路不完整。目标是以明确的持久提交边界、事务日志、确定性重放和统一客户端恢复契约,做到“已确认不丢、重复不重、重启可恢复、分支语义一致”。

2. 组件拆分

2.1 qim-store::MessageStore

职责:保存消息提交意图、MessageRecord、幂等状态与历史索引,并提供恢复扫描。

最小接口语义:

  • reserve(request) -> Fresh(intent) | Existing(intent):原子绑定 (tenant, sender, client_message_id)、唯一消息坐标与完整 canonical MessageRecord;重复不消耗新序号。持久 intent 必须足以在客户端永不重试时独立恢复。
  • put_record(intent):从 intent 自带的 canonical record 按稳定主键幂等写入,禁止调用方再传一份可能漂移的 payload。
  • mark_outbox_published(intent, log_position):记录已发布 Outbox 的可恢复证据。
  • mark_committed(intent):只有 record 与 outbox 均已成立后才能进入最终态。
  • load_commit(key) / scan_recoverable_commits(...):请求重试与后台恢复器共用,不把恢复责任推给客户端。
  • history_page(query):保真返回类型、custom_type、方向、边界、max_items 与 max_bytes。

实现:

  • main / store/redis:仅 standalone Redis 的事务/Lua 实现(AOF + noeviction);一期保留现有可运行后端,但消息正文、提交意图和 Outbox 证据必须具备同一语义。客户端提交/去重窗口固定为 2 小时,历史则按 retention_class 保留。
  • store/scylla:ScyllaDB LWT + 普通幂等写实现;初始化一律 IF NOT EXISTS,RF=3,关键读写 LOCAL_QUORUM,逐行 TTL/合规保留。

隐藏内容:Redis key、Lua 状态值、CQL schema、LWT applied 行、序号租约策略、桶布局。

2.2 qim-commit-log::CommitLog

职责:封装 Redpanda/Kafka API,使业务模块只操作领域事件和稳定位置。

接口分为两层:

  • OutboxPublisher:以 message_id 作为稳定 Kafka key 至少一次发布 CommittedEnvelope,返回持久 LogPosition。librdkafka 的幂等 producer 只吸收同一 producer 的协议级重试;应用在“日志成功、状态标记失败”后可再次发布同 key,必须由确定性 dispatch 与 MailboxStore 持久去重收敛,禁止宣称跨故障 exactly-once。
  • TransactionalFanout:以 read_committed 读取 Outbox,在同一事务内发布零到多条 GroupDispatch 并提交输入位点。

实现固定使用 rust-rdkafka。transactional.id、幂等 producer、init_transactions fencing、send_offsets_to_transaction、commit_transaction 与 abort/retry 全部藏在该深模块内;调用方不得直接拼 topic、header 或 offset commit。

2.3 qim-writer::MessageCommitter

职责:成为唯一允许生成 SEND_ACK 的组件。

提交状态机:

ABSENT
  │ reserve(固化 cmid/message_id/conversation_seq/成员快照)
  ▼
RESERVED(已持久保存完整 canonical MessageRecord)
  │ put_record(intent,自包含且幂等)
  ▼
RECORDED
  │ publish_outbox(message_id key,幂等/可检测重复)
  ▼
OUTBOXED
  │ mark_committed
  ▼
COMMITTED ──▶ SEND_ACK

请求命中任一中间态时从该态续做;命中 COMMITTED 时只回放同一结果。后台 CommitRecovery 扫描超时中间态并调用同一套 resume,避免出现一套在线逻辑、一套恢复逻辑。

直接会话收件人,或群会话的不可变 membership_version、成员数与目标分片集合,以及消息类型、正文引用/内联决定、created_at 都在 reserve 时固化;恢复不得改用“当前群成员”。Fanout 只能按该精确版本读取快照。

HLC 使用持久高水位或启动时租用新的 writer 身份。固定 writer_id 只有在成功恢复其 HLC checkpoint 后才允许接流量;无法证明单调时 readiness 失败。

2.4 新服务 qim-fanout

职责:消费 Outbox,按稳定路由规则拆成 MailboxShard dispatch。

  • 输入:CommittedEnvelope。
  • 输出:每个目标 shard 恰好一条 GroupDispatch,包含 dispatch_id、完整确定性输入、该分片的全部收件人块、原消息坐标与创建时间;分块只是同一 dispatch 的物化进度,不产生额外日志记录。
  • dispatch_id = blake3("qim.group-dispatch.v1\0" || message_id_be16 || target_shard_be4)[0..16],重试不变;字段均为无符号大端固定宽度编码。
  • 发布 dispatch 与提交 Outbox 位点必须在一个 Redpanda 事务中完成。
  • 不读写用户邮箱、不发 Socket、不修改 MessageRecord。

一期 1000 人群可以把不可变版本实现为 Redis 中的排序成员快照;后续 Bitmap/MemberSlotMap 优化保持相同的“按精确版本读取”接口,不改变 writer、mailbox 或客户端接口。快照保留期必须覆盖 Outbox、分发日志与 DispatchProgress 的最长恢复窗口。

2.5 qim-mailbox::MailboxMaterializer

职责:静态拥有 MailboxShard 分区,消费分发日志并物化邮箱、会话投影与在线 PushBatch。

  • mailbox_seq = encode(shard_epoch, dispatch_log_offset);同一 GroupDispatch 日志记录重放永远复用同一值,其全部 lane/chunk 共用该序号。
  • 每个 dispatch 按稳定 chunk 写入;条目先成功,进度后成功。支持原子后端时可同批提交。
  • DispatchProgress 记录 chunk 完成位;所有块完成后推进连续 W[lane],同时持久 W_floor[lane]。
  • 启动/接管从 min(W[lane]) + 1 seek,而非信任消费者组位点;重算低于 W_floor 时该 lane readiness 失败。
  • 日志消费位点只作性能提示,不作正确性真相。
  • 在线 push 是易失优化;邮箱与投影是持久事实。

2.6 Redis MailboxStore

职责:实现精确、有界、可恢复的一期邮箱存储。

关键决策:

  • ZSET score 只保存 48 位 offset 或其他精确整数;epoch 单独进入 key/元数据,完整复合序号使用固定宽度大端编码作为 member/游标比较依据。
  • 邮箱按 Entry.created_at 进入 30 天保留的日桶;RangeScan 对最多 32 个日桶执行有界 LIMIT + pipeline,再用小顶堆按完整 mailbox_seq 全局归并;损坏条目返回错误,不推进 covered。
  • 分页只在同一 mailbox_seq 事件组边界切分。
  • TruncateBefore 的物理裁剪与 user_trim_seq 读路径闭环;越界返回 CURSOR_EXPIRED。
  • DispatchProgress 只在 7 天 replay-safe 窗口后 GC;清理前必须仍可由分发日志与水位证明不再需要重放。
  • 当前只支持 standalone Redis,启动必须拒绝 Redis Cluster。GroupMembership 会跨槽,MessageStore 的全局 {commit} 会形成热槽,且 ConnectionManager 没有 Cluster 路由、MOVED/ASK 重试或拓扑刷新;局部 hash tag 不能补齐这些缺口,也不能伪造跨槽原子。
  • Redis 集成门禁使用隔离实例;连接失败即测试失败。

2.7 ScyllaDB MessageStore / MailboxStore

职责:在 store/scylla 分支实现相同契约。

关键决策:

  • keyspace 为 NetworkTopologyStrategy、RF=3;MessageRecord、邮箱、进度和 LWT 关键路径显式 LOCAL_QUORUM / LOCAL_SERIAL。
  • 所有“首次初始化”使用 INSERT ... IF NOT EXISTS 并解析 applied;禁止普通 upsert 初始化 allocator/trim watermark。
  • 不再把租到但未发出的整个段登记成水位空洞。邮箱序号来自 dispatch offset;会话序号使用可接管的单号/LWT 或持久租约游标,重复请求在分配前先命中 dedup。
  • MessageRecord 保存 message_type、custom_type、retention、DEK 等完整字段;表级 TTL=0,逐行 TTL;MessageStore 表禁止显式 DELETE。
  • earliest_available_conversation_seq 是权威元数据,不通过固定扫描 8 个桶猜测。
  • 所有 u64 -> bigint 经统一编码器断言 bit63 为零;热路径 prepared + shard-aware,范围查询有页界。

2.8 群权限边界

将现有混用的 GroupAdd 拆成两个语义:

  • 客户端 CreateGroup(actor, group_id, invitees):wire group_id 先以 strip_kind 归一,再由 group_conversation_id 生成唯一 canonical ConversationId;唯一性边界为 (tenant_id, canonical_conversation_id)。只允许创建不存在的群;服务端强制加入 actor、记录 owner/creator,整批校验上限并原子创建。客户端不能借该接口修改已有群。
  • 受信管理面 AdminSet/AddGroupMembers:仅供 qimctl/loadgen 或未来管理员鉴权服务,必须带明确内部身份;不从公网帧直达。

PullMembers、PullHistory、群发消息都在 writer 侧依据经过认证的 actor 再校验成员关系。网关负责传递身份,不负责代替存储侧授权判断。

2.9 qim-sdk::ClientRuntime

职责:统一连接生命周期、本地事务、同步/REBUILD、pending 对账与事件批量语义。

  • Handle 创建后由宿主显式发送一次 Connect;GUI、CLI、C FFI、UniFFI 通过共享 helper 建立同一行为。
  • LocalStore 的游标对象同时持久化 lane、shard_epoch、last_applied_seq;SQLite 迁移保持 pending,重置时游标和消息同边界处理。
  • CURSOR_EXPIRED 生产链路:SessionList → 每会话历史补齐 → finish_rebuild → 新游标/对账;不允许测试直接调用而生产不接线。
  • resolve_pending 保留 message_type / custom_type;历史只向 UI 广播实际新插入集合。
  • ACK 与 PONG 有 deadline;半开连接进入断开/重连,而不是无限 Online。
  • Session、History、Members 分页元数据进入 SDK 领域事件,宿主可以继续拉取而不猜测。
  • C FFI 使用内部事件拆包队列,批量事件逐项无丢失导出;UniFFI/Tauri 不复制状态机。

2.10 GUI / CLI 宿主

  • GUI 默认数据库位于 Tauri app_data_dir 下,按 endpoint/tenant/user 隔离;CLI 未提供 --db 时使用平台数据目录,而非内存库。测试可显式选择内存。
  • 登录只在用户操作后触发,不以默认 user=1 自动连接;事件 handler 通过 ref/领域事件读取当前会话,避免闭包捕获旧表单值。
  • 建群请求显示 pending,收到 MembersUpdated 或明确错误后才创建/打开会话。
  • 会话、房间入口使用语义按钮或具备键盘/ARIA 等价行为;连接状态使用可播报区域。
  • 服务端尚未支持的能力明确禁用或标记,不生成伪成功本地状态。

2.11 分支同步与发布

  • 所有共享协议、qim-commit-log、qim-fanout、SDK、测试向量先落 main。
  • store/redis 当前是 main 祖先且无独有差异,完成后按单向规则更新到 main 对应提交。
  • store/scylla 先吸收 main,再只改 qim-store/scylla、后端装配和迁移文档;不得反向覆盖 main。
  • 当前工作树含用户未提交改动,实施期间按文件所有权拆分,禁止 checkout/reset;生成产物单独说明。

3. 数据流

3.1 发送与可靠扩散

Client
  │ SEND_MESSAGE(cmid)
  ▼
Gateway ──authenticated actor──▶ ConversationWriter
                                  │
                                  ▼
                         MessageCommitter.commit
                                  │
                 ┌────────────────┼────────────────┐
                 ▼                ▼                ▼
            CommitIntent     MessageRecord    CommitLog.Outbox
                 │                                 │
                 └──────────── COMMITTED ◀─────────┘
                                  │
                                  └── SEND_ACK

CommitLog.Outbox --read_committed--> FanoutCoordinator
                                         │ one transaction
                                         ├── dispatch[MailboxShard A]
                                         ├── dispatch[MailboxShard B]
                                         └── commit outbox offset

dispatch partition/offset --> MailboxMaterializer
                                  │ stable mailbox_seq
                                  ├── UserMailboxEntry
                                  ├── SessionProjection
                                  ├── DispatchProgress / W / W_floor
                                  └── best-effort PushBatch --> Gateway --> Client

3.2 客户端重启恢复

stable app data path
        │ open/reopen
        ▼
SQLite LocalStore ──cursor(lane, epoch, seq), pending, messages, sessions──▶ ClientRuntime
        │                                                                  │
        │ local-first UI                                                   ├── Connect/Auth
        ▼                                                                  ├── PullMailbox to sync_to
GUI / CLI                                                                  ├── reconcile pending
                                                                           └── Online

CURSOR_EXPIRED --> authoritative SessionList --> per-conversation History
               --> persist rebuilt state/cursor --> reconcile pending --> Online

4. 技术选择与理由

选择 理由 优点 代价/约束
Redpanda + Kafka API 项目 PLAN 已锁定 CommitLog;需要事务、幂等 producer 与可重放分区日志 事务 consume-transform-produce;成熟运维/保留策略;Redis AOF 缺口可重放 新增基础设施与运维;本地/CI 必须提供真实 broker
rust-rdkafka Rust 生态中成熟的 librdkafka 绑定,支持事务 producer、read_committed 与 fencing 与 Redpanda 兼容;事务 API 完整;性能成熟 引入 C/librdkafka 构建链;需选择 vendored cmake-build 或系统库。项目优先可复现构建,采用 cmake-build,镜像体积略增
standalone Redis Lua + 精确编码 main/一期现有后端,需最小化迁移风险 延续现有部署;单实例原子与高吞吐 必须 AOF + noeviction 并绑定分发日志;Redis Cluster 未实现且启动拒绝,未来需独立状态机/迁移工作后再评估
ScyllaDB LWT + LOCAL_QUORUM 分支既定规模化后端 高吞吐、TTL/TWCS、按用户分区 LWT/一致性成本;需要 3 节点真实故障验证
SQLite SDK 已有实现,适合桌面/移动离线恢复 单文件、事务、跨平台、无需额外服务 宿主必须选稳定路径;迁移与并发访问需严格测试
Tauri app_data_dir 平台标准用户数据目录 无需用户手填;跨平台权限/路径正确 测试需注入临时目录,账号切换必须隔离文件

未选择:

  • “ACK 等到所有邮箱写完”:把大群扩散延迟和可用性耦合到发送成功,不符合已定义的提交语义。
  • 继续使用进程内 mpsc/TCP:没有崩溃恢复源,也无法可靠提交消费进度。
  • 仅用 Redis Stream 替代 CommitLog:偏离既定跨后端架构,且不能为 Redis 自身 AOF 丢失提供独立恢复域。
  • writer 本地 WAL:引入节点亲和、复制和接管问题,等同自研日志系统。
  • GUI 让用户填写数据库路径:把平台与隔离复杂度推给用户,默认路径仍会假持久化。

5. 复杂度估算

XL(至少一周级,按可独立验收的多个增量交付)。

原因:涉及协议与共享类型、两个服务端持久后端、一个新服务、消息提交时序、故障恢复、客户端状态机、四类宿主绑定、真实基础设施测试和三分支同步;任何单点绿色都不能替代端到端故障证明。

建议实施顺序:

  1. 契约、编码向量、群权限与客户端显性 P0。
  2. MessageStore 提交状态机与 Redis 实现。
  3. CommitLog/Fanout/Mailbox 日志消费与 Redis 故障恢复。
  4. SDK 磁盘重启、REBUILD、deadline、FFI 批量语义和 GUI 产品路径。
  5. 同步 Scylla 分支并修复 LWT/一致性/序号模型。
  6. 隔离集成、混沌、性能和分支门禁。

6. TestPlan 映射

子模块 Must Have
MessageStore + MessageCommitter M-COMMIT-01..05、M-HIST-01..02
CommitLog + FanoutCoordinator M-LOG-01..02
MailboxMaterializer M-MB-01..03
Redis MailboxStore M-REDIS-01..04
ScyllaDB 后端 M-SCYLLA-01..03
群权限 M-AUTH-01..02
SDK / SQLite M-CLIENT-01..02
GUI / CLI / 绑定 M-CLIENT-03..04
分支与发布 M-BRANCH-01、M-PERF-01

详细输入、故障点和断言见 testplan_20260823_message_persistence_client.md。

7. 复杂度分析

7.1 主要复杂度来源

  • 提交跨 MessageStore 与 CommitLog,底层不存在通用分布式事务。
  • fanout 与 mailbox 采用至少一次处理,正确性依赖稳定身份和同值覆盖。
  • Redis 与 ScyllaDB 的原子性、精度、TTL 和一致性模型不同。
  • 客户端同步包含实时推送、持久游标、历史补洞和 pending 对账四类输入。
  • 当前代码与文档存在“已经声明但未实现”的能力,常规单元测试会假绿。

7.2 降低变更放大

  • writer 只依赖 MessageCommitter,ACK 规则只有一处。
  • fanout 只依赖领域事件与 TransactionalFanout,topic/事务配置只有一处。
  • mailbox 只依赖 GroupDispatch 与 MailboxStore,不关心 Outbox 或 writer 重试。
  • Redis/Scylla 通过共享契约套件验证,不在业务服务中分散 cfg(feature) 分支。
  • GUI/FFI/UniFFI 复用 SDK runtime 和连接 helper,不各写一套同步逻辑。

7.3 降低认知负担

  • 用显式提交状态代替靠代码行顺序猜测崩溃语义。
  • 把“日志位置 ↔ mailbox_seq”固化为共享值对象和向量测试。
  • 将可恢复错误归一为领域错误,业务层无需理解 Redis MOVED、CQL applied 或 Kafka rebalance。
  • 每个模块只拥有一个事实源:MessageStore 管提交、CommitLog 管可重放输入、MailboxStore 管物化事实、LocalStore 管设备状态。

7.4 避免未知未知

  • 每个跨进程边界都有稳定 ID、持久进度、恢复起点和拒绝服务条件。
  • TestPlan 明确每个 kill 点与“不得推进”的坐标,测试不是仅断言返回成功。
  • 基础设施不可用时门禁失败而非跳过。
  • 分支矩阵显式记录共享与专属所有权,避免 Scylla 合并覆盖 Redis/客户端修复。

7.5 原则校验

原则 校验结果
深模块 MessageCommitter、CommitLog、ClientRuntime 的接口小,恢复/事务复杂度下沉
信息隐藏 key、CQL、topic、offset commit、平台路径不穿透到上层
抽象层次 提交、扩散、物化、同步各自表达不同领域能力,无透传式服务
内聚与分离 record 与 outbox 提交状态同属 committer;在线 push 与 durable mailbox 保持分离
错误处理 通过幂等与稳定坐标消除大量“是否已经执行”的分支;不可证明安全时 readiness 失败
命名 统一使用 CommitIntent、CommittedEnvelope、GroupDispatch、DispatchProgress、CursorIdentity

8. 测试策略

  1. 纯契约测试:序号编码、dispatch_id、事件边界切分、提交状态迁移、错误映射;Redis/Scylla/Memory 测同一向量。
  2. SQLite 重开测试:临时文件写入后 drop/reopen,验证会话、预览、消息类型、游标身份、已读和 pending。
  3. 隔离基础设施集成:Docker 启动独立 Redis、Redpanda、3 节点 Scylla;禁止连接或清空共享 127.0.0.1:6379。
  4. 故障注入:在 reserve、record、outbox、commit、dispatch chunk、watermark、client ACK/PONG 等边界 kill/断连,验证最终集合和坐标。
  5. SDK/绑定测试:真实 socket 建连;C FFI 批量事件全量导出;UniFFI/Tauri 构造后进入 Online;REBUILD 走生产事件循环。
  6. GUI 自动化:登录、发送、建群、历史分页、键盘操作、关闭重开与本地恢复;命令入队不作为成功。
  7. 分支矩阵:main/Redis 与 Scylla 各跑真实后端契约和同一 E2E;结果单独记录。
  8. 性能门禁:先断言实际到达率,再评估 P99;30 秒只算 smoke,180 秒/30 分钟/4 小时按文档独立标记,未执行不得声称通过。

9. 实施约束

  • 当前架构冻结后,跨模块共享类型只能经 contracts 目录和共享 crate 修改;不得由某个实现私自扩展语义。
  • 所有新增 Rust 依赖记录在 CLAUDE.md,并说明选择理由、替代方案和构建代价。
  • 修改后运行 cargo fmt --all、cargo clippy --workspace --all-targets --all-features -- -D warnings;GUI 运行 TypeScript 构建及项目可用的 lint/formatter。
  • 任何会 FLUSHALL、杀服务、修改 Redis 配置或覆盖分支的脚本只可在明确隔离环境执行。
  • 生成物与用户原有脏工作树分别列出,禁止用 reset/checkout 清理。