Browse Source

Merge pull request #22134 from abpframework/LocalDistributedEventBus

[rel-9.1] Changes to make LocalDistributedEventBus able to use outbox/inbox patterns.
pull/22136/head
liangshiwei 2 years ago
committed by GitHub
parent
commit
829485fff7
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 16
      framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/AbpAspNetCoreMvcDaprEventBusModule.cs
  2. 17
      framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/EventTypeWithEventHandlerFactories.cs
  3. 10
      framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs
  4. 30
      framework/src/Volo.Abp.EventBus.Dapr/Volo/Abp/EventBus/Dapr/DaprDistributedEventBus.cs
  5. 7
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs
  6. 17
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs
  7. 224
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs
  8. 13
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs
  9. 5
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs
  10. 6
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs

16
framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/AbpAspNetCoreMvcDaprEventBusModule.cs

@ -97,13 +97,21 @@ public class AbpAspNetCoreMvcDaprEventBusModule : AbpModule
if (IsAbpDaprEventData(data)) if (IsAbpDaprEventData(data))
{ {
var daprEventData = daprSerializer.Deserialize(data, typeof(AbpDaprEventData)).As<AbpDaprEventData>(); var daprEventData = daprSerializer.Deserialize(data, typeof(AbpDaprEventData)).As<AbpDaprEventData>();
var eventData = daprSerializer.Deserialize(daprEventData.JsonData, distributedEventBus.GetEventType(daprEventData.Topic)); var eventType = distributedEventBus.GetEventType(daprEventData.Topic);
await distributedEventBus.TriggerHandlersAsync(distributedEventBus.GetEventType(daprEventData.Topic), eventData, daprEventData.MessageId, daprEventData.CorrelationId); if (eventType != null)
{
var eventData = daprSerializer.Deserialize(daprEventData.JsonData, eventType);
await distributedEventBus.TriggerHandlersAsync(eventType, eventData, daprEventData.MessageId, daprEventData.CorrelationId);
}
} }
else else
{ {
var eventData = daprSerializer.Deserialize(data, distributedEventBus.GetEventType(topic!)); var eventType = distributedEventBus.GetEventType(topic);
await distributedEventBus.TriggerHandlersAsync(distributedEventBus.GetEventType(topic!), eventData); if (eventType != null)
{
var eventData = daprSerializer.Deserialize(data, eventType);
await distributedEventBus.TriggerHandlersAsync(eventType, eventData);
}
} }
httpContext.Response.StatusCode = 200; httpContext.Response.StatusCode = 200;

17
framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/EventTypeWithEventHandlerFactories.cs

@ -0,0 +1,17 @@
using System;
using System.Collections.Generic;
namespace Volo.Abp.EventBus;
public class EventTypeWithEventHandlerFactories
{
public Type EventType { get; }
public List<IEventHandlerFactory> EventHandlerFactories { get; }
public EventTypeWithEventHandlerFactories(Type eventType, List<IEventHandlerFactory> eventHandlerFactories)
{
EventType = eventType;
EventHandlerFactories = eventHandlerFactories;
}
}

10
framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs

@ -1,4 +1,5 @@
using System; using System;
using System.Collections.Generic;
namespace Volo.Abp.EventBus.Local; namespace Volo.Abp.EventBus.Local;
@ -8,11 +9,18 @@ namespace Volo.Abp.EventBus.Local;
public interface ILocalEventBus : IEventBus public interface ILocalEventBus : IEventBus
{ {
/// <summary> /// <summary>
/// Registers to an event. /// Registers to an event.
/// Same (given) instance of the handler is used for all event occurrences. /// Same (given) instance of the handler is used for all event occurrences.
/// </summary> /// </summary>
/// <typeparam name="TEvent">Event type</typeparam> /// <typeparam name="TEvent">Event type</typeparam>
/// <param name="handler">Object to handle the event</param> /// <param name="handler">Object to handle the event</param>
IDisposable Subscribe<TEvent>(ILocalEventHandler<TEvent> handler) IDisposable Subscribe<TEvent>(ILocalEventHandler<TEvent> handler)
where TEvent : class; where TEvent : class;
/// <summary>
/// Gets the list of event handler factories for the given event type.
/// </summary>
/// <param name="eventType">Event type</param>
/// <returns></returns>
List<EventTypeWithEventHandlerFactories> GetEventHandlerFactories(Type eventType);
} }

