跳转至

MailboxMaterializer 与 MailboxStore 契约

依赖:shared.md。邮箱是持久事实,在线 PushBatch 是在成功物化后产生的易失优化。

pub trait MailboxStore: Send + Sync {
    async fn append_batch(
        &self,
        batch: AppendMailboxBatch,
    ) -> Result<DurabilityReceipt, MailboxStoreError>;

    async fn range_scan(
        &self,
        query: MailboxRangeQuery,
    ) -> Result<MailboxPage, MailboxStoreError>;

    async fn truncate_before(
        &self,
        tenant_id: TenantId,
        user_id: UserId,
        before: MailboxSeq,
    ) -> Result<MailboxSeq, MailboxStoreError>;

    async fn dispatch_progress(
        &self,
        shard: MailboxShardId,
        dispatch_id: DispatchId,
    ) -> Result<Vec<DispatchProgress>, MailboxStoreError>;

    async fn lane_watermarks(
        &self,
        shard: MailboxShardId,
    ) -> Result<Vec<LaneWatermark>, MailboxStoreError>;

    async fn advance_watermark(
        &self,
        advance: WatermarkAdvance,
    ) -> Result<(), MailboxStoreError>;
}

pub trait MailboxMaterializer: Send + Sync {
    async fn materialize(
        &self,
        record: DispatchLogRecord,
    ) -> Result<Vec<PushBatch>, MailboxMaterializeError>;

    async fn prepare_takeover(
        &self,
        takeover: ShardTakeover,
    ) -> Result<Vec<ShardReadiness>, MailboxMaterializeError>;

    async fn pull_mailbox(
        &self,
        request: MailboxPullRequest,
    ) -> Result<MailboxPage, MailboxMaterializeError>;
}

前置条件

  • materialize 只处理静态归属分片的 DispatchLogRecord;其位置、epoch 与 GroupDispatch.target_shard 必须一致。
  • append_batch 的 entries 必须是同一 (dispatch_id, lane, chunk_id, mailbox_seq) 的完整事件组,且 progress 只描述该批已成功写入的块。
  • range_scan 的区间为 (after_seq, up_to_seq]。max_items 和 max_bytes 为软上限,不能切开同一 MailboxSeq 的事件组。
  • pull_mailbox 必须验证 MailboxPullRequest.actor、签名游标身份和范围中的 tenant/user 一致;存储层 range_scan 只接受已获授权的内部查询。

后置条件

  • MailboxSeq 仅由分发日志记录位置和当前分片 epoch 构造;同一 GroupDispatch 的全部 lane/chunk 共用该值,重复物化重用同一值,条目按稳定身份同值覆盖。
  • append_batch 成功后,条目与该块进度要么一同可恢复,要么遵守“条目先成功、进度后成功”的等价顺序。绝不允许先记录完成进度后写条目。
  • 仅在一个 lane 的全部相关块已持久化后,advance_watermark 才能单调推进,并同步持久化其 floor。物化失败、解码失败或缺块时不得推进。
  • prepare_takeover 从持久水位的最小值推导重放起点,并验证每 lane 的重算水位不低于 floor;未追平 lane 返回 LaneNotReady,不得提供该 lane 的认证或拉取服务。
  • truncate_before 返回的新裁剪下界,后续扫描越界必须返回 CursorExpired,不能把已裁剪区伪装为空洞。
  • 只有持久 append_batch 成功后,materialize 才能返回对应 PushBatch;Push 失败不影响邮箱恢复性。

错误

  • MailboxStoreError::CorruptEntry 必须使当前扫描失败且不推进 covered_through_seq。
  • MailboxMaterializeError::ReadinessFailed 表示不能证明水位安全,必须拒绝就绪而非依赖消费者位点继续服务。