profileShare

rasmusjy / splitapp-backend-microservices

Read-only snapshot

No repository description.

main default branch 501 files Expires Sep 13, 2026, 9:06 AM
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