From 2d84e52fcdb2175c98d2561b36d7cd7b0996970e Mon Sep 17 00:00:00 2001 From: Nuno Vieira Date: Tue, 21 Jan 2025 12:33:34 +0000 Subject: [PATCH 1/5] Changes to make LocalDistributedEventBus able to use outbox/inbox patterns. --- .../Distributed/LocalDistributedEventBus.cs | 197 +++++++++--------- 1 file changed, 102 insertions(+), 95 deletions(-) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs index d653c63bb7..79fb1add78 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs @@ -1,35 +1,87 @@ using System; +using System.Collections.Concurrent; +using System.Collections.Generic; using System.Reflection; +using System.Text.Json; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using Volo.Abp.Collections; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Local; +using System.Linq; +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; [Dependency(TryRegister = true)] [ExposeServices(typeof(IDistributedEventBus), typeof(LocalDistributedEventBus))] -public class LocalDistributedEventBus : IDistributedEventBus, ISingletonDependency +public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDependency { - private readonly ILocalEventBus _localEventBus; + protected ConcurrentDictionary> HandlerFactories { get; } - protected IServiceScopeFactory ServiceScopeFactory { get; } + protected ConcurrentDictionary EventTypes { get; } - protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } + public LocalDistributedEventBus(IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant, Volo.Abp.Uow.IUnitOfWorkManager unitOfWorkManager, IOptions abpDistributedEventBusOptions, + IGuidGenerator guidGenerator, IClock clock, IEventHandlerInvoker eventHandlerInvoker, ILocalEventBus localEventBus, ICorrelationIdProvider correlationIdProvider) + : base(serviceScopeFactory, currentTenant, unitOfWorkManager, abpDistributedEventBusOptions, guidGenerator, clock, eventHandlerInvoker, localEventBus, correlationIdProvider) + { + HandlerFactories = new ConcurrentDictionary>(); + EventTypes = new ConcurrentDictionary(); + Subscribe(abpDistributedEventBusOptions.Value.Handlers); + } - public LocalDistributedEventBus( - ILocalEventBus localEventBus, - IServiceScopeFactory serviceScopeFactory, - IOptions distributedEventBusOptions) + protected override Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData) { - _localEventBus = localEventBus; - ServiceScopeFactory = serviceScopeFactory; - AbpDistributedEventBusOptions = distributedEventBusOptions.Value; - Subscribe(distributedEventBusOptions.Value.Handlers); + EventTypes.GetOrAdd(eventName, eventType); + return base.OnAddToOutboxAsync(eventName, eventType, eventData); } + + public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) + { + var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); + if (eventType == null) + { + return; + } + + var eventData = JsonSerializer.Deserialize(incomingEvent.EventData, eventType); + + if (eventData == null) + { + return; + } + await LocalEventBus.PublishAsync(eventType, eventData); + + } + + public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) + { + var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName); + if (eventType == null) + return; + var eventData = JsonSerializer.Deserialize(outgoingEvent.EventData, eventType); + if (eventData == null) + { + return; + } + await LocalEventBus.PublishAsync(eventType, eventData); + } + + public async override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) + { + foreach (var outgoingEvent in outgoingEvents) + { + await PublishFromOutboxAsync(outgoingEvent, outboxConfig); + } + } + + public virtual void Subscribe(ITypeList handlers) { foreach (var handler in handlers) @@ -51,122 +103,77 @@ public class LocalDistributedEventBus : IDistributedEventBus, ISingletonDependen } } - /// - public virtual IDisposable Subscribe(IDistributedEventHandler handler) where TEvent : class - { - return Subscribe(typeof(TEvent), handler); - } - public IDisposable Subscribe(Func action) where TEvent : class + public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) { - return _localEventBus.Subscribe(action); + return LocalEventBus.Subscribe(eventType, factory); } - public IDisposable Subscribe(ILocalEventHandler handler) where TEvent : class + public override void Unsubscribe(Func action) { - return _localEventBus.Subscribe(handler); + LocalEventBus.Unsubscribe(action); } - public IDisposable Subscribe() where TEvent : class where THandler : IEventHandler, new() + public override void Unsubscribe(Type eventType, IEventHandler handler) { - return _localEventBus.Subscribe(); + LocalEventBus.Unsubscribe(eventType, handler); } - public IDisposable Subscribe(Type eventType, IEventHandler handler) + public override void Unsubscribe(Type eventType, IEventHandlerFactory factory) { - return _localEventBus.Subscribe(eventType, handler); + LocalEventBus.Unsubscribe(eventType, factory); } - public IDisposable Subscribe(IEventHandlerFactory factory) where TEvent : class + public override void UnsubscribeAll(Type eventType) { - return _localEventBus.Subscribe(factory); + LocalEventBus.UnsubscribeAll(eventType); } - public IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) + protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord) { - return _localEventBus.Subscribe(eventType, factory); + unitOfWork.AddOrReplaceDistributedEvent(eventRecord); } - public void Unsubscribe(Func action) where TEvent : class + protected override IEnumerable GetHandlerFactories(Type eventType) { - _localEventBus.Unsubscribe(action); - } + var handlerFactoryList = new List(); - public void Unsubscribe(ILocalEventHandler handler) where TEvent : class - { - _localEventBus.Unsubscribe(handler); - } - - public void Unsubscribe(Type eventType, IEventHandler handler) - { - _localEventBus.Unsubscribe(eventType, handler); - } - - public void Unsubscribe(IEventHandlerFactory factory) where TEvent : class - { - _localEventBus.Unsubscribe(factory); - } + foreach (var handlerFactory in HandlerFactories.Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key)) + ) + { + handlerFactoryList.Add( + new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value)); + } - public void Unsubscribe(Type eventType, IEventHandlerFactory factory) - { - _localEventBus.Unsubscribe(eventType, factory); + return handlerFactoryList.ToArray(); } - public void UnsubscribeAll() where TEvent : class + private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) { - _localEventBus.UnsubscribeAll(); - } + //Should trigger same type + if (handlerEventType == targetEventType) + { + return true; + } - public void UnsubscribeAll(Type eventType) - { - _localEventBus.UnsubscribeAll(eventType); - } + //Should trigger for inherited types + if (handlerEventType.IsAssignableFrom(targetEventType)) + { + return true; + } - public async Task PublishAsync(TEvent eventData, bool onUnitOfWorkComplete = true) - where TEvent : class - { - await PublishDistributedEventSentReceivedAsync(typeof(TEvent), eventData, onUnitOfWorkComplete); - await _localEventBus.PublishAsync(eventData, onUnitOfWorkComplete); + return false; } - public async Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true) - { - await PublishDistributedEventSentReceivedAsync(eventType, eventData, onUnitOfWorkComplete); - await _localEventBus.PublishAsync(eventType, eventData, onUnitOfWorkComplete); - } - public async Task PublishAsync(TEvent eventData, bool onUnitOfWorkComplete = true, bool useOutbox = true) where TEvent : class + protected async override Task PublishToEventBusAsync(Type eventType, object eventData) { - await PublishDistributedEventSentReceivedAsync(typeof(TEvent), eventData, onUnitOfWorkComplete); - await _localEventBus.PublishAsync(eventData, onUnitOfWorkComplete); - } + await LocalEventBus.PublishAsync(eventType, eventData); - public async Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true, bool useOutbox = true) - { - await PublishDistributedEventSentReceivedAsync(eventType, eventData, onUnitOfWorkComplete); - await _localEventBus.PublishAsync(eventType, eventData, onUnitOfWorkComplete); } - private async Task PublishDistributedEventSentReceivedAsync(Type eventType, object eventData, bool onUnitOfWorkComplete) + protected override byte[] Serialize(object eventData) { - if (eventType != typeof(DistributedEventSent)) - { - 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); - } + return JsonSerializer.SerializeToUtf8Bytes(eventData); } -} \ No newline at end of file +} From fbe3fe10bfa534eff15aadab7ea8e4926bb1e5a0 Mon Sep 17 00:00:00 2001 From: maliming Date: Wed, 29 Jan 2025 13:51:58 +0800 Subject: [PATCH 2/5] Refactor `LocalDistributedEventBus`. --- .../EventTypeWithEventHandlerFactories.cs | 17 ++ .../Volo/Abp/EventBus/Local/ILocalEventBus.cs | 10 +- .../Distributed/DistributedEventBusBase.cs | 6 +- .../Distributed/LocalDistributedEventBus.cs | 181 ++++++++++-------- .../Volo/Abp/EventBus/EventBusBase.cs | 13 -- .../Volo/Abp/EventBus/Local/LocalEventBus.cs | 5 + .../Abp/EventBus/Local/NullLocalEventBus.cs | 6 + 7 files changed, 145 insertions(+), 93 deletions(-) create mode 100644 framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/EventTypeWithEventHandlerFactories.cs diff --git a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/EventTypeWithEventHandlerFactories.cs b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/EventTypeWithEventHandlerFactories.cs new file mode 100644 index 0000000000..836d5cb486 --- /dev/null +++ b/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 EventHandlerFactories { get; } + + public EventTypeWithEventHandlerFactories(Type eventType, List eventHandlerFactories) + { + EventType = eventType; + EventHandlerFactories = eventHandlerFactories; + } +} diff --git a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs index 654895911d..e691b6c58c 100644 --- a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Generic; namespace Volo.Abp.EventBus.Local; @@ -8,11 +9,18 @@ namespace Volo.Abp.EventBus.Local; public interface ILocalEventBus : IEventBus { /// - /// Registers to an event. + /// Registers to an event. /// Same (given) instance of the handler is used for all event occurrences. /// /// Event type /// Object to handle the event IDisposable Subscribe(ILocalEventHandler handler) where TEvent : class; + + /// + /// Gets the list of event handler factories for the given event type. + /// + /// Event type + /// + List GetEventHandlerFactories(Type eventType); } diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs index ad69ed124b..d6b4ea51bd 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs +++ b/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); } - public async Task PublishAsync( + public virtual async Task PublishAsync( Type eventType, object eventData, bool onUnitOfWorkComplete = true, @@ -227,7 +227,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB { try { - await LocalEventBus.PublishAsync(distributedEvent); + await LocalEventBus.PublishAsync(distributedEvent, onUnitOfWorkComplete: false); } catch (Exception) { @@ -239,7 +239,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB { try { - await LocalEventBus.PublishAsync(distributedEvent); + await LocalEventBus.PublishAsync(distributedEvent, false); } catch (Exception) { diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs index 79fb1add78..430c933bac 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs @@ -1,7 +1,9 @@ using System; using System.Collections.Concurrent; using System.Collections.Generic; +using System.Linq; using System.Reflection; +using System.Text; using System.Text.Json; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; @@ -9,7 +11,6 @@ using Microsoft.Extensions.Options; using Volo.Abp.Collections; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Local; -using System.Linq; using Volo.Abp.Guids; using Volo.Abp.MultiTenancy; using Volo.Abp.Timing; @@ -22,66 +23,32 @@ namespace Volo.Abp.EventBus.Distributed; [ExposeServices(typeof(IDistributedEventBus), typeof(LocalDistributedEventBus))] public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDependency { - protected ConcurrentDictionary> HandlerFactories { get; } - protected ConcurrentDictionary EventTypes { get; } - public LocalDistributedEventBus(IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant, Volo.Abp.Uow.IUnitOfWorkManager unitOfWorkManager, IOptions abpDistributedEventBusOptions, - IGuidGenerator guidGenerator, IClock clock, IEventHandlerInvoker eventHandlerInvoker, ILocalEventBus localEventBus, ICorrelationIdProvider correlationIdProvider) - : base(serviceScopeFactory, currentTenant, unitOfWorkManager, abpDistributedEventBusOptions, guidGenerator, clock, eventHandlerInvoker, localEventBus, correlationIdProvider) + public LocalDistributedEventBus( + IServiceScopeFactory serviceScopeFactory, + ICurrentTenant currentTenant, + IUnitOfWorkManager unitOfWorkManager, + IOptions abpDistributedEventBusOptions, + IGuidGenerator guidGenerator, + IClock clock, + IEventHandlerInvoker eventHandlerInvoker, + ILocalEventBus localEventBus, + ICorrelationIdProvider correlationIdProvider) + : base(serviceScopeFactory, + currentTenant, + unitOfWorkManager, + abpDistributedEventBusOptions, + guidGenerator, + clock, + eventHandlerInvoker, + localEventBus, + correlationIdProvider) { - HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); Subscribe(abpDistributedEventBusOptions.Value.Handlers); } - protected override Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData) - { - EventTypes.GetOrAdd(eventName, eventType); - return base.OnAddToOutboxAsync(eventName, eventType, eventData); - } - - - public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) - { - var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); - if (eventType == null) - { - return; - } - - var eventData = JsonSerializer.Deserialize(incomingEvent.EventData, eventType); - - if (eventData == null) - { - return; - } - await LocalEventBus.PublishAsync(eventType, eventData); - - } - - public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) - { - var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName); - if (eventType == null) - return; - var eventData = JsonSerializer.Deserialize(outgoingEvent.EventData, eventType); - if (eventData == null) - { - return; - } - await LocalEventBus.PublishAsync(eventType, eventData); - } - - public async override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) - { - foreach (var outgoingEvent in outgoingEvents) - { - await PublishFromOutboxAsync(outgoingEvent, outboxConfig); - } - } - - public virtual void Subscribe(ITypeList handlers) { foreach (var handler in handlers) @@ -103,7 +70,6 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen } } - public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) { return LocalEventBus.Subscribe(eventType, factory); @@ -129,51 +95,114 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen LocalEventBus.UnsubscribeAll(eventType); } + public async override Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true, bool useOutbox = true) + { + if (onUnitOfWorkComplete && UnitOfWorkManager.Current != null) + { + AddToUnitOfWork( + UnitOfWorkManager.Current, + new UnitOfWorkEventRecord(eventType, eventData, EventOrderGenerator.GetNext(), useOutbox) + ); + return; + } + + if (useOutbox) + { + if (await AddToOutboxAsync(eventType, eventData)) + { + return; + } + } + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Direct, + EventName = EventNameAttribute.GetNameOrDefault(eventType), + EventData = eventData + }); + + await TriggerDistributedEventReceivedAsync(new DistributedEventReceived + { + Source = DistributedEventSource.Direct, + EventName = EventNameAttribute.GetNameOrDefault(eventType), + EventData = eventData + }); + + await PublishToEventBusAsync(eventType, eventData); + } + + protected async override Task PublishToEventBusAsync(Type eventType, object eventData) + { + await LocalEventBus.PublishAsync(eventType, eventData, false); + } + protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord) { unitOfWork.AddOrReplaceDistributedEvent(eventRecord); } - protected override IEnumerable GetHandlerFactories(Type eventType) + public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) { - var handlerFactoryList = new List(); + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); - foreach (var handlerFactory in HandlerFactories.Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key)) - ) + await TriggerDistributedEventReceivedAsync(new DistributedEventReceived { - handlerFactoryList.Add( - new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value)); - } + Source = DistributedEventSource.Direct, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); - return handlerFactoryList.ToArray(); + var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName)!; + var eventData = JsonSerializer.Deserialize(Encoding.UTF8.GetString(outgoingEvent.EventData), eventType)!; + await LocalEventBus.PublishAsync(eventType, eventData, false); } - private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) + public async override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) { - //Should trigger same type - if (handlerEventType == targetEventType) + foreach (var outgoingEvent in outgoingEvents) { - return true; + await PublishFromOutboxAsync(outgoingEvent, outboxConfig); } + } - //Should trigger for inherited types - if (handlerEventType.IsAssignableFrom(targetEventType)) + public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) + { + var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); + if (eventType == null) { - return true; + return; } - return false; + var eventData = JsonSerializer.Deserialize(incomingEvent.EventData, eventType); + var exceptions = new List(); + using (CorrelationIdProvider.Change(incomingEvent.GetCorrelationId())) + { + await TriggerHandlersFromInboxAsync(eventType, eventData!, exceptions, inboxConfig); + } + if (exceptions.Any()) + { + ThrowOriginalExceptions(eventType, exceptions); + } } - - protected async override Task PublishToEventBusAsync(Type eventType, object eventData) + protected override byte[] Serialize(object eventData) { - await LocalEventBus.PublishAsync(eventType, eventData); + return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(eventData)); + } + protected override Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData) + { + EventTypes.GetOrAdd(eventName, eventType); + return base.OnAddToOutboxAsync(eventName, eventType, eventData); } - protected override byte[] Serialize(object eventData) + protected override IEnumerable GetHandlerFactories(Type eventType) { - return JsonSerializer.SerializeToUtf8Bytes(eventData); + return LocalEventBus.GetEventHandlerFactories(eventType); } } diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs index 151a149281..231b4c7a0f 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/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 EventHandlerFactories { get; } - - public EventTypeWithEventHandlerFactories(Type eventType, List eventHandlerFactories) - { - EventType = eventType; - EventHandlerFactories = eventHandlerFactories; - } - } - // Reference from // https://blogs.msdn.microsoft.com/benwilli/2017/02/09/an-alternative-to-configureawaitfalse-everywhere/ protected struct SynchronizationContextRemover : INotifyCompletion diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs index 48f70ac8c1..7123ff340a 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs +++ b/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); } + public virtual List GetEventHandlerFactories(Type eventType) + { + return GetHandlerFactories(eventType).ToList(); + } + protected override IEnumerable GetHandlerFactories(Type eventType) { var handlerFactoryList = new List>(); diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs index 94c5f4ff83..3ffcd911ce 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Generic; using System.Threading.Tasks; namespace Volo.Abp.EventBus.Local; @@ -22,6 +23,11 @@ public sealed class NullLocalEventBus : ILocalEventBus return NullDisposable.Instance; } + public List GetEventHandlerFactories(Type eventType) + { + return new List(); + } + public IDisposable Subscribe() where TEvent : class where THandler : IEventHandler, new() { return NullDisposable.Instance; From 40c9f63371eaac195b1a79c29a1fd46aa44b113d Mon Sep 17 00:00:00 2001 From: maliming Date: Sun, 2 Feb 2025 18:06:59 +0800 Subject: [PATCH 3/5] Make sure `eventType` is not `null`. --- .../AbpAspNetCoreMvcDaprEventBusModule.cs | 16 +++++++--- .../EventBus/Dapr/DaprDistributedEventBus.cs | 30 ++++++++----------- .../Rebus/RebusDistributedEventBus.cs | 7 ++++- .../Distributed/LocalDistributedEventBus.cs | 9 +++++- 4 files changed, 38 insertions(+), 24 deletions(-) diff --git a/framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/AbpAspNetCoreMvcDaprEventBusModule.cs b/framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/AbpAspNetCoreMvcDaprEventBusModule.cs index a4a4f97f5d..fba9e12707 100644 --- a/framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/AbpAspNetCoreMvcDaprEventBusModule.cs +++ b/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)) { var daprEventData = daprSerializer.Deserialize(data, typeof(AbpDaprEventData)).As(); - var eventData = daprSerializer.Deserialize(daprEventData.JsonData, distributedEventBus.GetEventType(daprEventData.Topic)); - await distributedEventBus.TriggerHandlersAsync(distributedEventBus.GetEventType(daprEventData.Topic), eventData, daprEventData.MessageId, daprEventData.CorrelationId); + var eventType = distributedEventBus.GetEventType(daprEventData.Topic); + if (eventType != null) + { + var eventData = daprSerializer.Deserialize(daprEventData.JsonData, eventType); + await distributedEventBus.TriggerHandlersAsync(eventType, eventData, daprEventData.MessageId, daprEventData.CorrelationId); + } } else { - var eventData = daprSerializer.Deserialize(data, distributedEventBus.GetEventType(topic!)); - await distributedEventBus.TriggerHandlersAsync(distributedEventBus.GetEventType(topic!), eventData); + var eventType = distributedEventBus.GetEventType(topic); + if (eventType != null) + { + var eventData = daprSerializer.Deserialize(data, eventType); + await distributedEventBus.TriggerHandlersAsync(eventType, eventData); + } } httpContext.Response.StatusCode = 200; diff --git a/framework/src/Volo.Abp.EventBus.Dapr/Volo/Abp/EventBus/Dapr/DaprDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Dapr/Volo/Abp/EventBus/Dapr/DaprDistributedEventBus.cs index 2b0d2c7d0d..7c77340dda 100644 --- a/framework/src/Volo.Abp.EventBus.Dapr/Volo/Abp/EventBus/Dapr/DaprDistributedEventBus.cs +++ b/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) { - 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())) { @@ -168,21 +174,9 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend public async override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) { - var outgoingEventArray = outgoingEvents.ToArray(); - - foreach (var outgoingEvent in outgoingEventArray) + foreach (var outgoingEvent in outgoingEvents) { - await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, GetEventType(outgoingEvent.EventName)), outgoingEvent.Id, outgoingEvent.GetCorrelationId()); - - using (CorrelationIdProvider.Change(outgoingEvent.GetCorrelationId())) - { - await TriggerDistributedEventSentAsync(new DistributedEventSent() - { - Source = DistributedEventSource.Outbox, - EventName = outgoingEvent.EventName, - EventData = outgoingEvent.EventData - }); - } + await PublishFromOutboxAsync(outgoingEvent, outboxConfig); } } @@ -201,7 +195,7 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) { - var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); + var eventType = GetEventType(incomingEvent.EventName); if (eventType == null) { 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) diff --git a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs index 9e8398a495..7a3d79e7d8 100644 --- a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs @@ -250,7 +250,12 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen OutgoingEventInfo outgoingEvent, 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 headers = new Dictionary(); diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs index 430c933bac..cdf2849c4a 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs @@ -72,6 +72,8 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) { + var eventName = EventNameAttribute.GetNameOrDefault(eventType); + EventTypes.GetOrAdd(eventName, eventType); return LocalEventBus.Subscribe(eventType, factory); } @@ -157,7 +159,12 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen EventData = outgoingEvent.EventData }); - var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName)!; + var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName); + if (eventType == null) + { + return; + } + var eventData = JsonSerializer.Deserialize(Encoding.UTF8.GetString(outgoingEvent.EventData), eventType)!; await LocalEventBus.PublishAsync(eventType, eventData, false); } From d70255a25445b6e6745949f3444c2456a5d4c963 Mon Sep 17 00:00:00 2001 From: maliming Date: Sun, 2 Feb 2025 19:24:45 +0800 Subject: [PATCH 4/5] Add event to all Inbox/Outbox. --- .../EventBus/Distributed/DistributedEventBusBase.cs | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs index d6b4ea51bd..ac1e8c6565 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs @@ -117,6 +117,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB return false; } + var addedToOutbox = false; + foreach (var outboxConfig in AbpDistributedEventBusOptions.Outboxes.Values.OrderBy(x => x.Selector is null)) { if (outboxConfig.Selector == null || outboxConfig.Selector(eventType)) @@ -140,11 +142,11 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB } await eventOutbox.EnqueueAsync(outgoingEventInfo); - return true; + addedToOutbox = true; } } - return false; + return addedToOutbox; } protected virtual Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData) @@ -164,6 +166,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB return false; } + var addToInbox = false; + using (var scope = ServiceScopeFactory.CreateScope()) { 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!); await eventInbox.EnqueueAsync(incomingEventInfo); + addToInbox = true; } } } - return true; + return addToInbox; } protected abstract byte[] Serialize(object eventData); From 2a0f5e8741e6f15411e52de4c275f19353cf4387 Mon Sep 17 00:00:00 2001 From: maliming Date: Sun, 2 Feb 2025 20:23:14 +0800 Subject: [PATCH 5/5] Try to add message to inbox. --- .../EventBus/Distributed/LocalDistributedEventBus.cs | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs index cdf2849c4a..843fb4f8ea 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs @@ -5,6 +5,7 @@ using System.Linq; using System.Reflection; using System.Text; using System.Text.Json; +using System.Text.Unicode; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; @@ -135,6 +136,11 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen protected async override Task PublishToEventBusAsync(Type eventType, object eventData) { + if (await AddToInboxAsync(Guid.NewGuid().ToString(), EventNameAttribute.GetNameOrDefault(eventType), eventType, eventData, null)) + { + return; + } + await LocalEventBus.PublishAsync(eventType, eventData, false); } @@ -166,6 +172,11 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen } 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); }