Dreamine.Communication.Core 1.0.2
Dreamine.Communication.Core 통신 기능과 관련 API를 제공합니다.
로딩중...
검색중...
일치하는것 없음
InMemoryMessageBus.cs
이 파일의 문서화 페이지로 가기
1using System;
2using System.Collections.Concurrent;
3using System.Collections.Generic;
4using System.Threading;
5using System.Threading.Tasks;
6using Dreamine.Communication.Abstractions.Enums;
7using Dreamine.Communication.Abstractions.Interfaces;
8using Dreamine.Communication.Abstractions.Models;
9
11
20public sealed class InMemoryMessageBus : IMessageBus
21{
30 private readonly ConcurrentDictionary<string, List<Func<MessageEnvelope, CancellationToken, Task>>> _handlers = new();
31
40 public ConnectionState State { get; private set; } = ConnectionState.Disconnected;
41
50 public TransportKind Kind => TransportKind.InMemory;
51
84 public Task ConnectAsync(CancellationToken cancellationToken = default)
85 {
86 cancellationToken.ThrowIfCancellationRequested();
87 State = ConnectionState.Connected;
88 return Task.CompletedTask;
89 }
90
147 public async Task PublishAsync(
148 MessageEnvelope message,
149 CancellationToken cancellationToken = default)
150 {
151 ArgumentNullException.ThrowIfNull(message);
152
153 if (State != ConnectionState.Connected)
154 {
155 throw new InvalidOperationException("The message bus is not connected.");
156 }
157
158 if (!_handlers.TryGetValue(message.Route, out var handlers))
159 {
160 return;
161 }
162
163 List<Func<MessageEnvelope, CancellationToken, Task>> snapshot;
164
165 lock (handlers)
166 {
167 snapshot = new List<Func<MessageEnvelope, CancellationToken, Task>>(handlers);
168 }
169
170 foreach (var handler in snapshot)
171 {
172 cancellationToken.ThrowIfCancellationRequested();
173 await handler(message, cancellationToken).ConfigureAwait(false);
174 }
175 }
176
241 public Task SubscribeAsync(
242 string route,
243 Func<MessageEnvelope, CancellationToken, Task> handler,
244 CancellationToken cancellationToken = default)
245 {
246 ArgumentException.ThrowIfNullOrWhiteSpace(route);
247 ArgumentNullException.ThrowIfNull(handler);
248 cancellationToken.ThrowIfCancellationRequested();
249
250 var handlers = _handlers.GetOrAdd(
251 route,
252 _ => new List<Func<MessageEnvelope, CancellationToken, Task>>());
253
254 lock (handlers)
255 {
256 handlers.Add(handler);
257 }
258
259 return Task.CompletedTask;
260 }
261
294 public Task DisconnectAsync(CancellationToken cancellationToken = default)
295 {
296 cancellationToken.ThrowIfCancellationRequested();
297 State = ConnectionState.Disconnected;
298 return Task.CompletedTask;
299 }
300
317 public ValueTask DisposeAsync()
318 {
319 _handlers.Clear();
320 State = ConnectionState.Disconnected;
321 return ValueTask.CompletedTask;
322 }
323}
Task DisconnectAsync(CancellationToken cancellationToken=default)
async Task PublishAsync(MessageEnvelope message, CancellationToken cancellationToken=default)
readonly ConcurrentDictionary< string, List< Func< MessageEnvelope, CancellationToken, Task > > > _handlers
Task ConnectAsync(CancellationToken cancellationToken=default)
Task SubscribeAsync(string route, Func< MessageEnvelope, CancellationToken, Task > handler, CancellationToken cancellationToken=default)