diff --git a/Samples.sln b/Samples.sln index 2ea48508d..2913a709d 100644 --- a/Samples.sln +++ b/Samples.sln @@ -148,6 +148,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.ServiceBus.AzureServic EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.ServiceBus.RabbitMq", "src\servicebus\Elsa.ServiceBus.RabbitMq\Elsa.ServiceBus.RabbitMq.csproj", "{FC656D28-1AA3-462D-9B18-A10436E47C08}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Activities.AzureServiceBus", "src\activities\Elsa.Activities.AzureServiceBus\Elsa.Activities.AzureServiceBus.csproj", "{DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.AzureServiceBusWorker", "src\samples\Elsa.Samples.AzureServiceBusWorker\Elsa.Samples.AzureServiceBusWorker.csproj", "{1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -363,6 +367,14 @@ Global {FC656D28-1AA3-462D-9B18-A10436E47C08}.Debug|Any CPU.Build.0 = Debug|Any CPU {FC656D28-1AA3-462D-9B18-A10436E47C08}.Release|Any CPU.ActiveCfg = Release|Any CPU {FC656D28-1AA3-462D-9B18-A10436E47C08}.Release|Any CPU.Build.0 = Release|Any CPU + {DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Debug|Any CPU.Build.0 = Debug|Any CPU + {DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Release|Any CPU.ActiveCfg = Release|Any CPU + {DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Release|Any CPU.Build.0 = Release|Any CPU + {1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Debug|Any CPU.Build.0 = Debug|Any CPU + {1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Release|Any CPU.ActiveCfg = Release|Any CPU + {1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -433,6 +445,8 @@ Global {47EDFC14-BB99-44BF-A447-035D40D7347B} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925} {D3440E66-A871-4AB0-A1E0-8B3864ACA16A} = {47EDFC14-BB99-44BF-A447-035D40D7347B} {FC656D28-1AA3-462D-9B18-A10436E47C08} = {47EDFC14-BB99-44BF-A447-035D40D7347B} + {DD607BC1-F9D9-4685-BC8F-A4056C2C31FA} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} + {1D8B906A-5FCB-413C-8972-91D7A3EB8AE7} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158} diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived.cs new file mode 100644 index 000000000..696426a3f --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived.cs @@ -0,0 +1,37 @@ +using System; +using System.Text; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Serialization; +using Elsa.Services; +using Elsa.Services.Models; +using Microsoft.Azure.ServiceBus; + +namespace Elsa.Activities.AzureServiceBus.Activities +{ + [Trigger(Category = "Azure Service Bus", DisplayName = "Service Bus Message Received", Description = "Triggered when a message is received on the specified queue", Outcomes = new[] { OutcomeNames.Done })] + public class AzureServiceBusMessageReceived : Activity + { + private readonly IContentSerializer _serializer; + + public AzureServiceBusMessageReceived(IContentSerializer serializer) + { + _serializer = serializer; + } + + [ActivityProperty] public string QueueName { get; set; } = default!; + [ActivityProperty] public Type MessageType { get; set; } = default!; + + protected override IActivityExecutionResult OnExecute() => Suspend(); + + protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) + { + var message = (Message) context.Input!; + var bytes = message.Body; + var json = Encoding.UTF8.GetString(bytes); + var model = _serializer.Deserialize(json, MessageType); + + return Done(model); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage.cs new file mode 100644 index 000000000..8fe9f0a7d --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage.cs @@ -0,0 +1,43 @@ +using System.Text; +using System.Threading; +using System.Threading.Tasks; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Serialization; +using Elsa.Services; +using Elsa.Services.Models; +using Microsoft.Azure.ServiceBus; +using IServiceBusFactory = Elsa.Activities.AzureServiceBus.Services.IServiceBusFactory; + +namespace Elsa.Activities.AzureServiceBus.Activities +{ + [Trigger(Category = "Azure Service Bus", DisplayName = "Send Service Bus Message", Description = "Sends a message to the specified queue", Outcomes = new[] { OutcomeNames.Done })] + public class SendAzureServiceBusMessage : Activity + { + private readonly IServiceBusFactory _serviceBusFactory; + private readonly IContentSerializer _serializer; + + public SendAzureServiceBusMessage(IServiceBusFactory serviceBusFactory, IContentSerializer serializer) + { + _serviceBusFactory = serviceBusFactory; + _serializer = serializer; + } + + [ActivityProperty] public string QueueName { get; set; } = default!; + [ActivityProperty] public object Message { get; set; } = default!; + + protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) + { + var sender = await _serviceBusFactory.GetSenderAsync(QueueName, cancellationToken); + var json = _serializer.Serialize(Message); + var bytes = Encoding.UTF8.GetBytes(json); + var message = new Message(bytes); + + if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId)) + message.CorrelationId = context.WorkflowExecutionContext.CorrelationId; + + await sender.SendAsync(message); + return Done(); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj b/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj new file mode 100644 index 000000000..d064f2ff9 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Elsa.Activities.AzureServiceBus.csproj @@ -0,0 +1,18 @@ + + + + netstandard2.0 + latest + enable + + + + + + + + + + + + diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs new file mode 100644 index 000000000..c46d48006 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs @@ -0,0 +1,46 @@ +using System; +using Elsa.Activities.AzureServiceBus.Activities; +using Elsa.Activities.AzureServiceBus.Options; +using Elsa.Activities.AzureServiceBus.Services; +using Elsa.Activities.AzureServiceBus.StartupTasks; +using Elsa.Runtime; +using Microsoft.Azure.ServiceBus; +using Microsoft.Azure.ServiceBus.Management; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; + +namespace Elsa.Activities.AzureServiceBus.Extensions +{ + public static class ServiceCollectionExtensions + { + public static IServiceCollection AddAzureServiceBusActivities(this IServiceCollection services, Action? configure) + { + if (configure != null) + services.Configure(configure); + else + services.AddOptions(); + + return services + .AddSingleton(CreateServiceBusConnection) + .AddSingleton(CreateServiceBusManagementClient) + .AddSingleton() + .AddStartupTask() + .AddActivity() + .AddActivity(); + } + + private static ServiceBusConnection CreateServiceBusConnection(IServiceProvider serviceProvider) + { + var options = serviceProvider.GetRequiredService>().Value; + var connectionString = options.ConnectionString; + return new ServiceBusConnection(connectionString, RetryPolicy.Default); + } + + private static ManagementClient CreateServiceBusManagementClient(IServiceProvider serviceProvider) + { + var options = serviceProvider.GetRequiredService>().Value; + var connectionString = options.ConnectionString; + return new ManagementClient(connectionString); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Options/AzureServiceBusOptions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Options/AzureServiceBusOptions.cs new file mode 100644 index 000000000..457f33081 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Options/AzureServiceBusOptions.cs @@ -0,0 +1,7 @@ +namespace Elsa.Activities.AzureServiceBus.Options +{ + public class AzureServiceBusOptions + { + public string ConnectionString { get; set; } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusFactory.cs new file mode 100644 index 000000000..4d0e0dd55 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/IServiceBusFactory.cs @@ -0,0 +1,12 @@ +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Azure.ServiceBus.Core; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public interface IServiceBusFactory + { + Task GetSenderAsync(string queueName, CancellationToken cancellationToken = default); + Task GetReceiverAsync(string queueName, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs new file mode 100644 index 000000000..14089b2f9 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -0,0 +1,40 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Triggers; +using Elsa.Services; +using Microsoft.Azure.ServiceBus.Core; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public class QueueWorker + { + private readonly IMessageReceiver _messageReceiver; + private readonly IWorkflowScheduler _workflowScheduler; + + public QueueWorker(IMessageReceiver messageReceiver, IWorkflowScheduler workflowScheduler) + { + _messageReceiver = messageReceiver; + _workflowScheduler = workflowScheduler; + } + + public async Task StartAsync(CancellationToken cancellationToken) + { + var cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + await Task.Factory.StartNew(() => ReadQueueAsync(cancellationTokenSource.Token), cancellationToken); + } + + private async Task ReadQueueAsync(CancellationToken cancellationToken) + { + while (!cancellationToken.IsCancellationRequested) + { + var message = await _messageReceiver.ReceiveAsync(); + + if(message == null) + continue; + + await _workflowScheduler.TriggerWorkflowsAsync(x => x.QueueName == _messageReceiver.Path && (string.IsNullOrWhiteSpace(x.CorrelationId) || x.CorrelationId == message.CorrelationId), message, + message.CorrelationId, cancellationToken: cancellationToken); + } + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusFactory.cs new file mode 100644 index 000000000..bfec991c2 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ServiceBusFactory.cs @@ -0,0 +1,54 @@ +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Azure.ServiceBus; +using Microsoft.Azure.ServiceBus.Core; +using Microsoft.Azure.ServiceBus.Management; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public class ServiceBusFactory : IServiceBusFactory + { + private readonly ServiceBusConnection _connection; + private readonly ManagementClient _managementClient; + private readonly IDictionary _senders = new Dictionary(); + private readonly IDictionary _receivers = new Dictionary(); + + public ServiceBusFactory(ServiceBusConnection connection, ManagementClient managementClient) + { + _connection = connection; + _managementClient = managementClient; + } + + public async Task GetSenderAsync(string queueName, CancellationToken cancellationToken) + { + if (_senders.TryGetValue(queueName, out var messageSender)) + return messageSender; + + await EnsureQueueExistsAsync(queueName, cancellationToken); + + var newMessageSender = new MessageSender(_connection, queueName); + _senders.Add(queueName, newMessageSender); + return newMessageSender; + } + + public async Task GetReceiverAsync(string queueName, CancellationToken cancellationToken) + { + if (_receivers.TryGetValue(queueName, out var messageReceiver)) + return messageReceiver; + + await EnsureQueueExistsAsync(queueName, cancellationToken); + var newMessageReceiver = new MessageReceiver(_connection, queueName); + _receivers.Add(queueName, newMessageReceiver); + return newMessageReceiver; + } + + private async Task EnsureQueueExistsAsync(string queueName, CancellationToken cancellationToken) + { + if (await _managementClient.QueueExistsAsync(queueName, cancellationToken)) + return; + + await _managementClient.CreateQueueAsync(queueName, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs new file mode 100644 index 000000000..7f51ae4b0 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs @@ -0,0 +1,65 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Runtime.CompilerServices; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Activities; +using Elsa.Activities.AzureServiceBus.Services; +using Elsa.Services; +using Microsoft.Extensions.DependencyInjection; +using IServiceBusFactory = Elsa.Activities.AzureServiceBus.Services.IServiceBusFactory; + +namespace Elsa.Activities.AzureServiceBus.StartupTasks +{ + public class StartServiceBusQueues : IStartupTask + { + private readonly IWorkflowRegistry _workflowRegistry; + private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector; + private readonly IServiceBusFactory _serviceBusFactory; + private readonly IServiceProvider _serviceProvider; + + public StartServiceBusQueues(IWorkflowRegistry workflowRegistry, IWorkflowBlueprintReflector workflowBlueprintReflector, IServiceBusFactory serviceBusFactory, IServiceProvider serviceProvider) + { + _workflowRegistry = workflowRegistry; + _workflowBlueprintReflector = workflowBlueprintReflector; + _serviceBusFactory = serviceBusFactory; + _serviceProvider = serviceProvider; + } + + public async Task ExecuteAsync(CancellationToken cancellationToken = default) + { + var queueNames = await GetQueueNamesAsync(cancellationToken).ToListAsync(cancellationToken); + + foreach (var queueName in queueNames) + { + var receiver = await _serviceBusFactory.GetReceiverAsync(queueName, cancellationToken); + var worker = ActivatorUtilities.CreateInstance(_serviceProvider, receiver); + + await worker.StartAsync(cancellationToken); + } + } + + private async IAsyncEnumerable GetQueueNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + var workflows = await _workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken); + + var query = + from workflow in workflows + from activity in workflow.Activities + where activity.Type == nameof(AzureServiceBusMessageReceived) + select workflow; + + foreach (var workflow in query) + { + var workflowBlueprintWrapper = await _workflowBlueprintReflector.ReflectAsync(workflow, cancellationToken); + + foreach (var activity in workflowBlueprintWrapper.Filter()) + { + var queueName = await activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken); + yield return queueName; + } + } + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Triggers/MessageReceivedTrigger.cs b/src/activities/Elsa.Activities.AzureServiceBus/Triggers/MessageReceivedTrigger.cs new file mode 100644 index 000000000..602cfe0c0 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Triggers/MessageReceivedTrigger.cs @@ -0,0 +1,23 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Activities; +using Elsa.Triggers; + +namespace Elsa.Activities.AzureServiceBus.Triggers +{ + public class MessageReceivedTrigger : Trigger + { + public string QueueName { get; set; } = default!; + public string? CorrelationId { get; set; } + } + + public class MessageReceivedTriggerProvider : TriggerProvider + { + public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => + new MessageReceivedTrigger + { + QueueName = (await context.Activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken)), + CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId + }; + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Rebus/Activities/PublishMessage.cs b/src/activities/Elsa.Activities.Rebus/Activities/PublishRebusMessage.cs similarity index 89% rename from src/activities/Elsa.Activities.Rebus/Activities/PublishMessage.cs rename to src/activities/Elsa.Activities.Rebus/Activities/PublishRebusMessage.cs index e9a6d61f1..aebaa83f3 100644 --- a/src/activities/Elsa.Activities.Rebus/Activities/PublishMessage.cs +++ b/src/activities/Elsa.Activities.Rebus/Activities/PublishRebusMessage.cs @@ -9,11 +9,11 @@ using Elsa.Services.Models; namespace Elsa.Activities.Rebus { [Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })] - public class PublishMessage : Activity + public class PublishRebusMessage : Activity { private readonly IEventPublisher _eventPublisher; - public PublishMessage(IEventPublisher eventPublisher) + public PublishRebusMessage(IEventPublisher eventPublisher) { _eventPublisher = eventPublisher; } diff --git a/src/activities/Elsa.Activities.Rebus/Activities/MessageReceived.cs b/src/activities/Elsa.Activities.Rebus/Activities/RebusMessageReceived.cs similarity index 92% rename from src/activities/Elsa.Activities.Rebus/Activities/MessageReceived.cs rename to src/activities/Elsa.Activities.Rebus/Activities/RebusMessageReceived.cs index 38f337e65..9df44e58b 100644 --- a/src/activities/Elsa.Activities.Rebus/Activities/MessageReceived.cs +++ b/src/activities/Elsa.Activities.Rebus/Activities/RebusMessageReceived.cs @@ -7,7 +7,7 @@ using Elsa.Services.Models; namespace Elsa.Activities.Rebus { [Trigger(Category = "Rebus", Description = "Triggered when a message is received.", Outcomes = new[] { OutcomeNames.Done })] - public class MessageReceived : Activity + public class RebusMessageReceived : Activity { [ActivityProperty(Hint = "The type of message to receive.")] public Type MessageType { get; set; } = default!; diff --git a/src/activities/Elsa.Activities.Rebus/Activities/SendMessage.cs b/src/activities/Elsa.Activities.Rebus/Activities/SendRebusMessage.cs similarity index 90% rename from src/activities/Elsa.Activities.Rebus/Activities/SendMessage.cs rename to src/activities/Elsa.Activities.Rebus/Activities/SendRebusMessage.cs index c6010a9ff..d91cd3c9e 100644 --- a/src/activities/Elsa.Activities.Rebus/Activities/SendMessage.cs +++ b/src/activities/Elsa.Activities.Rebus/Activities/SendRebusMessage.cs @@ -9,11 +9,11 @@ using Elsa.Services.Models; namespace Elsa.Activities.Rebus { [Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })] - public class SendMessage : Activity + public class SendRebusMessage : Activity { private readonly ICommandSender _bus; - public SendMessage(ICommandSender bus) + public SendRebusMessage(ICommandSender bus) { _bus = bus; } diff --git a/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj b/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj index b4b502b80..ba48a6008 100644 --- a/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj +++ b/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj @@ -19,13 +19,13 @@ - + True - + diff --git a/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs index a74a7b7bd..abe5ac158 100644 --- a/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs @@ -22,9 +22,9 @@ namespace Elsa.Activities.Rebus.Extensions return services .AddTriggerProvider() .AddStartupTask(sp => ActivatorUtilities.CreateInstance(sp, (object)messageTypes)) - .AddActivity() - .AddActivity() - .AddActivity(); + .AddActivity() + .AddActivity() + .AddActivity(); } public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities(typeof(T)); diff --git a/src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs b/src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs index 8fde08bfd..7b3028b89 100644 --- a/src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs +++ b/src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs @@ -10,9 +10,9 @@ namespace Elsa.Activities.Rebus.Triggers public string? CorrelationId { get; set; } } - public class MessageReceivedTriggerProvider : TriggerProvider + public class MessageReceivedTriggerProvider : TriggerProvider { - public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => + public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => new MessageReceivedTrigger { MessageType = (await context.Activity.GetPropertyValueAsync(x => x.MessageType, cancellationToken)).Name, diff --git a/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs index f4c509aa5..1d27a5b01 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs @@ -26,16 +26,17 @@ namespace Elsa.Activities.Timers [ActivityProperty(Hint = "An instant in the future at which this activity should execute.")] public Instant Instant { get; set; } - public Instant ExecuteAt + public Instant? ExecuteAt { - get => GetState(); + get => GetState(); set => SetState(value); } protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { ExecuteAt = Instant; - return Suspend(); + var now = _clock.GetCurrentInstant(); + return ExecuteAt <= now ? Done() : Suspend(); } protected override IActivityExecutionResult OnResume() => Done(); diff --git a/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs b/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs index 2c3d31d8f..d8c1d00bc 100644 --- a/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs +++ b/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs @@ -1,6 +1,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Extensions; +using Elsa.Models; using Elsa.Services; using Elsa.Triggers; using NodaTime; @@ -32,18 +33,25 @@ namespace Elsa.Activities.Timers.Triggers public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) { // Only provide a trigger if the workflow hasn't executed already sometime in the past. - var workflowDefinitionId = context.ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint.Id; - var instanceCount = await _workflowInstanceManager.ListByDefinitionAsync(workflowDefinitionId, cancellationToken).Count(); + var activity = context.GetActivity(); var executeAt = context.Activity.GetState(x => x.ExecuteAt); + var workflowDefinitionId = context.ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint.Id; + var instanceCount = await _workflowInstanceManager.ListByDefinitionAsync(workflowDefinitionId, cancellationToken).Where(x => x.Status == WorkflowStatus.Finished).Count(); var now = _clock.GetCurrentInstant(); - + + if (executeAt == null) + { + var futureInstant = await activity.GetPropertyValueAsync(x => x.Instant, cancellationToken); + executeAt = futureInstant; + } + // If the configured instant lies in the past, and the workflow was already executed once, we don't trigger again. if (executeAt <= now && instanceCount > 0) return NullTrigger.Instance; return new InstantEventTrigger { - ExecuteAt = executeAt + ExecuteAt = executeAt.Value }; } } diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowBlueprintWrapperExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowBlueprintWrapperExtensions.cs new file mode 100644 index 000000000..242cb3369 --- /dev/null +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowBlueprintWrapperExtensions.cs @@ -0,0 +1,14 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using Elsa.Services; +using Elsa.Services.Models; + +namespace Elsa +{ + public static class WorkflowBlueprintWrapperExtensions + { + public static IEnumerable> Filter(this IWorkflowBlueprintWrapper workflowBlueprintWrapper, Func? predicate = default) where TActivity : IActivity => + workflowBlueprintWrapper.Activities.Where(x => x.ActivityBlueprint.Type == typeof(TActivity).Name).Select(x => x.As()); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowBlueprintReflector.cs b/src/core/Elsa.Abstractions/Services/IWorkflowBlueprintReflector.cs new file mode 100644 index 000000000..1531082ab --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/IWorkflowBlueprintReflector.cs @@ -0,0 +1,11 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services.Models; + +namespace Elsa.Services +{ + public interface IWorkflowBlueprintReflector + { + public Task ReflectAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/ActivityBlueprintWrapper.cs b/src/core/Elsa.Abstractions/Services/Models/ActivityBlueprintWrapper.cs new file mode 100644 index 000000000..b8b37102a --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/ActivityBlueprintWrapper.cs @@ -0,0 +1,41 @@ +using System; +using System.Linq.Expressions; +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Services.Models +{ + public class ActivityBlueprintWrapper : IActivityBlueprintWrapper + { + protected ActivityExecutionContext ActivityExecutionContext { get; } + + public ActivityBlueprintWrapper(ActivityExecutionContext activityExecutionContext) + { + ActivityExecutionContext = activityExecutionContext; + } + + public IActivityBlueprint ActivityBlueprint => ActivityExecutionContext.ActivityBlueprint; + + public IActivityBlueprintWrapper As() where TActivity : IActivity => new ActivityBlueprintWrapper(ActivityExecutionContext); + } + + public class ActivityBlueprintWrapper : ActivityBlueprintWrapper, IActivityBlueprintWrapper where TActivity : IActivity + { + public ActivityBlueprintWrapper(ActivityExecutionContext activityExecutionContext) : base(activityExecutionContext) + { + } + + public async ValueTask GetPropertyValueAsync(Expression> propertyExpression, CancellationToken cancellationToken = default) + { + var workflowBlueprint = ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint; + var activityId = ActivityExecutionContext.ActivityBlueprint.Id; + return await workflowBlueprint.GetActivityPropertyValue(activityId, propertyExpression, ActivityExecutionContext, cancellationToken); + } + + public T? GetState(Expression> propertyExpression) + { + var workflowBlueprint = ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint; + return workflowBlueprint.GetActivityState(propertyExpression, ActivityExecutionContext); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IActivityBlueprintWrapper.cs b/src/core/Elsa.Abstractions/Services/Models/IActivityBlueprintWrapper.cs new file mode 100644 index 000000000..d404d57ce --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IActivityBlueprintWrapper.cs @@ -0,0 +1,19 @@ +using System; +using System.Linq.Expressions; +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Services.Models +{ + public interface IActivityBlueprintWrapper + { + IActivityBlueprint ActivityBlueprint { get; } + IActivityBlueprintWrapper As() where TActivity : IActivity; + } + + public interface IActivityBlueprintWrapper : IActivityBlueprintWrapper where TActivity:IActivity + { + ValueTask GetPropertyValueAsync(Expression> propertyExpression, CancellationToken cancellationToken = default); + T? GetState(Expression> propertyExpression); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprintWrapper.cs b/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprintWrapper.cs new file mode 100644 index 000000000..11b044225 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprintWrapper.cs @@ -0,0 +1,10 @@ +using System.Collections.Generic; + +namespace Elsa.Services.Models +{ + public interface IWorkflowBlueprintWrapper + { + IWorkflowBlueprint WorkflowBlueprint { get; } + IEnumerable Activities { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprintWrapper.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprintWrapper.cs new file mode 100644 index 000000000..6d9bc80da --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprintWrapper.cs @@ -0,0 +1,31 @@ +using System.Collections.Generic; + +namespace Elsa.Services.Models +{ + public class WorkflowBlueprintWrapper : IWorkflowBlueprintWrapper + { + private readonly WorkflowExecutionContext _workflowExecutionContext; + + public WorkflowBlueprintWrapper(IWorkflowBlueprint workflowBlueprint, WorkflowExecutionContext workflowExecutionContext) + { + _workflowExecutionContext = workflowExecutionContext; + WorkflowBlueprint = workflowBlueprint; + Activities = GetActivityBlueprintWrappers(); + } + + public IWorkflowBlueprint WorkflowBlueprint { get; } + + public IEnumerable Activities { get; } + + private IEnumerable GetActivityBlueprintWrappers() + { + var activities = WorkflowBlueprint.Activities; + + foreach (var activity in activities) + { + var activityExecutionContext = new ActivityExecutionContext(_workflowExecutionContext, _workflowExecutionContext.ServiceProvider, activity); + yield return new ActivityBlueprintWrapper(activityExecutionContext); + } + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Triggers/ActivityBlueprintWrapper.cs b/src/core/Elsa.Abstractions/Triggers/ActivityBlueprintWrapper.cs deleted file mode 100644 index 8edf448e2..000000000 --- a/src/core/Elsa.Abstractions/Triggers/ActivityBlueprintWrapper.cs +++ /dev/null @@ -1,32 +0,0 @@ -using System; -using System.Linq.Expressions; -using System.Threading; -using System.Threading.Tasks; -using Elsa.Services; -using Elsa.Services.Models; - -namespace Elsa.Triggers -{ - public class ActivityBlueprintWrapper where TActivity : IActivity - { - private readonly ActivityExecutionContext _activityExecutionContext; - - public ActivityBlueprintWrapper(ActivityExecutionContext activityExecutionContext) - { - _activityExecutionContext = activityExecutionContext; - } - - public async ValueTask GetPropertyValueAsync(Expression> propertyExpression, CancellationToken cancellationToken = default) - { - var workflowBlueprint = _activityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint; - var activityId = _activityExecutionContext.ActivityBlueprint.Id; - return await workflowBlueprint.GetActivityPropertyValue(activityId, propertyExpression, _activityExecutionContext, cancellationToken); - } - - public T? GetState(Expression> propertyExpression) - { - var workflowBlueprint = _activityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint; - return workflowBlueprint.GetActivityState(propertyExpression, _activityExecutionContext); - } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs b/src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs index 155984cd2..ef117cf93 100644 --- a/src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs +++ b/src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs @@ -9,7 +9,7 @@ namespace Elsa.Triggers { public Type ForType() => typeof(T); public Type ForActivityType() => typeof(TActivity); - public virtual ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => new ValueTask(GetTrigger(context)); + public virtual ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => new(GetTrigger(context)); public virtual ITrigger GetTrigger(TriggerProviderContext context) => NullTrigger.Instance; async ValueTask ITriggerProvider.GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) diff --git a/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs b/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs index 927866bfe..0fbe55fa9 100644 --- a/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs +++ b/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs @@ -11,7 +11,7 @@ namespace Elsa.Triggers } public ActivityExecutionContext ActivityExecutionContext { get; } - public ActivityBlueprintWrapper GetActivity() where TActivity : IActivity => new(ActivityExecutionContext); + public IActivityBlueprintWrapper GetActivity() where TActivity : IActivity => new ActivityBlueprintWrapper(ActivityExecutionContext); } public class TriggerProviderContext : TriggerProviderContext where T:IActivity @@ -20,6 +20,6 @@ namespace Elsa.Triggers { } - public ActivityBlueprintWrapper Activity => GetActivity(); + public IActivityBlueprintWrapper Activity => GetActivity(); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index 16b977df6..d61372f73 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -100,6 +100,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/core/Elsa.Core/Services/ServiceBusFactory.cs b/src/core/Elsa.Core/Services/ServiceBusFactory.cs index dbae89da6..400a9206c 100644 --- a/src/core/Elsa.Core/Services/ServiceBusFactory.cs +++ b/src/core/Elsa.Core/Services/ServiceBusFactory.cs @@ -32,7 +32,7 @@ namespace Elsa.Services var configureContext = new ServiceBusEndpointConfigurationContext(configurer, queueName, map, _serviceProvider); _elsaOptions.ConfigureServiceBusEndpoint(configureContext); - + var newBus = configurer.Start(); _serviceBuses.Add(queueName, newBus); diff --git a/src/core/Elsa.Core/Services/WorkflowBlueprintReflector.cs b/src/core/Elsa.Core/Services/WorkflowBlueprintReflector.cs new file mode 100644 index 000000000..254afc675 --- /dev/null +++ b/src/core/Elsa.Core/Services/WorkflowBlueprintReflector.cs @@ -0,0 +1,26 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services.Models; + +namespace Elsa.Services +{ + public class WorkflowBlueprintReflector : IWorkflowBlueprintReflector + { + private readonly IWorkflowFactory _workflowFactory; + private readonly IServiceProvider _serviceProvider; + + public WorkflowBlueprintReflector(IWorkflowFactory workflowFactory, IServiceProvider serviceProvider) + { + _workflowFactory = workflowFactory; + _serviceProvider = serviceProvider; + } + + public async Task ReflectAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default) + { + var workflowInstance = await _workflowFactory.InstantiateAsync(workflowBlueprint, cancellationToken: cancellationToken); + var workflowExecutionContext = new WorkflowExecutionContext(_serviceProvider, workflowBlueprint, workflowInstance); + return new WorkflowBlueprintWrapper(workflowBlueprint, workflowExecutionContext); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Triggers/WorkflowSelector.cs b/src/core/Elsa.Core/Triggers/WorkflowSelector.cs index 20d545d44..ccf385ce7 100644 --- a/src/core/Elsa.Core/Triggers/WorkflowSelector.cs +++ b/src/core/Elsa.Core/Triggers/WorkflowSelector.cs @@ -19,6 +19,7 @@ namespace Elsa.Triggers private readonly IWorkflowFactory _workflowFactory; private readonly IWorkflowInstanceManager _workflowInstanceManager; private readonly IWorkflowContextManager _workflowContextManager; + private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector; private readonly IEnumerable _triggerProviders; private readonly IMemoryCache _memoryCache; private readonly IServiceProvider _serviceProvider; @@ -29,6 +30,7 @@ namespace Elsa.Triggers IWorkflowFactory workflowFactory, IWorkflowInstanceManager workflowInstanceManager, IWorkflowContextManager workflowContextManager, + IWorkflowBlueprintReflector workflowBlueprintReflector, IEnumerable triggerProviders, IMemoryCache memoryCache, IServiceProvider serviceProvider) @@ -37,6 +39,7 @@ namespace Elsa.Triggers _workflowFactory = workflowFactory; _workflowInstanceManager = workflowInstanceManager; _workflowContextManager = workflowContextManager; + _workflowBlueprintReflector = workflowBlueprintReflector; _triggerProviders = triggerProviders; _memoryCache = memoryCache; _serviceProvider = serviceProvider; diff --git a/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModelChangedEventArgs.cs b/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModelChangedEventArgs.cs index 38ff9d47f..ea652c0ae 100644 --- a/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModelChangedEventArgs.cs +++ b/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModelChangedEventArgs.cs @@ -1,7 +1,6 @@ using System; -using ElsaDashboard.Application.Models; -namespace ElsaDashboard.Application.Shared +namespace ElsaDashboard.Application.Models { public class WorkflowModelChangedEventArgs : EventArgs { diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/Elsa.Samples.AzureServiceBusWorker.csproj b/src/samples/Elsa.Samples.AzureServiceBusWorker/Elsa.Samples.AzureServiceBusWorker.csproj new file mode 100644 index 000000000..a96c9373d --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/Elsa.Samples.AzureServiceBusWorker.csproj @@ -0,0 +1,15 @@ + + + + net5.0 + + + + + + + + + + + diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/Messages/Greeting.cs b/src/samples/Elsa.Samples.AzureServiceBusWorker/Messages/Greeting.cs new file mode 100644 index 000000000..5a2d64e96 --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/Messages/Greeting.cs @@ -0,0 +1,10 @@ +namespace Elsa.Samples.AzureServiceBusWorker.Messages +{ + public class Greeting + { + public string From { get; set; } + public string To { get; set; } + public string Message { get; set; } + public override string ToString() => $"{From} says \"{Message}\" to {To}."; + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/Program.cs b/src/samples/Elsa.Samples.AzureServiceBusWorker/Program.cs new file mode 100644 index 000000000..9f4e9289a --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/Program.cs @@ -0,0 +1,30 @@ +using Elsa.Activities.AzureServiceBus.Extensions; +using Elsa.Samples.AzureServiceBusWorker.Workflows; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using NodaTime; +using YesSql.Provider.Sqlite; + +namespace Elsa.Samples.AzureServiceBusWorker +{ + public class Program + { + public static void Main(string[] args) + { + CreateHostBuilder(args).Build().Run(); + } + + public static IHostBuilder CreateHostBuilder(string[] args) => + Host.CreateDefaultBuilder(args) + .ConfigureServices((hostContext, services) => + { + services + .AddElsa() + .AddConsoleActivities() + .AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1)) + .AddAzureServiceBusActivities(options => options.ConnectionString = "Endpoint=sb://elsa-workflows-2.servicebus.windows.net/;SharedAccessKeyName=Elsa;SharedAccessKey=hAIa+fFuUbHi94y1Z0uO/2UTccjN/y4W0xvpaUd/cr4=") + .AddWorkflow() + .AddWorkflow(); + }); + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/Properties/launchSettings.json b/src/samples/Elsa.Samples.AzureServiceBusWorker/Properties/launchSettings.json new file mode 100644 index 000000000..a535fa630 --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/Properties/launchSettings.json @@ -0,0 +1,11 @@ +{ + "profiles": { + "Elsa.Samples.Rebus.AzureServiceBusWorker": { + "commandName": "Project", + "dotnetRunMessages": "true", + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + } + } + } +} diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/Workflows/ConsumerWorkflow.cs b/src/samples/Elsa.Samples.AzureServiceBusWorker/Workflows/ConsumerWorkflow.cs new file mode 100644 index 000000000..861a68822 --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/Workflows/ConsumerWorkflow.cs @@ -0,0 +1,21 @@ +using Elsa.Activities.AzureServiceBus.Activities; +using Elsa.Activities.Console; +using Elsa.Builders; +using Elsa.Samples.AzureServiceBusWorker.Messages; + +namespace Elsa.Samples.AzureServiceBusWorker.Workflows +{ + public class ConsumerWorkflow : IWorkflow + { + public void Build(IWorkflowBuilder workflow) + { + workflow + .StartWith(messageReceived => messageReceived.Set(x => x.QueueName, "greetings")) + .WriteLine(context => + { + var greeting = context.GetInput(); + return $"Received a greeting from {greeting.From}, saying \"{greeting.Message}\" to {greeting.To}!"; + }); + } + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/Workflows/ProducerWorkflow.cs b/src/samples/Elsa.Samples.AzureServiceBusWorker/Workflows/ProducerWorkflow.cs new file mode 100644 index 000000000..d479a5016 --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/Workflows/ProducerWorkflow.cs @@ -0,0 +1,48 @@ +using System; +using Elsa.Activities.AzureServiceBus.Activities; +using Elsa.Activities.Console; +using Elsa.Activities.Timers; +using Elsa.Builders; +using Elsa.Samples.AzureServiceBusWorker.Messages; +using NodaTime; + +namespace Elsa.Samples.AzureServiceBusWorker.Workflows +{ + public class ProducerWorkflow : IWorkflow + { + private readonly IClock _clock; + private readonly Random _random; + + public ProducerWorkflow(IClock clock) + { + _clock = clock; + _random = new Random(); + } + + public void Build(IWorkflowBuilder workflow) + { + workflow + .InstantEvent(_clock.GetCurrentInstant().Plus(Duration.FromSeconds(5))) + .WriteLine("Sending a random greeting to the \"greetings\" queue.") + .Then(sendMessage => sendMessage + .Set(x => x.Message, GetRandomGreeting) + .Set(x => x.QueueName, "greetings")); + } + + private Greeting GetRandomGreeting() + { + var names = new[] { "John", "Jill", "Julia", "Miriam", "Jack", "Bob" }; + var messages = new[] { "Hello!", "How do you do?", "Happy Monday!" }; + var from = _random.Next(0, names.Length); + var to = _random.Next(0, names.Length); + var message = _random.Next(0, messages.Length); + + return new Greeting + { + From = names[from], + To = names[to], + Message = messages[message] + }; + } + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/appsettings.Development.json b/src/samples/Elsa.Samples.AzureServiceBusWorker/appsettings.Development.json new file mode 100644 index 000000000..8983e0fc1 --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/appsettings.Development.json @@ -0,0 +1,9 @@ +{ + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + } +} diff --git a/src/samples/Elsa.Samples.AzureServiceBusWorker/appsettings.json b/src/samples/Elsa.Samples.AzureServiceBusWorker/appsettings.json new file mode 100644 index 000000000..8983e0fc1 --- /dev/null +++ b/src/samples/Elsa.Samples.AzureServiceBusWorker/appsettings.json @@ -0,0 +1,9 @@ +{ + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + } +} diff --git a/src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj b/src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj index aec835fa3..aa90ba825 100644 --- a/src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj +++ b/src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj @@ -7,7 +7,7 @@ - + diff --git a/src/samples/Elsa.Samples.RebusWorker/Program.cs b/src/samples/Elsa.Samples.RebusWorker/Program.cs index 5b1d704c0..04310c8f4 100644 --- a/src/samples/Elsa.Samples.RebusWorker/Program.cs +++ b/src/samples/Elsa.Samples.RebusWorker/Program.cs @@ -2,7 +2,6 @@ using Elsa.Activities.Rebus.Extensions; using Elsa.Samples.RebusWorker.Messages; using Elsa.Samples.RebusWorker.Workflows; using Elsa.ServiceBus.AzureServiceBus; -using Elsa.ServiceBus.RabbitMq.Extensions; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using NodaTime; @@ -23,15 +22,14 @@ namespace Elsa.Samples.RebusWorker .ConfigureServices((hostContext, services) => { services - .AddElsa(option => option - .UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared")) - .UseAzureServiceBus("Endpoint=sb://elsa-workflows.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=n4NBTw9eSX12AG5BdIkyxCRroJGvh+EMOOM8ypWxWrQ=", LogLevel.Debug)) - //.UseRabbitMq("amqp://localhost")) + .AddElsa(option => option.UseAzureServiceBus("Endpoint=sb://elsa-workflows-2.servicebus.windows.net/;SharedAccessKeyName=Elsa;SharedAccessKey=hAIa+fFuUbHi94y1Z0uO/2UTccjN/y4W0xvpaUd/cr4=", LogLevel.Debug)) + //.UseRabbitMq("amqp://localhost")) .AddConsoleActivities() .AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1)) .AddRebusActivities() - .AddWorkflow() - .AddWorkflow(); + //.AddHostedService() + .AddWorkflow() + .AddWorkflow(); }); } } \ No newline at end of file diff --git a/src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs b/src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs index dce443023..01a036495 100644 --- a/src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs +++ b/src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs @@ -10,7 +10,7 @@ namespace Elsa.Samples.RebusWorker.Workflows public void Build(IWorkflowBuilder workflow) { workflow - .StartWith(messageReceived => messageReceived.Set(x => x.MessageType, typeof(Greeting))) + .StartWith(messageReceived => messageReceived.Set(x => x.MessageType, typeof(Greeting))) .WriteLine(context => { var greeting = context.GetInput(); diff --git a/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs b/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs index b184235e1..f23857425 100644 --- a/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs +++ b/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs @@ -22,9 +22,9 @@ namespace Elsa.Samples.RebusWorker.Workflows public void Build(IWorkflowBuilder workflow) { workflow - .TimerEvent(Duration.FromSeconds(5)) + .InstantEvent(_clock.GetCurrentInstant().Plus(Duration.FromSeconds(5))) .WriteLine("Sending a random greeting to the \"greetings\" queue.") - .Then(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting)); + .Then(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting)); } private Greeting GetRandomGreeting() diff --git a/src/servicebus/Elsa.ServiceBus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs b/src/servicebus/Elsa.ServiceBus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs index 6db14b8f2..100c4acd3 100644 --- a/src/servicebus/Elsa.ServiceBus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs +++ b/src/servicebus/Elsa.ServiceBus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs @@ -1,6 +1,8 @@ -using Elsa.Extensions; +using System.Threading.Tasks; +using Elsa.Extensions; using Elsa.Services; using Microsoft.Azure.ServiceBus.Primitives; +using Microsoft.Azure.Services.AppAuthentication; using Rebus.Config; using Rebus.Logging; using Rebus.Routing.TypeBased; @@ -9,12 +11,12 @@ namespace Elsa.ServiceBus.AzureServiceBus { public static class ElsaOptionsExtensions { - public static ElsaOptions UseAzureServiceBus(this ElsaOptions elsaOptions, string connectionString, LogLevel logLevel = LogLevel.Info, ITokenProvider tokenProvider = default) + public static ElsaOptions UseAzureServiceBus(this ElsaOptions elsaOptions, string connectionString, LogLevel logLevel = LogLevel.Info, ITokenProvider? tokenProvider = default) { return elsaOptions.UseServiceBus(context => ConfigureAzureServiceBusEndpoint(context, connectionString, logLevel, tokenProvider)); } - private static void ConfigureAzureServiceBusEndpoint(ServiceBusEndpointConfigurationContext context, string connectionString, LogLevel logLevel, ITokenProvider tokenProvider) + private static void ConfigureAzureServiceBusEndpoint(ServiceBusEndpointConfigurationContext context, string connectionString, LogLevel logLevel, ITokenProvider? tokenProvider) { var queueName = context.QueueName;