/
vshmidt
/
masstransit
Обзор
Документация
Войти
/
vshmidt
/
masstransit
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
develop
src/MassTransit/Consumers/Configuration/MessageObserverConnector.cs
42 строки
2 KB
Chris Patterson
Integration of all external packages into MassTransit, split out into Abstractions, Middleware, and the core MassTransit assembly.
22 янв 2022, 18:24
22 янв 2022, 18:24
b5b4f50
Код
Авторство
О чём код?
namespace MassTransit.Configuration { using System; /// <summary> /// Connects a message handler to a pipe /// </summary> /// <typeparam name="TMessage"></typeparam> public class MessageObserverConnector<TMessage> : IObserverConnector<TMessage> where TMessage : class { public ConnectHandle ConnectObserver(IConsumePipeConnector consumePipe, IObserver<ConsumeContext<TMessage>> observer, params IFilter<ConsumeContext<TMessage>>[] filters) { IPipe<ConsumeContext<TMessage>> pipe = Pipe.New<ConsumeContext<TMessage>>(x => { foreach (IFilter<ConsumeContext<TMessage>> filter in filters) x.UseFilter(filter); x.AddPipeSpecification(new ObserverPipeSpecification<TMessage>(observer)); }); return consumePipe.ConnectConsumePipe(pipe); } public ConnectHandle ConnectRequestObserver(IRequestPipeConnector consumePipe, Guid requestId, IObserver<ConsumeContext<TMessage>> observer, params IFilter<ConsumeContext<TMessage>>[] filters) { IPipe<ConsumeContext<TMessage>> pipe = Pipe.New<ConsumeContext<TMessage>>(x => { foreach (IFilter<ConsumeContext<TMessage>> filter in filters) x.UseFilter(filter); x.AddPipeSpecification(new ObserverPipeSpecification<TMessage>(observer)); }); return consumePipe.ConnectRequestPipe(requestId, pipe); } } }