跳转至

数据流图

本页描述数据在进程、日志与存储之间如何流动、落在哪里、谁是权威。以现行一期实现为准,目标形态单独标注。

1. 系统总体数据流

flowchart LR
    subgraph 客户端
        SDK[qim-sdk<br/>SQLite 本地库]
    end
    subgraph 接入
        GW[qim-gateway<br/>ConnectionNode]
    end
    subgraph 提交
        WR[qim-writer<br/>ConversationWriter]
    end
    subgraph 日志["Redpanda"]
        OB[(Outbox topic<br/>canonical MessageRecord)]
        DL[(dispatch topic<br/>64 分区 = MailboxShard)]
    end
    subgraph 异步
        FO[qim-fanout<br/>FanoutCoordinator]
        AR[qim-archiver<br/>HistoryArchiver]
        MB[qim-mailbox<br/>MailboxNode]
        SS[qim-session<br/>SessionProjection]
    end
    subgraph 存储
        RC[(core Redis 分片<br/>提交事实 · 幂等 · 序号<br/>7 天热历史 · 群成员 · 会话集合 · 已读)]
        RM[(邮箱 Redis 分片<br/>个人邮箱 · DP · packed W)]
        SC[(ScyllaDB<br/>永久历史 · user_conversation_state<br/>可选邮箱后端)]
    end

    SDK <-->|TLS 二进制帧| GW
    GW -->|SendMessage / PullHistory<br/>PullMembers / CreateGroup| WR
    GW -->|PullMailbox| MB
    GW -->|SessionListReq / MarkRead| SS
    WR <--> RC
    WR -->|ACK 前 durable| OB
    OB --> FO
    OB --> AR
    FO -->|读冻结成员快照| RC
    FO --> DL
    DL --> MB
    DL --> SS
    MB <--> RM
    MB -->|PushBatch 尽力| GW
    MB -->|Projection 尽力| SS
    AR -->|批量 INSERT| SC
    AR -->|归档水位| RC
    SS <--> RC
    SS -->|会话首次登记| SC
    WR -->|超出热层下界的历史| SC

读法:

  • 唯一同步写路径是 GW → WR → core Redis + Outbox;其余全部由日志驱动、可重放。
  • Outbox 有两个独立消费者:FanoutCoordinator(投递)与 HistoryArchiver(归档),互不阻塞。
  • dispatch 日志也有两个独立消费者:MailboxNode(邮箱)与 SessionProjection(会话集合),各自持有位点。

2. 提交与幂等的数据落点

flowchart TD
    REQ[SEND_MESSAGE<br/>tenant, sender, cmid, body] --> ID[AuthenticatedSend<br/>reserve_identity_digest 36B]
    ID --> LUA{{reserve_and_record.lua<br/>同一 Lua 原子}}
    LUA --> K1["cmidown:{commit}…<br/>幂等归属,2 h"]
    LUA --> K2["msgcseq<br/>INCR 分配 conversation_seq"]
    LUA --> K3["msgcommit:{commit}…<br/>提交 hash:intent / 状态"]
    LUA --> K4["msgrecord / msghistory<br/>canonical 记录与历史索引"]
    K3 -->|COMMITTED| K3B["压缩为 ACK 坐标 listpack<br/>PEXPIREAT committed_at + 2h"]
    K4 --> OB[(Outbox)]
    K4 -->|保留期到期登记| GC["msghistorygc<br/>到期清理索引"]

{commit} 为提交桶(256 个)哈希标签,群键 grp/grpmeta/grpsnap/grpmember 与会话提交事实同桶,所以群发送的版本校验能与序号分配在同一 Lua 内完成。

3. 历史分层:热层与归档

flowchart LR
    WR[writer 提交] --> HOT[(Redis 热层<br/>msgrecord / msghistory)]
    WR --> OB[(Outbox)]
    OB --> AR[qim-archiver]
    AR -->|先记录后索引<br/>普通 INSERT,同值幂等<br/>按剩余保留期 USING TTL| SC[(ScyllaDB)]
    AR -->|全批确认 → 提交位点<br/>位点 ≥ H0 后推进到 T0| WM[msgarchive:watermark]
    WM --> TR[writer 每 5 s 裁剪]
    TR -->|只删 committed_at ≤ min 热窗口, 水位 − 恢复窗口 − 60 s| HOT
    TR -->|抬高| FL[msghistfloor 每会话热层下界]
    RD[TieredMessageStore 读] -->|seq ≤ 下界| SC
    RD -->|seq > 下界| HOT

归档器停摆时热层只涨不删;下界为 0 时读路径与纯 Redis 完全相同。

4. 分发:Outbox → dispatch

