-
-
Notifications
You must be signed in to change notification settings - Fork 242
/
InMemoryMessageBus.cs
80 lines (65 loc) · 2.8 KB
/
InMemoryMessageBus.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
using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
using Foundatio.Utility;
using Microsoft.Extensions.Logging;
namespace Foundatio.Messaging;
public class InMemoryMessageBus : MessageBusBase<InMemoryMessageBusOptions>
{
private readonly ConcurrentDictionary<string, long> _messageCounts = new();
private long _messagesSent;
public InMemoryMessageBus() : this(o => o) { }
public InMemoryMessageBus(InMemoryMessageBusOptions options) : base(options) { }
public InMemoryMessageBus(Builder<InMemoryMessageBusOptionsBuilder, InMemoryMessageBusOptions> config)
: this(config(new InMemoryMessageBusOptionsBuilder()).Build()) { }
public long MessagesSent => _messagesSent;
public long GetMessagesSent(Type messageType)
{
return _messageCounts.TryGetValue(GetMappedMessageType(messageType), out long count) ? count : 0;
}
public long GetMessagesSent<T>()
{
return _messageCounts.TryGetValue(GetMappedMessageType(typeof(T)), out long count) ? count : 0;
}
public void ResetMessagesSent()
{
Interlocked.Exchange(ref _messagesSent, 0);
_messageCounts.Clear();
}
protected override async Task PublishImplAsync(string messageType, object message, MessageOptions options, CancellationToken cancellationToken)
{
Interlocked.Increment(ref _messagesSent);
_messageCounts.AddOrUpdate(messageType, t => 1, (t, c) => c + 1);
var mappedType = GetMappedMessageType(messageType);
if (_subscribers.IsEmpty)
return;
bool isTraceLogLevelEnabled = _logger.IsEnabled(LogLevel.Trace);
if (options.DeliveryDelay.HasValue && options.DeliveryDelay.Value > TimeSpan.Zero)
{
if (isTraceLogLevelEnabled)
_logger.LogTrace("Schedule delayed message: {MessageType} ({Delay}ms)", messageType, options.DeliveryDelay.Value.TotalMilliseconds);
SendDelayedMessage(mappedType, message, options.DeliveryDelay.Value);
return;
}
byte[] body = SerializeMessageBody(messageType, message);
var messageData = new Message(body, DeserializeMessageBody)
{
CorrelationId = options.CorrelationId,
UniqueId = options.UniqueId,
Type = messageType,
ClrType = mappedType
};
foreach (var property in options.Properties)
messageData.Properties[property.Key] = property.Value;
try
{
await SendMessageToSubscribersAsync(messageData).AnyContext();
}
catch (Exception ex)
{
// swallow exceptions from subscriber handlers for the in memory bus
_logger.LogWarning(ex, "Error sending message to subscribers: {Message}", ex.Message);
}
}
}