diff --git a/docs/en/Distributed-Event-Bus-Azure-Integration.md b/docs/en/Distributed-Event-Bus-Azure-Integration.md new file mode 100644 index 0000000000..ca66cbca8a --- /dev/null +++ b/docs/en/Distributed-Event-Bus-Azure-Integration.md @@ -0,0 +1,133 @@ +# Distributed Event Bus Azure Integration + +> This document explains **how to configure the [Azure Service Bus](https://azure.microsoft.com/en-us/services/service-bus/)** as the distributed event bus provider. See the [distributed event bus document](Distributed-Event-Bus.md) to learn how to use the distributed event bus system + +## Installation + +Use the ABP CLI to add [Volo.Abp.EventBus.Azure](https://www.nuget.org/packages/Volo.Abp.EventBus.Azure) NuGet package to your project: + +* Install the [ABP CLI](https://docs.abp.io/en/abp/latest/CLI) if you haven't installed before. +* Open a command line (terminal) in the directory of the `.csproj` file you want to add the `Volo.Abp.EventBus.Azure` package. +* Run `abp add-package Volo.Abp.EventBus.Azure` command. + +If you want to do it manually, install the [Volo.Abp.EventBus.Azure](https://www.nuget.org/packages/Volo.Abp.EventBus.Azure) NuGet package to your project and add `[DependsOn(typeof(AbpEventBusAzureModule))]` to the [ABP module](Module-Development-Basics.md) class inside your project. + +## Configuration + +You can configure using the standard [configuration system](Configuration.md), like using the `appsettings.json` file, or using the [options](Options.md) classes. + +### `appsettings.json` file configuration + +This is the simplest way to configure the Azure Service Bus settings. It is also very strong since you can use any other configuration source (like environment variables) that is [supported by the AspNet Core](https://docs.microsoft.com/en-us/aspnet/core/fundamentals/configuration/). + +**Example: The minimal configuration to connect to Azure Service Bus Namespace with default configurations** + +````json +{ + "Azure": { + "ServiceBus": { + "Connections": { + "Default": { + "ConnectionString": "Endpoint=sb://sb-my-app.servicebus.windows.net/;SharedAccessKeyName={{Policy Name}};SharedAccessKey={};EntityPath=marketing-consent" + } + } + }, + "EventBus": { + "ConnectionName": "Default", + "SubscriberName": "MySubscriberName", + "TopicName": "MyTopicName" + } + } +} +```` + +* `MySubscriberName` is the name of this subscription, which is used as the **Subscriber** on the Azure Service Bus. +* `MyTopicName` is the **topic name**. + +See [the Azure Service Bus document](https://docs.microsoft.com/en-us/azure/service-bus-messaging/service-bus-queues-topics-subscriptions) to understand these options better. + +#### Connections + +If you need to connect to another Azure Service Bus Namespace the Default, you need to configure the connection properties. + +**Example: Declare two connections and use one of them for the event bus** + +````json +{ + "Azure": { + "ServiceBus": { + "Connections": { + "Default": { + "ConnectionString": "Endpoint=sb://sb-my-app.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey={{SharedAccessKey}}" + }, + "SecondConnection": { + "ConnectionString": "Endpoint=sb://sb-my-app.servicebus.windows.net/;SharedAccessKeyName={{Policy Name}};SharedAccessKey={{SharedAccessKey}}" + } + } + }, + "EventBus": { + "ConnectionName": "SecondConnection", + "SubscriberName": "MySubscriberName", + "TopicName": "MyTopicName" + } + } +} +```` + +This allows you to use multiple Azure Service Bus namespaces in your application, but select one of them for the event bus. + +You can use any of the [ServiceBusAdministrationClientOptions](https://docs.microsoft.com/en-us/dotnet/api/azure.messaging.servicebus.administration.servicebusadministrationclientoptions?view=azure-dotnet), [ServiceBusClientOptions](https://docs.microsoft.com/en-us/dotnet/api/azure.messaging.servicebus.servicebusclientoptions?view=azure-dotnet), [ServiceBusProcessorOptions](https://docs.microsoft.com/en-us/dotnet/api/azure.messaging.servicebus.servicebusprocessoroptions?view=azure-dotnet) properties for the connection. + +**Example: Specify the Admin, Client and Processor options** + +````json +{ + "Azure": { + "ServiceBus": { + "Connections": { + "Default": { + "ConnectionString": "Endpoint=sb://sb-my-app.servicebus.windows.net/;SharedAccessKeyName={{Policy Name}};SharedAccessKey={};EntityPath=marketing-consent", + "Admin": { + "Retry": { + "MaxRetries": 3 + } + }, + "Client": { + "RetryOptions": { + "MaxRetries": 1 + } + }, + "Processor": { + "AutoCompleteMessages": true, + "ReceiveMode": "ReceiveAndDelete" + } + } + } + }, + "EventBus": { + "ConnectionName": "Default", + "SubscriberName": "MySubscriberName", + "TopicName": "MyTopicName" + } + } +} +```` + +### The Options Classes + +`AbpAzureServiceBusOptions` and `AbpAzureEventBusOptions` classes can be used to configure the connection strings and event bus options for Azure Service Bus. + +You can configure this options inside the `ConfigureServices` of your [module](Module-Development-Basics.md). + +**Example: Configure the connection** + +````csharp +Configure(options => +{ + options.Connections.Default.ConnectionString = "Endpoint=sb://sb-my-app.servicebus.windows.net/;SharedAccessKeyName={{Policy Name}};SharedAccessKey={}"; + options.Connections.Default.Admin.Retry.MaxRetries = 3; + options.Connections.Default.Client.RetryOptions.MaxRetries = 1; +}); +```` + +Using these options classes can be combined with the `appsettings.json` way. Configuring an option property in the code overrides the value in the configuration file. diff --git a/docs/en/Distributed-Event-Bus.md b/docs/en/Distributed-Event-Bus.md index 05ef16a59e..1e8549d8f0 100644 --- a/docs/en/Distributed-Event-Bus.md +++ b/docs/en/Distributed-Event-Bus.md @@ -7,6 +7,7 @@ Distributed Event bus system allows to **publish** and **subscribe** to events t Distributed event bus system provides an **abstraction** that can be implemented by any vendor/provider. There are four providers implemented out of the box: * `LocalDistributedEventBus` is the default implementation that implements the distributed event bus to work as in-process. Yes! The **default implementation works just like the [local event bus](Local-Event-Bus.md)**, if you don't configure a real distributed provider. +* `AzureDistributedEventBus` implements the distributed event bus with the [Azure Service Bus](https://azure.microsoft.com/en-us/services/service-bus/). See the [Azure Service Bus integration document](Distributed-Event-Bus-Azure-Integration.md) to learn how to configure it. * `RabbitMqDistributedEventBus` implements the distributed event bus with the [RabbitMQ](https://www.rabbitmq.com/). See the [RabbitMQ integration document](Distributed-Event-Bus-RabbitMQ-Integration.md) to learn how to configure it. * `KafkaDistributedEventBus` implements the distributed event bus with the [Kafka](https://kafka.apache.org/). See the [Kafka integration document](Distributed-Event-Bus-Kafka-Integration.md) to learn how to configure it. * `RebusDistributedEventBus` implements the distributed event bus with the [Rebus](http://mookid.dk/category/rebus/). See the [Rebus integration document](Distributed-Event-Bus-Rebus-Integration.md) to learn how to configure it. diff --git a/docs/en/docs-nav.json b/docs/en/docs-nav.json index 76e3f70d1c..6bffa3d5ec 100644 --- a/docs/en/docs-nav.json +++ b/docs/en/docs-nav.json @@ -328,6 +328,10 @@ "text": "Distributed Event Bus", "path": "Distributed-Event-Bus.md", "items": [ + { + "text": "Azure Service Bus Integration", + "path": "Distributed-Event-Bus-Azure-Integration.md" + }, { "text": "RabbitMQ Integration", "path": "Distributed-Event-Bus-RabbitMQ-Integration.md" diff --git a/framework/Volo.Abp.sln b/framework/Volo.Abp.sln index 6ea15635e9..309be85ea1 100644 --- a/framework/Volo.Abp.sln +++ b/framework/Volo.Abp.sln @@ -369,6 +369,9 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Volo.Abp.AspNetCore.Compone EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Volo.Abp.AspNetCore.Mvc.UI.Bundling.Abstractions", "src\Volo.Abp.AspNetCore.Mvc.UI.Bundling.Abstractions\Volo.Abp.AspNetCore.Mvc.UI.Bundling.Abstractions.csproj", "{E9CE58DB-0789-4D18-8B63-474F7D7B14B4}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.AzureServiceBus", "src\Volo.Abp.AzureServiceBus\Volo.Abp.AzureServiceBus.csproj", "{808EC18E-C8CC-4F5C-82B6-984EADBBF85D}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.EventBus.Azure", "src\Volo.Abp.EventBus.Azure\Volo.Abp.EventBus.Azure.csproj", "{FB27F78E-F10E-4810-9B8E-BCD67DCFC8A2}" Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.Authorization.Abstractions", "src\Volo.Abp.Authorization.Abstractions\Volo.Abp.Authorization.Abstractions.csproj", "{87B0C2A8-FE95-4779-8B9C-2181AA52B3FA}" EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.TextTemplating.Core", "src\Volo.Abp.TextTemplating.Core\Volo.Abp.TextTemplating.Core.csproj", "{184E859A-282D-44D7-B8E9-FEA874644013}" @@ -1124,6 +1127,14 @@ Global {E9CE58DB-0789-4D18-8B63-474F7D7B14B4}.Debug|Any CPU.Build.0 = Debug|Any CPU {E9CE58DB-0789-4D18-8B63-474F7D7B14B4}.Release|Any CPU.ActiveCfg = Release|Any CPU {E9CE58DB-0789-4D18-8B63-474F7D7B14B4}.Release|Any CPU.Build.0 = Release|Any CPU + {808EC18E-C8CC-4F5C-82B6-984EADBBF85D}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {808EC18E-C8CC-4F5C-82B6-984EADBBF85D}.Debug|Any CPU.Build.0 = Debug|Any CPU + {808EC18E-C8CC-4F5C-82B6-984EADBBF85D}.Release|Any CPU.ActiveCfg = Release|Any CPU + {808EC18E-C8CC-4F5C-82B6-984EADBBF85D}.Release|Any CPU.Build.0 = Release|Any CPU + {FB27F78E-F10E-4810-9B8E-BCD67DCFC8A2}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {FB27F78E-F10E-4810-9B8E-BCD67DCFC8A2}.Debug|Any CPU.Build.0 = Debug|Any CPU + {FB27F78E-F10E-4810-9B8E-BCD67DCFC8A2}.Release|Any CPU.ActiveCfg = Release|Any CPU + {FB27F78E-F10E-4810-9B8E-BCD67DCFC8A2}.Release|Any CPU.Build.0 = Release|Any CPU {87B0C2A8-FE95-4779-8B9C-2181AA52B3FA}.Debug|Any CPU.ActiveCfg = Debug|Any CPU {87B0C2A8-FE95-4779-8B9C-2181AA52B3FA}.Debug|Any CPU.Build.0 = Debug|Any CPU {87B0C2A8-FE95-4779-8B9C-2181AA52B3FA}.Release|Any CPU.ActiveCfg = Release|Any CPU @@ -1362,6 +1373,8 @@ Global {863C18F9-2407-49F9-9ADC-F6229AF3B385} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} {B4B6B7DE-9798-4007-B1DF-7EE7929E392A} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} {E9CE58DB-0789-4D18-8B63-474F7D7B14B4} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} + {808EC18E-C8CC-4F5C-82B6-984EADBBF85D} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} + {FB27F78E-F10E-4810-9B8E-BCD67DCFC8A2} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} {87B0C2A8-FE95-4779-8B9C-2181AA52B3FA} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} {184E859A-282D-44D7-B8E9-FEA874644013} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} {228723E6-FA6D-406B-B8F8-C9BCC547AF8E} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6} diff --git a/framework/src/Volo.Abp.AzureServiceBus/FodyWeavers.xml b/framework/src/Volo.Abp.AzureServiceBus/FodyWeavers.xml new file mode 100644 index 0000000000..be0de3a908 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/FodyWeavers.xml @@ -0,0 +1,3 @@ + + + \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/FodyWeavers.xsd b/framework/src/Volo.Abp.AzureServiceBus/FodyWeavers.xsd new file mode 100644 index 0000000000..3f3946e282 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/FodyWeavers.xsd @@ -0,0 +1,30 @@ + + + + + + + + + + + + + + + 'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed. + + + + + A comma-separated list of error codes that can be safely ignored in assembly verification. + + + + + 'false' to turn off automatic generation of the XML Schema file. + + + + + \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo.Abp.AzureServiceBus.csproj b/framework/src/Volo.Abp.AzureServiceBus/Volo.Abp.AzureServiceBus.csproj new file mode 100644 index 0000000000..51cab2b9bf --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo.Abp.AzureServiceBus.csproj @@ -0,0 +1,23 @@ + + + + + + + netstandard2.0 + Volo.Abp.AzureServiceBus + Volo.Abp.AzureServiceBus + $(AssetTargetFallback);portable-net45+win8+wp8+wpa81; + false + false + false + + + + + + + + + + diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AbpAzureServiceBusModule.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AbpAzureServiceBusModule.cs new file mode 100644 index 0000000000..e9d1c20cd3 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AbpAzureServiceBusModule.cs @@ -0,0 +1,20 @@ +using Microsoft.Extensions.DependencyInjection; +using Volo.Abp.Json; +using Volo.Abp.Modularity; +using Volo.Abp.Threading; + +namespace Volo.Abp.AzureServiceBus +{ + [DependsOn( + typeof(AbpJsonModule), + typeof(AbpThreadingModule) + )] + public class AbpAzureServiceBusModule : AbpModule + { + public override void ConfigureServices(ServiceConfigurationContext context) + { + var configuration = context.Services.GetConfiguration(); + Configure(configuration.GetSection("Azure:ServiceBus")); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AbpAzureServiceBusOptions.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AbpAzureServiceBusOptions.cs new file mode 100644 index 0000000000..7672502d28 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AbpAzureServiceBusOptions.cs @@ -0,0 +1,12 @@ +namespace Volo.Abp.AzureServiceBus +{ + public class AbpAzureServiceBusOptions + { + public AzureServiceBusConnections Connections { get; } + + public AbpAzureServiceBusOptions() + { + Connections = new AzureServiceBusConnections(); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusConnections.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusConnections.cs new file mode 100644 index 0000000000..f8b0683d96 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusConnections.cs @@ -0,0 +1,31 @@ +using System; +using System.Collections.Generic; +using JetBrains.Annotations; + +namespace Volo.Abp.AzureServiceBus +{ + [Serializable] + public class AzureServiceBusConnections : Dictionary + { + public const string DefaultConnectionName = "Default"; + + [NotNull] + public ClientConfig Default + { + get => this[DefaultConnectionName]; + set => this[DefaultConnectionName] = Check.NotNull(value, nameof(value)); + } + + public AzureServiceBusConnections() + { + Default = new ClientConfig(); + } + + public ClientConfig GetOrDefault(string connectionName) + { + return TryGetValue(connectionName, out var connectionFactory) + ? connectionFactory + : Default; + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusMessageConsumer.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusMessageConsumer.cs new file mode 100644 index 0000000000..6eeb8f7693 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusMessageConsumer.cs @@ -0,0 +1,104 @@ +using System; +using System.Collections.Concurrent; +using System.Threading; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; +using JetBrains.Annotations; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using Volo.Abp.DependencyInjection; +using Volo.Abp.ExceptionHandling; +using Volo.Abp.Threading; + +namespace Volo.Abp.AzureServiceBus +{ + public class AzureServiceBusMessageConsumer : IAzureServiceBusMessageConsumer, ITransientDependency + { + public ILogger Logger { get; set; } + + private readonly IExceptionNotifier _exceptionNotifier; + private readonly IProcessorPool _processorPool; + private readonly ConcurrentBag> _callbacks; + private string _connectionName; + private string _subscriptionName; + private string _topicName; + + public AzureServiceBusMessageConsumer( + IExceptionNotifier exceptionNotifier, + IProcessorPool processorPool) + { + _exceptionNotifier = exceptionNotifier; + _processorPool = processorPool; + Logger = NullLogger.Instance; + _callbacks = new ConcurrentBag>(); + } + + public virtual void Initialize( + [NotNull] string topicName, + [NotNull] string subscriptionName, + string connectionName) + { + Check.NotNull(topicName, nameof(topicName)); + Check.NotNull(subscriptionName, nameof(subscriptionName)); + + _topicName = topicName; + _connectionName = connectionName ?? AzureServiceBusConnections.DefaultConnectionName; + _subscriptionName = subscriptionName; + StartProcessing(); + } + + public void OnMessageReceived(Func callback) + { + _callbacks.Add(callback); + } + + protected virtual void StartProcessing() + { + Task.Factory.StartNew(function: async () => + { + var serviceBusProcessor = await _processorPool.GetAsync(_subscriptionName, _topicName, _connectionName); + serviceBusProcessor.ProcessErrorAsync += HandleIncomingError; + serviceBusProcessor.ProcessMessageAsync += HandleIncomingMessage; + + if (!serviceBusProcessor.IsProcessing) + { + await serviceBusProcessor.StartProcessingAsync(); + } + + while (true) + { + Thread.Sleep(1000); + } + }, TaskCreationOptions.LongRunning); + } + + protected virtual async Task HandleIncomingMessage(ProcessMessageEventArgs args) + { + try + { + foreach (var callback in _callbacks) + { + await callback(args.Message); + } + + await args.CompleteMessageAsync(args.Message); + } + catch (Exception exception) + { + await HandleError(exception); + } + } + + protected virtual async Task HandleIncomingError(ProcessErrorEventArgs args) + { + await HandleError(args.Exception); + } + + protected virtual async Task HandleError(Exception exception) + { + Logger.LogException(exception); + await _exceptionNotifier.NotifyAsync(exception); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusMessageConsumerFactory.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusMessageConsumerFactory.cs new file mode 100644 index 0000000000..d373de43cf --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/AzureServiceBusMessageConsumerFactory.cs @@ -0,0 +1,29 @@ +using System; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.AzureServiceBus +{ + public class AzureServiceBusMessageConsumerFactory : IAzureServiceBusMessageConsumerFactory, ISingletonDependency, IDisposable + { + protected IServiceScope ServiceScope { get; } + + public AzureServiceBusMessageConsumerFactory(IServiceScopeFactory serviceScopeFactory) + { + ServiceScope = serviceScopeFactory.CreateScope(); + } + + public IAzureServiceBusMessageConsumer CreateMessageConsumer(string topicName, string subscriptionName, string connectionName) + { + var processor = ServiceScope.ServiceProvider.GetRequiredService(); + processor.Initialize(topicName, subscriptionName, connectionName); + return processor; + } + + public void Dispose() + { + ServiceScope?.Dispose(); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ClientConfig.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ClientConfig.cs new file mode 100644 index 0000000000..9a4c8c2856 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ClientConfig.cs @@ -0,0 +1,16 @@ +using Azure.Messaging.ServiceBus; +using Azure.Messaging.ServiceBus.Administration; + +namespace Volo.Abp.AzureServiceBus +{ + public class ClientConfig + { + public string ConnectionString { get; set; } + + public ServiceBusAdministrationClientOptions Admin { get; set; } = new(); + + public ServiceBusClientOptions Client { get; set; } = new(); + + public ServiceBusProcessorOptions Processor { get; set; } = new(); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ConnectionPool.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ConnectionPool.cs new file mode 100644 index 0000000000..4f38250dec --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ConnectionPool.cs @@ -0,0 +1,88 @@ +using System; +using System.Collections.Concurrent; +using System.Linq; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; +using Azure.Messaging.ServiceBus.Administration; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.AzureServiceBus +{ + public class ConnectionPool : IConnectionPool, ISingletonDependency + { + public ILogger Logger { get; set; } + + private bool _isDisposed; + private readonly AbpAzureServiceBusOptions _options; + private readonly ConcurrentDictionary> _clients; + private readonly ConcurrentDictionary> _adminClients; + + public ConnectionPool(IOptions options) + { + _options = options.Value; + _clients = new ConcurrentDictionary>(); + _adminClients = new ConcurrentDictionary>(); + Logger = new NullLogger(); + } + + public ServiceBusClient GetClient(string connectionName) + { + connectionName ??= AzureServiceBusConnections.DefaultConnectionName; + return _clients.GetOrAdd( + connectionName, new Lazy(() => + { + var config = _options.Connections.GetOrDefault(connectionName); + return new ServiceBusClient(config.ConnectionString, config.Client); + }) + ).Value; + } + + public ServiceBusAdministrationClient GetAdministrationClient(string connectionName) + { + connectionName ??= AzureServiceBusConnections.DefaultConnectionName; + return _adminClients.GetOrAdd( + connectionName, new Lazy(() => + { + var config = _options.Connections.GetOrDefault(connectionName); + return new ServiceBusAdministrationClient(config.ConnectionString); + }) + ).Value; + } + + public async ValueTask DisposeAsync() + { + if (_isDisposed) + { + return; + } + + _isDisposed = true; + if (!_clients.Any()) + { + Logger.LogDebug($"Disposed connection pool with no connection in the pool."); + return; + } + + Logger.LogInformation($"Disposing connection pool ({_clients.Count} connections)."); + + foreach (var connection in _clients.Values) + { + await connection.Value.DisposeAsync(); + } + + _clients.Clear(); + + if (!_adminClients.Any()) + { + Logger.LogDebug($"Disposed admin connection pool with no admin connection in the pool."); + return; + } + + Logger.LogInformation($"Disposing admin connection pool ({_adminClients.Count} admin connections)."); + _adminClients.Clear(); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusMessageConsumer.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusMessageConsumer.cs new file mode 100644 index 0000000000..7a47e42e40 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusMessageConsumer.cs @@ -0,0 +1,11 @@ +using System; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; + +namespace Volo.Abp.AzureServiceBus +{ + public interface IAzureServiceBusMessageConsumer + { + void OnMessageReceived(Func callback); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusMessageConsumerFactory.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusMessageConsumerFactory.cs new file mode 100644 index 0000000000..0f9f20f0aa --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusMessageConsumerFactory.cs @@ -0,0 +1,21 @@ +using System.Threading.Tasks; + +namespace Volo.Abp.AzureServiceBus +{ + public interface IAzureServiceBusMessageConsumerFactory + { + /// + /// Creates a new . + /// Avoid to create too many consumers since they are + /// not disposed until end of the application. + /// + /// + /// + /// + /// + IAzureServiceBusMessageConsumer CreateMessageConsumer( + string topicName, + string subscriptionName, + string connectionName); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusSerializer.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusSerializer.cs new file mode 100644 index 0000000000..61fc9e5ca0 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IAzureServiceBusSerializer.cs @@ -0,0 +1,13 @@ +using System; + +namespace Volo.Abp.AzureServiceBus +{ + public interface IAzureServiceBusSerializer + { + byte[] Serialize(object obj); + + object Deserialize(byte[] value, Type type); + + T Deserialize(byte[] value); + } +} diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IConnectionPool.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IConnectionPool.cs new file mode 100644 index 0000000000..43ed188022 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IConnectionPool.cs @@ -0,0 +1,13 @@ +using System; +using Azure.Messaging.ServiceBus; +using Azure.Messaging.ServiceBus.Administration; + +namespace Volo.Abp.AzureServiceBus +{ + public interface IConnectionPool : IAsyncDisposable + { + ServiceBusClient GetClient(string connectionName); + + ServiceBusAdministrationClient GetAdministrationClient(string connectionName); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IProcessorPool.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IProcessorPool.cs new file mode 100644 index 0000000000..338bee393c --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IProcessorPool.cs @@ -0,0 +1,11 @@ +using System; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; + +namespace Volo.Abp.AzureServiceBus +{ + public interface IProcessorPool : IAsyncDisposable + { + Task GetAsync(string subscriptionName, string topicName, string connectionName); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IPublisherPool.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IPublisherPool.cs new file mode 100644 index 0000000000..9ae141a252 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/IPublisherPool.cs @@ -0,0 +1,11 @@ +using System; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; + +namespace Volo.Abp.AzureServiceBus +{ + public interface IPublisherPool : IAsyncDisposable + { + Task GetAsync(string topicName, string connectionName); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ProcessorPool.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ProcessorPool.cs new file mode 100644 index 0000000000..daa804e930 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ProcessorPool.cs @@ -0,0 +1,82 @@ +using System; +using System.Collections.Concurrent; +using System.Linq; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.AzureServiceBus +{ + public class ProcessorPool : IProcessorPool, ISingletonDependency + { + public ILogger Logger { get; set; } + + private bool _isDisposed; + private readonly AbpAzureServiceBusOptions _options; + private readonly IConnectionPool _connectionPool; + private readonly ConcurrentDictionary> _processors; + + public ProcessorPool( + IOptions options, + IConnectionPool connectionPool) + { + _options = options.Value; + _connectionPool = connectionPool; + _processors = new ConcurrentDictionary>(); + Logger = new NullLogger(); + } + + public async Task GetAsync(string subscriptionName, string topicName, string connectionName) + { + var admin = _connectionPool.GetAdministrationClient(connectionName); + await admin.SetupSubscriptionAsync(topicName, subscriptionName); + + return _processors.GetOrAdd( + $"{topicName}-{subscriptionName}", new Lazy(() => + { + var config = _options.Connections.GetOrDefault(connectionName); + var client = _connectionPool.GetClient(connectionName); + return client.CreateProcessor(topicName, subscriptionName, config.Processor); + }) + ).Value; + } + + public async ValueTask DisposeAsync() + { + if (_isDisposed) + { + return; + } + + _isDisposed = true; + if (!_processors.Any()) + { + Logger.LogDebug($"Disposed processor pool with no processors in the pool."); + return; + } + + Logger.LogInformation($"Disposing processor pool ({_processors.Count} processors)."); + + foreach (var item in _processors.Values) + { + var processor = item.Value; + if (processor.IsProcessing) + { + await processor.StopProcessingAsync(); + } + + if (!processor.IsClosed) + { + await processor.CloseAsync(); + } + + await processor.DisposeAsync(); + } + + _processors.Clear(); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/PublisherPool.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/PublisherPool.cs new file mode 100644 index 0000000000..57008e07cb --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/PublisherPool.cs @@ -0,0 +1,66 @@ +using System; +using System.Collections.Concurrent; +using System.Linq; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.AzureServiceBus +{ + public class PublisherPool : IPublisherPool, ISingletonDependency + { + public ILogger Logger { get; set; } + + private bool _isDisposed; + private readonly IConnectionPool _connectionPool; + private readonly ConcurrentDictionary> _publishers; + + public PublisherPool(IConnectionPool connectionPool) + { + _connectionPool = connectionPool; + _publishers = new ConcurrentDictionary>(); + Logger = new NullLogger(); + } + + public async Task GetAsync(string topicName, string connectionName) + { + var admin = _connectionPool.GetAdministrationClient(connectionName); + await admin.SetupTopicAsync(topicName); + + return _publishers.GetOrAdd( + topicName, new Lazy(() => + { + var client = _connectionPool.GetClient(connectionName); + return client.CreateSender(topicName); + }) + ).Value; + } + + public async ValueTask DisposeAsync() + { + if (_isDisposed) + { + return; + } + + _isDisposed = true; + if (!_publishers.Any()) + { + Logger.LogDebug($"Disposed publisher pool with no publisher in the pool."); + return; + } + + Logger.LogInformation($"Disposing publisher pool ({_publishers.Count} publishers)."); + + foreach (var publisher in _publishers.Values) + { + await publisher.Value.CloseAsync(); + await publisher.Value.DisposeAsync(); + } + + _publishers.Clear(); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ServiceBusAdministrationClientExtensions.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ServiceBusAdministrationClientExtensions.cs new file mode 100644 index 0000000000..c6a09a155b --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/ServiceBusAdministrationClientExtensions.cs @@ -0,0 +1,25 @@ +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus.Administration; + +namespace Volo.Abp.AzureServiceBus +{ + public static class ServiceBusAdministrationClientExtensions + { + public static async Task SetupTopicAsync(this ServiceBusAdministrationClient client, string topicName) + { + if (!await client.TopicExistsAsync(topicName)) + { + await client.CreateTopicAsync(topicName); + } + } + + public static async Task SetupSubscriptionAsync(this ServiceBusAdministrationClient client, string topicName, string subscriptionName) + { + await client.SetupTopicAsync(topicName); + if (!await client.SubscriptionExistsAsync(topicName, subscriptionName)) + { + await client.CreateSubscriptionAsync(topicName, subscriptionName); + } + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/Utf8JsonAzureServiceBusSerializer.cs b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/Utf8JsonAzureServiceBusSerializer.cs new file mode 100644 index 0000000000..d8d6fb82a2 --- /dev/null +++ b/framework/src/Volo.Abp.AzureServiceBus/Volo/Abp/AzureServiceBus/Utf8JsonAzureServiceBusSerializer.cs @@ -0,0 +1,32 @@ +using System; +using System.Text; +using Volo.Abp.DependencyInjection; +using Volo.Abp.Json; + +namespace Volo.Abp.AzureServiceBus +{ + public class Utf8JsonAzureServiceBusSerializer : IAzureServiceBusSerializer, ITransientDependency + { + private readonly IJsonSerializer _jsonSerializer; + + public Utf8JsonAzureServiceBusSerializer(IJsonSerializer jsonSerializer) + { + _jsonSerializer = jsonSerializer; + } + + public byte[] Serialize(object obj) + { + return Encoding.UTF8.GetBytes(_jsonSerializer.Serialize(obj)); + } + + public object Deserialize(byte[] value, Type type) + { + return _jsonSerializer.Deserialize(type, Encoding.UTF8.GetString(value)); + } + + public T Deserialize(byte[] value) + { + return _jsonSerializer.Deserialize(Encoding.UTF8.GetString(value)); + } + } +} diff --git a/framework/src/Volo.Abp.EventBus.Azure/FodyWeavers.xml b/framework/src/Volo.Abp.EventBus.Azure/FodyWeavers.xml new file mode 100644 index 0000000000..be0de3a908 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Azure/FodyWeavers.xml @@ -0,0 +1,3 @@ + + + \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Azure/FodyWeavers.xsd b/framework/src/Volo.Abp.EventBus.Azure/FodyWeavers.xsd new file mode 100644 index 0000000000..3f3946e282 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Azure/FodyWeavers.xsd @@ -0,0 +1,30 @@ + + + + + + + + + + + + + + + 'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed. + + + + + A comma-separated list of error codes that can be safely ignored in assembly verification. + + + + + 'false' to turn off automatic generation of the XML Schema file. + + + + + \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Azure/Volo.Abp.EventBus.Azure.csproj b/framework/src/Volo.Abp.EventBus.Azure/Volo.Abp.EventBus.Azure.csproj new file mode 100644 index 0000000000..2ffd853cdc --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Azure/Volo.Abp.EventBus.Azure.csproj @@ -0,0 +1,22 @@ + + + + + + + netstandard2.0 + Volo.Abp.EventBus.Azure + Volo.Abp.EventBus.Azure + $(AssetTargetFallback);portable-net45+win8+wp8+wpa81; + false + false + false + + + + + + + + + diff --git a/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AbpAzureEventBusOptions.cs b/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AbpAzureEventBusOptions.cs new file mode 100644 index 0000000000..cb56de424a --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AbpAzureEventBusOptions.cs @@ -0,0 +1,11 @@ +namespace Volo.Abp.EventBus.Azure +{ + public class AbpAzureEventBusOptions + { + public string ConnectionName { get; set; } + + public string SubscriberName { get; set; } + + public string TopicName { get; set; } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AbpEventBusAzureModule.cs b/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AbpEventBusAzureModule.cs new file mode 100644 index 0000000000..08feabd366 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AbpEventBusAzureModule.cs @@ -0,0 +1,28 @@ +using Microsoft.Extensions.DependencyInjection; +using Volo.Abp.AzureServiceBus; +using Volo.Abp.Modularity; + +namespace Volo.Abp.EventBus.Azure +{ + [DependsOn( + typeof(AbpEventBusModule), + typeof(AbpAzureServiceBusModule) + )] + public class AbpEventBusAzureModule : AbpModule + { + public override void ConfigureServices(ServiceConfigurationContext context) + { + var configuration = context.Services.GetConfiguration(); + + Configure(configuration.GetSection("Azure:EventBus")); + } + + public override void OnApplicationInitialization(ApplicationInitializationContext context) + { + context + .ServiceProvider + .GetRequiredService() + .Initialize(); + } + } +} \ No newline at end of file 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 new file mode 100644 index 0000000000..a719e0a8d5 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Azure/Volo/Abp/EventBus/Azure/AzureDistributedEventBus.cs @@ -0,0 +1,234 @@ +using System; +using System.Collections.Concurrent; +using System.Collections.Generic; +using System.Linq; +using System.Threading.Tasks; +using Azure.Messaging.ServiceBus; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Volo.Abp.DependencyInjection; +using Volo.Abp.EventBus.Distributed; +using Volo.Abp.AzureServiceBus; +using Volo.Abp.Guids; +using Volo.Abp.MultiTenancy; +using Volo.Abp.Threading; +using Volo.Abp.Timing; +using Volo.Abp.Uow; + +namespace Volo.Abp.EventBus.Azure +{ + [Dependency(ReplaceServices = true)] + [ExposeServices(typeof(IDistributedEventBus), typeof(AzureDistributedEventBus))] + public class AzureDistributedEventBus : DistributedEventBusBase, ISingletonDependency + { + private readonly AbpAzureEventBusOptions _options; + private readonly IAzureServiceBusMessageConsumerFactory _messageConsumerFactory; + private readonly IPublisherPool _publisherPool; + private readonly IAzureServiceBusSerializer _serializer; + private readonly ConcurrentDictionary> _handlerFactories; + private readonly ConcurrentDictionary _eventTypes; + private IAzureServiceBusMessageConsumer _consumer; + + public AzureDistributedEventBus( + IServiceScopeFactory serviceScopeFactory, + ICurrentTenant currentTenant, + IUnitOfWorkManager unitOfWorkManager, + IOptions abpDistributedEventBusOptions, + IGuidGenerator guidGenerator, + IClock clock, + IOptions abpAzureEventBusOptions, + IAzureServiceBusSerializer serializer, + IAzureServiceBusMessageConsumerFactory messageConsumerFactory, + IPublisherPool publisherPool) + : base(serviceScopeFactory, + currentTenant, + unitOfWorkManager, + abpDistributedEventBusOptions, + guidGenerator, + clock) + { + _options = abpAzureEventBusOptions.Value; + _serializer = serializer; + _messageConsumerFactory = messageConsumerFactory; + _publisherPool = publisherPool; + _handlerFactories = new ConcurrentDictionary>(); + _eventTypes = new ConcurrentDictionary(); + } + + public void Initialize() + { + _consumer = _messageConsumerFactory.CreateMessageConsumer( + _options.TopicName, + _options.SubscriberName, + _options.ConnectionName); + + _consumer.OnMessageReceived(ProcessEventAsync); + SubscribeHandlers(AbpDistributedEventBusOptions.Handlers); + } + + private async Task ProcessEventAsync(ServiceBusReceivedMessage message) + { + var eventName = message.Subject; + var eventType = _eventTypes.GetOrDefault(eventName); + if (eventType == null) + { + return; + } + + if (await AddToInboxAsync(message.MessageId, eventName, eventType, message.Body.ToArray())) + { + return; + } + + var eventData = _serializer.Deserialize(message.Body.ToArray(), eventType); + + await TriggerHandlersAsync(eventType, eventData); + } + + public override async Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) + { + await PublishAsync(outgoingEvent.EventName, outgoingEvent.EventData); + } + + public override async Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) + { + var eventType = _eventTypes.GetOrDefault(incomingEvent.EventName); + if (eventType == null) + { + return; + } + + var eventData = _serializer.Deserialize(incomingEvent.EventData, eventType); + var exceptions = new List(); + await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); + if (exceptions.Any()) + { + ThrowOriginalExceptions(eventType, exceptions); + } + } + + protected override byte[] Serialize(object eventData) + { + return _serializer.Serialize(eventData); + } + + public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) + { + var handlerFactories = GetOrCreateHandlerFactories(eventType); + + if (factory.IsInFactories(handlerFactories)) + { + return NullDisposable.Instance; + } + + handlerFactories.Add(factory); + + return new EventHandlerFactoryUnregistrar(this, eventType, factory); + } + + public override void Unsubscribe(Func action) + { + Check.NotNull(action, nameof(action)); + + GetOrCreateHandlerFactories(typeof(TEvent)) + .Locking(factories => + { + factories.RemoveAll( + factory => + { + var singleInstanceFactory = factory as SingleInstanceHandlerFactory; + if (singleInstanceFactory == null) + { + return false; + } + + var actionHandler = singleInstanceFactory.HandlerInstance as ActionEventHandler; + if (actionHandler == null) + { + return false; + } + + return actionHandler.Action == action; + }); + }); + } + + public override void Unsubscribe(Type eventType, IEventHandler handler) + { + GetOrCreateHandlerFactories(eventType) + .Locking(factories => + { + factories.RemoveAll( + factory => + factory is SingleInstanceHandlerFactory handlerFactory && + handlerFactory.HandlerInstance == handler + ); + }); + } + + public override void Unsubscribe(Type eventType, IEventHandlerFactory factory) + { + GetOrCreateHandlerFactories(eventType) + .Locking(factories => factories.Remove(factory)); + } + + public override void UnsubscribeAll(Type eventType) + { + GetOrCreateHandlerFactories(eventType) + .Locking(factories => factories.Clear()); + } + + protected override async Task PublishToEventBusAsync(Type eventType, object eventData) + { + await PublishAsync(eventType, eventData); + } + + protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord) + { + unitOfWork.AddOrReplaceDistributedEvent(eventRecord); + } + + protected virtual async Task PublishAsync(string eventName, object eventData) + { + var body = _serializer.Serialize(eventData); + + var message = new ServiceBusMessage(body) + { + Subject = eventName + }; + + var publisher = await _publisherPool.GetAsync( + _options.TopicName, + _options.ConnectionName); + + await publisher.SendMessageAsync(message); + } + + protected override IEnumerable GetHandlerFactories(Type eventType) + { + return _handlerFactories + .Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key)) + .Select(handlerFactory => + new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value)) + .ToArray(); + } + + private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) + { + return handlerEventType == targetEventType || handlerEventType.IsAssignableFrom(targetEventType); + } + + private List GetOrCreateHandlerFactories(Type eventType) + { + return _handlerFactories.GetOrAdd( + eventType, + type => + { + var eventName = EventNameAttribute.GetNameOrDefault(type); + _eventTypes[eventName] = type; + return new List(); + } + ); + } + } +} diff --git a/nupkg/common.ps1 b/nupkg/common.ps1 index 714da2465c..7d48c1ac1d 100644 --- a/nupkg/common.ps1 +++ b/nupkg/common.ps1 @@ -65,6 +65,7 @@ $projects = ( "framework/src/Volo.Abp.Autofac", "framework/src/Volo.Abp.Autofac.WebAssembly", "framework/src/Volo.Abp.AutoMapper", + "framework/src/Volo.Abp.AzureServiceBus", "framework/src/Volo.Abp.BackgroundJobs.Abstractions", "framework/src/Volo.Abp.BackgroundJobs", "framework/src/Volo.Abp.BackgroundJobs.HangFire", @@ -106,6 +107,7 @@ $projects = ( "framework/src/Volo.Abp.EventBus.RabbitMQ", "framework/src/Volo.Abp.EventBus.Kafka", "framework/src/Volo.Abp.EventBus.Rebus", + "framework/src/Volo.Abp.EventBus.Azure", "framework/src/Volo.Abp.ExceptionHandling", "framework/src/Volo.Abp.Features", "framework/src/Volo.Abp.FluentValidation",