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¶
- 原实现把 dedup 首次写成最终态,并在历史与可靠分发前返回
SEND_ACK;新设计只允许CommittedMessage生成 ACK,未完成状态可由请求或后台恢复器续做。 - 原实现把 writer 到 mailbox 的一次 TCP 写成功当成可靠分发;新设计以持久 Outbox 和事务分发日志作为唯一恢复源。
- 原实现每次 dispatch 重试重新分配邮箱序号;新设计令日志 offset 成为稳定邮箱 offset,消除重试重复和“只剩 pending 号、没有 payload”的永久空洞。
- 原 Redis 扫描按日期拼接、用完整
u64作浮点 score;新设计把精确序号编码与跨桶有界归并封装在 RedisMailboxStore内。 - 原客户端各宿主存在未 Connect、默认内存库、游标身份未持久化、REBUILD 未接线和批量事件丢失;新设计统一由
ClientRuntime契约驱动,宿主只提供稳定路径和事件呈现。 - 原
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]) + 1seek,而非信任消费者组位点;重算低于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):wiregroup_id先以strip_kind归一,再由group_conversation_id生成唯一 canonicalConversationId;唯一性边界为(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(至少一周级,按可独立验收的多个增量交付)。
原因:涉及协议与共享类型、两个服务端持久后端、一个新服务、消息提交时序、故障恢复、客户端状态机、四类宿主绑定、真实基础设施测试和三分支同步;任何单点绿色都不能替代端到端故障证明。
建议实施顺序:
- 契约、编码向量、群权限与客户端显性 P0。
- MessageStore 提交状态机与 Redis 实现。
- CommitLog/Fanout/Mailbox 日志消费与 Redis 故障恢复。
- SDK 磁盘重启、REBUILD、deadline、FFI 批量语义和 GUI 产品路径。
- 同步 Scylla 分支并修复 LWT/一致性/序号模型。
- 隔离集成、混沌、性能和分支门禁。
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. 测试策略¶
- 纯契约测试:序号编码、dispatch_id、事件边界切分、提交状态迁移、错误映射;Redis/Scylla/Memory 测同一向量。
- SQLite 重开测试:临时文件写入后 drop/reopen,验证会话、预览、消息类型、游标身份、已读和 pending。
- 隔离基础设施集成:Docker 启动独立 Redis、Redpanda、3 节点 Scylla;禁止连接或清空共享
127.0.0.1:6379。 - 故障注入:在 reserve、record、outbox、commit、dispatch chunk、watermark、client ACK/PONG 等边界 kill/断连,验证最终集合和坐标。
- SDK/绑定测试:真实 socket 建连;C FFI 批量事件全量导出;UniFFI/Tauri 构造后进入 Online;REBUILD 走生产事件循环。
- GUI 自动化:登录、发送、建群、历史分页、键盘操作、关闭重开与本地恢复;命令入队不作为成功。
- 分支矩阵:main/Redis 与 Scylla 各跑真实后端契约和同一 E2E;结果单独记录。
- 性能门禁:先断言实际到达率,再评估 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 清理。