Dreamine.Communication.Core 1.0.2
Dreamine.Communication.Core 통신 기능과 관련 API를 제공합니다.
로딩중...
검색중...
일치하는것 없음
TransportMessageBusAdapter.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 TransportMessageBusAdapter : IMessageBus, IDisposable
21{
30 private readonly IMessageTransport _transport;
39 private readonly ConcurrentDictionary<string, List<Func<MessageEnvelope, CancellationToken, Task>>> _handlers = new();
40
49 public Action<Exception, string>? OnHandlerError { get; set; }
50
75 public TransportMessageBusAdapter(IMessageTransport transport)
76 {
77 _transport = transport ?? throw new ArgumentNullException(nameof(transport));
78 _transport.MessageReceived += OnMessageReceived;
79 }
80
89 public ConnectionState State => _transport.State;
90
99 public TransportKind Kind => _transport.Kind;
100
125 public Task ConnectAsync(CancellationToken cancellationToken = default)
126 {
127 return _transport.ConnectAsync(cancellationToken);
128 }
129
170 public Task PublishAsync(
171 MessageEnvelope message,
172 CancellationToken cancellationToken = default)
173 {
174 ArgumentNullException.ThrowIfNull(message);
175 return _transport.SendAsync(message, cancellationToken);
176 }
177
242 public Task SubscribeAsync(
243 string route,
244 Func<MessageEnvelope, CancellationToken, Task> handler,
245 CancellationToken cancellationToken = default)
246 {
247 ArgumentException.ThrowIfNullOrWhiteSpace(route);
248 ArgumentNullException.ThrowIfNull(handler);
249 cancellationToken.ThrowIfCancellationRequested();
250
251 var handlers = _handlers.GetOrAdd(
252 route,
253 _ => new List<Func<MessageEnvelope, CancellationToken, Task>>());
254
255 lock (handlers)
256 {
257 handlers.Add(handler);
258 }
259
260 return Task.CompletedTask;
261 }
262
287 public Task DisconnectAsync(CancellationToken cancellationToken = default)
288 {
289 return _transport.DisconnectAsync(cancellationToken);
290 }
291
300 public void Dispose()
301 {
302 _transport.MessageReceived -= OnMessageReceived;
303 _ = _transport.DisposeAsync().AsTask();
304 }
305
322 public async ValueTask DisposeAsync()
323 {
324 _transport.MessageReceived -= OnMessageReceived;
325 await _transport.DisposeAsync().ConfigureAwait(false);
326 }
327
360 private async void OnMessageReceived(object? sender, MessageEnvelope message)
361 {
362 try
363 {
364 if (!_handlers.TryGetValue(message.Route, out var handlers))
365 {
366 return;
367 }
368
369 List<Func<MessageEnvelope, CancellationToken, Task>> snapshot;
370
371 lock (handlers)
372 {
373 snapshot = new List<Func<MessageEnvelope, CancellationToken, Task>>(handlers);
374 }
375
376 foreach (var handler in snapshot)
377 {
378 await handler(message, CancellationToken.None).ConfigureAwait(false);
379 }
380 }
381 catch (Exception ex)
382 {
383 // Event dispatchers cannot propagate async-void exceptions.
384 // Delegate to caller-supplied handler so failures are observable.
385 // 이벤트 디스패처는 async void 예외를 전파할 수 없습니다.
386 // 호출자가 OnHandlerError를 설정하면 실패를 관찰할 수 있습니다.
387 OnHandlerError?.Invoke(ex, message.Route);
388 }
389 }
390}
async void OnMessageReceived(object? sender, MessageEnvelope message)
Task PublishAsync(MessageEnvelope message, CancellationToken cancellationToken=default)
Task DisconnectAsync(CancellationToken cancellationToken=default)
Task ConnectAsync(CancellationToken cancellationToken=default)
readonly ConcurrentDictionary< string, List< Func< MessageEnvelope, CancellationToken, Task > > > _handlers
Task SubscribeAsync(string route, Func< MessageEnvelope, CancellationToken, Task > handler, CancellationToken cancellationToken=default)