30
framework/src/Volo.Abp.EventBus.Dapr/Volo/Abp/EventBus/Dapr/DaprDistributedEventBus.cs

@ -153,7 +153,13 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend
public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig)
{ {
await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, GetEventType(outgoingEvent.EventName)), outgoingEvent.Id, outgoingEvent.GetCorrelationId()); var eventType = GetEventType(outgoingEvent.EventName);
if (eventType == null)
{
return;
}
await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, eventType), outgoingEvent.Id, outgoingEvent.GetCorrelationId());
using (CorrelationIdProvider.Change(outgoingEvent.GetCorrelationId())) using (CorrelationIdProvider.Change(outgoingEvent.GetCorrelationId()))
{ {
@ -168,21 +174,9 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend
public async override Task PublishManyFromOutboxAsync(IEnumerable<OutgoingEventInfo> outgoingEvents, OutboxConfig outboxConfig) public async override Task PublishManyFromOutboxAsync(IEnumerable<OutgoingEventInfo> outgoingEvents, OutboxConfig outboxConfig)
{ {
var outgoingEventArray = outgoingEvents.ToArray(); foreach (var outgoingEvent in outgoingEvents)
foreach (var outgoingEvent in outgoingEventArray)
{ {
await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, GetEventType(outgoingEvent.EventName)), outgoingEvent.Id, outgoingEvent.GetCorrelationId()); await PublishFromOutboxAsync(outgoingEvent, outboxConfig);
using (CorrelationIdProvider.Change(outgoingEvent.GetCorrelationId()))
{
await TriggerDistributedEventSentAsync(new DistributedEventSent()
{
Source = DistributedEventSource.Outbox,
EventName = outgoingEvent.EventName,
EventData = outgoingEvent.EventData
});
}
} }
} }
@ -201,7 +195,7 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend
public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig)
{ {
var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); var eventType = GetEventType(incomingEvent.EventName);
if (eventType == null) if (eventType == null)
{ {
return; return;
@ -243,9 +237,9 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend
); );
} }
public Type GetEventType(string eventName) public Type? GetEventType(string eventName)
{ {
return EventTypes.GetOrDefault(eventName)!; return EventTypes.GetOrDefault(eventName);
} }
protected virtual async Task PublishToDaprAsync(Type eventType, object eventData, Guid? messageId = null, string? correlationId = null) protected virtual async Task PublishToDaprAsync(Type eventType, object eventData, Guid? messageId = null, string? correlationId = null)

7
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs

@ -250,7 +250,12 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen
OutgoingEventInfo outgoingEvent, OutgoingEventInfo outgoingEvent,
OutboxConfig outboxConfig) OutboxConfig outboxConfig)
{ {
var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName)!; var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName);
if (eventType == null)
{
return;
}
var eventData = Serializer.Deserialize(outgoingEvent.EventData, eventType); var eventData = Serializer.Deserialize(outgoingEvent.EventData, eventType);
var headers = new Dictionary<string, string>(); var headers = new Dictionary<string, string>();

17
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs

