跳转至

业务流程图

本页用时序图描述现行(一期日志驱动)实现的主要业务流程。机制细节与编号以 消息处理流程 和 PLAN 为准;标注"目标形态"的部分尚未实现。

参与者简称:C 客户端 SDK · GW ConnectionNode(qim-gateway)· WR ConversationWriter(qim-writer)· R core Redis(MessageStore 热层)· OB Redpanda Outbox · FO FanoutCoordinator(qim-fanout)· DL Redpanda dispatch 日志 · MB MailboxNode(qim-mailbox)· MR 邮箱 Redis · SS SessionProjection(qim-session)

1. 单聊发送:从 SEND_MESSAGE 到对端收到

sequenceDiagram
    autonumber
    participant A as C(发送者)
    participant GW as GW
    participant WR as WR
    participant R as core Redis
    participant OB as Outbox
    participant FO as FO
    participant DL as dispatch 日志
    participant MB as MB
    participant MR as 邮箱 Redis
    participant B as C(接收者)
    A->>A: 落库 local_pending(先落库再发)
    A->>GW: SEND_MESSAGE{cmid, conversation_id, body}
    GW->>GW: 校验 cmid/cid 非零 16 字节、custom_type ≤ 64B
    GW->>WR: SendMessage RPC(按 conversation_id 会合哈希亲和)
    WR->>R: reserve_and_record.lua<br/>幂等检查 + INCR cseq + 写 intent 与历史
    R-->>WR: Fresh(RECORDED)或 Existing(原提交坐标)
    WR->>OB: 发布 canonical MessageRecord(acks=all, key=message_id)
    OB-->>WR: durable
    WR->>R: 标记 OUTBOXED → COMMITTED<br/>(提交 hash 压缩为 ACK 坐标、2h 过期)
    WR-->>GW: CommittedMessage
    GW-->>A: SEND_ACK{message_id, conversation_seq}
    A->>A: local_pending → 已发送
    Note over FO,DL: 以下异步,不在 ACK 路径
    FO->>OB: 批量消费(默认 512 条)
    FO->>DL: 每目标 shard 一条 dispatch<br/>(单聊:收件人 + 发送者)
    DL-->>FO: 整批 durable 后同步提交 source offset
    MB->>DL: assign() 静态消费 partition = shard
    MB->>MR: begin → append 邮箱条目 → finish(推进 64 lane W)
    MB-->>GW: PushBatch(durable 之后,尽力而为)
    GW-->>B: PUSH_EVENTS(不推进游标)
    B->>GW: PULL_MAILBOX{after_seq = cursor}
    GW->>MB: PullMailbox RPC
    MB-->>B: MAILBOX_BATCH(≤ W 的连续前缀)
    B->>B: SQLite 原子落库后推进 cursor、去重

要点:

  • SEND_ACK 只由 COMMITTED 生成;ACK 健康不代表投递健康,Outbox → dispatch → 物化 → 拉取必须独立观测。
  • 发送者自己也写邮箱条目,且只有发送者的条目带 client_message_id,用作重连对账锚点(ADR-0008)。
  • PUSH 只是"有新数据"的信号,游标只能由 MAILBOX_BATCH 连续推进。

2. writer 提交状态机与重试

stateDiagram-v2
    [*] --> RESERVED: 预留序号 + 持久 CommitIntent<br/>(含完整 MessageRecord)
    RESERVED --> RECORDED: 写 MessageRecord 与历史索引<br/>(一期与预留合并为一个 Lua)
    RECORDED --> OUTBOXED: Outbox 发布成功
    OUTBOXED --> COMMITTED: 标记完成
    COMMITTED --> [*]: 回放 SEND_ACK
    RESERVED --> RESERVED: 崩溃 → 恢复器续做
    RECORDED --> RECORDED: 崩溃 → 恢复器续做
    OUTBOXED --> OUTBOXED: 标记失败 → 允许重复发布<br/>由 dispatch_id + digest 收敛
