2024-10-30 17:49:05 +08:00

46 lines
2.0 KiB
C#

using System.Security.Cryptography.X509Certificates;
using JiShe.CollectBus.Common.Enums;
using JiShe.CollectBus.MongoDB;
using JiShe.CollectBus.Protocol.Contracts.Models;
using MassTransit;
using Microsoft.Extensions.Logging;
using TouchSocket.Sockets;
namespace JiShe.CollectBus.RabbitMQ.Consumers
{
public class MessageIssuedConsumer : IConsumer<MessageIssuedEvent>
{
private readonly ILogger<MessageIssuedEvent> _logger;
private readonly ITcpService _tcpService;
public readonly IMongoRepository<MessageReceivedHeartbeatEvent> _mongoHeartbeatRepository;
public readonly IMongoRepository<MessageReceivedLoginEvent> _mongoLoginRepository;
public MessageIssuedConsumer(ILogger<MessageIssuedEvent> logger, ITcpService tcpService, IMongoRepository<MessageReceivedHeartbeatEvent> mongoHeartbeatRepository, IMongoRepository<MessageReceivedLoginEvent> mongoLoginRepository)
{
_logger = logger;
_tcpService = tcpService;
_mongoHeartbeatRepository = mongoHeartbeatRepository;
_mongoLoginRepository = mongoLoginRepository;
}
public async Task Consume(ConsumeContext<MessageIssuedEvent> context)
{
switch (context.Message.Type)
{
case IssuedEventType.Heartbeat:
await _mongoHeartbeatRepository.UpdateAsync(a => a.MessageId == context.Message.MessageId, b =>new MessageReceivedHeartbeatEvent { IsAck = true,AckTime=DateTime.Now});
break;
case IssuedEventType.Login:
await _mongoLoginRepository.UpdateAsync(a => a.MessageId == context.Message.MessageId, b => new MessageReceivedLoginEvent { IsAck = true, AckTime = DateTime.Now });
break;
case IssuedEventType.Data:
break;
default:
throw new ArgumentOutOfRangeException();
}
await _tcpService.SendAsync(context.Message.ClientId, context.Message.Message);
}
}
}