@ -62,7 +62,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
return PublishAsync(typeof(TEvent), eventData, onUnitOfWorkComplete, useOutbox); return PublishAsync(typeof(TEvent), eventData, onUnitOfWorkComplete, useOutbox);
} }
public async Task PublishAsync( public virtual async Task PublishAsync(
Type eventType, Type eventType,
object eventData, object eventData,
bool onUnitOfWorkComplete = true, bool onUnitOfWorkComplete = true,
@ -117,6 +117,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
return false; return false;
} }
var addedToOutbox = false;
foreach (var outboxConfig in AbpDistributedEventBusOptions.Outboxes.Values.OrderBy(x => x.Selector is null)) foreach (var outboxConfig in AbpDistributedEventBusOptions.Outboxes.Values.OrderBy(x => x.Selector is null))
{ {
if (outboxConfig.Selector == null || outboxConfig.Selector(eventType)) if (outboxConfig.Selector == null || outboxConfig.Selector(eventType))
@ -140,11 +142,11 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
} }
await eventOutbox.EnqueueAsync(outgoingEventInfo); await eventOutbox.EnqueueAsync(outgoingEventInfo);
return true; addedToOutbox = true;
} }
} }
return false; return addedToOutbox;
} }
protected virtual Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData) protected virtual Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData)
@ -164,6 +166,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
return false; return false;
} }
var addToInbox = false;
using (var scope = ServiceScopeFactory.CreateScope()) using (var scope = ServiceScopeFactory.CreateScope())
{ {
foreach (var inboxConfig in AbpDistributedEventBusOptions.Inboxes.Values.OrderBy(x => x.EventSelector is null)) foreach (var inboxConfig in AbpDistributedEventBusOptions.Inboxes.Values.OrderBy(x => x.EventSelector is null))
@ -190,11 +194,12 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
); );
incomingEventInfo.SetCorrelationId(correlationId!); incomingEventInfo.SetCorrelationId(correlationId!);
await eventInbox.EnqueueAsync(incomingEventInfo); await eventInbox.EnqueueAsync(incomingEventInfo);
addToInbox = true;
} }
} }
} }
return true; return addToInbox;
} }
protected abstract byte[] Serialize(object eventData); protected abstract byte[] Serialize(object eventData);
@ -227,7 +232,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
{ {
try try
{ {
await LocalEventBus.PublishAsync(distributedEvent); await LocalEventBus.PublishAsync(distributedEvent, onUnitOfWorkComplete: false);
} }
catch (Exception) catch (Exception)
{ {
@ -239,7 +244,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
{ {
try try
{ {
await LocalEventBus.PublishAsync(distributedEvent); await LocalEventBus.PublishAsync(distributedEvent, false);
} }
catch (Exception) catch (Exception)
{ {

224
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs

@ -1,33 +1,53 @@
using System; using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
using System.Reflection; using System.Reflection;
using System.Text;
using System.Text.Json;
using System.Text.Unicode;
using System.Threading.Tasks; using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
using Volo.Abp.Collections; using Volo.Abp.Collections;
using Volo.Abp.DependencyInjection; using Volo.Abp.DependencyInjection;
using Volo.Abp.EventBus.Local; using Volo.Abp.EventBus.Local;
using Volo.Abp.Guids;
using Volo.Abp.MultiTenancy;
using Volo.Abp.Timing;
using Volo.Abp.Tracing;
using Volo.Abp.Uow;
namespace Volo.Abp.EventBus.Distributed; namespace Volo.Abp.EventBus.Distributed;
[Dependency(TryRegister = true)] [Dependency(TryRegister = true)]
[ExposeServices(typeof(IDistributedEventBus), typeof(LocalDistributedEventBus))] [ExposeServices(typeof(IDistributedEventBus), typeof(LocalDistributedEventBus))]
public class LocalDistributedEventBus : IDistributedEventBus, ISingletonDependency public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDependency
{ {
private readonly ILocalEventBus _localEventBus; protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected IServiceScopeFactory ServiceScopeFactory { get; }
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; }
public LocalDistributedEventBus( public LocalDistributedEventBus(
ILocalEventBus localEventBus,
IServiceScopeFactory serviceScopeFactory, IServiceScopeFactory serviceScopeFactory,
IOptions<AbpDistributedEventBusOptions> distributedEventBusOptions) ICurrentTenant currentTenant,
IUnitOfWorkManager unitOfWorkManager,
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions,
IGuidGenerator guidGenerator,
IClock clock,
IEventHandlerInvoker eventHandlerInvoker,
ILocalEventBus localEventBus,
ICorrelationIdProvider correlationIdProvider)
: base(serviceScopeFactory,
currentTenant,
unitOfWorkManager,
abpDistributedEventBusOptions,
guidGenerator,
clock,
eventHandlerInvoker,
localEventBus,
correlationIdProvider)
{ {
_localEventBus = localEventBus; EventTypes = new ConcurrentDictionary<string, Type>();
ServiceScopeFactory = serviceScopeFactory; Subscribe(abpDistributedEventBusOptions.Value.Handlers);
AbpDistributedEventBusOptions = distributedEventBusOptions.Value;
Subscribe(distributedEventBusOptions.Value.Handlers);
} }
public virtual void Subscribe(ITypeList<IEventHandler> handlers) public virtual void Subscribe(ITypeList<IEventHandler> handlers)
@ -51,122 +71,156 @@ public class LocalDistributedEventBus : IDistributedEventBus, ISingletonDependen
} }
} }
/// <inheritdoc/> public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory)
public virtual IDisposable Subscribe<TEvent>(IDistributedEventHandler<TEvent> handler) where TEvent : class
{ {
return Subscribe(typeof(TEvent), handler); var eventName = EventNameAttribute.GetNameOrDefault(eventType);
EventTypes.GetOrAdd(eventName, eventType);
return LocalEventBus.Subscribe(eventType, factory);
} }
public IDisposable Subscribe<TEvent>(Func<TEvent, Task> action) where TEvent : class public override void Unsubscribe<TEvent>(Func<TEvent, Task> action)
{ {
return _localEventBus.Subscribe(action); LocalEventBus.Unsubscribe(action);
} }
public IDisposable Subscribe<TEvent>(ILocalEventHandler<TEvent> handler) where TEvent : class public override void Unsubscribe(Type eventType, IEventHandler handler)
{ {
return _localEventBus.Subscribe(handler); LocalEventBus.Unsubscribe(eventType, handler);
} }
public IDisposable Subscribe<TEvent, THandler>() where TEvent : class where THandler : IEventHandler, new() public override void Unsubscribe(Type eventType, IEventHandlerFactory factory)
{ {
return _localEventBus.Subscribe<TEvent, THandler>(); LocalEventBus.Unsubscribe(eventType, factory);
} }
public IDisposable Subscribe(Type eventType, IEventHandler handler) public override void UnsubscribeAll(Type eventType)
{ {
return _localEventBus.Subscribe(eventType, handler); LocalEventBus.UnsubscribeAll(eventType);
} }
public IDisposable Subscribe<TEvent>(IEventHandlerFactory factory) where TEvent : class public async override Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true, bool useOutbox = true)
{ {
return _localEventBus.Subscribe<TEvent>(factory); if (onUnitOfWorkComplete && UnitOfWorkManager.Current != null)
} {
AddToUnitOfWork(
UnitOfWorkManager.Current,
new UnitOfWorkEventRecord(eventType, eventData, EventOrderGenerator.GetNext(), useOutbox)
);
return;
}
public IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) if (useOutbox)
{ {
return _localEventBus.Subscribe(eventType, factory); if (await AddToOutboxAsync(eventType, eventData))
} {
return;
}
}
public void Unsubscribe<TEvent>(Func<TEvent, Task> action) where TEvent : class await TriggerDistributedEventSentAsync(new DistributedEventSent()
{ {
_localEventBus.Unsubscribe(action); Source = DistributedEventSource.Direct,
} EventName = EventNameAttribute.GetNameOrDefault(eventType),
EventData = eventData
});
public void Unsubscribe<TEvent>(ILocalEventHandler<TEvent> handler) where TEvent : class await TriggerDistributedEventReceivedAsync(new DistributedEventReceived
{ {
_localEventBus.Unsubscribe(handler); Source = DistributedEventSource.Direct,
} EventName = EventNameAttribute.GetNameOrDefault(eventType),
EventData = eventData
});
public void Unsubscribe(Type eventType, IEventHandler handler) await PublishToEventBusAsync(eventType, eventData);
{
_localEventBus.Unsubscribe(eventType, handler);
} }
public void Unsubscribe<TEvent>(IEventHandlerFactory factory) where TEvent : class protected async override Task PublishToEventBusAsync(Type eventType, object eventData)
{ {
_localEventBus.Unsubscribe<TEvent>(factory); if (await AddToInboxAsync(Guid.NewGuid().ToString(), EventNameAttribute.GetNameOrDefault(eventType), eventType, eventData, null))
} {
return;
}
public void Unsubscribe(Type eventType, IEventHandlerFactory factory) await LocalEventBus.PublishAsync(eventType, eventData, false);
{
_localEventBus.Unsubscribe(eventType, factory);
} }
public void UnsubscribeAll<TEvent>() where TEvent : class protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord)
{ {
_localEventBus.UnsubscribeAll<TEvent>(); unitOfWork.AddOrReplaceDistributedEvent(eventRecord);
} }
public void UnsubscribeAll(Type eventType) public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig)
{ {
_localEventBus.UnsubscribeAll(eventType); await TriggerDistributedEventSentAsync(new DistributedEventSent()
{
Source = DistributedEventSource.Outbox,
EventName = outgoingEvent.EventName,
EventData = outgoingEvent.EventData
});
await TriggerDistributedEventReceivedAsync(new DistributedEventReceived
{
Source = DistributedEventSource.Direct,
EventName = outgoingEvent.EventName,
EventData = outgoingEvent.EventData
});
var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName);
if (eventType == null)
{
return;
}
var eventData = JsonSerializer.Deserialize(Encoding.UTF8.GetString(outgoingEvent.EventData), eventType)!;
if (await AddToInboxAsync(Guid.NewGuid().ToString(), outgoingEvent.EventName, eventType, eventData, null))
{
return;
}
await LocalEventBus.PublishAsync(eventType, eventData, false);
} }
public async Task PublishAsync<TEvent>(TEvent eventData, bool onUnitOfWorkComplete = true) public async override Task PublishManyFromOutboxAsync(IEnumerable<OutgoingEventInfo> outgoingEvents, OutboxConfig outboxConfig)
where TEvent : class
{ {
await PublishDistributedEventSentReceivedAsync(typeof(TEvent), eventData, onUnitOfWorkComplete); foreach (var outgoingEvent in outgoingEvents)
await _localEventBus.PublishAsync(eventData, onUnitOfWorkComplete); {
await PublishFromOutboxAsync(outgoingEvent, outboxConfig);
}
} }
public async Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true) public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig)
{ {
await PublishDistributedEventSentReceivedAsync(eventType, eventData, onUnitOfWorkComplete); var eventType = EventTypes.GetOrDefault(incomingEvent.EventName);
await _localEventBus.PublishAsync(eventType, eventData, onUnitOfWorkComplete); if (eventType == null)
{
return;
}
var eventData = JsonSerializer.Deserialize(incomingEvent.EventData, eventType);
var exceptions = new List<Exception>();
using (CorrelationIdProvider.Change(incomingEvent.GetCorrelationId()))
{
await TriggerHandlersFromInboxAsync(eventType, eventData!, exceptions, inboxConfig);
}
if (exceptions.Any())
{
ThrowOriginalExceptions(eventType, exceptions);
}
} }
public async Task PublishAsync<TEvent>(TEvent eventData, bool onUnitOfWorkComplete = true, bool useOutbox = true) where TEvent : class protected override byte[] Serialize(object eventData)
{ {
await PublishDistributedEventSentReceivedAsync(typeof(TEvent), eventData, onUnitOfWorkComplete); return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(eventData));
await _localEventBus.PublishAsync(eventData, onUnitOfWorkComplete);
} }
public async Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true, bool useOutbox = true) protected override Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData)
{ {
await PublishDistributedEventSentReceivedAsync(eventType, eventData, onUnitOfWorkComplete); EventTypes.GetOrAdd(eventName, eventType);
await _localEventBus.PublishAsync(eventType, eventData, onUnitOfWorkComplete); return base.OnAddToOutboxAsync(eventName, eventType, eventData);
} }
private async Task PublishDistributedEventSentReceivedAsync(Type eventType, object eventData, bool onUnitOfWorkComplete) protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType)
{ {
if (eventType != typeof(DistributedEventSent)) return LocalEventBus.GetEventHandlerFactories(eventType);
{
await _localEventBus.PublishAsync(new DistributedEventSent
{
Source = DistributedEventSource.Direct,
EventName = EventNameAttribute.GetNameOrDefault(eventType),
EventData = eventData
}, onUnitOfWorkComplete);
}
if (eventType != typeof(DistributedEventReceived))
{
await _localEventBus.PublishAsync(new DistributedEventReceived
{
Source = DistributedEventSource.Direct,
EventName = EventNameAttribute.GetNameOrDefault(eventType),
EventData = eventData
}, onUnitOfWorkComplete);
}
} }
} }

