using Microsoft.Extensions.Logging; using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ; using QYZH.InteractiveMagazine.IService; using QYZH.InteractiveMagazine.Models.Common; using QYZH.InteractiveMagazine.Models.Dto; using QYZH.InteractiveMagazine.Models.Dto.RabbitMQ; using QYZH.InteractiveMagazine.Models.Entity; using QYZH.InteractiveMagazine.Models.Enum; using QYZH.InteractiveMagazine.Repository; using System.Text.Json; namespace QYZH.InteractiveMagazine.Service; /// /// MQ消息发布服务。 /// public class MessagePublishService( IRabbitMQService rabbitMQService, ILogger logger) : BaseRepository, IMessagePublishService { /// /// 可靠发布消息,默认写入Outbox。 /// public async Task PublishAsync(MessagePublishInput input, CancellationToken cancellationToken = default) { ValidateInput(input); var message = BuildOutboxMessage(input); await Context.Insertable(message).ExecuteCommandAsync(cancellationToken); return new MessagePublishResult { Success = true, OutboxIds = [message.Id], Message = "消息已写入Outbox" }; } /// /// 批量可靠发布消息,默认写入Outbox。 /// public async Task PublishBatchAsync(IEnumerable> inputs, CancellationToken cancellationToken = default) { var inputList = inputs.ToList(); BusinessException.ThrowIf(inputList.Count == 0, "消息发布列表不能为空", ResultCode.BAD_REQUEST); inputList.ForEach(ValidateInput); var messages = inputList.Select(BuildOutboxMessage).ToList(); await Context.Insertable(messages).ExecuteCommandAsync(cancellationToken); return new MessagePublishResult { Success = true, OutboxIds = messages.Select(x => x.Id).ToList(), Message = "消息已批量写入Outbox" }; } /// /// 直接发布消息,不写Outbox。 /// public async Task PublishDirectAsync(MessagePublishInput input, CancellationToken cancellationToken = default) { ValidateInput(input); var sent = await rabbitMQService.SendAsync(new RabbitMQSendParam { Exchange = input.Exchange, Queue = input.Queue, RoutingKey = input.RoutingKey, Data = input.Data! }, cancellationToken); if (!sent) { logger.LogError("MQ直接发布失败,Exchange: {Exchange}, Queue: {Queue}, RoutingKey: {RoutingKey}, BusinessType: {BusinessType}, BusinessId: {BusinessId}", input.Exchange, input.Queue, input.RoutingKey, input.BusinessType, input.BusinessId); throw new BusinessException("MQ消息发送失败", ResultCode.GLOBAL_ERROR); } return new MessagePublishResult { Success = true, Message = "消息已直接发送" }; } private static MessageOutbox BuildOutboxMessage(MessagePublishInput input) { var now = DateTime.Now; return new MessageOutbox { Exchange = input.Exchange, Queue = input.Queue, RoutingKey = input.RoutingKey, Payload = JsonSerializer.Serialize(input.Data), BusinessType = input.BusinessType, BusinessId = input.BusinessId, Status = (int)MessageOutboxStatusEnum.Pending, CreatedBy = "System", CreatedAt = now, UpdatedBy = "System", UpdatedAt = now }; } private static void ValidateInput(MessagePublishInput input) { BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.Exchange), "MQ交换机不能为空", ResultCode.BAD_REQUEST); BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.Queue), "MQ队列不能为空", ResultCode.BAD_REQUEST); BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.RoutingKey), "MQ路由键不能为空", ResultCode.BAD_REQUEST); BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.BusinessType), "MQ业务类型不能为空", ResultCode.BAD_REQUEST); BusinessException.ThrowIf(input.BusinessId <= 0, "MQ业务ID无效", ResultCode.BAD_REQUEST); BusinessException.ThrowIf(input.Data == null, "MQ消息数据不能为空", ResultCode.BAD_REQUEST); } }