From aae1e330ba04baf90e06a9ae7aa86d2c626e66ad Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?SAL=C4=B0H=20=C3=96ZKARA?= Date: Wed, 4 Mar 2026 12:42:15 +0300 Subject: [PATCH] Anonymous subscriptions and deserialization fixes Add and centralize Subscribe/Unsubscribe(string, IEventHandlerFactory) implementations for Azure and Kafka to avoid duplicate anonymous handler registrations (checks IsInFactories / returns NullDisposable or skips adding). Switch Kafka anonymous payload handling from JsonElement to the generic Serializer.Deserialize to preserve original types. Refactor RabbitMQ handler resolution to include anonymous handler factories by matching event names and return handler list early when concrete event type is found. Update DistributedEventBusBase to use ResolveEventForPublishing to obtain event name and data together, and ensure GetEventData is applied at the correct point when processing incoming events. --- .../Azure/AzureDistributedEventBus.cs | 38 ++++++++-------- .../Kafka/KafkaDistributedEventBus.cs | 44 +++++++++---------- .../RabbitMq/RabbitMqDistributedEventBus.cs | 16 ++++--- .../Distributed/DistributedEventBusBase.cs | 5 ++- 4 files changed, 55 insertions(+), 48 deletions(-) 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 8cf42cc796..cb85a7bf41 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 @@ -221,6 +221,20 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen return new EventHandlerFactoryUnregistrar(this, eventType, factory); } + + public override IDisposable Subscribe(string eventName, IEventHandlerFactory handler) + { + var handlerFactories = GetOrCreateAnonymousHandlerFactories(eventName); + + if (handler.IsInFactories(handlerFactories)) + { + return NullDisposable.Instance; + } + + handlerFactories.Add(handler); + + return new AnonymousEventHandlerFactoryUnregistrar(this, eventName, handler); + } public override void Unsubscribe(Func action) { @@ -267,6 +281,12 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen 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) { @@ -349,24 +369,6 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen 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(); 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 17d4bccc12..0be4e22bf6 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 @@ -94,8 +94,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen } else if (AnonymousHandlerFactories.ContainsKey(eventName)) { - var jsonElement = JsonSerializer.Deserialize(message.Value); - eventData = new AnonymousEventData(eventName, jsonElement); + var element = Serializer.Deserialize(message.Value); + eventData = new AnonymousEventData(eventName, element); eventType = typeof(AnonymousEventData); } else @@ -127,6 +127,19 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen return new EventHandlerFactoryUnregistrar(this, eventType, factory); } + + 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(Func action) @@ -175,6 +188,11 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen { 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) @@ -318,8 +336,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen } else if (AnonymousHandlerFactories.ContainsKey(incomingEvent.EventName)) { - var jsonElement = JsonSerializer.Deserialize(incomingEvent.EventData); - eventData = new AnonymousEventData(incomingEvent.EventName, jsonElement); + var element = Serializer.Deserialize(incomingEvent.EventData); + eventData = new AnonymousEventData(incomingEvent.EventName, element); eventType = typeof(AnonymousEventData); } else @@ -410,24 +428,6 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen 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(); 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 b34e97d62b..31207da4ab 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 @@ -458,16 +458,20 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis } ); } - + 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))) + { + handlerFactoryList.Add(new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value)); + } - foreach (var handlerFactory in - HandlerFactories.Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key))) + foreach (var handlerFactory in AnonymousHandlerFactories.Where(aehf => eventNames.Contains(aehf.Key))) { - handlerFactoryList.Add( - new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value)); + handlerFactoryList.Add(new EventTypeWithEventHandlerFactories(typeof(AnonymousEventData), handlerFactory.Value)); } return handlerFactoryList.ToArray(); @@ -490,7 +494,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis var eventType = GetEventTypeByEventName(eventName); if (eventType != null) { - result.AddRange(GetHandlerFactories(eventType)); + return GetHandlerFactories(eventType); } foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName)) 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 ff9d587769..9dca5a7466 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 @@ -124,8 +124,7 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB if (outboxConfig.Selector == null || outboxConfig.Selector(eventType)) { var eventOutbox = (IEventOutbox)unitOfWork.ServiceProvider.GetRequiredService(outboxConfig.ImplementationType); - var eventName = GetEventName(eventType, eventData); - eventData = GetEventData(eventData); + (var eventName, eventData) = ResolveEventForPublishing(eventType, eventData); await OnAddToOutboxAsync(eventName, eventType, eventData); @@ -186,6 +185,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB continue; } } + + eventData = GetEventData(eventData); var incomingEventInfo = new IncomingEventInfo( GuidGenerator.Create(),