跳转至

MessageStore 与 MessageCommitter 契约

依赖:shared.md。MessageStore 是提交事实源;MessageCommitter 是唯一能把成功提交映射为 SEND_ACK 的边界。

pub trait MessageStore: Send + Sync {
    async fn reserve(
        &self,
        request: &AuthenticatedSend,
        delivery: FrozenDeliveryInput,
        candidate: ReserveCandidate,
    ) -> Result<ReserveResult, MessageStoreError>;

    async fn put_record(
        &self,
        intent: &CommitIntent,
    ) -> Result<(), MessageStoreError>;

    async fn mark_outbox_published(
        &self,
        intent: &CommitIntent,
        publication: OutboxPublication,
    ) -> Result<(), MessageStoreError>;

    async fn mark_published_and_committed(
        &self,
        intent: &CommitIntent,
        publication: OutboxPublication,
    ) -> Result<CommittedMessage, MessageStoreError>;

    async fn mark_committed(
        &self,
        intent: &CommitIntent,
    ) -> Result<CommittedMessage, MessageStoreError>;

    async fn load_commit(
        &self,
        key: &CommitKey,
        conversation_id: ConversationId,
    ) -> Result<Option<LoadedCommit>, MessageStoreError>;

    async fn load_commit_for_request(
        &self,
        request: &AuthenticatedSend,
    ) -> Result<Option<LoadedCommit>, MessageStoreError>;

    async fn scan_recoverable_commits(
        &self,
        older_than: TimestampMillis,
        max_items: u32,
    ) -> Result<Vec<CommitIntent>, MessageStoreError>;

    async fn history_page(
        &self,
        query: HistoryQuery,
    ) -> Result<HistoryPage, MessageStoreError>;

    async fn readiness(&self) -> Result<(), MessageStoreError>;
}

pub trait MessageCommitter: Send + Sync {
    async fn commit(
        &self,
        request: &AuthenticatedSend,
        delivery: FrozenDeliveryInput,
        candidate: ReserveCandidate,
    ) -> Result<CommittedMessage, CommitError>;

    async fn resume(&self, key: &CommitKey) -> Result<CommittedMessage, CommitError>;

    async fn resume_existing(
        &self,
        request: &AuthenticatedSend,
    ) -> Result<Option<CommittedMessage>, CommitError>;

    async fn recover_due(
        &self,
        older_than: TimestampMillis,
        max_items: u32,
    ) -> Result<u32, CommitError>;

    async fn readiness(&self) -> Result<(), CommitError>;
}

前置条件

  • commit 的 AuthenticatedSend 已携带网关认证身份;写入者仍须完成发送者成员准入、内容与配额检查。无效 ClientMessageId 必须返回 CommitError::InvalidClientMessageId。
  • reserve 之前固定 FrozenDeliveryInput 与完整 canonical MessageRecord;不得用恢复时的当前成员关系或客户端再次上传的 payload 替换它。
  • load_commit_for_request 必须比较持久的规范请求身份;相同 (tenant,sender,cmid) 但 device、conversation、content、mention 或 reply 任一不同均返回冲突,绝不能回放旧 ACK。
  • 已提交的消息只承诺回放 SEND_ACK 所需的坐标:reserve 命中返回 ReserveResult::Committed, load_commit* 返回 LoadedCommit::Committed。后端可在进入 COMMITTED 时丢弃记录本体(此时它已在 会话历史与 Outbox 中);Redis 把提交 hash 压缩为身份摘要 + 坐标 + publication(2026-09-28, 100 B 正文每条 2,640 B → 552 B 且保持 listpack),升级前的完整格式照常读取。持久请求身份可存 域分隔摘要(AuthenticatedSend::reserve_identity_digest,头 QIMD);读到升级前的完整身份(头 QIMC)时逐字节比较。
  • put_record 只能写 CommitIntent.record;history_page 必须先对 query.actor 执行会话成员授权。
  • readiness 必须真实访问生产依赖;参数为零的本地早退、lazy client/producer 构造成功都不是就绪证据。

后置条件

  • 同一 CommitKey 的首次和重复 reserve 返回相同坐标,且重复请求不消耗新的消息或会话序号。
  • reserve 成功持久化的 intent 单独包含恢复后续步骤所需的全部消息事实;客户端永不重试也不能形成不可恢复 RESERVED。
  • put_record 与 mark_outbox_published 可安全重试并收敛为同一值;mark_committed 只有在两者均已有可恢复证据时才能返回 CommittedMessage。
  • mark_published_and_committed 接受 durable OutboxPublication,并在成功时返回该冻结 intent 的 CommittedMessage。Memory 与 Redis 必须将 Recorded + publication → Committed 作为一个后端原子迁移:同值重试回放原 ACK,异 publication、缺少 record 或错误前态均不得留下部分写入;成功后必须移除恢复候选并登记应有的 history expiry。
  • 后端没有此原子能力时,默认实现依次调用 mark_outbox_published 与 mark_committed。若第二步失败,必须保留 Outboxed 及 publication 作为可恢复事实,随后由 resume 的 Outboxed 分支调用 mark_committed。Scylla 当前采用该默认两步条件写,不宣称单次 LWT 或跨表原子提交。
  • MessageCommitter 在 Recorded 分支取得 Outbox durable ack 并经过 M-COMMIT-04 故障点后调用 mark_published_and_committed;已持久化的 Outboxed 仍只调用 mark_committed。调用响应丢失或并发推进时,必须重读事实源并只向前收敛。
  • commit 或 resume 命中中间态必须从持久状态继续;命中 Committed 只回放原 CommittedMessage。CommitRecovery 必须调用同一 resume 语义。
  • 成功 CommittedMessage 是且仅是生成 SEND_ACK 的依据;它不表示邮箱物化、推送或已读。
  • 历史页保留 MessageType、custom_type、正文引用、方向、锚点、字节上限和可用边界;历史空洞不能被解释为丢失。

错误

  • 存储边界只暴露 MessageStoreError,不得透出具体后端条件写、键、表或脚本错误。
  • 提交边界只暴露 CommitError;Recoverable { key } 表示可由同一幂等键继续,绝不能伪造成功 ACK。