Browse Source

Add anonymous event support to EventBus

Introduce AnonymousEventData and add support for anonymous (name-based) events across the event bus implementations. Adds string-based Subscribe/Unsubscribe APIs, anonymous handler factories, and handling in distributed providers (Azure, Dapr, Kafka, RabbitMQ, Rebus) and local buses. Update EventBusBase and DistributedEventBusBase to resolve event names/data (GetEventName/GetEventData/ResolveEventForPublishing) and route/serialize/deserialize anonymous payloads. Also add AnonymousEventHandlerFactoryUnregistrar and minimal NullDistributedEventBus implementations, plus tests for anonymous local events.
pull/25023/head
SALİH ÖZKARA 7 months ago
parent
commit
58e0c2f609
  1. 103
      framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/AnonymousEventData.cs
  2. 4
      framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/IEventBus.cs
  3. 2
      framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs
  4. 98
      framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AzureDistributedEventBus.cs
  5. 66
      framework/src/Volo.Abp.EventBus.Dapr/Volo/Abp/EventBus/Dapr/DaprDistributedEventBus.cs
  6. 105
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs
  7. 100
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  8. 78
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs
  9. 42
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs
  10. 67
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs
  11. 17
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/NullDistributedEventBus.cs
  12. 63
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs
  13. 19
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventHandlerFactoryUnregistrar.cs
  14. 99
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs
  15. 20
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs
  16. 179
      framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus_Test.cs
  17. 162
      framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/LocalEventBus_Anonymous_Test.cs

103
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<T>()
{
if (Data is T typedData)
{
return typedData;
}
return GetJsonElement().Deserialize<T>()
?? 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<string, object?>();
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!;
}
}
}

4
framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/IEventBus.cs

@ -52,6 +52,8 @@ public interface IEventBus
/// <param name="eventType">Event type</param>
/// <param name="handler">Object to handle the event</param>
IDisposable Subscribe(Type eventType, IEventHandler handler);
IDisposable Subscribe(string eventName, IEventHandlerFactory handler);
/// <summary>
/// Registers to an event.
@ -106,6 +108,8 @@ public interface IEventBus
/// <param name="eventType">Event type</param>
/// <param name="factory">Factory object that is registered before</param>
void Unsubscribe(Type eventType, IEventHandlerFactory factory);
void Unsubscribe(string eventName, IEventHandlerFactory factory);
/// <summary>
/// Unregisters all event handlers of given event type.

2
framework/src/Volo.Abp.EventBus.Abstractions/Volo/Abp/EventBus/Local/ILocalEventBus.cs

@ -23,4 +23,6 @@ public interface ILocalEventBus : IEventBus
/// <param name="eventType">Event type</param>
/// <returns></returns>
List<EventTypeWithEventHandlerFactories> GetEventHandlerFactories(Type eventType);
List<EventTypeWithEventHandlerFactories> GetAnonymousEventHandlerFactories(string eventName);
}

