using Dpz.Core.MessageQueue.Models;

namespace Dpz.Core.MessageQueue.Abstractions;

/// <summary>
/// 消息 Outbox 存储接口,用于持久化消息投递状态,实现发布和消费的数据库兜底。
/// </summary>
public interface IMessageOutboxStore
{
    /// <summary>
    /// 创建一条 Outbox 记录,在消息发布前调用,初始状态为"待发布"。
    /// </summary>
    /// <param name="messageId">消息唯一标识</param>
    /// <param name="messageType">消息类型的完全限定名</param>
    /// <param name="exchange">目标 Exchange 名称</param>
    /// <param name="routingKey">路由键</param>
    /// <param name="payload">JSON 序列化的消息体</param>
    /// <param name="source">消息来源标识,可为空</param>
    /// <param name="cancellationToken">取消令牌</param>
    Task CreateAsync(
        string messageId,
        string messageType,
        string exchange,
        string routingKey,
        string payload,
        string? source,
        CancellationToken cancellationToken = default
    );

    /// <summary>
    /// 将消息标记为已成功发布到 RabbitMQ,累加发布尝试次数,状态变为"已发布",
    /// 并记录本次尝试时间和发布时间。
    /// </summary>
    /// <param name="messageId">消息唯一标识</param>
    /// <param name="cancellationToken">取消令牌</param>
    Task MarkSentAsync(string messageId, CancellationToken cancellationToken = default);

    /// <summary>
    /// 将消息标记为发布失败,累加发布尝试次数并按指数退避计算下次重试时间,
    /// 状态变为"发布失败"。
    /// </summary>
    /// <param name="messageId">消息唯一标识</param>
    /// <param name="error">失败原因描述</param>
    /// <param name="cancellationToken">取消令牌</param>
    Task MarkPublishFailedAsync(
        string messageId,
        string error,
        CancellationToken cancellationToken = default
    );

    /// <summary>
    /// 将消息标记为已成功消费,累加消费尝试次数,状态变为"已消费",
    /// 并记录本次尝试时间和消费时间。
    /// </summary>
    /// <param name="messageId">消息唯一标识</param>
    /// <param name="cancellationToken">取消令牌</param>
    Task MarkConsumedAsync(string messageId, CancellationToken cancellationToken = default);

    /// <summary>
    /// 记录一次消费失败尝试,累加消费尝试次数并记录错误,但不改变消息当前状态。
    /// 用于消费者仍会立即重发消息体的短重试路径。
    /// </summary>
    /// <param name="messageId">消息唯一标识</param>
    /// <param name="error">失败原因描述</param>
    /// <param name="cancellationToken">取消令牌</param>
    Task RecordConsumeAttemptFailedAsync(
        string messageId,
        string error,
        CancellationToken cancellationToken = default
    );

    /// <summary>
    /// 将消息标记为消费失败,累加消费尝试次数并按指数退避计算下次重试时间,
    /// 状态变为"消费失败"。
    /// </summary>
    /// <param name="messageId">消息唯一标识</param>
    /// <param name="error">失败原因描述</param>
    /// <param name="cancellationToken">取消令牌</param>
    Task MarkConsumeFailedAsync(
        string messageId,
        string error,
        CancellationToken cancellationToken = default
    );

    /// <summary>
    /// 获取当前可进行发布重试的 Outbox 条目(状态为"发布失败"且退避时间已到)。
    /// </summary>
    /// <param name="batchSize">每批最大条数</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>待重试的 Outbox 条目列表</returns>
    Task<IReadOnlyList<MessageOutboxEntry>> GetPendingPublishRetryAsync(
        int batchSize,
        CancellationToken cancellationToken = default
    );

    /// <summary>
    /// 获取当前可进行消费重试的 Outbox 条目(状态为"消费失败"且退避时间已到)。
    /// </summary>
    /// <param name="batchSize">每批最大条数</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>待重试的 Outbox 条目列表</returns>
    Task<IReadOnlyList<MessageOutboxEntry>> GetPendingConsumeRetryAsync(
        int batchSize,
        CancellationToken cancellationToken = default
    );
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

代码解释:消息 Outbox 存储接口

这是一个用于实现 Outbox 模式 的接口定义,Outbox 模式是分布式系统中保证消息可靠投递的常用模式。

核心目的

该接口定义了消息持久化存储的标准操作,用于:

  1. 持久化消息状态:在数据库中记录消息的生命周期
  2. 实现消息可靠投递:防止消息丢失
  3. 支持重试机制:通过数据库兜底实现发布和消费的重试

主要方法说明

1. CreateAsync - 创建消息记录

Task CreateAsync(string messageId, string messageType, ...)
  • 在消息发布前创建一条记录
  • 初始状态为"待发布"
  • 记录消息的元数据(ID、类型、路由信息、消息体等)

2. MarkSentAsync - 标记发布成功

  • 消息成功发送到 RabbitMQ 后调用
  • 状态变更为"已发布"
  • 累加发布尝试次数

3. MarkPublishFailedAsync - 标记发布失败

  • 发布失败时调用
  • 记录失败原因
  • 使用指数退避算法计算下次重试时间

4. MarkConsumedAsync - 标记消费成功

  • 消费者成功处理消息后调用
  • 状态变更为"已消费"
  • 累加消费尝试次数

5. RecordConsumeAttemptFailedAsync - 记录消费尝试失败

  • 记录一次失败尝试,但不改变状态
  • 用于短重试路径(立即重试的场景)

6. MarkConsumeFailedAsync - 标记消费失败

  • 消费失败且需要延迟重试时调用
  • 状态变更为"消费失败"
  • 使用指数退避算法计算下次重试时间

7. GetPendingPublishRetryAsync - 获取待重试的发布消息

  • 查询状态为"发布失败"且退避时间已到的消息
  • 支持批量获取(batchSize)

8. GetPendingConsumeRetryAsync - 获取待重试的消费消息

  • 查询状态为"消费失败"且退避时间已到的消息
  • 支持批量获取

设计亮点

  1. 完整的生命周期管理:覆盖从创建到发布、消费的完整流程
  2. 区分短重试和长重试RecordConsumeAttemptFailedAsync vs MarkConsumeFailedAsync
  3. 指数退避:失败重试采用指数退避算法,避免系统压力
  4. 异步设计:所有方法都是异步的,支持高并发场景
  5. 取消令牌支持:所有方法都支持 CancellationToken,便于优雅关闭

典型使用场景

这个接口通常配合数据库事务使用,确保业务操作和消息发布的原子性,是实现最终一致性的关键组件。

评论加载中...