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