98
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<Type, List<IEventHandlerFactory>> HandlerFactories { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, List<IEventHandlerFactory>> 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<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>();
AnonymousHandlerFactories = new ConcurrentDictionary<string, List<IEventHandlerFactory>>();
}
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<object>(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<object>(incomingEvent.EventData);
eventData = new AnonymousEventData(incomingEvent.EventName, element);
eventType = typeof(AnonymousEventData);
}
else
{
return;
}
var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType);
var exceptions = new List<Exception>();
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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{
var result = new List<EventTypeWithEventHandlerFactories>();
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<IEventHandlerFactory> GetOrCreateAnonymousHandlerFactories(string eventName)
{
return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List<IEventHandlerFactory>());
}
private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType)
{
return handlerEventType == targetEventType || handlerEventType.IsAssignableFrom(targetEventType);

66
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<Type, List<IEventHandlerFactory>> HandlerFactories { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, List<IEventHandlerFactory>> AnonymousHandlerFactories { get; }
public DaprDistributedEventBus(
IServiceScopeFactory serviceScopeFactory,
@ -58,6 +59,7 @@ public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDepend
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>();
AnonymousHandlerFactories = new ConcurrentDictionary<string, List<IEventHandlerFactory>>();
}
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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{
var result = new List<EventTypeWithEventHandlerFactories>();
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<IEventHandlerFactory> GetOrCreateAnonymousHandlerFactories(string eventName)
{
return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List<IEventHandlerFactory>());
}
public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig)
{
var eventType = GetEventType(outgoingEvent.EventName);

105
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<Type, List<IEventHandlerFactory>> HandlerFactories { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, List<IEventHandlerFactory>> AnonymousHandlerFactories { get; }
protected IKafkaMessageConsumer Consumer { get; private set; } = default!;
public KafkaDistributedEventBus(
@ -63,6 +65,7 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>();
AnonymousHandlerFactories = new ConcurrentDictionary<string, List<IEventHandlerFactory>>();
}
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<JsonElement>(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<JsonElement>(incomingEvent.EventData);
eventData = new AnonymousEventData(incomingEvent.EventName, jsonElement);
eventType = typeof(AnonymousEventData);
}
else
{
return;
}
var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType);
var exceptions = new List<Exception>();
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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{
var result = new List<EventTypeWithEventHandlerFactories>();
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<IEventHandlerFactory> GetOrCreateAnonymousHandlerFactories(string eventName)
{
return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List<IEventHandlerFactory>());
}
private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType)
{
//Should trigger same type

100
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<IEventHandlerFactory> may not be thread-safe!
protected ConcurrentDictionary<Type, List<IEventHandlerFactory>> HandlerFactories { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, List<IEventHandlerFactory>> 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<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>();
AnonymousHandlerFactories = new ConcurrentDictionary<string, List<IEventHandlerFactory>>();
}
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<JsonElement>(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<JsonElement>(incomingEvent.EventData);
eventData = new AnonymousEventData(incomingEvent.EventName, jsonElement);
eventType = typeof(AnonymousEventData);
}
else
{
return;
}
var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType);
var exceptions = new List<Exception>();
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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{
var result = new List<EventTypeWithEventHandlerFactories>();
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<IEventHandlerFactory> GetOrCreateAnonymousHandlerFactories(string eventName)
{
return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List<IEventHandlerFactory>());
}
private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType)
{
//Should trigger same type

78
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<IEventHandlerFactory> may not be thread-safe!
protected ConcurrentDictionary<Type, List<IEventHandlerFactory>> HandlerFactories { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, List<IEventHandlerFactory>> AnonymousHandlerFactories { get; }
protected AbpRebusEventBusOptions AbpRebusEventBusOptions { get; }
public RebusDistributedEventBus(
@ -63,6 +64,7 @@ public class RebusDistributedEventBus : DistributedEventBusBase, ISingletonDepen
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>();
AnonymousHandlerFactories = new ConcurrentDictionary<string, List<IEventHandlerFactory>>();
}
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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{
var result = new List<EventTypeWithEventHandlerFactories>();
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<IEventHandlerFactory> GetOrCreateAnonymousHandlerFactories(string eventName)
{
return AnonymousHandlerFactories.GetOrAdd(eventName, _ => new List<IEventHandlerFactory>());
}
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<Exception>();
using (CorrelationIdProvider.Change(incomingEvent.GetCorrelationId()))
{

42
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));
}
}

67
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<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, bool> AnonymousEventNames { get; }
public LocalDistributedEventBus(
IServiceScopeFactory serviceScopeFactory,
ICurrentTenant currentTenant,
@ -47,6 +48,7 @@ public class LocalDistributedEventBus : DistributedEventBusBase, ISingletonDepen
correlationIdProvider)
{
EventTypes = new ConcurrentDictionary<string, Type>();
AnonymousEventNames = new ConcurrentDictionary<string, bool>();
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<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{
return LocalEventBus.GetAnonymousEventHandlerFactories(eventName);
}
protected override Type? GetEventTypeByEventName(string eventName)
{
return EventTypes.GetOrDefault(eventName);
}
}

