跳转至

CommitLog 与 FanoutCoordinator 契约(ADR-0012 无 Kafka 事务)

依赖:shared.md。此边界封装可靠 Outbox、已提交读取和 “全部 dispatch durable 后才提交 source offset”的顺序;调用方不操作消息中间件配置或位点提交。

实现中的 TransactionalFanout / FanoutBatchTransaction 是 ADR-0012 之前留下的名称, 不表示当前使用 Kafka transaction。现行实现允许“输出已 durable、source checkpoint 失败”后 整批重放,并由确定性 dispatch_id + payload_digest 持久收敛。

pub trait OutboxPublisher: Send + Sync {
    async fn publish(
        &self,
        envelope: CommittedEnvelope,
    ) -> Result<LogPosition, CommitLogError>;
}

pub struct FanoutBatchTransaction {
    pub transactions: Vec<FanoutTransaction>,
}

pub trait TransactionalFanout: Send + Sync {
    async fn read_committed(
        &self,
        after: Option<LogPosition>,
    ) -> Result<Option<(LogPosition, CommittedEnvelope)>, CommitLogError>;

    async fn read_committed_batch(
        &self,
        after: Option<LogPosition>,
        max_records: usize,
    ) -> Result<Vec<(LogPosition, CommittedEnvelope)>, CommitLogError>;

    async fn publish_dispatches_and_ack_source(
        &self,
        transaction: FanoutTransaction,
    ) -> Result<(), CommitLogError>;

    async fn publish_dispatches_and_ack_batch(
        &self,
        batch: FanoutBatchTransaction,
    ) -> Result<(), CommitLogError>;
}

pub trait FanoutCoordinator: Send + Sync {
    async fn plan(
        &self,
        source: LogPosition,
        envelope: CommittedEnvelope,
    ) -> Result<FanoutTransaction, FanoutError>;

    async fn fanout_once(
        &self,
        source: LogPosition,
        envelope: CommittedEnvelope,
    ) -> Result<(), FanoutError>;
}

前置条件

  • OutboxPublisher::publish 的输入必须来自已写入 CommitIntent.record 的同一 CommittedEnvelope;若 envelope 同时携带 record,其字节必须与 intent 内记录完全一致。稳定身份与 Kafka key 是其 MessageId,但应用级重试允许产生同 key 重复记录;逻辑幂等由确定性 dispatch 与 MailboxStore 持久收敛。
  • read_committed 只交付已完成的 Outbox 记录。plan 的直接收件人或群精确成员版本、正文内联决定和创建时间均取自固化信封;群成员只能按该版本读取不可变快照。
  • 调用 publish_dispatches_and_ack_source/batch 前,每个 FanoutTransaction.source 必须是本批 已读取且连续的 Outbox 位置,全部 dispatch 均须有确定性身份和目标分片。

后置条件

  • publish 成功返回的 LogPosition 是 MessageStore::mark_outbox_published 所需的持久证据。
  • 对相同输入,plan 输出的 dispatch 数量、DispatchId、目标分片和内容完全相同;每个目标 MailboxShard 恰好一条 GroupDispatch,收件人分块不产生额外 dispatch。不得读取当前成员版本或生成随机身份。
  • publish_dispatches_and_ack_source/batch 必须先等待本批全部 dispatch 的 durable delivery, 再同步提交批末 source 的下一 offset。若 dispatch 已 durable 而 checkpoint 失败,源记录仍可重试, 重复输出是允许的;若任一 dispatch 未取得 durable 证据,source checkpoint 必须保持不动。 成功或重放时都不得为同一 source 产生不同的 dispatch 集合。
  • FanoutCoordinator 不读写 MailboxEntry、不写 Socket、不修改 MessageRecord,也不以进程内队列的成功替代持久日志。

错误

  • CommitLogError::PublicationNotDurable 或历史命名的 TransactionAborted 返回时,调用方不得推进 source 位置;后者不表示发生了 Kafka transaction abort。
  • FanoutError::MembershipSnapshotUnavailable 与 DeterminismViolation 均不可降级为空 dispatch;应保留源记录以待恢复。