-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathMessageBusRPCListener.cs
94 lines (74 loc) · 3.28 KB
/
MessageBusRPCListener.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
using System.Text;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Pact.Core.Extensions;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
namespace Pact.RabbitMQ;
/// <summary>
/// A basic class used to listen for RPC events
/// </summary>
/// <typeparam name="T1">Type of object received in messages</typeparam>
/// <typeparam name="T2">Type of object sent as a response</typeparam>
public abstract class MessageBusRPCListener<T1, T2> : IMessageBusListener where T1 : class where T2 : class
{
protected MessageBusRPCListener(ILoggerFactory loggerFactory, IServiceProvider services)
{
_services = services;
Logger = loggerFactory.CreateLogger(GetType().Name);
}
private readonly IServiceProvider _services;
public readonly ILogger Logger;
public abstract string Name { get; }
public abstract string Key { get; }
public abstract string Exchange { get; }
/// <summary>
/// Message processing handler (RPC)
/// </summary>
/// <param name="services">A scoped service provider for this message</param>
/// <param name="message">The message object</param>
/// <returns>The result object of the message call</returns>
public abstract Task<T2> ProcessMessage(IServiceProvider services, T1 message);
public Task Setup(IMessageBusClient client)
{
try
{
//Create a disposable queue
var queueName = client.Channel.QueueDeclare(Name, true, false, false).QueueName;
client.Channel.BasicQos(0, 1, false);
//Bind it to the exchange
client.Channel.QueueBind(queueName, Exchange, Key);
var consumer = new EventingBasicConsumer(client.Channel);
consumer.Received += async (model, ea) =>
{
try
{
using var scope = _services.CreateScope();
var messageJson = Encoding.UTF8.GetString(ea.Body.ToArray());
var message = messageJson.FromJson<T1>();
var response = await ProcessMessage(scope.ServiceProvider, message);
var responseJson = response.ToJson();
var responseData = Encoding.UTF8.GetBytes(responseJson);
//Send result
var replyProps = client.Channel.CreateBasicProperties();
replyProps.CorrelationId = ea.BasicProperties.CorrelationId;
client.Channel.BasicPublish("", ea.BasicProperties.ReplyTo, replyProps, responseData);
client.Channel.BasicAck(ea.DeliveryTag, false);
}
catch (Exception exc)
{
var body = ea.Body.ToArray();
Logger.LogError(exc, "RabbitMQ Message Failed => {name} => {message}", Name, Encoding.UTF8.GetString(body));
client.Channel.BasicNack(ea.DeliveryTag, false, false);
client.SendError(body, ea.RoutingKey);
}
};
client.Channel.BasicConsume(queueName, false, consumer);
}
catch (Exception exc)
{
Logger.LogError(exc, "Error setting up message bus listener (rpc) {name}", Name);
}
return Task.CompletedTask;
}
}