17
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<TEvent>(Func<TEvent, Task> 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<TEvent>(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<TEvent>() where TEvent : class
{

63
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);
/// <inheritdoc/>
public virtual IDisposable Subscribe<TEvent>(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);
/// <inheritdoc/>
public virtual void UnsubscribeAll<TEvent>() 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<EventTypeWithEventHandlerFactories> 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<Exception> exceptions)
{
if (exceptions.Count == 1)
@ -200,6 +241,10 @@ public abstract class EventBusBase : IEventBus
protected abstract IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType);
protected abstract IEnumerable<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName);
protected abstract Type? GetEventTypeByEventName(string eventName);
protected virtual async Task TriggerHandlerAsync(IEventHandlerFactory asyncHandlerFactory, Type eventType,
object eventData, List<Exception> exceptions, InboxConfig? inboxConfig = null)
{

19
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);
}
}

99
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<Type, List<IEventHandlerFactory>> HandlerFactories { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected ConcurrentDictionary<string, List<IEventHandlerFactory>> AnonymousEventHandlerFactories { get; }
public LocalEventBus(
IOptions<AbpLocalEventBusOptions> options,
@ -45,6 +47,7 @@ public class LocalEventBus : EventBusBase, ILocalEventBus, ISingletonDependency
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>();
AnonymousEventHandlerFactories = new ConcurrentDictionary<string, List<IEventHandlerFactory>>();
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);
}
/// <inheritdoc/>
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));
}
/// <inheritdoc/>
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<EventTypeWithEventHandlerFactories> GetAnonymousEventHandlerFactories(string eventName)
{
return GetAnonymousHandlerFactories(eventName).ToList();
}
protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType)
{
var handlerFactoryList = new List<Tuple<IEventHandlerFactory, Type, int>>();
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<IEventHandlerFactory, Type, int>(
factory,
typeof(AnonymousEventData),
ReflectionHelper.GetAttributesOfMemberOrDeclaringType<LocalEventHandlerOrderAttribute>(factory.GetHandler().EventHandler.GetType()).FirstOrDefault()?.Order ?? 0));
}
}
return handlerFactoryList.OrderBy(x => x.Item3).Select(x => new EventTypeWithEventHandlerFactories(x.Item2, new List<IEventHandlerFactory> {x.Item1})).ToArray();
}
protected override IEnumerable<EventTypeWithEventHandlerFactories> GetAnonymousHandlerFactories(string eventName)
{
var eventType = EventTypes.GetOrDefault(eventName);
if (eventType != null)
{
return GetHandlerFactories(eventType);
}
var handlerFactoryList = new List<Tuple<IEventHandlerFactory, Type, int>>();
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<IEventHandlerFactory, Type, int>(
factory,
typeof(AnonymousEventData),
ReflectionHelper
.GetAttributesOfMemberOrDeclaringType<LocalEventHandlerOrderAttribute>(handlerType)
.FirstOrDefault()?.Order ?? 0));
}
}
return handlerFactoryList.OrderBy(x => x.Item3).Select(x =>
new EventTypeWithEventHandlerFactories(x.Item2, new List<IEventHandlerFactory> { x.Item1 })).ToArray();
}
protected override Type? GetEventTypeByEventName(string eventName)
{
return EventTypes.GetOrDefault(eventName);
}
private List<IEventHandlerFactory> GetOrCreateHandlerFactories(Type eventType)
{
return HandlerFactories.GetOrAdd(eventType, (type) => new List<IEventHandlerFactory>());
}
private List<IEventHandlerFactory> GetOrCreateAnonymousHandlerFactories(string eventName)
{
return AnonymousEventHandlerFactories.GetOrAdd(eventName, (name) => new List<IEventHandlerFactory>());
}
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;