flowchart TD
    OB[(Outbox 单分区)] -->|批 ≤ 512| PL[逐条确定性规划]
    PL --> D1{会话类型}
    D1 -- 单聊 --> T1[收件人与发送者所在 shard]
    D1 -- 群 --> T2[冻结 membership_version<br/>的 target_shards]
    T1 --> EN[全部 dispatch enqueue]
    T2 --> EN
    EN --> WAIT[等待全部 durable delivery]
    WAIT --> CM[同步提交批末 source offset]
    WAIT -. 任一失败 .-> RP[不推进位点,整批重放<br/>dispatch_id 确定性收敛]
    EN --> DL[(dispatch 分区 = shard<br/>offset 低 48 位 = mailbox_seq)]

5. 邮箱物化与连续水位

flowchart TD
    DL[(dispatch 分区 s)] -->|assign 静态指派<br/>起点 min W + 1| CON[中心消费者]
    CON -->|队满只暂停该分区| Q[shard s FIFO worker<br/>容量 128,同分区严格串行]
    Q --> P[全局 permit 64]
    P --> BEGIN[begin_dispatch<br/>DP + 稀疏 manifest]
    BEGIN --> APP[全部 lane/chunk/user 并发 append<br/>UserMailboxEntry 按天分桶 zset]
    APP --> FIN[finish_dispatch<br/>一条 Lua:final completion + 64 lane packed W]
    FIN --> OK[W 推进到该 record 的 observed offset]
    OK --> PUSH[PushBatch → gateway(易失)]
    OK --> PRJ[Projection → session(try_send,满则丢)]
    APP -. 任一失败 .-> DIRTY[零 finish,不推 W<br/>MailboxDirty 停止消费,重启重放]
    BEGIN -. 同 dispatch_id 异 digest .-> STOP[停止服务]

Fresh 单收件人走融合 Lua:DP、事件组、V4 proof、完成索引与 packed W 一次原子写入。Scylla 邮箱后端用分片前沿(mailbox_shard_frontier)与写时间戳单调写代替 LWT,并按页推进(ADR-0022)。

6. 会话投影数据流

flowchart LR
    DL[(dispatch 日志)] --> SC[qim-session durable 消费者]
    MB[MailboxNode Projection<br/>低延迟、可丢] --> LUA
    SC --> LUA{{共用登记 Lua}}
    LUA --> Z["convs:{user} zset<br/>ZADD GT 分值 = created_at"]
    LUA -->|首次进入集合| PQ[session_projection:ucs_pending]
    PQ --> BG[后台批量] --> UCS[(Scylla user_conversation_state)]
    SC -->|record 全部收件人写完后 CAS| CK[checkpoint<br/>绑定 cluster_id + topic_id + 分区数]
    CK -. 越窗 .-> EP[INCR 投影纪元 → 按用户惰性重建]
    Z --> SNAP[内存快照 5 min<br/>脏集合增量刷新]
    SNAP --> LIST[会话列表 + 定义式未读]

7. 客户端本地数据流

flowchart TD
    UI[宿主:GUI / desktop / C FFI / UniFFI] -->|Command| DRV[SDK driver 状态机]
    DRV -->|send_message| PEND[(local_pending<br/>立即落库)]
    DRV -->|帧| NET[TLS 连接]
    NET -->|MAILBOX_BATCH| TX[SQLite 事务<br/>消息 + cursor 原子提交]
    NET -->|PUSH_EVENTS| TXP[SQLite 落库<br/>不推进 cursor]
    TX --> DEDUP[提交成功后才更新内存去重<br/>与解除 pending]
    TXP --> DEDUP
    DEDUP -->|ConversationsUpdated 先于 MessagesAdded| UI
    CMID[BEGIN IMMEDIATE 预留 ordinal 号段] --> H[BLAKE3 device_id + ordinal → cmid]
    H --> PEND

本地库路径:规范化 endpoint 的 BLAKE3 scope + 显式 tenant + user,禁止 token 入路径。

8. 数据保留与回收

gantt
    title 各类数据在 Redis / Scylla 中的存活期
    dateFormat X
    axisFormat %s 天
    section Redis
    提交辅助状态(COMMITTED 后) :0, 1
    幂等归属 cmidown(2 h)       :0, 1
    个人邮箱条目                  :0, 7
    DispatchProgress              :0, 7
    历史热窗口                    :0, 7
    section ScyllaDB
    默认历史(可配 1..3650 天)   :0, 30
    ephemeral 历史                :0, 1
    邮箱 dispatch 身份(Scylla 后端) :0, 31

compliance_hold 永不到期;tenant_custom 当前失败闭合为不到期;时间轴的 1 天格代表 ≤ 1 天。