13
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs

@ -242,19 +242,6 @@ public abstract class EventBusBase : IEventBus
}; };
} }
protected class EventTypeWithEventHandlerFactories
{
public Type EventType { get; }
public List<IEventHandlerFactory> EventHandlerFactories { get; }
public EventTypeWithEventHandlerFactories(Type eventType, List<IEventHandlerFactory> eventHandlerFactories)
{
EventType = eventType;
EventHandlerFactories = eventHandlerFactories;
}
}
// Reference from // Reference from
// https://blogs.msdn.microsoft.com/benwilli/2017/02/09/an-alternative-to-configureawaitfalse-everywhere/ // https://blogs.msdn.microsoft.com/benwilli/2017/02/09/an-alternative-to-configureawaitfalse-everywhere/
protected struct SynchronizationContextRemover : INotifyCompletion protected struct SynchronizationContextRemover : INotifyCompletion

5
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs

@ -136,6 +136,11 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency
await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData); await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData);
} }
public virtual List<EventTypeWithEventHandlerFactories> GetEventHandlerFactories(Type eventType)
{
return GetHandlerFactories(eventType).ToList();
}
protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType) protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType)
{ {
var handlerFactoryList = new List<Tuple<IEventHandlerFactory, Type, int>>(); var handlerFactoryList = new List<Tuple<IEventHandlerFactory, Type, int>>();

6
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs

@ -1,4 +1,5 @@
using System; using System;
using System.Collections.Generic;
using System.Threading.Tasks; using System.Threading.Tasks;
namespace Volo.Abp.EventBus.Local; namespace Volo.Abp.EventBus.Local;
@ -22,6 +23,11 @@ public sealed class NullLocalEventBus : ILocalEventBus
return NullDisposable.Instance; return NullDisposable.Instance;
} }
public List<EventTypeWithEventHandlerFactories> GetEventHandlerFactories(Type eventType)
{
return new List<EventTypeWithEventHandlerFactories>();
}
public IDisposable Subscribe<TEvent, THandler>() where TEvent : class where THandler : IEventHandler, new() public IDisposable Subscribe<TEvent, THandler>() where TEvent : class where THandler : IEventHandler, new()
{ {
return NullDisposable.Instance; return NullDisposable.Instance;

Loading…
Cancel
Save