20
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<TEvent>(Func<TEvent, Task> action) where TEvent : class
{
return NullDisposable.Instance;
@ -28,6 +33,11 @@ public sealed class NullLocalEventBus : ILocalEventBus
return new List<EventTypeWithEventHandlerFactories>();
}
public List<EventTypeWithEventHandlerFactories> GetAnonymousEventHandlerFactories(string eventName)
{
return new List<EventTypeWithEventHandlerFactories>();
}
public IDisposable Subscribe<TEvent, THandler>() 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<TEvent>(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<TEvent>() where TEvent : class
{

179
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<MySimpleEventData, MySimpleDistributedTransientEventHandler>();
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
await DistributedEventBus.PublishAsync(eventName, new MySimpleEventData(1));
await DistributedEventBus.PublishAsync(eventName, new Dictionary<string, object>()
{
{"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<AnonymousEventData>(async (d) =>
{
handleCount++;
await Task.CompletedTask;
})));
await DistributedEventBus.PublishAsync("MyEvent", new MySimpleEventData(1));
await DistributedEventBus.PublishAsync("MyEvent", new Dictionary<string, object>()
{
{"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<AnonymousEventData>(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<string, object>()
{
{"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<MySimpleEventData, MySimpleDistributedTransientEventHandler>();
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new MySimpleEventData(1)));
await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new Dictionary<string, object>()
{
{"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<MySimpleEventData, MySimpleDistributedTransientEventHandler>();
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
var anonymousHandleCount = 0;
DistributedEventBus.Subscribe(eventName, new SingleInstanceHandlerFactory(new ActionEventHandler<AnonymousEventData>(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<MySimpleEventData, MySimpleDistributedTransientEventHandler>();
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
var anonymousHandleCount = 0;
DistributedEventBus.Subscribe(eventName, new SingleInstanceHandlerFactory(new ActionEventHandler<AnonymousEventData>(async (d) =>
{
anonymousHandleCount++;
await Task.CompletedTask;
})));
await DistributedEventBus.PublishAsync(new MySimpleEventData(1));
await DistributedEventBus.PublishAsync(new AnonymousEventData(eventName, new Dictionary<string, object>()
{
{"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<AnonymousEventData>(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<AbpException>(() =>
DistributedEventBus.PublishAsync("MyEvent", new { Value = 2 }));
Assert.Equal(1, handleCount);
}
[Fact]
public async Task Should_Throw_For_Unknown_Event_Name()
{
await Assert.ThrowsAsync<AbpException>(() =>
DistributedEventBus.PublishAsync("NonExistentEvent", new { Value = 1 }));
}
[Fact]
public async Task Should_Convert_AnonymousEventData_To_Typed_Object()
{
MySimpleEventData? receivedData = null;
DistributedEventBus.Subscribe<MySimpleEventData>(async (data) =>
{
receivedData = data;
await Task.CompletedTask;
});
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
await DistributedEventBus.PublishAsync(eventName, new { Value = 42 });
receivedData.ShouldNotBeNull();
receivedData.Value.ShouldBe(42);
}
[Fact]
public async Task Should_Change_TenantId_If_EventData_Is_MultiTenant()
{

162
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<AnonymousEventData>(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<MySimpleEventData>(async (data) =>
{
handleCount++;
await Task.CompletedTask;
});
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
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<MySimpleEventData>(async (data) =>
{
receivedData = data;
await Task.CompletedTask;
});
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
await LocalEventBus.PublishAsync(eventName, new Dictionary<string, object>
{
{ "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<MySimpleEventData>(async (data) =>
{
typedHandleCount++;
await Task.CompletedTask;
});
var eventName = EventNameAttribute.GetNameOrDefault<MySimpleEventData>();
LocalEventBus.Subscribe(eventName,
new SingleInstanceHandlerFactory(new ActionEventHandler<AnonymousEventData>(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<AnonymousEventData>(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<AbpException>(() =>
LocalEventBus.PublishAsync("TestEvent", new { Value = 2 }));
handleCount.ShouldBe(1);
}
[Fact]
public async Task Should_Throw_For_Unknown_Event_Name()
{
await Assert.ThrowsAsync<AbpException>(() =>
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<AnonymousEventData>(async (d) =>
{
receivedData = d.ConvertToTypedObject();
await Task.CompletedTask;
})));
await LocalEventBus.PublishAsync("TestEvent", new { Name = "Hello", Count = 42 });
receivedData.ShouldNotBeNull();
var dict = receivedData.ShouldBeOfType<Dictionary<string, object?>>();
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<AnonymousEventData>(async (d) =>
{
receivedData = d.ConvertToTypedObject<MySimpleEventData>();
await Task.CompletedTask;
})));
await LocalEventBus.PublishAsync("TestEvent", new MySimpleEventData(99));
receivedData.ShouldNotBeNull();
receivedData.Value.ShouldBe(99);
}
}
Loading…
Cancel
Save