diff --git a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/AnonymousEventData.cs b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/AnonymousEventData.cs new file mode 100644 index 0000000000..c4b5d1c3fd --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/AnonymousEventData.cs @@ -0,0 +1,103 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text.Json; + +namespace Volo.Abp.EventBus; + +[Serializable] +public class AnonymousEventData +{ + public string EventName { get; } + public object Data { get; } + + private JsonElement? _cachedJsonElement; + + public AnonymousEventData(string eventName, object data) + { + EventName = eventName; + Data = data; + } + + public object ConvertToTypedObject() + { + return ConvertElement(GetJsonElement()); + } + + public T ConvertToTypedObject() + { + if (Data is T typedData) + { + return typedData; + } + + return GetJsonElement().Deserialize() + ?? throw new InvalidOperationException($"Failed to deserialize AnonymousEventData to {typeof(T).FullName}."); + } + + public object ConvertToTypedObject(Type type) + { + if (type.IsInstanceOfType(Data)) + { + return Data; + } + + return GetJsonElement().Deserialize(type) + ?? throw new InvalidOperationException($"Failed to deserialize AnonymousEventData to {type.FullName}."); + } + + private JsonElement GetJsonElement() + { + if (_cachedJsonElement.HasValue) + { + return _cachedJsonElement.Value; + } + + if (Data is JsonElement existingElement) + { + _cachedJsonElement = existingElement; + return existingElement; + } + + _cachedJsonElement = JsonSerializer.SerializeToElement(Data); + return _cachedJsonElement.Value; + } + + private static object ConvertElement(JsonElement element) + { + switch (element.ValueKind) + { + case JsonValueKind.Object: + { + var obj = new Dictionary(); + foreach (var property in element.EnumerateObject()) + { + obj[property.Name] = property.Value.ValueKind == JsonValueKind.Null + ? null + : ConvertElement(property.Value); + } + return obj; + } + case JsonValueKind.Array: + return element.EnumerateArray() + .Select(item => item.ValueKind == JsonValueKind.Null ? null : (object?)ConvertElement(item)) + .ToList(); + case JsonValueKind.String: + return element.GetString()!; + case JsonValueKind.Number when element.TryGetInt64(out var longValue): + return longValue; + case JsonValueKind.Number when element.TryGetDecimal(out var decimalValue): + return decimalValue; + case JsonValueKind.Number when element.TryGetDouble(out var doubleValue): + return doubleValue; + case JsonValueKind.True: + return true; + case JsonValueKind.False: + return false; + case JsonValueKind.Null: + case JsonValueKind.Undefined: + default: + return null!; + } + } +} diff --git a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/IEventBus.cs b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/IEventBus.cs index 938fd4d97b..fa17528750 100644 --- a/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/IEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/IEventBus.cs @@ -52,6 +52,8 @@ public interface IEventBus /// Event type /// Object to handle the event IDisposable Subscribe(Type eventType, IEventHandler handler); + + IDisposable Subscribe(string eventName, IEventHandlerFactory handler); /// /// Registers to an event. @@ -106,6 +108,8 @@ public interface IEventBus /// Event type /// Factory object that is registered before void Unsubscribe(Type eventType, IEventHandlerFactory factory); + + void Unsubscribe(string eventName, IEventHandlerFactory factory); /// /// Unregisters all event handlers of given event type. 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 e691b6c58c..d816d61c14 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 @@ -23,4 +23,6 @@ public interface ILocalEventBus : IEventBus /// Event type /// List GetEventHandlerFactories(Type eventType); + + List GetAnonymousEventHandlerFactories(string eventName); } 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 da997f23f4..8cf42cc796 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 @@ -2,6 +2,7 @@ using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; +using System.Text.Json; using System.Threading.Tasks; using Azure.Messaging.ServiceBus; using Microsoft.Extensions.DependencyInjection; @@ -29,6 +30,7 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen protected IAzureServiceBusSerializer Serializer { get; } protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } + protected ConcurrentDictionary> AnonymousHandlerFactories { get; } protected IAzureServiceBusMessageConsumer Consumer { get; private set; } = default!; public AzureDistributedEventBus( @@ -61,6 +63,7 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen PublisherPool = publisherPool; HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); + AnonymousHandlerFactories = new ConcurrentDictionary>(); } public void Initialize() @@ -81,14 +84,25 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen { return; } + var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + object eventData; + + if (eventType != null) + { + eventData = Serializer.Deserialize(message.Body.ToArray(), eventType); + } + else if (AnonymousHandlerFactories.ContainsKey(eventName)) + { + var data = Serializer.Deserialize(message.Body.ToArray()); + eventData = new AnonymousEventData(eventName, data); + eventType = typeof(AnonymousEventData); + } + else { return; } - var eventData = Serializer.Deserialize(message.Body.ToArray(), eventType); - if (await AddToInboxAsync(message.MessageId, eventName, eventType, eventData, message.CorrelationId)) { return; @@ -162,12 +176,22 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) { var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); - if (eventType == null) + object eventData; + + if (eventType != null) + { + eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); + } + else if (AnonymousHandlerFactories.ContainsKey(incomingEvent.EventName)) + { + var element = Serializer.Deserialize(incomingEvent.EventData); + eventData = new AnonymousEventData(incomingEvent.EventName, element); + eventType = typeof(AnonymousEventData); + } + else { return; } - - var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); using (CorrelationIdProvider.Change(incomingEvent.GetCorrelationId())) { @@ -253,17 +277,25 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen public override Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) { var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + var anonymousEventData = eventData as AnonymousEventData ?? new AnonymousEventData(eventName, eventData); + + if (eventType != null) + { + return PublishAsync(eventType, anonymousEventData.ConvertToTypedObject(eventType), onUnitOfWorkComplete); + } + + if (AnonymousHandlerFactories.ContainsKey(eventName)) { - throw new AbpException($"Unknown event name: {eventName}"); + return PublishAsync(typeof(AnonymousEventData), anonymousEventData, onUnitOfWorkComplete); } - return PublishAsync(eventType, eventData, onUnitOfWorkComplete); + throw new AbpException($"Unknown event name: {eventName}"); } protected async override Task PublishToEventBusAsync(Type eventType, object eventData) { - await PublishAsync(EventNameAttribute.GetNameOrDefault(eventType), eventData); + var (eventName, resolvedData) = ResolveEventForPublishing(eventType, eventData); + await PublishAsync(eventName, resolvedData); } protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord) @@ -312,6 +344,52 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen .ToArray(); } + protected override Type? GetEventTypeByEventName(string eventName) + { + return EventTypes.GetOrDefault(eventName); + } + + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => + { + if (!handler.IsInFactories(factories)) + { + factories.Add(handler); + } + }); + + return new AnonymousEventHandlerFactoryUnregistrar(this, eventName, handler); + } + + public override void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => factories.Remove(factory)); + } + + protected override IEnumerable GetAnonymousHandlerFactories(string eventName) + { + var result = new List(); + + var eventType = GetEventTypeByEventName(eventName); + if (eventType != null) + { + result.AddRange(GetHandlerFactories(eventType)); + } + + foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName)) + { + result.Add(new EventTypeWithEventHandlerFactories(typeof(AnonymousEventData), handlerFactory.Value)); + } + + return result; + } + + private List GetOrCreateAnonymousHandlerFactories(string eventName) + { + return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List()); + } + private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) { return handlerEventType == targetEventType || handlerEventType.IsAssignableFrom(targetEventType); 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 548467d1c5..f036514115 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 @@ -1,4 +1,4 @@ -using System; +using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; @@ -28,6 +28,7 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } + protected ConcurrentDictionary> AnonymousHandlerFactories { get; } public DaprDistributedEventBus( IServiceScopeFactory serviceScopeFactory, @@ -58,6 +59,7 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); + AnonymousHandlerFactories = new ConcurrentDictionary>(); } public void Initialize() @@ -132,17 +134,25 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend public override Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) { var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + var anonymousEventData = eventData as AnonymousEventData ?? new AnonymousEventData(eventName, eventData); + + if (eventType != null) + { + return PublishAsync(eventType, anonymousEventData.ConvertToTypedObject(eventType), onUnitOfWorkComplete); + } + + if (AnonymousHandlerFactories.ContainsKey(eventName)) { - throw new AbpException($"Unknown event name: {eventName}"); + return PublishAsync(typeof(AnonymousEventData), new AnonymousEventData(eventName, eventData), onUnitOfWorkComplete); } - return PublishAsync(eventType, eventData, onUnitOfWorkComplete); + throw new AbpException($"Unknown event name: {eventName}"); } protected async override Task PublishToEventBusAsync(Type eventType, object eventData) { - await PublishToDaprAsync(eventType, eventData, null, CorrelationIdProvider.Get()); + var (eventName, resolvedData) = ResolveEventForPublishing(eventType, eventData); + await PublishToDaprAsync(eventName, resolvedData, null, CorrelationIdProvider.Get()); } protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord) @@ -162,6 +172,52 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend return handlerFactoryList.ToArray(); } + protected override Type? GetEventTypeByEventName(string eventName) + { + return EventTypes.GetOrDefault(eventName); + } + + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => + { + if (!handler.IsInFactories(factories)) + { + factories.Add(handler); + } + }); + + return new AnonymousEventHandlerFactoryUnregistrar(this, eventName, handler); + } + + public override void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => factories.Remove(factory)); + } + + protected override IEnumerable GetAnonymousHandlerFactories(string eventName) + { + var result = new List(); + + var eventType = GetEventTypeByEventName(eventName); + if (eventType != null) + { + result.AddRange(GetHandlerFactories(eventType)); + } + + foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName)) + { + result.Add(new EventTypeWithEventHandlerFactories(typeof(AnonymousEventData), handlerFactory.Value)); + } + + return result; + } + + private List GetOrCreateAnonymousHandlerFactories(string eventName) + { + return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List()); + } + public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) { var eventType = GetEventType(outgoingEvent.EventName); 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 5909614fb3..17d4bccc12 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 @@ -1,7 +1,8 @@ -using System; +using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; +using System.Text.Json; using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.DependencyInjection; @@ -29,6 +30,7 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen protected IProducerPool ProducerPool { get; } protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } + protected ConcurrentDictionary> AnonymousHandlerFactories { get; } protected IKafkaMessageConsumer Consumer { get; private set; } = default!; public KafkaDistributedEventBus( @@ -63,6 +65,7 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); + AnonymousHandlerFactories = new ConcurrentDictionary>(); } public void Initialize() @@ -80,14 +83,25 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen { var eventName = message.Key; var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) - { - return; - } var messageId = message.GetMessageId(); - var eventData = Serializer.Deserialize(message.Value, eventType); var correlationId = message.GetCorrelationId(); + object eventData; + + if (eventType != null) + { + eventData = Serializer.Deserialize(message.Value, eventType); + } + else if (AnonymousHandlerFactories.ContainsKey(eventName)) + { + var jsonElement = JsonSerializer.Deserialize(message.Value); + eventData = new AnonymousEventData(eventName, jsonElement); + eventType = typeof(AnonymousEventData); + } + else + { + return; + } if (await AddToInboxAsync(messageId, eventName, eventType, eventData, correlationId)) { @@ -171,12 +185,19 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen public override Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) { var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + var anonymousEventData = eventData as AnonymousEventData ?? new AnonymousEventData(eventName, eventData); + + if (eventType != null) { - throw new AbpException($"Unknown event name: {eventName}"); + return PublishAsync(eventType, anonymousEventData.ConvertToTypedObject(eventType), onUnitOfWorkComplete); } - return PublishAsync(eventType, eventData, onUnitOfWorkComplete); + if (AnonymousHandlerFactories.ContainsKey(eventName)) + { + return PublishAsync(typeof(AnonymousEventData), new AnonymousEventData(eventName, eventData), onUnitOfWorkComplete); + } + + throw new AbpException($"Unknown event name: {eventName}"); } protected override async Task PublishToEventBusAsync(Type eventType, object eventData) @@ -289,12 +310,22 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen InboxConfig inboxConfig) { var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); - if (eventType == null) + object eventData; + + if (eventType != null) + { + eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); + } + else if (AnonymousHandlerFactories.ContainsKey(incomingEvent.EventName)) + { + var jsonElement = JsonSerializer.Deserialize(incomingEvent.EventData); + eventData = new AnonymousEventData(incomingEvent.EventName, jsonElement); + eventType = typeof(AnonymousEventData); + } + else { return; } - - var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); using (CorrelationIdProvider.Change(incomingEvent.GetCorrelationId())) { @@ -313,8 +344,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen private async Task PublishAsync(string topicName, Type eventType, object eventData, Headers headers) { - var eventName = EventNameAttribute.GetNameOrDefault(eventType); - var body = Serializer.Serialize(eventData); + var (eventName, resolvedData) = ResolveEventForPublishing(eventType, eventData); + var body = Serializer.Serialize(resolvedData); var result = await PublishAsync(topicName, eventName, body, headers); if (result.Status != PersistenceStatus.Persisted) @@ -374,6 +405,52 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen return handlerFactoryList.ToArray(); } + protected override Type? GetEventTypeByEventName(string eventName) + { + return EventTypes.GetOrDefault(eventName); + } + + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => + { + if (!handler.IsInFactories(factories)) + { + factories.Add(handler); + } + }); + + return new AnonymousEventHandlerFactoryUnregistrar(this, eventName, handler); + } + + public override void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => factories.Remove(factory)); + } + + protected override IEnumerable GetAnonymousHandlerFactories(string eventName) + { + var result = new List(); + + var eventType = GetEventTypeByEventName(eventName); + if (eventType != null) + { + result.AddRange(GetHandlerFactories(eventType)); + } + + foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName)) + { + result.Add(new EventTypeWithEventHandlerFactories(typeof(AnonymousEventData), handlerFactory.Value)); + } + + return result; + } + + private List GetOrCreateAnonymousHandlerFactories(string eventName) + { + return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List()); + } + private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) { //Should trigger same type 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 c58e40a1ec..928c99762d 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 @@ -1,7 +1,8 @@ -using System; +using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; +using System.Text.Json; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; @@ -33,6 +34,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis //TODO: Accessing to the List may not be thread-safe! protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } + protected ConcurrentDictionary> AnonymousHandlerFactories { get; } protected IRabbitMqMessageConsumerFactory MessageConsumerFactory { get; } protected IRabbitMqMessageConsumer Consumer { get; private set; } = default!; @@ -70,6 +72,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); + AnonymousHandlerFactories = new ConcurrentDictionary>(); } public virtual void Initialize() @@ -101,13 +104,23 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis { var eventName = ea.RoutingKey; var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + object eventData; + + if (eventType != null) + { + eventData = Serializer.Deserialize(ea.Body.ToArray(), eventType); + } + else if (AnonymousHandlerFactories.ContainsKey(eventName)) + { + var jsonElement = JsonSerializer.Deserialize(ea.Body.ToArray()); + eventData = new AnonymousEventData(eventName, jsonElement); + eventType = typeof(AnonymousEventData); + } + else { return; } - var eventData = Serializer.Deserialize(ea.Body.ToArray(), eventType); - var correlationId = ea.BasicProperties.CorrelationId; if (await AddToInboxAsync(ea.BasicProperties.MessageId, eventName, eventType, eventData, correlationId)) { @@ -196,12 +209,19 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis public override Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) { var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + var anonymousEventData = eventData as AnonymousEventData ?? new AnonymousEventData(eventName, eventData); + + if (eventType != null) { - throw new AbpException($"Unknown event name: {eventName}"); + return PublishAsync(eventType, anonymousEventData.ConvertToTypedObject(eventType), onUnitOfWorkComplete); } - return PublishAsync(eventType, eventData, onUnitOfWorkComplete); + if (AnonymousHandlerFactories.ContainsKey(eventName)) + { + return PublishAsync(typeof(AnonymousEventData), new AnonymousEventData(eventName, eventData), onUnitOfWorkComplete); + } + + throw new AbpException($"Unknown event name: {eventName}"); } protected async override Task PublishToEventBusAsync(Type eventType, object eventData) @@ -267,12 +287,22 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis InboxConfig inboxConfig) { var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); - if (eventType == null) + object eventData; + + if (eventType != null) + { + eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); + } + else if (AnonymousHandlerFactories.ContainsKey(incomingEvent.EventName)) + { + var jsonElement = JsonSerializer.Deserialize(incomingEvent.EventData); + eventData = new AnonymousEventData(incomingEvent.EventName, jsonElement); + eventType = typeof(AnonymousEventData); + } + else { return; } - - var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); using (CorrelationIdProvider.Change(incomingEvent.GetCorrelationId())) { @@ -296,8 +326,8 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis Guid? eventId = null, string? correlationId = null) { - var eventName = EventNameAttribute.GetNameOrDefault(eventType); - var body = Serializer.Serialize(eventData); + var (eventName, resolvedData) = ResolveEventForPublishing(eventType, eventData); + var body = Serializer.Serialize(resolvedData); return PublishAsync(eventName, body, headersArguments, eventId, correlationId); } @@ -424,6 +454,52 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis return handlerFactoryList.ToArray(); } + protected override Type? GetEventTypeByEventName(string eventName) + { + return EventTypes.GetOrDefault(eventName); + } + + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => + { + if (!handler.IsInFactories(factories)) + { + factories.Add(handler); + } + }); + + return new AnonymousEventHandlerFactoryUnregistrar(this, eventName, handler); + } + + public override void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => factories.Remove(factory)); + } + + protected override IEnumerable GetAnonymousHandlerFactories(string eventName) + { + var result = new List(); + + var eventType = GetEventTypeByEventName(eventName); + if (eventType != null) + { + result.AddRange(GetHandlerFactories(eventType)); + } + + foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName)) + { + result.Add(new EventTypeWithEventHandlerFactories(typeof(AnonymousEventData), handlerFactory.Value)); + } + + return result; + } + + private List GetOrCreateAnonymousHandlerFactories(string eventName) + { + return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List()); + } + private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) { //Should trigger same type 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 a90e166289..a1f9e0116e 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 @@ -1,4 +1,4 @@ -using System; +using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; @@ -31,6 +31,7 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen //TODO: Accessing to the List may not be thread-safe! protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } + protected ConcurrentDictionary> AnonymousHandlerFactories { get; } protected AbpRebusEventBusOptions AbpRebusEventBusOptions { get; } public RebusDistributedEventBus( @@ -63,6 +64,7 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); + AnonymousHandlerFactories = new ConcurrentDictionary>(); } public void Initialize() @@ -163,12 +165,19 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen public override Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) { var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + var anonymousEventData = eventData as AnonymousEventData ?? new AnonymousEventData(eventName, eventData); + + if (eventType != null) + { + return PublishAsync(eventType, anonymousEventData.ConvertToTypedObject(eventType), onUnitOfWorkComplete); + } + + if (AnonymousHandlerFactories.ContainsKey(eventName)) { - throw new AbpException($"Unknown event name: {eventName}"); + return PublishAsync(typeof(AnonymousEventData), new AnonymousEventData(eventName, eventData), onUnitOfWorkComplete); } - return PublishAsync(eventType, eventData, onUnitOfWorkComplete); + throw new AbpException($"Unknown event name: {eventName}"); } protected async override Task PublishToEventBusAsync(Type eventType, object eventData) @@ -240,6 +249,52 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen return handlerFactoryList.ToArray(); } + protected override Type? GetEventTypeByEventName(string eventName) + { + return EventTypes.GetOrDefault(eventName); + } + + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => + { + if (!handler.IsInFactories(factories)) + { + factories.Add(handler); + } + }); + + return new AnonymousEventHandlerFactoryUnregistrar(this, eventName, handler); + } + + public override void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => factories.Remove(factory)); + } + + protected override IEnumerable GetAnonymousHandlerFactories(string eventName) + { + var result = new List(); + + var eventType = GetEventTypeByEventName(eventName); + if (eventType != null) + { + result.AddRange(GetHandlerFactories(eventType)); + } + + foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName)) + { + result.Add(new EventTypeWithEventHandlerFactories(typeof(AnonymousEventData), handlerFactory.Value)); + } + + return result; + } + + private List GetOrCreateAnonymousHandlerFactories(string eventName) + { + return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List()); + } + private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) { //Should trigger same type @@ -317,12 +372,21 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen InboxConfig inboxConfig) { var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); - if (eventType == null) + object eventData; + + if (eventType != null) + { + eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); + } + else if (AnonymousHandlerFactories.ContainsKey(incomingEvent.EventName)) + { + eventData = new AnonymousEventData(incomingEvent.EventName, Serializer.Deserialize(incomingEvent.EventData, typeof(object))); + eventType = typeof(AnonymousEventData); + } + else { return; } - - var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); using (CorrelationIdProvider.Change(incomingEvent.GetCorrelationId())) { 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 3668193c9b..ff9d587769 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 @@ -90,8 +90,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB await TriggerDistributedEventSentAsync(new DistributedEventSent() { Source = DistributedEventSource.Direct, - EventName = EventNameAttribute.GetNameOrDefault(eventType), - EventData = eventData + EventName = GetEventName(eventType, eventData), + EventData = GetEventData(eventData) }); } @@ -124,7 +124,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB if (outboxConfig.Selector == null || outboxConfig.Selector(eventType)) { var eventOutbox = (IEventOutbox)unitOfWork.ServiceProvider.GetRequiredService(outboxConfig.ImplementationType); - var eventName = EventNameAttribute.GetNameOrDefault(eventType); + var eventName = GetEventName(eventType, eventData); + eventData = GetEventData(eventData); await OnAddToOutboxAsync(eventName, eventType, eventData); @@ -181,8 +182,6 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB { if (await eventInbox.ExistsByMessageIdAsync(messageId!)) { - // Message already exists in the inbox, no need to add again. - // This can happen in case of retries from the sender side. addToInbox = true; continue; } @@ -212,8 +211,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB await TriggerDistributedEventReceivedAsync(new DistributedEventReceived { Source = DistributedEventSource.Direct, - EventName = EventNameAttribute.GetNameOrDefault(eventType), - EventData = eventData + EventName = GetEventName(eventType, eventData), + EventData = GetEventData(eventData) }); await TriggerHandlersAsync(eventType, eventData); @@ -224,8 +223,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB await TriggerDistributedEventReceivedAsync(new DistributedEventReceived { Source = DistributedEventSource.Inbox, - EventName = EventNameAttribute.GetNameOrDefault(eventType), - EventData = eventData + EventName = GetEventName(eventType, eventData), + EventData = GetEventData(eventData) }); await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); @@ -254,4 +253,29 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB // ignored } } + + protected virtual string GetEventName(Type eventType, object eventData) + { + if (eventData is AnonymousEventData anonymousEventData) + { + return anonymousEventData.EventName; + } + + return EventNameAttribute.GetNameOrDefault(eventType); + } + + protected virtual object GetEventData(object eventData) + { + if (eventData is AnonymousEventData anonymousEventData) + { + return anonymousEventData.ConvertToTypedObject(); + } + + return eventData; + } + + protected virtual (string EventName, object EventData) ResolveEventForPublishing(Type eventType, object eventData) + { + return (GetEventName(eventType, eventData), GetEventData(eventData)); + } } 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 abcd2dbd05..3cb012b58b 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,7 +5,6 @@ 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; @@ -26,6 +25,8 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen { protected ConcurrentDictionary EventTypes { get; } + protected ConcurrentDictionary AnonymousEventNames { get; } + public LocalDistributedEventBus( IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant, @@ -47,6 +48,7 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen correlationIdProvider) { EventTypes = new ConcurrentDictionary(); + AnonymousEventNames = new ConcurrentDictionary(); Subscribe(abpDistributedEventBusOptions.Value.Handlers); } @@ -71,6 +73,12 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen } } + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + AnonymousEventNames.GetOrAdd(eventName, true); + return LocalEventBus.Subscribe(eventName, handler); + } + public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) { var eventName = EventNameAttribute.GetNameOrDefault(eventType); @@ -93,6 +101,11 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen LocalEventBus.Unsubscribe(eventType, factory); } + public override void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + LocalEventBus.Unsubscribe(eventName, factory); + } + public override void UnsubscribeAll(Type eventType) { LocalEventBus.UnsubscribeAll(eventType); @@ -120,15 +133,15 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen await TriggerDistributedEventSentAsync(new DistributedEventSent() { Source = DistributedEventSource.Direct, - EventName = EventNameAttribute.GetNameOrDefault(eventType), - EventData = eventData + EventName = GetEventName(eventType, eventData), + EventData = GetEventData(eventData) }); await TriggerDistributedEventReceivedAsync(new DistributedEventReceived { Source = DistributedEventSource.Direct, - EventName = EventNameAttribute.GetNameOrDefault(eventType), - EventData = eventData + EventName = GetEventName(eventType, eventData), + EventData = GetEventData(eventData) }); await PublishToEventBusAsync(eventType, eventData); @@ -137,12 +150,21 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen public override Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) { var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + + var anonymousEventData = eventData as AnonymousEventData ?? new AnonymousEventData(eventName, eventData); + + if (eventType != null) + { + return PublishAsync(eventType, anonymousEventData.ConvertToTypedObject(eventType), onUnitOfWorkComplete); + } + + var isAnonymous = AnonymousEventNames.ContainsKey(eventName); + if (!isAnonymous) { throw new AbpException($"Unknown event name: {eventName}"); } - return PublishAsync(eventType, eventData, onUnitOfWorkComplete); + return PublishAsync(typeof(AnonymousEventData), new AnonymousEventData(eventName, eventData), onUnitOfWorkComplete); } protected async override Task PublishToEventBusAsync(Type eventType, object eventData) @@ -179,7 +201,13 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName); if (eventType == null) { - return; + var isAnonymous = AnonymousEventNames.ContainsKey(outgoingEvent.EventName); + if (!isAnonymous) + { + return; + } + + eventType = typeof(AnonymousEventData); } var eventData = JsonSerializer.Deserialize(Encoding.UTF8.GetString(outgoingEvent.EventData), eventType)!; @@ -204,7 +232,13 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); if (eventType == null) { - return; + var isAnonymous = AnonymousEventNames.ContainsKey(incomingEvent.EventName); + if (!isAnonymous) + { + return; + } + + eventType = typeof(AnonymousEventData); } var eventData = JsonSerializer.Deserialize(incomingEvent.EventData, eventType); @@ -226,7 +260,10 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen protected override Task OnAddToOutboxAsync(string eventName, Type eventType, object eventData) { - EventTypes.GetOrAdd(eventName, eventType); + if (eventType != typeof(AnonymousEventData)) + { + EventTypes.GetOrAdd(eventName, eventType); + } return base.OnAddToOutboxAsync(eventName, eventType, eventData); } @@ -234,4 +271,14 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen { return LocalEventBus.GetEventHandlerFactories(eventType); } + + protected override IEnumerable GetAnonymousHandlerFactories(string eventName) + { + return LocalEventBus.GetAnonymousEventHandlerFactories(eventName); + } + + protected override Type? GetEventTypeByEventName(string eventName) + { + return EventTypes.GetOrDefault(eventName); + } } diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/NullDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/NullDistributedEventBus.cs index bf64839863..12a4b3235b 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/NullDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/NullDistributedEventBus.cs @@ -1,4 +1,4 @@ -using System; +using System; using System.Threading.Tasks; namespace Volo.Abp.EventBus.Distributed; @@ -12,6 +12,11 @@ public sealed class NullDistributedEventBus : IDistributedEventBus } + public Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) + { + return Task.CompletedTask; + } + public IDisposable Subscribe(Func action) where TEvent : class { return NullDisposable.Instance; @@ -32,6 +37,11 @@ public sealed class NullDistributedEventBus : IDistributedEventBus return NullDisposable.Instance; } + public IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + return NullDisposable.Instance; + } + public IDisposable Subscribe(IEventHandlerFactory factory) where TEvent : class { return NullDisposable.Instance; @@ -67,6 +77,11 @@ public sealed class NullDistributedEventBus : IDistributedEventBus } + public void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + + } + public void UnsubscribeAll() where TEvent : class { 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 42d35ad314..b5283f6208 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -57,6 +57,8 @@ public abstract class EventBusBase : IEventBus return Subscribe(eventType, new SingleInstanceHandlerFactory(handler)); } + public abstract IDisposable Subscribe(string eventName, IEventHandlerFactory handler); + /// public virtual IDisposable Subscribe(IEventHandlerFactory factory) where TEvent : class { @@ -83,6 +85,8 @@ public abstract class EventBusBase : IEventBus public abstract void Unsubscribe(Type eventType, IEventHandlerFactory factory); + public abstract void Unsubscribe(string eventName, IEventHandlerFactory factory); + /// public virtual void UnsubscribeAll() where TEvent : class { @@ -139,31 +143,68 @@ public abstract class EventBusBase : IEventBus { await new SynchronizationContextRemover(); - foreach (var handlerFactories in GetHandlerFactories(eventType).ToList()) + var (handlerFactoriesList, actualEventType) = ResolveHandlerFactories(eventType, eventData); + + foreach (var handlerFactories in handlerFactoriesList) { foreach (var handlerFactory in handlerFactories.EventHandlerFactories.ToList()) { - await TriggerHandlerAsync(handlerFactory, handlerFactories.EventType, eventData, exceptions, inboxConfig); + var resolvedEventData = ResolveEventDataForHandler(eventData, eventType, handlerFactories.EventType); + await TriggerHandlerAsync(handlerFactory, handlerFactories.EventType, resolvedEventData, exceptions, inboxConfig); } } - //Implements generic argument inheritance. See IEventDataWithInheritableGenericArgument - if (eventType.GetTypeInfo().IsGenericType && - eventType.GetGenericArguments().Length == 1 && - typeof(IEventDataWithInheritableGenericArgument).IsAssignableFrom(eventType)) + if (actualEventType != null && + actualEventType.GetTypeInfo().IsGenericType && + actualEventType.GetGenericArguments().Length == 1 && + typeof(IEventDataWithInheritableGenericArgument).IsAssignableFrom(actualEventType)) { - var genericArg = eventType.GetGenericArguments()[0]; + var resolvedEventData = eventData is AnonymousEventData aed + ? aed.ConvertToTypedObject(actualEventType) + : eventData; + + var genericArg = actualEventType.GetGenericArguments()[0]; var baseArg = genericArg.GetTypeInfo().BaseType; if (baseArg != null) { - var baseEventType = eventType.GetGenericTypeDefinition().MakeGenericType(baseArg); - var constructorArgs = ((IEventDataWithInheritableGenericArgument)eventData).GetConstructorArgs(); + var baseEventType = actualEventType.GetGenericTypeDefinition().MakeGenericType(baseArg); + var constructorArgs = ((IEventDataWithInheritableGenericArgument)resolvedEventData).GetConstructorArgs(); var baseEventData = Activator.CreateInstance(baseEventType, constructorArgs)!; await PublishToEventBusAsync(baseEventType, baseEventData); } } } + protected virtual (List Factories, Type? ActualEventType) ResolveHandlerFactories( + Type eventType, + object eventData) + { + if (eventData is AnonymousEventData anonymousEventData) + { + return ( + GetAnonymousHandlerFactories(anonymousEventData.EventName).ToList(), + GetEventTypeByEventName(anonymousEventData.EventName) + ); + } + + return (GetHandlerFactories(eventType).ToList(), eventType); + } + + protected virtual object ResolveEventDataForHandler(object eventData, Type sourceEventType, Type handlerEventType) + { + if (eventData is AnonymousEventData anonymousEventData && handlerEventType != typeof(AnonymousEventData)) + { + return anonymousEventData.ConvertToTypedObject(handlerEventType); + } + + if (handlerEventType == typeof(AnonymousEventData) && eventData is not AnonymousEventData) + { + return new AnonymousEventData(EventNameAttribute.GetNameOrDefault(sourceEventType), eventData); + } + + return eventData; + } + protected void ThrowOriginalExceptions(Type eventType, List exceptions) { if (exceptions.Count == 1) @@ -200,6 +241,10 @@ public abstract class EventBusBase : IEventBus protected abstract IEnumerable GetHandlerFactories(Type eventType); + protected abstract IEnumerable GetAnonymousHandlerFactories(string eventName); + + protected abstract Type? GetEventTypeByEventName(string eventName); + protected virtual async Task TriggerHandlerAsync(IEventHandlerFactory asyncHandlerFactory, Type eventType, object eventData, List exceptions, InboxConfig? inboxConfig = null) { diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventHandlerFactoryUnregistrar.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventHandlerFactoryUnregistrar.cs index f94d7a4ed1..4c58978b75 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventHandlerFactoryUnregistrar.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventHandlerFactoryUnregistrar.cs @@ -23,3 +23,22 @@ public class EventHandlerFactoryUnregistrar : IDisposable _eventBus.Unsubscribe(_eventType, _factory); } } + +public class AnonymousEventHandlerFactoryUnregistrar : IDisposable +{ + private readonly IEventBus _eventBus; + private readonly string _eventName; + private readonly IEventHandlerFactory _factory; + + public AnonymousEventHandlerFactoryUnregistrar(IEventBus eventBus, string eventName, IEventHandlerFactory factory) + { + _eventBus = eventBus; + _eventName = eventName; + _factory = factory; + } + + public void Dispose() + { + _eventBus.Unsubscribe(_eventName, _factory); + } +} \ No newline at end of file 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 8631eaf5cf..0d6fbace63 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 @@ -29,8 +29,10 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency protected AbpLocalEventBusOptions Options { get; } protected ConcurrentDictionary> HandlerFactories { get; } - - protected ConcurrentDictionary EventTypes { get; } + + protected ConcurrentDictionary EventTypes { get; } + + protected ConcurrentDictionary> AnonymousEventHandlerFactories { get; } public LocalEventBus( IOptions options, @@ -45,6 +47,7 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); + AnonymousEventHandlerFactories = new ConcurrentDictionary>(); SubscribeHandlers(Options.Handlers); } @@ -54,11 +57,24 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency return Subscribe(typeof(TEvent), handler); } + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => + { + if (!handler.IsInFactories(factories)) + { + factories.Add(handler); + } + }); + + return new AnonymousEventHandlerFactoryUnregistrar(this, eventName, handler); + } + /// public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) { EventTypes.GetOrAdd(EventNameAttribute.GetNameOrDefault(eventType), eventType); - + GetOrCreateHandlerFactories(eventType) .Locking(factories => { @@ -120,6 +136,11 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Remove(factory)); } + public override void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => factories.Remove(factory)); + } + /// public override void UnsubscribeAll(Type eventType) { @@ -129,12 +150,21 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency public override Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) { var eventType = EventTypes.GetOrDefault(eventName); - if (eventType == null) + + var anonymousEventData = eventData as AnonymousEventData ?? new AnonymousEventData(eventName, eventData); + + if (eventType != null) + { + return PublishAsync(eventType, anonymousEventData.ConvertToTypedObject(eventType), onUnitOfWorkComplete); + } + + var isAnonymous = AnonymousEventHandlerFactories.ContainsKey(eventName); + if (!isAnonymous) { throw new AbpException($"Unknown event name: {eventName}"); } - return PublishAsync(eventType, eventData, onUnitOfWorkComplete); + return PublishAsync(typeof(AnonymousEventData), anonymousEventData, onUnitOfWorkComplete); } protected override async Task PublishToEventBusAsync(Type eventType, object eventData) @@ -157,9 +187,16 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency return GetHandlerFactories(eventType).ToList(); } + public virtual List GetAnonymousEventHandlerFactories(string eventName) + { + return GetAnonymousHandlerFactories(eventName).ToList(); + } + protected override IEnumerable GetHandlerFactories(Type eventType) { var handlerFactoryList = new List>(); + var eventNames = EventTypes.Where(x => ShouldTriggerEventForHandler(eventType, x.Value)).Select(x => x.Key).ToList(); + foreach (var handlerFactory in HandlerFactories.Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key))) { foreach (var factory in handlerFactory.Value) @@ -171,23 +208,71 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency } } + foreach (var handlerFactory in AnonymousEventHandlerFactories.Where(aehf => eventNames.Contains(aehf.Key))) + { + foreach (var factory in handlerFactory.Value) + { + handlerFactoryList.Add(new Tuple( + factory, + typeof(AnonymousEventData), + ReflectionHelper.GetAttributesOfMemberOrDeclaringType(factory.GetHandler().EventHandler.GetType()).FirstOrDefault()?.Order ?? 0)); + } + } + return handlerFactoryList.OrderBy(x => x.Item3).Select(x => new EventTypeWithEventHandlerFactories(x.Item2, new List {x.Item1})).ToArray(); } + protected override IEnumerable GetAnonymousHandlerFactories(string eventName) + { + var eventType = EventTypes.GetOrDefault(eventName); + if (eventType != null) + { + return GetHandlerFactories(eventType); + } + + var handlerFactoryList = new List>(); + + foreach (var handlerFactory in AnonymousEventHandlerFactories.Where(aehf => aehf.Key == eventName)) + { + foreach (var factory in handlerFactory.Value) + { + using var handler = factory.GetHandler(); + var handlerType = handler.EventHandler.GetType(); + handlerFactoryList.Add(new Tuple( + factory, + typeof(AnonymousEventData), + ReflectionHelper + .GetAttributesOfMemberOrDeclaringType(handlerType) + .FirstOrDefault()?.Order ?? 0)); + } + } + + return handlerFactoryList.OrderBy(x => x.Item3).Select(x => + new EventTypeWithEventHandlerFactories(x.Item2, new List { x.Item1 })).ToArray(); + } + + protected override Type? GetEventTypeByEventName(string eventName) + { + return EventTypes.GetOrDefault(eventName); + } + private List GetOrCreateHandlerFactories(Type eventType) { return HandlerFactories.GetOrAdd(eventType, (type) => new List()); } + private List GetOrCreateAnonymousHandlerFactories(string eventName) + { + return AnonymousEventHandlerFactories.GetOrAdd(eventName, (name) => new List()); + } + private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) { - //Should trigger same type if (handlerEventType == targetEventType) { return true; } - //Should trigger for inherited types if (handlerEventType.IsAssignableFrom(targetEventType)) { return true; 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 3ffcd911ce..76e61a263e 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 @@ -13,6 +13,11 @@ public sealed class NullLocalEventBus : ILocalEventBus } + public Task PublishAsync(string eventName, object eventData, bool onUnitOfWorkComplete = true) + { + return Task.CompletedTask; + } + public IDisposable Subscribe(Func action) where TEvent : class { return NullDisposable.Instance; @@ -28,6 +33,11 @@ public sealed class NullLocalEventBus : ILocalEventBus return new List(); } + public List GetAnonymousEventHandlerFactories(string eventName) + { + return new List(); + } + public IDisposable Subscribe() where TEvent : class where THandler : IEventHandler, new() { return NullDisposable.Instance; @@ -38,6 +48,11 @@ public sealed class NullLocalEventBus : ILocalEventBus return NullDisposable.Instance; } + public IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + return NullDisposable.Instance; + } + public IDisposable Subscribe(IEventHandlerFactory factory) where TEvent : class { return NullDisposable.Instance; @@ -73,6 +88,11 @@ public sealed class NullLocalEventBus : ILocalEventBus } + public void Unsubscribe(string eventName, IEventHandlerFactory factory) + { + + } + public void UnsubscribeAll() where TEvent : class { 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 88c515f8dc..c306b6c37d 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,4 +1,5 @@ using System; +using System.Collections.Generic; using System.Threading.Tasks; using Shouldly; using Volo.Abp.Domain.Entities.Events.Distributed; @@ -23,6 +24,184 @@ public class LocalDistributedEventBus_Test : LocalDistributedEventBusTestBase Assert.Equal(3, MySimpleDistributedTransientEventHandler.DisposeCount); } + [Fact] + public async Task Should_Handle_Typed_Handler_When_Published_With_EventName() + { + DistributedEventBus.Subscribe(); + + var eventName = EventNameAttribute.GetNameOrDefault(); + await DistributedEventBus.PublishAsync(eventName, new MySimpleEventData(1)); + await DistributedEventBus.PublishAsync(eventName, new Dictionary() + { + {"Value", 2} + }); + await DistributedEventBus.PublishAsync(eventName, new { Value = 3 }); + + Assert.Equal(3, MySimpleDistributedTransientEventHandler.HandleCount); + Assert.Equal(3, MySimpleDistributedTransientEventHandler.DisposeCount); + } + + [Fact] + public async Task Should_Handle_Anonymous_Handler_When_Published_With_EventName() + { + var handleCount = 0; + + DistributedEventBus.Subscribe("MyEvent", + new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + handleCount++; + await Task.CompletedTask; + }))); + + await DistributedEventBus.PublishAsync("MyEvent", new MySimpleEventData(1)); + await DistributedEventBus.PublishAsync("MyEvent", new Dictionary() + { + {"Value", 2} + }); + await DistributedEventBus.PublishAsync("MyEvent", new { Value = 3 }); + await DistributedEventBus.PublishAsync("MyEvent", new[] { 1, 2, 3 }); + + Assert.Equal(4, handleCount); + } + + [Fact] + public async Task Should_Handle_Anonymous_Handler_When_Published_With_AnonymousEventData() + { + var handleCount = 0; + + DistributedEventBus.Subscribe("MyEvent", + new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + handleCount++; + d.ConvertToTypedObject().ShouldNotBeNull(); + await Task.CompletedTask; + }))); + + await DistributedEventBus.PublishAsync(new AnonymousEventData("MyEvent", new MySimpleEventData(1))); + await DistributedEventBus.PublishAsync(new AnonymousEventData("MyEvent", new Dictionary() + { + {"Value", 2} + })); + await DistributedEventBus.PublishAsync(new AnonymousEventData("MyEvent", new { Value = 3 })); + + Assert.Equal(3, handleCount); + } + + [Fact] + public async Task Should_Handle_Typed_Handler_When_Published_With_AnonymousEventData() + { + DistributedEventBus.Subscribe(); + + var eventName = EventNameAttribute.GetNameOrDefault(); + + await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new MySimpleEventData(1))); + await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new Dictionary() + { + {"Value", 2} + })); + await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new { Value = 3 })); + + Assert.Equal(3, MySimpleDistributedTransientEventHandler.HandleCount); + } + + [Fact] + public async Task Should_Trigger_Both_Typed_And_Anonymous_Handlers_For_Typed_Event() + { + DistributedEventBus.Subscribe(); + + var eventName = EventNameAttribute.GetNameOrDefault(); + + var anonymousHandleCount = 0; + + DistributedEventBus.Subscribe(eventName, new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + anonymousHandleCount++; + await Task.CompletedTask; + }))); + + await DistributedEventBus.PublishAsync(new MySimpleEventData(1)); + await DistributedEventBus.PublishAsync(new MySimpleEventData(2)); + await DistributedEventBus.PublishAsync(new MySimpleEventData(3)); + + Assert.Equal(3, MySimpleDistributedTransientEventHandler.HandleCount); + Assert.Equal(3, anonymousHandleCount); + } + + [Fact] + public async Task Should_Trigger_Both_Handlers_For_Mixed_Typed_And_Anonymous_Publish() + { + DistributedEventBus.Subscribe(); + + var eventName = EventNameAttribute.GetNameOrDefault(); + + var anonymousHandleCount = 0; + + DistributedEventBus.Subscribe(eventName, new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + anonymousHandleCount++; + await Task.CompletedTask; + }))); + + await DistributedEventBus.PublishAsync(new MySimpleEventData(1)); + await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new Dictionary() + { + {"Value", 2} + })); + await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new { Value = 3 })); + + Assert.Equal(3, MySimpleDistributedTransientEventHandler.HandleCount); + Assert.Equal(3, anonymousHandleCount); + } + + [Fact] + public async Task Should_Unsubscribe_Anonymous_Handler() + { + var handleCount = 0; + + var handler = new ActionEventHandler(async (d) => + { + handleCount++; + await Task.CompletedTask; + }); + var factory = new SingleInstanceHandlerFactory(handler); + + var disposable = DistributedEventBus.Subscribe("MyEvent", factory); + + await DistributedEventBus.PublishAsync("MyEvent", new { Value = 1 }); + Assert.Equal(1, handleCount); + + disposable.Dispose(); + + await Assert.ThrowsAsync(() => + DistributedEventBus.PublishAsync("MyEvent", new { Value = 2 })); + Assert.Equal(1, handleCount); + } + + [Fact] + public async Task Should_Throw_For_Unknown_Event_Name() + { + await Assert.ThrowsAsync(() => + DistributedEventBus.PublishAsync("NonExistentEvent", new { Value = 1 })); + } + + [Fact] + public async Task Should_Convert_AnonymousEventData_To_Typed_Object() + { + MySimpleEventData? receivedData = null; + + DistributedEventBus.Subscribe(async (data) => + { + receivedData = data; + await Task.CompletedTask; + }); + + var eventName = EventNameAttribute.GetNameOrDefault(); + await DistributedEventBus.PublishAsync(eventName, new { Value = 42 }); + + receivedData.ShouldNotBeNull(); + receivedData.Value.ShouldBe(42); + } + [Fact] public async Task Should_Change_TenantId_If_EventData_Is_MultiTenant() { diff --git a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/LocalEventBus_Anonymous_Test.cs b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/LocalEventBus_Anonymous_Test.cs new file mode 100644 index 0000000000..a0b3e12e61 --- /dev/null +++ b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/LocalEventBus_Anonymous_Test.cs @@ -0,0 +1,162 @@ +using System.Collections.Generic; +using System.Threading.Tasks; +using Shouldly; +using Xunit; + +namespace Volo.Abp.EventBus.Local; + +public class LocalEventBus_Anonymous_Test : EventBusTestBase +{ + [Fact] + public async Task Should_Handle_Anonymous_Handler_With_EventName() + { + var handleCount = 0; + + LocalEventBus.Subscribe("TestEvent", + new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + handleCount++; + d.EventName.ShouldBe("TestEvent"); + await Task.CompletedTask; + }))); + + await LocalEventBus.PublishAsync("TestEvent", new { Value = 1 }); + await LocalEventBus.PublishAsync("TestEvent", new { Value = 2 }); + + handleCount.ShouldBe(2); + } + + [Fact] + public async Task Should_Handle_Typed_Handler_When_Published_With_EventName() + { + var handleCount = 0; + + LocalEventBus.Subscribe(async (data) => + { + handleCount++; + await Task.CompletedTask; + }); + + var eventName = EventNameAttribute.GetNameOrDefault(); + await LocalEventBus.PublishAsync(eventName, new MySimpleEventData(42)); + + handleCount.ShouldBe(1); + } + + [Fact] + public async Task Should_Convert_Dictionary_To_Typed_Handler() + { + MySimpleEventData? receivedData = null; + + LocalEventBus.Subscribe(async (data) => + { + receivedData = data; + await Task.CompletedTask; + }); + + var eventName = EventNameAttribute.GetNameOrDefault(); + await LocalEventBus.PublishAsync(eventName, new Dictionary + { + { "Value", 42 } + }); + + receivedData.ShouldNotBeNull(); + receivedData.Value.ShouldBe(42); + } + + [Fact] + public async Task Should_Trigger_Both_Typed_And_Anonymous_Handlers() + { + var typedHandleCount = 0; + var anonymousHandleCount = 0; + + LocalEventBus.Subscribe(async (data) => + { + typedHandleCount++; + await Task.CompletedTask; + }); + + var eventName = EventNameAttribute.GetNameOrDefault(); + + LocalEventBus.Subscribe(eventName, + new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + anonymousHandleCount++; + await Task.CompletedTask; + }))); + + await LocalEventBus.PublishAsync(new MySimpleEventData(1)); + + typedHandleCount.ShouldBe(1); + anonymousHandleCount.ShouldBe(1); + } + + [Fact] + public async Task Should_Unsubscribe_Anonymous_Handler() + { + var handleCount = 0; + + var handler = new ActionEventHandler(async (d) => + { + handleCount++; + await Task.CompletedTask; + }); + var factory = new SingleInstanceHandlerFactory(handler); + + var disposable = LocalEventBus.Subscribe("TestEvent", factory); + + await LocalEventBus.PublishAsync("TestEvent", new { Value = 1 }); + handleCount.ShouldBe(1); + + disposable.Dispose(); + + await Assert.ThrowsAsync(() => + LocalEventBus.PublishAsync("TestEvent", new { Value = 2 })); + handleCount.ShouldBe(1); + } + + [Fact] + public async Task Should_Throw_For_Unknown_Event_Name() + { + await Assert.ThrowsAsync(() => + LocalEventBus.PublishAsync("NonExistentEvent", new { Value = 1 })); + } + + [Fact] + public async Task Should_ConvertToTypedObject_In_Anonymous_Handler() + { + object? receivedData = null; + + LocalEventBus.Subscribe("TestEvent", + new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + receivedData = d.ConvertToTypedObject(); + await Task.CompletedTask; + }))); + + await LocalEventBus.PublishAsync("TestEvent", new { Name = "Hello", Count = 42 }); + + receivedData.ShouldNotBeNull(); + var dict = receivedData.ShouldBeOfType>(); + dict["Name"].ShouldBe("Hello"); + dict["Count"].ShouldBe(42L); + } + + [Fact] + public async Task Should_ConvertToTypedObject_Generic_In_Anonymous_Handler() + { + MySimpleEventData? receivedData = null; + + LocalEventBus.Subscribe("TestEvent", + new SingleInstanceHandlerFactory(new ActionEventHandler(async (d) => + { + receivedData = d.ConvertToTypedObject(); + await Task.CompletedTask; + }))); + + await LocalEventBus.PublishAsync("TestEvent", new MySimpleEventData(99)); + + receivedData.ShouldNotBeNull(); + receivedData.Value.ShouldBe(99); + } +}