Browse Source

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<object> 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.
pull/25023/head
SALİH ÖZKARA 7 months ago
parent
commit
aae1e330ba
  1. 38
      framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AzureDistributedEventBus.cs
  2. 44
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs
  3. 16
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  4. 5
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs

38
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); 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<TEvent>(Func<TEvent, Task> action) public override void Unsubscribe<TEvent>(Func<TEvent, Task> action)
{ {
@ -267,6 +281,12 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen
GetOrCreateHandlerFactories(eventType) GetOrCreateHandlerFactories(eventType)
.Locking(factories => factories.Remove(factory)); .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) public override void UnsubscribeAll(Type eventType)
{ {
@ -349,24 +369,6 @@ public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDepen
return EventTypes.GetOrDefault(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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName) protected override IEnumerable<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{ {
var result = new List<EventTypeWithEventHandlerFactories>(); var result = new List<EventTypeWithEventHandlerFactories>();

44
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)) else if (AnonymousHandlerFactories.ContainsKey(eventName))
{ {
var jsonElement = JsonSerializer.Deserialize<JsonElement>(message.Value); var element = Serializer.Deserialize<object>(message.Value);
eventData = new AnonymousEventData(eventName, jsonElement); eventData = new AnonymousEventData(eventName, element);
eventType = typeof(AnonymousEventData); eventType = typeof(AnonymousEventData);
} }
else else
@ -127,6 +127,19 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen
return new EventHandlerFactoryUnregistrar(this, eventType, factory); 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);
}
/// <inheritdoc/> /// <inheritdoc/>
public override void Unsubscribe<TEvent>(Func<TEvent, Task> action) public override void Unsubscribe<TEvent>(Func<TEvent, Task> action)
@ -175,6 +188,11 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen
{ {
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Remove(factory)); GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Remove(factory));
} }
public override void Unsubscribe(string eventName, IEventHandlerFactory factory)
{
GetOrCreateAnonymousHandlerFactories(eventName).Locking(factories => factories.Remove(factory));
}
/// <inheritdoc/> /// <inheritdoc/>
public override void UnsubscribeAll(Type eventType) public override void UnsubscribeAll(Type eventType)
@ -318,8 +336,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen
} }
else if (AnonymousHandlerFactories.ContainsKey(incomingEvent.EventName)) else if (AnonymousHandlerFactories.ContainsKey(incomingEvent.EventName))
{ {
var jsonElement = JsonSerializer.Deserialize<JsonElement>(incomingEvent.EventData); var element = Serializer.Deserialize<object>(incomingEvent.EventData);
eventData = new AnonymousEventData(incomingEvent.EventName, jsonElement); eventData = new AnonymousEventData(incomingEvent.EventName, element);
eventType = typeof(AnonymousEventData); eventType = typeof(AnonymousEventData);
} }
else else
@ -410,24 +428,6 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen
return EventTypes.GetOrDefault(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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName) protected override IEnumerable<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{ {
var result = new List<EventTypeWithEventHandlerFactories>(); var result = new List<EventTypeWithEventHandlerFactories>();

16
framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs

@ -458,16 +458,20 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis
} }
); );
} }
protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType) protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType)
{ {
var handlerFactoryList = new List<EventTypeWithEventHandlerFactories>(); var handlerFactoryList = new List<EventTypeWithEventHandlerFactories>();
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 foreach (var handlerFactory in AnonymousHandlerFactories.Where(aehf => eventNames.Contains(aehf.Key)))
HandlerFactories.Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key)))
{ {
handlerFactoryList.Add( handlerFactoryList.Add(new EventTypeWithEventHandlerFactories(typeof(AnonymousEventData), handlerFactory.Value));
new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value));
} }
return handlerFactoryList.ToArray(); return handlerFactoryList.ToArray();
@ -490,7 +494,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis
var eventType = GetEventTypeByEventName(eventName); var eventType = GetEventTypeByEventName(eventName);
if (eventType != null) if (eventType != null)
{ {
result.AddRange(GetHandlerFactories(eventType)); return GetHandlerFactories(eventType);
} }
foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName)) foreach (var handlerFactory in AnonymousHandlerFactories.Where(hf => hf.Key == eventName))

5
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)) if (outboxConfig.Selector == null || outboxConfig.Selector(eventType))
{ {
var eventOutbox = (IEventOutbox)unitOfWork.ServiceProvider.GetRequiredService(outboxConfig.ImplementationType); var eventOutbox = (IEventOutbox)unitOfWork.ServiceProvider.GetRequiredService(outboxConfig.ImplementationType);
var eventName = GetEventName(eventType, eventData); (var eventName, eventData) = ResolveEventForPublishing(eventType, eventData);
eventData = GetEventData(eventData);
await OnAddToOutboxAsync(eventName, eventType, eventData); await OnAddToOutboxAsync(eventName, eventType, eventData);
@ -186,6 +185,8 @@ public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventB
continue; continue;
} }
} }
eventData = GetEventData(eventData);
var incomingEventInfo = new IncomingEventInfo( var incomingEventInfo = new IncomingEventInfo(
GuidGenerator.Create(), GuidGenerator.Create(),

Loading…
Cancel
Save