RabbitMqConsumerHostedService.cs
6,533 bytes
| 1 | using Microsoft.Extensions.DependencyInjection; |
|---|---|
| 2 | using Microsoft.Extensions.Hosting; |
| 3 | using Microsoft.Extensions.Logging; |
| 4 | using Microsoft.Extensions.Options; |
| 5 | using RabbitMQ.Client; |
| 6 | using RabbitMQ.Client.Events; |
| 7 | |
| 8 | namespace SplitApp.Shared.Messaging.Internal; |
| 9 | |
| 10 | /// <summary> |
| 11 | /// Server-side of the bus. On start, declares one durable queue per registered handler, |
| 12 | /// binds it to the right exchange + routing key, and starts an async consumer. |
| 13 | /// For RPC handlers, publishes the response back to the caller's ReplyTo queue with |
| 14 | /// the original CorrelationId. Acks on success; nacks (no requeue) on handler exception. |
| 15 | /// </summary> |
| 16 | public class RabbitMqConsumerHostedService : BackgroundService |
| 17 | { |
| 18 | private readonly RabbitMqConnectionProvider _connections; |
| 19 | private readonly HandlerRegistry _registry; |
| 20 | private readonly IServiceScopeFactory _scopeFactory; |
| 21 | private readonly MessagingOptions _options; |
| 22 | private readonly ILogger<RabbitMqConsumerHostedService> _logger; |
| 23 | |
| 24 | private readonly List<IChannel> _consumerChannels = new(); |
| 25 | |
| 26 | public RabbitMqConsumerHostedService( |
| 27 | RabbitMqConnectionProvider connections, |
| 28 | HandlerRegistry registry, |
| 29 | IServiceScopeFactory scopeFactory, |
| 30 | IOptions<MessagingOptions> options, |
| 31 | ILogger<RabbitMqConsumerHostedService> logger) |
| 32 | { |
| 33 | _connections = connections; |
| 34 | _registry = registry; |
| 35 | _scopeFactory = scopeFactory; |
| 36 | _options = options.Value; |
| 37 | _logger = logger; |
| 38 | } |
| 39 | |
| 40 | protected override async Task ExecuteAsync(CancellationToken stoppingToken) |
| 41 | { |
| 42 | if (_registry.All.Count == 0) |
| 43 | { |
| 44 | _logger.LogInformation("No integration handlers registered; consumer service idle"); |
| 45 | return; |
| 46 | } |
| 47 | |
| 48 | var conn = await _connections.GetConnectionAsync(stoppingToken); |
| 49 | |
| 50 | foreach (var registration in _registry.All) |
| 51 | { |
| 52 | await SubscribeAsync(conn, registration, stoppingToken); |
| 53 | } |
| 54 | |
| 55 | try |
| 56 | { |
| 57 | await Task.Delay(Timeout.Infinite, stoppingToken); |
| 58 | } |
| 59 | catch (OperationCanceledException) |
| 60 | { |
| 61 | // Expected on shutdown |
| 62 | } |
| 63 | } |
| 64 | |
| 65 | private async Task SubscribeAsync(IConnection conn, MessageHandlerRegistration reg, CancellationToken ct) |
| 66 | { |
| 67 | var channel = await conn.CreateChannelAsync(cancellationToken: ct); |
| 68 | _consumerChannels.Add(channel); |
| 69 | |
| 70 | await channel.ExchangeDeclareAsync( |
| 71 | exchange: reg.Exchange, |
| 72 | type: reg.Exchange == RabbitMqTopology.EventsExchange |
| 73 | ? RabbitMqTopology.EventsExchangeType |
| 74 | : RabbitMqTopology.RequestsExchangeType, |
| 75 | durable: true, |
| 76 | autoDelete: false, |
| 77 | cancellationToken: ct); |
| 78 | |
| 79 | var queueName = RabbitMqTopology.HandlerQueueName(_options.ServiceName, reg.MessageType.Name); |
| 80 | await channel.QueueDeclareAsync( |
| 81 | queue: queueName, |
| 82 | durable: true, |
| 83 | exclusive: false, |
| 84 | autoDelete: false, |
| 85 | cancellationToken: ct); |
| 86 | |
| 87 | await channel.QueueBindAsync( |
| 88 | queue: queueName, |
| 89 | exchange: reg.Exchange, |
| 90 | routingKey: reg.RoutingKey, |
| 91 | cancellationToken: ct); |
| 92 | |
| 93 | // Process one message at a time per consumer (predictable handler concurrency). |
| 94 | await channel.BasicQosAsync(prefetchSize: 0, prefetchCount: 1, global: false, cancellationToken: ct); |
| 95 | |
| 96 | var consumer = new AsyncEventingBasicConsumer(channel); |
| 97 | consumer.ReceivedAsync += (sender, ea) => OnReceivedAsync(channel, reg, ea, ct); |
| 98 | |
| 99 | await channel.BasicConsumeAsync( |
| 100 | queue: queueName, |
| 101 | autoAck: false, |
| 102 | consumer: consumer, |
| 103 | cancellationToken: ct); |
| 104 | |
| 105 | _logger.LogInformation( |
| 106 | "Subscribed: queue={Queue} exchange={Exchange} routingKey={RoutingKey} handler={Handler}", |
| 107 | queueName, reg.Exchange, reg.RoutingKey, reg.HandlerType.Name); |
| 108 | } |
| 109 | |
| 110 | private async Task OnReceivedAsync( |
| 111 | IChannel channel, |
| 112 | MessageHandlerRegistration reg, |
| 113 | BasicDeliverEventArgs ea, |
| 114 | CancellationToken stoppingToken) |
| 115 | { |
| 116 | try |
| 117 | { |
| 118 | using var scope = _scopeFactory.CreateScope(); |
| 119 | var responseBytes = await reg.Dispatcher(ea.Body, scope.ServiceProvider, stoppingToken); |
| 120 | |
| 121 | if (reg.Kind == MessageKind.Request && responseBytes is not null) |
| 122 | { |
| 123 | var replyTo = ea.BasicProperties.ReplyTo; |
| 124 | var correlationId = ea.BasicProperties.CorrelationId; |
| 125 | if (string.IsNullOrEmpty(replyTo) || string.IsNullOrEmpty(correlationId)) |
| 126 | { |
| 127 | _logger.LogWarning( |
| 128 | "Request {Type} missing ReplyTo/CorrelationId; dropping response", reg.MessageType.Name); |
| 129 | } |
| 130 | else |
| 131 | { |
| 132 | var replyProps = new BasicProperties |
| 133 | { |
| 134 | CorrelationId = correlationId, |
| 135 | ContentType = "application/json", |
| 136 | }; |
| 137 | await channel.BasicPublishAsync( |
| 138 | exchange: "", |
| 139 | routingKey: replyTo, |
| 140 | mandatory: false, |
| 141 | basicProperties: replyProps, |
| 142 | body: responseBytes, |
| 143 | cancellationToken: stoppingToken); |
| 144 | } |
| 145 | } |
| 146 | |
| 147 | await channel.BasicAckAsync(ea.DeliveryTag, multiple: false); |
| 148 | } |
| 149 | catch (Exception ex) |
| 150 | { |
| 151 | _logger.LogError(ex, "Handler {Handler} failed for {Type}", reg.HandlerType.Name, reg.MessageType.Name); |
| 152 | try |
| 153 | { |
| 154 | await channel.BasicNackAsync(ea.DeliveryTag, multiple: false, requeue: false); |
| 155 | } |
| 156 | catch (Exception nackEx) |
| 157 | { |
| 158 | _logger.LogError(nackEx, "Nack failed for {Type}", reg.MessageType.Name); |
| 159 | } |
| 160 | } |
| 161 | } |
| 162 | |
| 163 | public override async Task StopAsync(CancellationToken cancellationToken) |
| 164 | { |
| 165 | await base.StopAsync(cancellationToken); |
| 166 | |
| 167 | foreach (var ch in _consumerChannels) |
| 168 | { |
| 169 | try |
| 170 | { |
| 171 | await ch.CloseAsync(cancellationToken); |
| 172 | await ch.DisposeAsync(); |
| 173 | } |
| 174 | catch (Exception ex) |
| 175 | { |
| 176 | _logger.LogWarning(ex, "Error closing consumer channel"); |
| 177 | } |
| 178 | } |
| 179 | _consumerChannels.Clear(); |
| 180 | } |
| 181 | } |
| 182 | |