From 88609d5a62e32d6029ec7dc77ee7e8f1f253d814 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 21 Feb 2022 16:13:55 +0100 Subject: [PATCH] Incremental work on Azure Service Bus module --- Elsa.sln | 14 ++++ .../Elsa.Core/Attributes/ActivityAttribute.cs | 6 +- .../Models/ActivityExecutionContext.cs | 4 +- .../Elsa.Formatting/Contracts/IFormatter.cs | 10 +++ .../Elsa.Formatting/Elsa.Formatting.csproj | 9 +++ .../Formatters/JsonFormatter.cs | 20 +++++ .../Activities/Send.cs | 74 +++++++++++++++++++ .../Contracts/IQueueProvider.cs | 11 +++ .../Contracts/IServiceBusInitializer.cs | 12 +++ .../Contracts/ISubscriptionProvider.cs | 11 +++ .../Contracts/ITopicProvider.cs | 11 +++ .../Elsa.Modules.AzureServiceBus.csproj | 23 ++++++ .../Extensions/ServiceCollectionExtensions.cs | 48 ++++++++++++ .../CreateQueuesTopicsAndSubscriptions.cs | 15 ++++ .../Models/QueueDefinition.cs | 9 +++ .../Models/SubscriptionDefinition.cs | 10 +++ .../Models/TopicDefinition.cs | 9 +++ .../Options/AzureServiceBusOptions.cs | 11 +++ ...rationQueueTopicAndSubscriptionProvider.cs | 18 +++++ .../Services/ServiceBusInitializer.cs | 63 ++++++++++++++++ .../Elsa.Modules.Quartz/AssemblyInfo.cs | 1 - .../Elsa.Samples.Web1.csproj | 17 +++-- .../aspnet/Elsa.Samples.Web1/Program.cs | 8 +- .../Workflows/AzureServiceBusWorkflow.cs | 19 +++++ .../aspnet/Elsa.Samples.Web1/appsettings.json | 11 ++- .../console/Elsa.Samples.Console1/Program.cs | 2 +- 26 files changed, 431 insertions(+), 15 deletions(-) create mode 100644 src/core/Elsa.Formatting/Contracts/IFormatter.cs create mode 100644 src/core/Elsa.Formatting/Elsa.Formatting.csproj create mode 100644 src/core/Elsa.Formatting/Formatters/JsonFormatter.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Activities/Send.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Contracts/IQueueProvider.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Contracts/IServiceBusInitializer.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Contracts/ISubscriptionProvider.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Contracts/ITopicProvider.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Elsa.Modules.AzureServiceBus.csproj create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/HostedServices/CreateQueuesTopicsAndSubscriptions.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Models/QueueDefinition.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Models/SubscriptionDefinition.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Models/TopicDefinition.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Options/AzureServiceBusOptions.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Providers/ConfigurationQueueTopicAndSubscriptionProvider.cs create mode 100644 src/modules/Elsa.Modules.AzureServiceBus/Services/ServiceBusInitializer.cs delete mode 100644 src/modules/Elsa.Modules.Quartz/AssemblyInfo.cs create mode 100644 src/samples/aspnet/Elsa.Samples.Web1/Workflows/AzureServiceBusWorkflow.cs diff --git a/Elsa.sln b/Elsa.sln index 986a26735..3d8134ca8 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -82,6 +82,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Modules.Http", "src\mo EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Modules.Scheduling", "src\modules\Elsa.Modules.Scheduling\Elsa.Modules.Scheduling.csproj", "{ACD65CE5-3CC2-47B1-BFAC-72443D764F6E}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Modules.AzureServiceBus", "src\modules\Elsa.Modules.AzureServiceBus\Elsa.Modules.AzureServiceBus.csproj", "{24C7095B-5CDD-4369-9F34-5565A019B195}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Formatting", "src\core\Elsa.Formatting\Elsa.Formatting.csproj", "{49716A83-239C-4913-BC11-E379ED2F676E}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -176,6 +180,14 @@ Global {ACD65CE5-3CC2-47B1-BFAC-72443D764F6E}.Debug|Any CPU.Build.0 = Debug|Any CPU {ACD65CE5-3CC2-47B1-BFAC-72443D764F6E}.Release|Any CPU.ActiveCfg = Release|Any CPU {ACD65CE5-3CC2-47B1-BFAC-72443D764F6E}.Release|Any CPU.Build.0 = Release|Any CPU + {24C7095B-5CDD-4369-9F34-5565A019B195}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {24C7095B-5CDD-4369-9F34-5565A019B195}.Debug|Any CPU.Build.0 = Debug|Any CPU + {24C7095B-5CDD-4369-9F34-5565A019B195}.Release|Any CPU.ActiveCfg = Release|Any CPU + {24C7095B-5CDD-4369-9F34-5565A019B195}.Release|Any CPU.Build.0 = Release|Any CPU + {49716A83-239C-4913-BC11-E379ED2F676E}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {49716A83-239C-4913-BC11-E379ED2F676E}.Debug|Any CPU.Build.0 = Debug|Any CPU + {49716A83-239C-4913-BC11-E379ED2F676E}.Release|Any CPU.ActiveCfg = Release|Any CPU + {49716A83-239C-4913-BC11-E379ED2F676E}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(NestedProjects) = preSolution {155227F0-A33B-40AA-A4B4-06F813EB921B} = {61017E64-6D00-49CB-9E81-5002DC8F7D5F} @@ -212,5 +224,7 @@ Global {50697882-E38B-4CE7-B041-8009C18490E0} = {C6658DE0-2B2F-47F0-BB61-2CA66D435C09} {82E26BF9-5F3A-4365-899A-AB1FFD54AA45} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} {ACD65CE5-3CC2-47B1-BFAC-72443D764F6E} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} + {24C7095B-5CDD-4369-9F34-5565A019B195} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} + {49716A83-239C-4913-BC11-E379ED2F676E} = {C6658DE0-2B2F-47F0-BB61-2CA66D435C09} EndGlobalSection EndGlobal diff --git a/src/core/Elsa.Core/Attributes/ActivityAttribute.cs b/src/core/Elsa.Core/Attributes/ActivityAttribute.cs index 379069559..434fe631a 100644 --- a/src/core/Elsa.Core/Attributes/ActivityAttribute.cs +++ b/src/core/Elsa.Core/Attributes/ActivityAttribute.cs @@ -3,10 +3,14 @@ namespace Elsa.Attributes; [AttributeUsage(AttributeTargets.Class)] public class ActivityAttribute : Attribute { - public ActivityAttribute(string? typeName = default) + public ActivityAttribute(string? typeName = default, string? description = default, string? category = default) { TypeName = typeName; + Description = description; + Category = category; } public string? TypeName { get; } + public string? Description { get; } + public string? Category { get; } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Models/ActivityExecutionContext.cs b/src/core/Elsa.Core/Models/ActivityExecutionContext.cs index bdda0a21e..280f25e90 100644 --- a/src/core/Elsa.Core/Models/ActivityExecutionContext.cs +++ b/src/core/Elsa.Core/Models/ActivityExecutionContext.cs @@ -59,7 +59,7 @@ public class ActivityExecutionContext public void SetBookmark(object? payload, ExecuteActivityDelegate? callback = default) { var hasher = GetRequiredService(); - + var identityGenerator = GetRequiredService(); var payloadSerializer = GetRequiredService(); var payloadJson = payload != null ? payloadSerializer.Serialize(payload) : default; @@ -87,7 +87,7 @@ public class ActivityExecutionContext } public T GetRequiredService() where T : notnull => WorkflowExecutionContext.GetRequiredService(); - public T? Get(Input input) => Get(input.LocationReference); + public T? Get(Input? input) => input == null ? default : Get(input.LocationReference); public object? Get(RegisterLocationReference locationReference) { diff --git a/src/core/Elsa.Formatting/Contracts/IFormatter.cs b/src/core/Elsa.Formatting/Contracts/IFormatter.cs new file mode 100644 index 000000000..e9002cf10 --- /dev/null +++ b/src/core/Elsa.Formatting/Contracts/IFormatter.cs @@ -0,0 +1,10 @@ +namespace Elsa.Formatting.Contracts; + +/// +/// Represents a formatter that can serialize an object to a string and deserialize a string into an object. +/// +public interface IFormatter +{ + ValueTask ToStringAsync(object body, CancellationToken cancellationToken = default); + ValueTask FromStringAsync(string data, Type? returnType, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/core/Elsa.Formatting/Elsa.Formatting.csproj b/src/core/Elsa.Formatting/Elsa.Formatting.csproj new file mode 100644 index 000000000..eb2460e91 --- /dev/null +++ b/src/core/Elsa.Formatting/Elsa.Formatting.csproj @@ -0,0 +1,9 @@ + + + + net6.0 + enable + enable + + + diff --git a/src/core/Elsa.Formatting/Formatters/JsonFormatter.cs b/src/core/Elsa.Formatting/Formatters/JsonFormatter.cs new file mode 100644 index 000000000..98687c3f1 --- /dev/null +++ b/src/core/Elsa.Formatting/Formatters/JsonFormatter.cs @@ -0,0 +1,20 @@ +using System.Text.Json; +using Elsa.Formatting.Contracts; + +namespace Elsa.Formatting.Formatters; + +public class JsonFormatter : IFormatter +{ + public ValueTask ToStringAsync(object body, CancellationToken cancellationToken) + { + var json = JsonSerializer.Serialize(body); + return ValueTask.FromResult(json); + } + + public ValueTask FromStringAsync(string data, Type? returnType, CancellationToken cancellationToken) + { + var options = new JsonSerializerOptions(); + var value = returnType != null ? JsonSerializer.Deserialize(data, returnType, options)! : JsonSerializer.Deserialize(data, options)!; + return ValueTask.FromResult(value); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Activities/Send.cs b/src/modules/Elsa.Modules.AzureServiceBus/Activities/Send.cs new file mode 100644 index 000000000..8ebefcab4 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Activities/Send.cs @@ -0,0 +1,74 @@ +using Azure.Messaging.ServiceBus; +using Elsa.Attributes; +using Elsa.Formatting.Contracts; +using Elsa.Formatting.Formatters; +using Elsa.Models; +using Microsoft.Extensions.Logging; + +namespace Elsa.Modules.AzureServiceBus.Activities; + +[Activity("Azure.ServiceBus.Send", "Send a message to a queue or topic", "Azure Service Bus")] +public class Send : Activity +{ + public Input MessageBody { get; set; } = default!; + public Input QueueOrTopic { get; set; } = default!; + public Input? ContentType { get; set; } + public Input? Subject { get; set; } + public Input? CorrelationId { get; set; } + public IFormatter? Formatter { get; set; } + + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) + { + var queueOrTopic = context.Get(QueueOrTopic); + var messageBody = context.Get(MessageBody); + + if (!ValidatePreconditions(context, queueOrTopic, messageBody)) + return; + + var cancellationToken = context.CancellationToken; + var serializedMessageBody = await SerializeMessageBodyAsync(messageBody!, cancellationToken); + + var message = new ServiceBusMessage(serializedMessageBody) + { + ContentType = context.Get(ContentType), + Subject = context.Get(Subject), + CorrelationId = context.Get(CorrelationId) + + // TODO: Maybe expose additional members. + }; + + var client = context.GetRequiredService(); + + await using var sender = client.CreateSender(queueOrTopic); + await sender.SendMessageAsync(message, cancellationToken); + } + + private async ValueTask SerializeMessageBodyAsync(object value, CancellationToken cancellationToken) + { + if (value is string s) return BinaryData.FromString(s); + + var formatter = Formatter ?? new JsonFormatter(); + var data = await formatter.ToStringAsync(value, cancellationToken); + + return BinaryData.FromString(data); + } + + private static bool ValidatePreconditions(ActivityExecutionContext context, string? queueOrTopic, object? messageBody) + { + var logger = context.GetRequiredService>(); + + if (string.IsNullOrWhiteSpace(queueOrTopic)) + { + logger.LogWarning("Can't send a message because no queue or topic was specified"); + return false; + } + + if (messageBody == null) + { + logger.LogWarning("Can't send a message because no message body was specified"); + return false; + } + + return true; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IQueueProvider.cs b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IQueueProvider.cs new file mode 100644 index 000000000..b783a0aae --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IQueueProvider.cs @@ -0,0 +1,11 @@ +using Elsa.Modules.AzureServiceBus.Models; + +namespace Elsa.Modules.AzureServiceBus.Contracts; + +/// +/// Provides queue definitions to the system. +/// +public interface IQueueProvider +{ + ValueTask> GetQueuesAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IServiceBusInitializer.cs b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IServiceBusInitializer.cs new file mode 100644 index 000000000..71dc9f040 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IServiceBusInitializer.cs @@ -0,0 +1,12 @@ +namespace Elsa.Modules.AzureServiceBus.Contracts; + +/// +/// Creates queues, topics and subscriptions provided by , and implementations. +/// +public interface IServiceBusInitializer +{ + /// + /// Creates queues, topics and subscriptions provided by , and implementations. + /// + Task InitializeAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Contracts/ISubscriptionProvider.cs b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/ISubscriptionProvider.cs new file mode 100644 index 000000000..abe7783fc --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/ISubscriptionProvider.cs @@ -0,0 +1,11 @@ +using Elsa.Modules.AzureServiceBus.Models; + +namespace Elsa.Modules.AzureServiceBus.Contracts; + +/// +/// Provides subscription definitions to the system. +/// +public interface ISubscriptionProvider +{ + ValueTask> GetSubscriptionsAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Contracts/ITopicProvider.cs b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/ITopicProvider.cs new file mode 100644 index 000000000..c844ab58a --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/ITopicProvider.cs @@ -0,0 +1,11 @@ +using Elsa.Modules.AzureServiceBus.Models; + +namespace Elsa.Modules.AzureServiceBus.Contracts; + +/// +/// Provides topic definitions to the system. +/// +public interface ITopicProvider +{ + ValueTask> GetTopicsAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Elsa.Modules.AzureServiceBus.csproj b/src/modules/Elsa.Modules.AzureServiceBus/Elsa.Modules.AzureServiceBus.csproj new file mode 100644 index 000000000..c56344153 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Elsa.Modules.AzureServiceBus.csproj @@ -0,0 +1,23 @@ + + + + net6.0 + enable + enable + + + + + + + + + + + + + + + + + diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs b/src/modules/Elsa.Modules.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs new file mode 100644 index 000000000..722830144 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs @@ -0,0 +1,48 @@ +using Azure.Messaging.ServiceBus; +using Azure.Messaging.ServiceBus.Administration; +using Elsa.Modules.AzureServiceBus.Contracts; +using Elsa.Modules.AzureServiceBus.HostedServices; +using Elsa.Modules.AzureServiceBus.Options; +using Elsa.Modules.AzureServiceBus.Providers; +using Elsa.Modules.AzureServiceBus.Services; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; + +namespace Elsa.Modules.AzureServiceBus.Extensions; + +public static class ServiceCollectionExtensions +{ + /// + /// Register required services for the Azure Service Bus module. + /// + /// A value indicating whether or not queues, topics and subscriptions should be created at application startup. + public static IServiceCollection AddAzureServiceBusServices(this IServiceCollection services, Action configure, bool autoCreateQueuesTopicsAndSubscriptions = true) + { + services.Configure(configure); + + services + .AddSingleton(CreateServiceBusManagementClient) + .AddSingleton(CreateServiceBusClient) + .AddSingleton() + .AddTransient() + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton(sp => sp.GetRequiredService()); + + if (autoCreateQueuesTopicsAndSubscriptions) + services.AddHostedService(); + + return services; + } + + private static ServiceBusClient CreateServiceBusClient(IServiceProvider serviceProvider) => new(GetConnectionString(serviceProvider)); + private static ServiceBusAdministrationClient CreateServiceBusManagementClient(IServiceProvider serviceProvider) => new(GetConnectionString(serviceProvider)); + + private static string GetConnectionString(IServiceProvider serviceProvider) + { + var options = serviceProvider.GetRequiredService>().Value; + var configuration = serviceProvider.GetRequiredService(); + return configuration.GetConnectionString(options.ConnectionStringOrName) ?? options.ConnectionStringOrName; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/HostedServices/CreateQueuesTopicsAndSubscriptions.cs b/src/modules/Elsa.Modules.AzureServiceBus/HostedServices/CreateQueuesTopicsAndSubscriptions.cs new file mode 100644 index 000000000..07a55096d --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/HostedServices/CreateQueuesTopicsAndSubscriptions.cs @@ -0,0 +1,15 @@ +using Elsa.Modules.AzureServiceBus.Contracts; +using Microsoft.Extensions.Hosting; + +namespace Elsa.Modules.AzureServiceBus.HostedServices; + +/// +/// A blocking hosted service that creates queues, topics and subscriptions. +/// +public class CreateQueuesTopicsAndSubscriptions : IHostedService +{ + private readonly IServiceBusInitializer _serviceBusInitializer; + public CreateQueuesTopicsAndSubscriptions(IServiceBusInitializer serviceBusInitializer) => _serviceBusInitializer = serviceBusInitializer; + public Task StartAsync(CancellationToken cancellationToken) => _serviceBusInitializer.InitializeAsync(cancellationToken); + public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Models/QueueDefinition.cs b/src/modules/Elsa.Modules.AzureServiceBus/Models/QueueDefinition.cs new file mode 100644 index 000000000..4f7d13d20 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Models/QueueDefinition.cs @@ -0,0 +1,9 @@ +namespace Elsa.Modules.AzureServiceBus.Models; + +/// +/// Represents a queue that is available to the system. +/// +public class QueueDefinition +{ + public string Name { get; set; } = default!; +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Models/SubscriptionDefinition.cs b/src/modules/Elsa.Modules.AzureServiceBus/Models/SubscriptionDefinition.cs new file mode 100644 index 000000000..f43d18aba --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Models/SubscriptionDefinition.cs @@ -0,0 +1,10 @@ +namespace Elsa.Modules.AzureServiceBus.Models; + +/// +/// Represents a topic subscription that is available to the system. +/// +public class SubscriptionDefinition +{ + public string Name { get; set; } = default!; + public string Topic { get; set; } = default!; +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Models/TopicDefinition.cs b/src/modules/Elsa.Modules.AzureServiceBus/Models/TopicDefinition.cs new file mode 100644 index 000000000..1abe9e312 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Models/TopicDefinition.cs @@ -0,0 +1,9 @@ +namespace Elsa.Modules.AzureServiceBus.Models; + +/// +/// Represents a topic that is available to the system. +/// +public class TopicDefinition +{ + public string Name { get; set; } = default!; +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Options/AzureServiceBusOptions.cs b/src/modules/Elsa.Modules.AzureServiceBus/Options/AzureServiceBusOptions.cs new file mode 100644 index 000000000..50279af84 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Options/AzureServiceBusOptions.cs @@ -0,0 +1,11 @@ +using Elsa.Modules.AzureServiceBus.Models; + +namespace Elsa.Modules.AzureServiceBus.Options; + +public class AzureServiceBusOptions +{ + public string ConnectionStringOrName { get; set; } = default!; + public ICollection Queues { get; set; } = new List(); + public ICollection Topics { get; set; } = new List(); + public ICollection Subscriptions { get; set; } = new List(); +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Providers/ConfigurationQueueTopicAndSubscriptionProvider.cs b/src/modules/Elsa.Modules.AzureServiceBus/Providers/ConfigurationQueueTopicAndSubscriptionProvider.cs new file mode 100644 index 000000000..a921ba534 --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Providers/ConfigurationQueueTopicAndSubscriptionProvider.cs @@ -0,0 +1,18 @@ +using Elsa.Modules.AzureServiceBus.Contracts; +using Elsa.Modules.AzureServiceBus.Models; +using Elsa.Modules.AzureServiceBus.Options; +using Microsoft.Extensions.Options; + +namespace Elsa.Modules.AzureServiceBus.Providers; + +/// +/// Represents a queue provider that reads queue definitions from configuration. +/// +public class ConfigurationQueueTopicAndSubscriptionProvider : IQueueProvider, ITopicProvider, ISubscriptionProvider +{ + private readonly AzureServiceBusOptions _options; + public ConfigurationQueueTopicAndSubscriptionProvider(IOptions options) => _options = options.Value; + public ValueTask> GetQueuesAsync(CancellationToken cancellationToken) => new(_options.Queues); + public ValueTask> GetTopicsAsync(CancellationToken cancellationToken) => new(_options.Topics); + public ValueTask> GetSubscriptionsAsync(CancellationToken cancellationToken) => new(_options.Subscriptions); +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Services/ServiceBusInitializer.cs b/src/modules/Elsa.Modules.AzureServiceBus/Services/ServiceBusInitializer.cs new file mode 100644 index 000000000..ddb9d7a6a --- /dev/null +++ b/src/modules/Elsa.Modules.AzureServiceBus/Services/ServiceBusInitializer.cs @@ -0,0 +1,63 @@ +using Azure.Messaging.ServiceBus.Administration; +using Elsa.Modules.AzureServiceBus.Contracts; + +namespace Elsa.Modules.AzureServiceBus.Services; + +public class ServiceBusInitializer : IServiceBusInitializer +{ + private readonly ServiceBusAdministrationClient _serviceBusAdministrationClient; + private readonly IReadOnlyCollection _queueProviders; + private readonly IReadOnlyCollection _topicProviders; + private readonly IReadOnlyCollection _subscriptionProviders; + + public ServiceBusInitializer( + ServiceBusAdministrationClient serviceBusAdministrationClient, + IEnumerable queueProviders, + IEnumerable topicProviders, + IEnumerable subscriptionProviders) + { + _serviceBusAdministrationClient = serviceBusAdministrationClient; + _queueProviders = queueProviders.ToList(); + _topicProviders = topicProviders.ToList(); + _subscriptionProviders = subscriptionProviders.ToList(); + } + + public async Task InitializeAsync(CancellationToken cancellationToken = default) + { + var tasks = new[] { CreateQueuesAsync(cancellationToken), CreateTopicsAsync(cancellationToken), CreateSubscriptionsAsync(cancellationToken) }; + await Task.WhenAll(tasks); + } + + private async Task CreateQueuesAsync(CancellationToken cancellationToken) + { + var definitions = (await Task.WhenAll(_queueProviders.Select(async x => await x.GetQueuesAsync(cancellationToken)))).SelectMany(x => x); + var parallelOptions = new ParallelOptions { CancellationToken = cancellationToken, MaxDegreeOfParallelism = 5 }; + await Parallel.ForEachAsync(definitions, parallelOptions, async (definition, ct) => + { + if (!await _serviceBusAdministrationClient.QueueExistsAsync(definition.Name, ct)) + await _serviceBusAdministrationClient.CreateQueueAsync(definition.Name, ct); + }); + } + + private async Task CreateTopicsAsync(CancellationToken cancellationToken) + { + var definitions = (await Task.WhenAll(_topicProviders.Select(async x => await x.GetTopicsAsync(cancellationToken)))).SelectMany(x => x); + var parallelOptions = new ParallelOptions { CancellationToken = cancellationToken, MaxDegreeOfParallelism = 5 }; + await Parallel.ForEachAsync(definitions, parallelOptions, async (definition, ct) => + { + if (!await _serviceBusAdministrationClient.TopicExistsAsync(definition.Name, ct)) + await _serviceBusAdministrationClient.CreateTopicAsync(definition.Name, ct); + }); + } + + private async Task CreateSubscriptionsAsync(CancellationToken cancellationToken) + { + var definitions = (await Task.WhenAll(_subscriptionProviders.Select(async x => await x.GetSubscriptionsAsync(cancellationToken)))).SelectMany(x => x); + var parallelOptions = new ParallelOptions { CancellationToken = cancellationToken, MaxDegreeOfParallelism = 5 }; + await Parallel.ForEachAsync(definitions, parallelOptions, async (definition, ct) => + { + if (!await _serviceBusAdministrationClient.SubscriptionExistsAsync(definition.Topic, definition.Name, ct)) + await _serviceBusAdministrationClient.CreateSubscriptionAsync(definition.Topic, definition.Name, ct); + }); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.Quartz/AssemblyInfo.cs b/src/modules/Elsa.Modules.Quartz/AssemblyInfo.cs deleted file mode 100644 index 95465252a..000000000 --- a/src/modules/Elsa.Modules.Quartz/AssemblyInfo.cs +++ /dev/null @@ -1 +0,0 @@ -[assembly: CLSCompliant(true)] \ No newline at end of file diff --git a/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj b/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj index 7d1e095c9..f67f9f6ca 100644 --- a/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj +++ b/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj @@ -6,14 +6,15 @@ - - - - - - - - + + + + + + + + + diff --git a/src/samples/aspnet/Elsa.Samples.Web1/Program.cs b/src/samples/aspnet/Elsa.Samples.Web1/Program.cs index 54ccd0a40..2a5817ed3 100644 --- a/src/samples/aspnet/Elsa.Samples.Web1/Program.cs +++ b/src/samples/aspnet/Elsa.Samples.Web1/Program.cs @@ -6,6 +6,8 @@ using Elsa.Extensions; using Elsa.Management.Contracts; using Elsa.Management.Extensions; using Elsa.Mediator.Extensions; +using Elsa.Modules.AzureServiceBus.Activities; +using Elsa.Modules.AzureServiceBus.Extensions; using Elsa.Modules.Http; using Elsa.Modules.Http.Extensions; using Elsa.Modules.Quartz.Services; @@ -24,12 +26,14 @@ using Elsa.Scheduling.Extensions; using Elsa.Scripting.JavaScript.Extensions; using Elsa.Scripting.Liquid.Extensions; using Microsoft.AspNetCore.Builder; +using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; var builder = WebApplication.CreateBuilder(args); var services = builder.Services; +var configuration = builder.Configuration; // Add services. services @@ -42,12 +46,13 @@ services .AddScheduling(new QuartzSchedulingServiceProvider()) .AddHttpActivityServices() .AddSchedulingActivities() - + .AddAzureServiceBusServices(options => configuration.GetSection("AzureServiceBus").Bind(options)) .ConfigureWorkflowRuntime(options => { options.Workflows.Add("HelloWorldWorkflow", new HelloWorldWorkflow()); options.Workflows.Add("HttpWorkflow", new HttpWorkflow()); options.Workflows.Add("ForkedHttpWorkflow", new ForkedHttpWorkflow()); + options.Workflows.Add("AzureServiceBusWorkflow", new AzureServiceBusWorkflow()); options.Workflows.Add(nameof(CompositeActivitiesWorkflow), new CompositeActivitiesWorkflow()); }); @@ -65,6 +70,7 @@ services .AddActivity() .AddActivity() .AddActivity() + .AddActivity() ; // Register available triggers. diff --git a/src/samples/aspnet/Elsa.Samples.Web1/Workflows/AzureServiceBusWorkflow.cs b/src/samples/aspnet/Elsa.Samples.Web1/Workflows/AzureServiceBusWorkflow.cs new file mode 100644 index 000000000..356ee2633 --- /dev/null +++ b/src/samples/aspnet/Elsa.Samples.Web1/Workflows/AzureServiceBusWorkflow.cs @@ -0,0 +1,19 @@ +using Elsa.Contracts; +using Elsa.Models; +using Elsa.Modules.AzureServiceBus.Activities; +using Elsa.Runtime.Contracts; + +namespace Elsa.Samples.Web1.Workflows; + +public class AzureServiceBusWorkflow : IWorkflow +{ + public void Build(IWorkflowDefinitionBuilder workflow) + { + workflow.WithRoot(new Send + { + QueueOrTopic = new Input("inbox"), + MessageBody = new Input(new { Subject = "Greetings", Message = "Hello World!" }), + ContentType = new Input("application/json") + }); + } +} \ No newline at end of file diff --git a/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json b/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json index 778fcef3d..4d330878b 100644 --- a/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json +++ b/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json @@ -7,5 +7,14 @@ "Microsoft.EntityFrameworkCore.Database.Command": "Warning" } }, - "AllowedHosts": "*" + "AllowedHosts": "*", + "ConnectionStrings": { + "AzureServiceBus": "" + }, + "AzureServiceBus": { + "ConnectionStringOrName": "AzureServiceBus", + "Queues": [{ + "Name": "inbox" + }] + } } diff --git a/src/samples/console/Elsa.Samples.Console1/Program.cs b/src/samples/console/Elsa.Samples.Console1/Program.cs index 99bde066f..69946ee3d 100644 --- a/src/samples/console/Elsa.Samples.Console1/Program.cs +++ b/src/samples/console/Elsa.Samples.Console1/Program.cs @@ -50,7 +50,7 @@ class Program var workflow13 = new Func(BlockingParallelForEachWorkflow.Create); var workflow14 = new Func(FlowchartWorkflow.Create); - var workflowFactory = workflow11; + var workflowFactory = workflow1; var workflowGraph = workflowFactory(); var workflow = Workflow.FromActivity(workflowGraph);