flowchart TD
    S[收到 SEND_MESSAGE] --> L{load_commit_for_request<br/>tenant, sender, cmid}
    L -- 不存在 --> N[新提交:走完整状态机]
    L -- 存在且身份摘要一致 --> E{状态}
    L -- 同 cmid 不同内容 --> X[拒绝]
    E -- COMMITTED --> ACK[用当前 request_id + 原提交坐标回 ACK]
    E -- R/D/O 中间态 --> RES[在新成员关系/审核之前续做首次 intent]
    RES --> ACK
    N --> ACK

3. 群消息发送(冻结成员版本)

sequenceDiagram
    autonumber
    participant C as C
    participant GW as GW
    participant WR as WR
    participant R as core Redis(群键与提交同桶)
    participant OB as Outbox
    participant FO as FO
    participant DL as dispatch 日志
    participant MB as MB(各 shard)
    C->>GW: SEND_MESSAGE(群会话)
    GW->>WR: SendMessage(注入 actor_user_id)
    WR->>R: 读当前成员版本 V 并冻结受众
    WR->>R: reserve_and_record.lua<br/>INCR 前比较群 meta 版本 == V
    alt 版本已变(有人进出群)
        R-->>WR: V → MembershipChanged(不消耗序号)
        WR->>R: 重新冻结,有界重试 3 次
    else 版本一致
        R-->>WR: RECORDED(固化 membership_version)
    end
    WR->>OB: 发布 MessageRecord(含冻结版本与 target shards)
    WR-->>C: SEND_ACK
    FO->>OB: 消费
    FO->>FO: 按精确版本读成员快照<br/>(进程内 FIFO 缓存,版本不可变故恒正确)
    FO->>DL: 每个目标 MailboxShard 恰好一条 GroupDispatch
    MB->>DL: 各 shard 消费自己的分区
    MB->>MB: 本 shard 成员按 lane/chunk 并发 append,完成后统一推进 W

重放时禁止改用当前成员;dispatch_id 由 message_id 与 shard 确定性派生。

4. 登录与离线同步

sequenceDiagram
    autonumber
    participant C as C
    participant GW as GW
    participant MB as MB
    participant SS as SS
    C->>GW: AUTH{access_token, device_id, mailbox_cursor}
    GW->>GW: 同 uid 旧连接 → KICKED{REPLACED}(顶号)
    alt 游标签名无效或越界
        GW-->>C: ERROR{CURSOR_INVALID}
    else 游标早于有效裁剪点
        GW-->>C: ERROR{CURSOR_EXPIRED} → 进入 REBUILD
    else 正常
        GW-->>C: AUTH_OK{sync_to_seq = 当前 lane 水位快照, has_offline, ...}
    end
    loop 直到覆盖 sync_to_seq(在途 ≤ 4)
        C->>GW: PULL_MAILBOX{after_seq, up_to_seq = sync_to_seq}
        GW->>MB: PullMailbox
        MB-->>C: MAILBOX_BATCH{entries, covered_through_seq, has_more}
        C->>C: SQLite 原子落库 → 推进 cursor
    end
    C->>GW: SYNC_COMPLETE{sync_to_seq}
    GW-->>C: ONLINE_READY(解除登录屏障)
    C->>GW: PULL_SESSION_LIST(登录序列自动拉一次)
    GW->>SS: SessionListReq
    SS-->>C: SESSION_LIST_BATCH(含 read_conversation_seq)
    C->>C: 已读水位只进不退,本地更高时回补 MARK_READ
    C->>C: 对 local_pending 先查邮箱对账锚点,未命中才按原 cmid 重发

新设备(零游标)读回保留窗口内尚存的全部邮箱条目(7 天),更早的走 PULL_HISTORY(ADR-0023)。

5. 在线实时推送与合并拉取

sequenceDiagram
    participant MB as MB
    participant GW as GW
    participant SDK as SDK 状态机
    participant DB as 本地 SQLite
    MB->>GW: PushBatch(W 推进之后)
    GW->>SDK: PUSH_EVENTS
    SDK->>DB: 低延迟落库(按 message_id 去重)
    Note over SDK: 不推进 cursor
    alt 没有在飞的拉取
        SDK->>GW: PULL_MAILBOX(单飞)
    else 已有在飞拉取
        SDK->>SDK: 只置 dirty
    end
    GW-->>SDK: MAILBOX_BATCH 末页
    SDK->>DB: 提交并推进 cursor
    SDK->>SDK: dirty 则续拉,PONG 作兜底

6. 游标过期与 REBUILD

flowchart TD
    A[收到 CURSOR_EXPIRED] --> B[保留本地消息,不清库]
    B --> C[以空游标重新 AUTH]
    C --> D[AUTH_OK → SYNC_COMPLETE → ONLINE_READY]
    D --> E[PULL_SESSION_LIST 取权威会话集合]
    E --> F[逐会话 anchor = 本地最大连续 conversation_seq]
    F --> G[PULL_HISTORY newer 直到 has_more=false]
    G --> H[按 message_id 去重合并]
    H --> I{latest - base ≤ 200 ?}
    I -- 是 --> J[按定义式精确重算未读]
    I -- 否 --> K[未读饱和 200,unread_exact=false]
    J --> L[游标换发为本次 sync_to_seq]
    K --> L

7. 会话列表、未读与已读

sequenceDiagram
    autonumber
    participant C as C
    participant GW as GW
    participant SS as SS
    participant R as core Redis
    C->>GW: PULL_SESSION_LIST
    GW->>SS: SessionListReq
    alt 有冻结快照且无脏会话
        SS-->>C: 快照分页(keyset)
    else 首次或快照失效
        SS->>R: ZREVRANGE convs:{user} 0 -1(全量会话集合)
        SS->>R: 两轮 pipeline 读 msghistory/msgrecord → 会话头
        SS->>R: HMGET read:{user}
        SS->>SS: 定义式未读:seq > read ∧ sender ≠ 自己<br/>排除 Control/Recalled/Deleted,窗口 200
        SS-->>C: SESSION_LIST_BATCH
    else 仅部分会话脏
        SS->>R: 只对脏集合重算并换 revision
        SS-->>C: SESSION_LIST_BATCH
    end
    C->>GW: MARK_READ{conversation_id, read_seq}
    GW->>SS: MarkRead
    SS->>R: read:{user} 只进不退
    SS->>SS: 把该会话记脏(下次首页增量刷新)

未读只在读路径算、不在写路径累加;客户端用同一定义式本地算未读,同步完成前是下界。

8. 建群、入群与退群(ADR-0021)

sequenceDiagram
    autonumber
    participant C as C / qimctl
    participant WR as WR
    participant R as core Redis
    participant OB as Outbox
    C->>WR: CREATE_GROUP / GroupAdd / LEAVE_GROUP / GroupRemove
    WR->>WR: 按当前版本过滤真实变更对象(空则 NoOp)
    WR->>WR: cmid = f(tenant, group, op, actor, V, users)
    WR->>R: membership_change.lua<br/>校验版本 → INCR 得 J → 改成员/区间 → 版本 V+1 → 写记录
    Note right of R: 入群 joined_at = J−1<br/>退群 left_at = J<br/>建群初始成员 joined_at = 0
    WR->>OB: 与普通消息同一状态机发布 CONTROL 消息<br/>custom_type = qim.membership.v1
    Note over OB: 变更后成员 ≤ 500 广播,否则只发给变更对象

区间外的历史、未读与会话头一律不可见;退群者仍可读退群点之前的历史。

9. 聊天室

flowchart LR
    C1[客户端 ROOM_JOIN] --> RW[RoomWriter 进程内环形缓冲]
    C2[发言 SEND_MESSAGE<br/>room_conversation_id] --> RW
    RW -->|分配 room_seq| BC[按 ConnectionShard 实时广播]
    BC --> C3[在房成员]
    RW -->|短期回放,可截断| C4[新加入者<br/>replay_truncated 提示]

聊天室不进个人邮箱、不计会话未读、不改会话列表;状态可丢失。