Incremental work on Azure Service Bus module
This commit is contained in:
parent
34f851d8c7
commit
88609d5a62
14
Elsa.sln
14
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
|
||||
|
|
|
|||
|
|
@ -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; }
|
||||
}
|
||||
|
|
@ -59,7 +59,7 @@ public class ActivityExecutionContext
|
|||
public void SetBookmark(object? payload, ExecuteActivityDelegate? callback = default)
|
||||
{
|
||||
var hasher = GetRequiredService<IHasher>();
|
||||
|
||||
|
||||
var identityGenerator = GetRequiredService<IIdentityGenerator>();
|
||||
var payloadSerializer = GetRequiredService<IPayloadSerializer>();
|
||||
var payloadJson = payload != null ? payloadSerializer.Serialize(payload) : default;
|
||||
|
|
@ -87,7 +87,7 @@ public class ActivityExecutionContext
|
|||
}
|
||||
|
||||
public T GetRequiredService<T>() where T : notnull => WorkflowExecutionContext.GetRequiredService<T>();
|
||||
public T? Get<T>(Input<T> input) => Get<T>(input.LocationReference);
|
||||
public T? Get<T>(Input<T>? input) => input == null ? default : Get<T>(input.LocationReference);
|
||||
|
||||
public object? Get(RegisterLocationReference locationReference)
|
||||
{
|
||||
|
|
|
|||
10
src/core/Elsa.Formatting/Contracts/IFormatter.cs
Normal file
10
src/core/Elsa.Formatting/Contracts/IFormatter.cs
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
namespace Elsa.Formatting.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a formatter that can serialize an object to a string and deserialize a string into an object.
|
||||
/// </summary>
|
||||
public interface IFormatter
|
||||
{
|
||||
ValueTask<string> ToStringAsync(object body, CancellationToken cancellationToken = default);
|
||||
ValueTask<object> FromStringAsync(string data, Type? returnType, CancellationToken cancellationToken = default);
|
||||
}
|
||||
9
src/core/Elsa.Formatting/Elsa.Formatting.csproj
Normal file
9
src/core/Elsa.Formatting/Elsa.Formatting.csproj
Normal file
|
|
@ -0,0 +1,9 @@
|
|||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>net6.0</TargetFramework>
|
||||
<ImplicitUsings>enable</ImplicitUsings>
|
||||
<Nullable>enable</Nullable>
|
||||
</PropertyGroup>
|
||||
|
||||
</Project>
|
||||
20
src/core/Elsa.Formatting/Formatters/JsonFormatter.cs
Normal file
20
src/core/Elsa.Formatting/Formatters/JsonFormatter.cs
Normal file
|
|
@ -0,0 +1,20 @@
|
|||
using System.Text.Json;
|
||||
using Elsa.Formatting.Contracts;
|
||||
|
||||
namespace Elsa.Formatting.Formatters;
|
||||
|
||||
public class JsonFormatter : IFormatter
|
||||
{
|
||||
public ValueTask<string> ToStringAsync(object body, CancellationToken cancellationToken)
|
||||
{
|
||||
var json = JsonSerializer.Serialize(body);
|
||||
return ValueTask.FromResult<string>(json);
|
||||
}
|
||||
|
||||
public ValueTask<object> FromStringAsync(string data, Type? returnType, CancellationToken cancellationToken)
|
||||
{
|
||||
var options = new JsonSerializerOptions();
|
||||
var value = returnType != null ? JsonSerializer.Deserialize(data, returnType, options)! : JsonSerializer.Deserialize<object>(data, options)!;
|
||||
return ValueTask.FromResult(value);
|
||||
}
|
||||
}
|
||||
74
src/modules/Elsa.Modules.AzureServiceBus/Activities/Send.cs
Normal file
74
src/modules/Elsa.Modules.AzureServiceBus/Activities/Send.cs
Normal file
|
|
@ -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<object> MessageBody { get; set; } = default!;
|
||||
public Input<string> QueueOrTopic { get; set; } = default!;
|
||||
public Input<string>? ContentType { get; set; }
|
||||
public Input<string>? Subject { get; set; }
|
||||
public Input<string>? 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<ServiceBusClient>();
|
||||
|
||||
await using var sender = client.CreateSender(queueOrTopic);
|
||||
await sender.SendMessageAsync(message, cancellationToken);
|
||||
}
|
||||
|
||||
private async ValueTask<BinaryData> 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<ILogger<Send>>();
|
||||
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,11 @@
|
|||
using Elsa.Modules.AzureServiceBus.Models;
|
||||
|
||||
namespace Elsa.Modules.AzureServiceBus.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Provides queue definitions to the system.
|
||||
/// </summary>
|
||||
public interface IQueueProvider
|
||||
{
|
||||
ValueTask<ICollection<QueueDefinition>> GetQueuesAsync(CancellationToken cancellationToken);
|
||||
}
|
||||
|
|
@ -0,0 +1,12 @@
|
|||
namespace Elsa.Modules.AzureServiceBus.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Creates queues, topics and subscriptions provided by <see cref="IQueueProvider"/>, <see cref="ITopicProvider"/> and <see cref="ISubscriptionProvider"/> implementations.
|
||||
/// </summary>
|
||||
public interface IServiceBusInitializer
|
||||
{
|
||||
/// <summary>
|
||||
/// Creates queues, topics and subscriptions provided by <see cref="IQueueProvider"/>, <see cref="ITopicProvider"/> and <see cref="ISubscriptionProvider"/> implementations.
|
||||
/// </summary>
|
||||
Task InitializeAsync(CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -0,0 +1,11 @@
|
|||
using Elsa.Modules.AzureServiceBus.Models;
|
||||
|
||||
namespace Elsa.Modules.AzureServiceBus.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Provides subscription definitions to the system.
|
||||
/// </summary>
|
||||
public interface ISubscriptionProvider
|
||||
{
|
||||
ValueTask<ICollection<SubscriptionDefinition>> GetSubscriptionsAsync(CancellationToken cancellationToken);
|
||||
}
|
||||
|
|
@ -0,0 +1,11 @@
|
|||
using Elsa.Modules.AzureServiceBus.Models;
|
||||
|
||||
namespace Elsa.Modules.AzureServiceBus.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Provides topic definitions to the system.
|
||||
/// </summary>
|
||||
public interface ITopicProvider
|
||||
{
|
||||
ValueTask<ICollection<TopicDefinition>> GetTopicsAsync(CancellationToken cancellationToken);
|
||||
}
|
||||
|
|
@ -0,0 +1,23 @@
|
|||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>net6.0</TargetFramework>
|
||||
<ImplicitUsings>enable</ImplicitUsings>
|
||||
<Nullable>enable</Nullable>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Azure.Messaging.ServiceBus" Version="7.5.1" />
|
||||
<PackageReference Include="Microsoft.Azure.Management.ServiceBus.Fluent" Version="1.38.0" />
|
||||
<PackageReference Include="Microsoft.Extensions.Configuration.Abstractions" Version="6.0.0" />
|
||||
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions" Version="6.0.0" />
|
||||
<PackageReference Include="System.Linq.Async" Version="6.0.1" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
|
||||
<ProjectReference Include="..\..\core\Elsa.Formatting\Elsa.Formatting.csproj" />
|
||||
<ProjectReference Include="..\..\core\Elsa.Mediator\Elsa.Mediator.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
@ -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
|
||||
{
|
||||
/// <summary>
|
||||
/// Register required services for the Azure Service Bus module.
|
||||
/// </summary>
|
||||
/// <param name="autoCreateQueuesTopicsAndSubscriptions">A value indicating whether or not queues, topics and subscriptions should be created at application startup.</param>
|
||||
public static IServiceCollection AddAzureServiceBusServices(this IServiceCollection services, Action<AzureServiceBusOptions> configure, bool autoCreateQueuesTopicsAndSubscriptions = true)
|
||||
{
|
||||
services.Configure(configure);
|
||||
|
||||
services
|
||||
.AddSingleton(CreateServiceBusManagementClient)
|
||||
.AddSingleton(CreateServiceBusClient)
|
||||
.AddSingleton<ConfigurationQueueTopicAndSubscriptionProvider>()
|
||||
.AddTransient<IServiceBusInitializer, ServiceBusInitializer>()
|
||||
.AddSingleton<IQueueProvider>(sp => sp.GetRequiredService<ConfigurationQueueTopicAndSubscriptionProvider>())
|
||||
.AddSingleton<ITopicProvider>(sp => sp.GetRequiredService<ConfigurationQueueTopicAndSubscriptionProvider>())
|
||||
.AddSingleton<ISubscriptionProvider>(sp => sp.GetRequiredService<ConfigurationQueueTopicAndSubscriptionProvider>());
|
||||
|
||||
if (autoCreateQueuesTopicsAndSubscriptions)
|
||||
services.AddHostedService<CreateQueuesTopicsAndSubscriptions>();
|
||||
|
||||
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<IOptions<AzureServiceBusOptions>>().Value;
|
||||
var configuration = serviceProvider.GetRequiredService<IConfiguration>();
|
||||
return configuration.GetConnectionString(options.ConnectionStringOrName) ?? options.ConnectionStringOrName;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,15 @@
|
|||
using Elsa.Modules.AzureServiceBus.Contracts;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
|
||||
namespace Elsa.Modules.AzureServiceBus.HostedServices;
|
||||
|
||||
/// <summary>
|
||||
/// A blocking hosted service that creates queues, topics and subscriptions.
|
||||
/// </summary>
|
||||
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;
|
||||
}
|
||||
|
|
@ -0,0 +1,9 @@
|
|||
namespace Elsa.Modules.AzureServiceBus.Models;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a queue that is available to the system.
|
||||
/// </summary>
|
||||
public class QueueDefinition
|
||||
{
|
||||
public string Name { get; set; } = default!;
|
||||
}
|
||||
|
|
@ -0,0 +1,10 @@
|
|||
namespace Elsa.Modules.AzureServiceBus.Models;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a topic subscription that is available to the system.
|
||||
/// </summary>
|
||||
public class SubscriptionDefinition
|
||||
{
|
||||
public string Name { get; set; } = default!;
|
||||
public string Topic { get; set; } = default!;
|
||||
}
|
||||
|
|
@ -0,0 +1,9 @@
|
|||
namespace Elsa.Modules.AzureServiceBus.Models;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a topic that is available to the system.
|
||||
/// </summary>
|
||||
public class TopicDefinition
|
||||
{
|
||||
public string Name { get; set; } = default!;
|
||||
}
|
||||
|
|
@ -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<QueueDefinition> Queues { get; set; } = new List<QueueDefinition>();
|
||||
public ICollection<TopicDefinition> Topics { get; set; } = new List<TopicDefinition>();
|
||||
public ICollection<SubscriptionDefinition> Subscriptions { get; set; } = new List<SubscriptionDefinition>();
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a queue provider that reads queue definitions from configuration.
|
||||
/// </summary>
|
||||
public class ConfigurationQueueTopicAndSubscriptionProvider : IQueueProvider, ITopicProvider, ISubscriptionProvider
|
||||
{
|
||||
private readonly AzureServiceBusOptions _options;
|
||||
public ConfigurationQueueTopicAndSubscriptionProvider(IOptions<AzureServiceBusOptions> options) => _options = options.Value;
|
||||
public ValueTask<ICollection<QueueDefinition>> GetQueuesAsync(CancellationToken cancellationToken) => new(_options.Queues);
|
||||
public ValueTask<ICollection<TopicDefinition>> GetTopicsAsync(CancellationToken cancellationToken) => new(_options.Topics);
|
||||
public ValueTask<ICollection<SubscriptionDefinition>> GetSubscriptionsAsync(CancellationToken cancellationToken) => new(_options.Subscriptions);
|
||||
}
|
||||
|
|
@ -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<IQueueProvider> _queueProviders;
|
||||
private readonly IReadOnlyCollection<ITopicProvider> _topicProviders;
|
||||
private readonly IReadOnlyCollection<ISubscriptionProvider> _subscriptionProviders;
|
||||
|
||||
public ServiceBusInitializer(
|
||||
ServiceBusAdministrationClient serviceBusAdministrationClient,
|
||||
IEnumerable<IQueueProvider> queueProviders,
|
||||
IEnumerable<ITopicProvider> topicProviders,
|
||||
IEnumerable<ISubscriptionProvider> 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);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -1 +0,0 @@
|
|||
[assembly: CLSCompliant(true)]
|
||||
|
|
@ -6,14 +6,15 @@
|
|||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\..\modules\Elsa.Modules.Http\Elsa.Modules.Http.csproj"/>
|
||||
<ProjectReference Include="..\..\..\modules\Elsa.Modules.Scheduling\Elsa.Modules.Scheduling.csproj"/>
|
||||
<ProjectReference Include="..\..\..\modules\Elsa.Modules.Quartz\Elsa.Modules.Quartz.csproj"/>
|
||||
<ProjectReference Include="..\..\..\api\Elsa.Api\Elsa.Api.csproj"/>
|
||||
<ProjectReference Include="..\..\..\core\Elsa\Elsa.csproj"/>
|
||||
<ProjectReference Include="..\..\..\persistence\Elsa.Persistence.EntityFrameworkCore.Sqlite\Elsa.Persistence.EntityFrameworkCore.Sqlite.csproj"/>
|
||||
<ProjectReference Include="..\..\..\runtime\Elsa.Runtime.ProtoActor\Elsa.Runtime.ProtoActor.csproj"/>
|
||||
<ProjectReference Include="..\..\..\scripting\Elsa.Scripting.Liquid\Elsa.Scripting.Liquid.csproj"/>
|
||||
<ProjectReference Include="..\..\..\modules\Elsa.Modules.AzureServiceBus\Elsa.Modules.AzureServiceBus.csproj" />
|
||||
<ProjectReference Include="..\..\..\modules\Elsa.Modules.Http\Elsa.Modules.Http.csproj" />
|
||||
<ProjectReference Include="..\..\..\modules\Elsa.Modules.Scheduling\Elsa.Modules.Scheduling.csproj" />
|
||||
<ProjectReference Include="..\..\..\modules\Elsa.Modules.Quartz\Elsa.Modules.Quartz.csproj" />
|
||||
<ProjectReference Include="..\..\..\api\Elsa.Api\Elsa.Api.csproj" />
|
||||
<ProjectReference Include="..\..\..\core\Elsa\Elsa.csproj" />
|
||||
<ProjectReference Include="..\..\..\persistence\Elsa.Persistence.EntityFrameworkCore.Sqlite\Elsa.Persistence.EntityFrameworkCore.Sqlite.csproj" />
|
||||
<ProjectReference Include="..\..\..\runtime\Elsa.Runtime.ProtoActor\Elsa.Runtime.ProtoActor.csproj" />
|
||||
<ProjectReference Include="..\..\..\scripting\Elsa.Scripting.Liquid\Elsa.Scripting.Liquid.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
|
|||
|
|
@ -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<Delay>()
|
||||
.AddActivity<ForEach>()
|
||||
.AddActivity<Switch>()
|
||||
.AddActivity<Send>()
|
||||
;
|
||||
|
||||
// Register available triggers.
|
||||
|
|
|
|||
|
|
@ -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<string>("inbox"),
|
||||
MessageBody = new Input<object>(new { Subject = "Greetings", Message = "Hello World!" }),
|
||||
ContentType = new Input<string>("application/json")
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -7,5 +7,14 @@
|
|||
"Microsoft.EntityFrameworkCore.Database.Command": "Warning"
|
||||
}
|
||||
},
|
||||
"AllowedHosts": "*"
|
||||
"AllowedHosts": "*",
|
||||
"ConnectionStrings": {
|
||||
"AzureServiceBus": ""
|
||||
},
|
||||
"AzureServiceBus": {
|
||||
"ConnectionStringOrName": "AzureServiceBus",
|
||||
"Queues": [{
|
||||
"Name": "inbox"
|
||||
}]
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -50,7 +50,7 @@ class Program
|
|||
var workflow13 = new Func<IActivity>(BlockingParallelForEachWorkflow.Create);
|
||||
var workflow14 = new Func<IActivity>(FlowchartWorkflow.Create);
|
||||
|
||||
var workflowFactory = workflow11;
|
||||
var workflowFactory = workflow1;
|
||||
var workflowGraph = workflowFactory();
|
||||
var workflow = Workflow.FromActivity(workflowGraph);
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue