From 3d3e024f8fa9220f939f9655edb53208805c35ff Mon Sep 17 00:00:00 2001 From: maliming Date: Thu, 20 Apr 2023 17:50:03 +0800 Subject: [PATCH] Inform when a new distributed message received or sent. Resolve #16321 --- .../AbpAspNetCoreMvcDaprEventsController.cs | 2 +- .../Distributed/DistributedEventReceived.cs | 10 +++ .../Distributed/DistributedEventSent.cs | 10 +++ .../Distributed/DistributedEventSource.cs | 10 +++ .../Azure/AzureDistributedEventBus.cs | 28 +++++++-- .../EventBus/Dapr/DaprDistributedEventBus.cs | 25 +++++++- .../Kafka/KafkaDistributedEventBus.cs | 43 ++++++++----- .../RabbitMq/RabbitMqDistributedEventBus.cs | 34 +++++++--- .../Rebus/RebusDistributedEventBus.cs | 31 +++++++--- .../Distributed/DistributedEventBusBase.cs | 62 ++++++++++++++++++- .../Distributed/LocalDistributedEventBus.cs | 28 ++++++++- .../Volo/Abp/EventBus/EventBusBase.cs | 8 ++- .../Volo/Abp/EventBus/Local/LocalEventBus.cs | 42 +++++++++++++ .../MongoDB/AbpMongoDbDateTimeSerializer.cs | 2 +- .../Distributed/DistributedEventHandles.cs | 22 +++++++ .../LocalDistributedEventBus_Test.cs | 44 ++++++++++++- 16 files changed, 353 insertions(+), 48 deletions(-) create mode 100644 framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventReceived.cs create mode 100644 framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSent.cs create mode 100644 framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSource.cs create mode 100644 framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/DistributedEventHandles.cs diff --git a/framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/Controllers/AbpAspNetCoreMvcDaprEventsController.cs b/framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/Controllers/AbpAspNetCoreMvcDaprEventsController.cs index 92cf41db90..d4d4585013 100644 --- a/framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/Controllers/AbpAspNetCoreMvcDaprEventsController.cs +++ b/framework/src/Volo.Abp.AspNetCore.Mvc.Dapr.EventBus/Volo/Abp/AspNetCore/Mvc/Dapr/EventBus/Controllers/AbpAspNetCoreMvcDaprEventsController.cs @@ -32,7 +32,7 @@ public class AbpAspNetCoreMvcDaprEventsController : AbpController var distributedEventBus = HttpContext.RequestServices.GetRequiredService(); var eventData = daprSerializer.Deserialize(data, distributedEventBus.GetEventType(topic)); - await distributedEventBus.TriggerHandlersAsync(distributedEventBus.GetEventType(topic), eventData); + await distributedEventBus.DaprTriggerHandlersDirectAsync(distributedEventBus.GetEventType(topic), eventData); return Ok(); } } diff --git a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventReceived.cs b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventReceived.cs new file mode 100644 index 0000000000..42036d85ce --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventReceived.cs @@ -0,0 +1,10 @@ +namespace Volo.Abp.EventBus.Distributed; + +public class DistributedEventReceived +{ + public DistributedEventSource Source { get; set; } + + public string EventName { get; set; } + + public object EventData { get; set; } +} diff --git a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSent.cs b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSent.cs new file mode 100644 index 0000000000..990068e0cc --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSent.cs @@ -0,0 +1,10 @@ +namespace Volo.Abp.EventBus.Distributed; + +public class DistributedEventSent +{ + public DistributedEventSource Source { get; set; } + + public string EventName { get; set; } + + public object EventData { get; set; } +} diff --git a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSource.cs b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSource.cs new file mode 100644 index 0000000000..a1cea8eef7 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Distributed/DistributedEventSource.cs @@ -0,0 +1,10 @@ +namespace Volo.Abp.EventBus.Distributed; + +public enum DistributedEventSource +{ + Direct, + + Inbox, + + Outbox +} diff --git a/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AzureDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AzureDistributedEventBus.cs index 78e26e7c32..0c45215c8d 100644 --- a/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AzureDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AzureDistributedEventBus.cs @@ -9,6 +9,7 @@ using Microsoft.Extensions.Options; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; using Volo.Abp.AzureServiceBus; +using Volo.Abp.EventBus.Local; using Volo.Abp.Guids; using Volo.Abp.MultiTenancy; using Volo.Abp.Threading; @@ -40,14 +41,16 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen IAzureServiceBusSerializer serializer, IAzureServiceBusMessageConsumerFactory messageConsumerFactory, IPublisherPool publisherPool, - IEventHandlerInvoker eventHandlerInvoker) + IEventHandlerInvoker eventHandlerInvoker, + ILocalEventBus localEventBus) : base(serviceScopeFactory, currentTenant, unitOfWorkManager, abpDistributedEventBusOptions, guidGenerator, clock, - eventHandlerInvoker) + eventHandlerInvoker, + localEventBus) { _options = abpAzureEventBusOptions.Value; _serializer = serializer; @@ -88,12 +91,19 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen var eventData = _serializer.Deserialize(message.Body.ToArray(), eventType); - await TriggerHandlersAsync(eventType, eventData); + await TriggerHandlersDirectAsync(eventType, eventData); } public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) { await PublishAsync(outgoingEvent.EventName, outgoingEvent.EventData, outgoingEvent.Id); + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } public async override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) @@ -123,6 +133,16 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen } await publisher.SendMessagesAsync(messageBatch); + + foreach (var outgoingEvent in outgoingEventArray) + { + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); + } } public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) @@ -135,7 +155,7 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen var eventData = _serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); - await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); + await TriggerHandlersFromInboxAsync(eventType, eventData, exceptions, inboxConfig); if (exceptions.Any()) { ThrowOriginalExceptions(eventType, exceptions); 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 a00127ac8c..f479370c8e 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 @@ -8,6 +8,7 @@ using Microsoft.Extensions.Options; using Volo.Abp.Dapr; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; +using Volo.Abp.EventBus.Local; using Volo.Abp.Guids; using Volo.Abp.MultiTenancy; using Volo.Abp.Threading; @@ -37,8 +38,9 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend IEventHandlerInvoker eventHandlerInvoker, IDaprSerializer serializer, IOptions daprEventBusOptions, - IAbpDaprClientFactory daprClientFactory) - : base(serviceScopeFactory, currentTenant, unitOfWorkManager, abpDistributedEventBusOptions, guidGenerator, clock, eventHandlerInvoker) + IAbpDaprClientFactory daprClientFactory, + ILocalEventBus localEventBus) + : base(serviceScopeFactory, currentTenant, unitOfWorkManager, abpDistributedEventBusOptions, guidGenerator, clock, eventHandlerInvoker, localEventBus) { Serializer = serializer; DaprEventBusOptions = daprEventBusOptions.Value; @@ -142,6 +144,12 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend public override async Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) { await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, GetEventType(outgoingEvent.EventName))); + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } public override async Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) @@ -151,9 +159,20 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend foreach (var outgoingEvent in outgoingEventArray) { await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, GetEventType(outgoingEvent.EventName))); + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } } + public virtual async Task DaprTriggerHandlersDirectAsync(Type eventType, object eventData) + { + await TriggerHandlersDirectAsync(eventType, eventData); + } + public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) { var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); @@ -164,7 +183,7 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); - await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); + await TriggerHandlersFromInboxAsync(eventType, eventData, exceptions, inboxConfig); if (exceptions.Any()) { ThrowOriginalExceptions(eventType, exceptions); diff --git a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs index aa79c81f50..f66f7cf60f 100644 --- a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs @@ -8,6 +8,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; +using Volo.Abp.EventBus.Local; using Volo.Abp.Guids; using Volo.Abp.Kafka; using Volo.Abp.MultiTenancy; @@ -40,7 +41,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen IProducerPool producerPool, IGuidGenerator guidGenerator, IClock clock, - IEventHandlerInvoker eventHandlerInvoker) + IEventHandlerInvoker eventHandlerInvoker, + ILocalEventBus localEventBus) : base( serviceScopeFactory, currentTenant, @@ -48,7 +50,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen abpDistributedEventBusOptions, guidGenerator, clock, - eventHandlerInvoker) + eventHandlerInvoker, + localEventBus) { AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; MessageConsumerFactory = messageConsumerFactory; @@ -88,7 +91,7 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen var eventData = Serializer.Deserialize(message.Value, eventType); - await TriggerHandlersAsync(eventType, eventData); + await TriggerHandlersDirectAsync(eventType, eventData); } public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) @@ -177,11 +180,11 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen unitOfWork.AddOrReplaceDistributedEvent(eventRecord); } - public override Task PublishFromOutboxAsync( + public override async Task PublishFromOutboxAsync( OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) { - return PublishAsync( + await PublishAsync( AbpKafkaEventBusOptions.TopicName, outgoingEvent.EventName, outgoingEvent.EventData, @@ -190,21 +193,28 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen { "messageId", System.Text.Encoding.UTF8.GetBytes(outgoingEvent.Id.ToString("N")) } } ); + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } - public override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) + public async override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) { var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); var outgoingEventArray = outgoingEvents.ToArray(); - + foreach (var outgoingEvent in outgoingEventArray) { var messageId = outgoingEvent.Id.ToString("N"); - var headers = new Headers + var headers = new Headers { { "messageId", System.Text.Encoding.UTF8.GetBytes(messageId)} }; - + producer.Produce( AbpKafkaEventBusOptions.TopicName, new Message @@ -213,9 +223,14 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen Value = outgoingEvent.EventData, Headers = headers }); + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } - - return Task.CompletedTask; } public async override Task ProcessFromInboxAsync( @@ -230,7 +245,7 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); - await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); + await TriggerHandlersFromInboxAsync(eventType, eventData, exceptions, inboxConfig); if (exceptions.Any()) { ThrowOriginalExceptions(eventType, exceptions); @@ -251,9 +266,9 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen } private Task> PublishAsync( - string topicName, + string topicName, string eventName, - byte[] body, + byte[] body, Headers headers) { var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs index 7e352ca530..e085cc2150 100644 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs @@ -9,6 +9,7 @@ using RabbitMQ.Client; using RabbitMQ.Client.Events; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; +using Volo.Abp.EventBus.Local; using Volo.Abp.Guids; using Volo.Abp.MultiTenancy; using Volo.Abp.RabbitMQ; @@ -47,7 +48,8 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe IUnitOfWorkManager unitOfWorkManager, IGuidGenerator guidGenerator, IClock clock, - IEventHandlerInvoker eventHandlerInvoker) + IEventHandlerInvoker eventHandlerInvoker, + ILocalEventBus localEventBus) : base( serviceScopeFactory, currentTenant, @@ -55,7 +57,8 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe distributedEventBusOptions, guidGenerator, clock, - eventHandlerInvoker) + eventHandlerInvoker, + localEventBus) { ConnectionPool = connectionPool; Serializer = serializer; @@ -107,7 +110,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe var eventData = Serializer.Deserialize(eventBytes, eventType); - await TriggerHandlersAsync(eventType, eventData); + await TriggerHandlersDirectAsync(eventType, eventData); } public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) @@ -193,11 +196,17 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe unitOfWork.AddOrReplaceDistributedEvent(eventRecord); } - public override Task PublishFromOutboxAsync( + public override async Task PublishFromOutboxAsync( OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) { - return PublishAsync(outgoingEvent.EventName, outgoingEvent.EventData, null, eventId: outgoingEvent.Id); + await PublishAsync(outgoingEvent.EventName, outgoingEvent.EventData, null, eventId: outgoingEvent.Id); + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } public async override Task PublishManyFromOutboxAsync( @@ -213,10 +222,17 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe { await PublishAsync( channel, - outgoingEvent.EventName, - outgoingEvent.EventData, + outgoingEvent.EventName, + outgoingEvent.EventData, properties: null, eventId: outgoingEvent.Id); + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } channel.WaitForConfirmsOrDie(); @@ -235,7 +251,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); - await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); + await TriggerHandlersFromInboxAsync(eventType, eventData, exceptions, inboxConfig); if (exceptions.Any()) { ThrowOriginalExceptions(eventType, exceptions); @@ -249,7 +265,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe public Task PublishAsync( Type eventType, - object eventData, + object eventData, IBasicProperties properties, Dictionary headersArguments = 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 bfe763d960..6540ef242a 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 @@ -11,6 +11,7 @@ using Rebus.Pipeline; using Rebus.Transport; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; +using Volo.Abp.EventBus.Local; using Volo.Abp.Guids; using Volo.Abp.MultiTenancy; using Volo.Abp.Threading; @@ -41,7 +42,8 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen IRebusSerializer serializer, IGuidGenerator guidGenerator, IClock clock, - IEventHandlerInvoker eventHandlerInvoker) : + IEventHandlerInvoker eventHandlerInvoker, + ILocalEventBus localEventBus) : base( serviceScopeFactory, currentTenant, @@ -49,7 +51,8 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen abpDistributedEventBusOptions, guidGenerator, clock, - eventHandlerInvoker) + eventHandlerInvoker, + localEventBus) { Rebus = rebus; Serializer = serializer; @@ -147,7 +150,7 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen return; } - await TriggerHandlersAsync(eventType, eventData); + await TriggerHandlersDirectAsync(eventType, eventData); } protected async override Task PublishToEventBusAsync(Type eventType, object eventData) @@ -225,14 +228,21 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen return false; } - public override Task PublishFromOutboxAsync( + public override async Task PublishFromOutboxAsync( OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) { var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName); var eventData = Serializer.Deserialize(outgoingEvent.EventData, eventType); - return PublishAsync(eventType, eventData, eventId: outgoingEvent.Id); + await PublishAsync(eventType, eventData, eventId: outgoingEvent.Id); + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } public async override Task PublishManyFromOutboxAsync(IEnumerable outgoingEvents, OutboxConfig outboxConfig) @@ -244,8 +254,15 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen foreach (var outgoingEvent in outgoingEventArray) { await PublishFromOutboxAsync(outgoingEvent, outboxConfig); + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Outbox, + EventName = outgoingEvent.EventName, + EventData = outgoingEvent.EventData + }); } - + await scope.CompleteAsync(); } } @@ -262,7 +279,7 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); - await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); + await TriggerHandlersFromInboxAsync(eventType, eventData, exceptions, inboxConfig); if (exceptions.Any()) { ThrowOriginalExceptions(eventType, exceptions); 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 5da66e03cd..d35377db2d 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 @@ -4,6 +4,7 @@ using System.Linq; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; +using Volo.Abp.EventBus.Local; using Volo.Abp.Guids; using Volo.Abp.MultiTenancy; using Volo.Abp.Timing; @@ -16,6 +17,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB protected IGuidGenerator GuidGenerator { get; } protected IClock Clock { get; } protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } + protected ILocalEventBus LocalEventBus { get; } protected DistributedEventBusBase( IServiceScopeFactory serviceScopeFactory, @@ -24,8 +26,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB IOptions abpDistributedEventBusOptions, IGuidGenerator guidGenerator, IClock clock, - IEventHandlerInvoker eventHandlerInvoker - ) : base( + IEventHandlerInvoker eventHandlerInvoker, + ILocalEventBus localEventBus) : base( serviceScopeFactory, currentTenant, unitOfWorkManager, @@ -34,6 +36,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB GuidGenerator = guidGenerator; Clock = clock; AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; + LocalEventBus = localEventBus; } public IDisposable Subscribe(IDistributedEventHandler handler) where TEvent : class @@ -79,6 +82,13 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB } await PublishToEventBusAsync(eventType, eventData); + + await TriggerDistributedEventSentAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Direct, + EventName = EventNameAttribute.GetNameOrDefault(eventType), + EventData = eventData + }); } public abstract Task PublishFromOutboxAsync( @@ -170,4 +180,52 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB } protected abstract byte[] Serialize(object eventData); + + protected virtual async Task TriggerHandlersDirectAsync(Type eventType, object eventData) + { + await TriggerHandlersAsync(eventType, eventData); + + await TriggerDistributedEventReceivedAsync(new DistributedEventReceived + { + Source = DistributedEventSource.Direct, + EventName = EventNameAttribute.GetNameOrDefault(eventType), + EventData = eventData + }); + } + + protected virtual async Task TriggerHandlersFromInboxAsync(Type eventType, object eventData, List exceptions, InboxConfig inboxConfig = null) + { + await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); + + await TriggerDistributedEventReceivedAsync(new DistributedEventReceived + { + Source = DistributedEventSource.Inbox, + EventName = EventNameAttribute.GetNameOrDefault(eventType), + EventData = eventData + }); + } + + public virtual async Task TriggerDistributedEventSentAsync(DistributedEventSent distributedEvent) + { + try + { + await LocalEventBus.PublishAsync(distributedEvent); + } + catch (Exception _) + { + // ignored + } + } + + public virtual async Task TriggerDistributedEventReceivedAsync(DistributedEventReceived distributedEvent) + { + try + { + await LocalEventBus.PublishAsync(distributedEvent); + } + catch (Exception _) + { + // ignored + } + } } 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 64f375130b..b63133a388 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 @@ -28,6 +28,30 @@ public class LocalDistributedEventBus : IDistributedEventBus, ISingletonDependen ServiceScopeFactory = serviceScopeFactory; AbpDistributedEventBusOptions = distributedEventBusOptions.Value; Subscribe(distributedEventBusOptions.Value.Handlers); + + // For unit testing + if (localEventBus is LocalEventBus eventBus) + { + eventBus.OnEventHandleInvoking = async (eventType, eventData) => + { + await localEventBus.PublishAsync(new DistributedEventReceived() + { + Source = DistributedEventSource.Direct, + EventName = EventNameAttribute.GetNameOrDefault(eventType), + EventData = eventData + }, onUnitOfWorkComplete: false); + }; + + eventBus.OnPublishing = async (eventType, eventData) => + { + await localEventBus.PublishAsync(new DistributedEventSent() + { + Source = DistributedEventSource.Direct, + EventName = EventNameAttribute.GetNameOrDefault(eventType), + EventData = eventData + }, onUnitOfWorkComplete: false); + }; + } } public virtual void Subscribe(ITypeList handlers) @@ -132,7 +156,7 @@ public class LocalDistributedEventBus : IDistributedEventBus, ISingletonDependen { return _localEventBus.PublishAsync(eventType, eventData, onUnitOfWorkComplete); } - + public Task PublishAsync(TEvent eventData, bool onUnitOfWorkComplete = true, bool useOutbox = true) where TEvent : class { return _localEventBus.PublishAsync(eventData, onUnitOfWorkComplete); @@ -142,4 +166,4 @@ public class LocalDistributedEventBus : IDistributedEventBus, ISingletonDependen { return _localEventBus.PublishAsync(eventType, eventData, onUnitOfWorkComplete); } -} \ No newline at end of file +} 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 6e48b27e5f..516a0008d7 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -9,7 +9,6 @@ using Microsoft.Extensions.DependencyInjection; using Volo.Abp.Collections; using Volo.Abp.EventBus.Distributed; using Volo.Abp.MultiTenancy; -using Volo.Abp.Reflection; using Volo.Abp.Uow; namespace Volo.Abp.EventBus; @@ -214,7 +213,7 @@ public abstract class EventBusBase : IEventBus using (CurrentTenant.Change(GetEventDataTenantId(eventData))) { - await EventHandlerInvoker.InvokeAsync(eventHandlerWrapper.EventHandler, eventData, eventType); + await InvokeEventHandlerAsync(eventHandlerWrapper.EventHandler, eventData, eventType); } } catch (TargetInvocationException ex) @@ -228,6 +227,11 @@ public abstract class EventBusBase : IEventBus } } + protected virtual Task InvokeEventHandlerAsync(IEventHandler eventHandler, object eventData, Type eventType) + { + return EventHandlerInvoker.InvokeAsync(eventHandler, eventData, eventType); + } + protected virtual Guid? GetEventDataTenantId(object eventData) { return eventData switch 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 3f7dbbc851..ce98cecf5f 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 @@ -8,6 +8,7 @@ using System.Linq; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Volo.Abp.DependencyInjection; +using Volo.Abp.EventBus.Distributed; using Volo.Abp.MultiTenancy; using Volo.Abp.Threading; using Volo.Abp.Uow; @@ -135,6 +136,33 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData); } + // Internal for unit testing + internal Func OnPublishing { get; set; } + + // For unit testing + public async override Task PublishAsync( + Type eventType, + object eventData, + bool onUnitOfWorkComplete = true) + { + if (onUnitOfWorkComplete && UnitOfWorkManager.Current != null) + { + AddToUnitOfWork( + UnitOfWorkManager.Current, + new UnitOfWorkEventRecord(eventType, eventData, EventOrderGenerator.GetNext()) + ); + return; + } + + // For unit testing + if (OnPublishing != null && eventType != typeof(DistributedEventSent) && eventType != typeof(DistributedEventReceived)) + { + await OnPublishing(eventType, eventData); + } + + await PublishToEventBusAsync(eventType, eventData); + } + protected override IEnumerable GetHandlerFactories(Type eventType) { var handlerFactoryList = new List(); @@ -168,4 +196,18 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency return false; } + + // Internal for unit testing + internal Func OnEventHandleInvoking { get; set; } + + // Internal for unit testing + protected async override Task InvokeEventHandlerAsync(IEventHandler eventHandler, object eventData, Type eventType) + { + if (OnEventHandleInvoking != null && eventType != typeof(DistributedEventSent) && eventType != typeof(DistributedEventReceived)) + { + await OnEventHandleInvoking(eventType, eventData); + } + + await base.InvokeEventHandlerAsync(eventHandler, eventData, eventType); + } } diff --git a/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/AbpMongoDbDateTimeSerializer.cs b/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/AbpMongoDbDateTimeSerializer.cs index f8619602cb..c1d8b8f514 100644 --- a/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/AbpMongoDbDateTimeSerializer.cs +++ b/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/AbpMongoDbDateTimeSerializer.cs @@ -34,7 +34,7 @@ public class AbpMongoDbDateTimeSerializer : DateTimeSerializer return (dateTime - BsonConstants.UnixEpoch).Ticks / 10000L; } - // For unit testing. + // For unit testing internal void SetDateTimeKind(DateTimeKind dateTimeKind) { DateTimeKind = dateTimeKind; diff --git a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/DistributedEventHandles.cs b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/DistributedEventHandles.cs new file mode 100644 index 0000000000..eb1f1a52df --- /dev/null +++ b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/DistributedEventHandles.cs @@ -0,0 +1,22 @@ +using System.Threading.Tasks; + +namespace Volo.Abp.EventBus.Distributed; + +public class DistributedEventHandles : ILocalEventHandler, ILocalEventHandler +{ + public static int SentCount { get; set; } + + public static int ReceivedCount { get; set; } + + public Task HandleEventAsync(DistributedEventSent eventData) + { + SentCount++; + return Task.CompletedTask; + } + + public Task HandleEventAsync(DistributedEventReceived eventData) + { + ReceivedCount++; + return Task.CompletedTask; + } +} diff --git a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus_Test.cs b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus_Test.cs index c239df38c9..bbb8838cf6 100644 --- a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus_Test.cs +++ b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus_Test.cs @@ -1,7 +1,8 @@ using System; using System.Threading.Tasks; using Volo.Abp.Domain.Entities.Events.Distributed; -using Volo.Abp.MultiTenancy; +using Volo.Abp.EventBus.Local; +using Volo.Abp.Uow; using Xunit; namespace Volo.Abp.EventBus.Distributed; @@ -49,7 +50,7 @@ public class LocalDistributedEventBus_Test : LocalDistributedEventBusTestBase public async Task Should_Get_TenantId_From_EventEto_Extra_Property() { var tenantId = Guid.NewGuid(); - + DistributedEventBus.Subscribe(GetRequiredService()); await DistributedEventBus.PublishAsync(new MySimpleEto @@ -59,7 +60,44 @@ public class LocalDistributedEventBus_Test : LocalDistributedEventBusTestBase {"TenantId", tenantId.ToString()} } }); - + Assert.Equal(tenantId, MySimpleDistributedSingleInstanceEventHandler.TenantId); } + + [Fact] + public async Task DistributedEventSentAndReceived_Test() + { + GetRequiredService().Subscribe(); + GetRequiredService().Subscribe(); + + DistributedEventBus.Subscribe(); + + using (var uow = GetRequiredService().Begin()) + { + await DistributedEventBus.PublishAsync(new MyEventDate(), onUnitOfWorkComplete: false); + + Assert.Equal(1, DistributedEventHandles.SentCount); + Assert.Equal(1, DistributedEventHandles.ReceivedCount); + + await DistributedEventBus.PublishAsync(new MyEventDate(), onUnitOfWorkComplete: true); + + await uow.CompleteAsync(); + + Assert.Equal(2, DistributedEventHandles.SentCount); + Assert.Equal(2, DistributedEventHandles.ReceivedCount); + } + } + + class MyEventDate + { + + } + + class MyEventHandle : IDistributedEventHandler + { + public Task HandleEventAsync(MyEventDate eventData) + { + return Task.CompletedTask; + } + } }