From 418cc85ba30a7857c296194513918c737f50615d Mon Sep 17 00:00:00 2001 From: jdevillard Date: Sun, 14 Mar 2021 12:10:31 +0100 Subject: [PATCH 1/2] Add Azure Servicebus Topic Activities (Sender & Receiver) (#749) * Add Service Bus Topic Sender & Receiver Activities * Ability to manage simple "string" message * Rename MessageReceived to MessageQueueReceived and SendMessage to SendQueueMessage * Missing renaming bookmark , fix TopicWorker Co-authored-by: jdevillard --- ...viceBusMessageReceivedBuilderExtensions.cs | 16 -- ...zureServiceBusMessageReceivedExtensions.cs | 23 --- ...=> AzureServiceBusQueueMessageReceived.cs} | 6 +- ...usQueueMessageReceivedBuilderExtensions.cs | 16 ++ ...erviceBusQueueMessageReceivedExtensions.cs | 23 +++ .../AzureServiceBusTopicMessageReceived.cs | 37 ++++ ...usTopicMessageReceivedBuilderExtensions.cs | 16 ++ ...erviceBusTopicMessageReceivedExtensions.cs | 29 +++ ...viceBusMessageReceivedBuilderExtensions.cs | 27 --- ...zureServiceBusMessageReceivedExtensions.cs | 22 --- ....cs => SendAzureServiceBusQueueMessage.cs} | 9 +- ...ServiceBusQueueMessageBuilderExtensions.cs | 27 +++ ...ndAzureServiceBusQueueMessageExtensions.cs | 22 +++ .../SendAzureServiceBusTopicMessage.cs | 41 +++++ ...ServiceBusTopicMessageBuilderExtensions.cs | 27 +++ ...ndAzureServiceBusTopicMessageExtensions.cs | 22 +++ ...ark.cs => QueueMessageReceivedBookmark.cs} | 14 +- .../Bookmarks/TopicMessageReceivedBookmark.cs | 39 ++++ .../Extensions/MessageBodyExtensions.cs | 22 ++- .../Extensions/ServiceCollectionExtensions.cs | 16 +- .../Services/ITopicMessageReceiverFactory.cs | 11 ++ .../Services/ITopicMessageSenderFactory.cs | 11 ++ .../Services/MessageBusFactory.cs | 64 ++++++- .../Services/QueueWorker.cs | 8 +- .../Services/TopicWorker.cs | 166 ++++++++++++++++++ .../StartupTasks/StartServiceBusQueues.cs | 4 +- .../StartServiceBusSubscription.cs | 64 +++++++ .../Workflows/ConsumerWorkflow.cs | 2 +- .../Workflows/ProducerWorkflow.cs | 2 +- 29 files changed, 669 insertions(+), 117 deletions(-) delete mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedBuilderExtensions.cs delete mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedExtensions.cs rename src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/{AzureServiceBusMessageReceived.cs => AzureServiceBusQueueMessageReceived.cs} (89%) create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedBuilderExtensions.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedExtensions.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceived.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedBuilderExtensions.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedExtensions.cs delete mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedBuilderExtensions.cs delete mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedExtensions.cs rename src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/{SendAzureServiceBusMessage.cs => SendAzureServiceBusQueueMessage.cs} (81%) create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageBuilderExtensions.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageExtensions.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageBuilderExtensions.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageExtensions.cs rename src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/{MessageReceivedBookmark.cs => QueueMessageReceivedBookmark.cs} (58%) create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/TopicMessageReceivedBookmark.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs create mode 100644 src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusSubscription.cs diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedBuilderExtensions.cs deleted file mode 100644 index e73b6eb54..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedBuilderExtensions.cs +++ /dev/null @@ -1,16 +0,0 @@ -using System; -using System.Runtime.CompilerServices; -using Elsa.Builders; - -// ReSharper disable ExplicitCallerInfoArgument -namespace Elsa.Activities.AzureServiceBus -{ - public static class AzureServiceBusMessageReceivedBuilderExtensions - { - public static IActivityBuilder MessageReceived(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => - builder.Then(setup, null, lineNumber, sourceFile); - - public static IActivityBuilder MessageReceived(this IBuilder builder, string queueName, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => - builder.MessageReceived(setup => setup.WithQueueName(queueName).WithMessageType(), lineNumber, sourceFile); - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedExtensions.cs deleted file mode 100644 index f091195ce..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceivedExtensions.cs +++ /dev/null @@ -1,23 +0,0 @@ -using System; -using System.Threading.Tasks; -using Elsa.Builders; -using Elsa.Services.Models; - -// ReSharper disable once CheckNamespace -namespace Elsa.Activities.AzureServiceBus -{ - public static class AzureServiceBusMessageReceivedExtensions - { - public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, string value) => messageReceived.Set(x => x.QueueName, value!); - - public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.MessageType, value!); - public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); - public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); - public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!); - public static ISetupActivity WithMessageType(this ISetupActivity messageReceived) => messageReceived.WithMessageType(typeof(T)); - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceived.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceived.cs similarity index 89% rename from src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceived.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceived.cs index 5843ccf7d..57a358c12 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusMessageReceived.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceived.cs @@ -1,4 +1,4 @@ -using System; +using System; using Elsa.Activities.AzureServiceBus.Extensions; using Elsa.Activities.AzureServiceBus.Models; using Elsa.ActivityResults; @@ -10,11 +10,11 @@ using Elsa.Services.Models; namespace Elsa.Activities.AzureServiceBus { [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 + public class AzureServiceBusQueueMessageReceived : Activity { private readonly IContentSerializer _serializer; - public AzureServiceBusMessageReceived(IContentSerializer serializer) + public AzureServiceBusQueueMessageReceived(IContentSerializer serializer) { _serializer = serializer; } diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedBuilderExtensions.cs new file mode 100644 index 000000000..b9cfabc6d --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedBuilderExtensions.cs @@ -0,0 +1,16 @@ +using System; +using System.Runtime.CompilerServices; +using Elsa.Builders; + +// ReSharper disable ExplicitCallerInfoArgument +namespace Elsa.Activities.AzureServiceBus +{ + public static class AzureServiceBusQueueMessageReceivedBuilderExtensions + { + public static IActivityBuilder MessageQueueReceived(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.Then(setup, null, lineNumber, sourceFile); + + public static IActivityBuilder MessageQueueReceived(this IBuilder builder, string queueName, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.MessageQueueReceived(setup => setup.WithQueueName(queueName).WithMessageType(), lineNumber, sourceFile); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedExtensions.cs new file mode 100644 index 000000000..ec702ca4a --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusQueueMessageReceivedExtensions.cs @@ -0,0 +1,23 @@ +using System; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.AzureServiceBus +{ + public static class AzureServiceBusQueueMessageReceivedExtensions + { + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity messageReceived, string value) => messageReceived.Set(x => x.QueueName, value!); + + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived) => messageReceived.WithMessageType(typeof(T)); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceived.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceived.cs new file mode 100644 index 000000000..1009297d1 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceived.cs @@ -0,0 +1,37 @@ +using System; +using Elsa.Activities.AzureServiceBus.Extensions; +using Elsa.Activities.AzureServiceBus.Models; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Serialization; +using Elsa.Services; +using Elsa.Services.Models; + +namespace Elsa.Activities.AzureServiceBus +{ + [Trigger(Category = "Azure Service Bus", DisplayName = "Service Bus Topic Message Received", Description = "Triggered when a message is received on the specified topic/subscription", Outcomes = new[] { OutcomeNames.Done })] + public class AzureServiceBusTopicMessageReceived : Activity + { + private readonly IContentSerializer _serializer; + + public AzureServiceBusTopicMessageReceived(IContentSerializer serializer) + { + _serializer = serializer; + } + + [ActivityProperty] public string TopicName { get; set; } = default!; + [ActivityProperty] public string SubscriptionName { get; set; } = default!; + [ActivityProperty] public Type MessageType { get; set; } = default!; + + protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => context.WorkflowExecutionContext.IsFirstPass ? ExecuteInternal(context) : Suspend(); + protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) => ExecuteInternal(context); + + private IActivityExecutionResult ExecuteInternal(ActivityExecutionContext context) + { + var message = (MessageModel)context.Input!; + var model = message.ReadBody(MessageType, _serializer); + + return Done(model); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedBuilderExtensions.cs new file mode 100644 index 000000000..5e3ef1131 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedBuilderExtensions.cs @@ -0,0 +1,16 @@ +using System; +using System.Runtime.CompilerServices; +using Elsa.Builders; + +// ReSharper disable ExplicitCallerInfoArgument +namespace Elsa.Activities.AzureServiceBus +{ + public static class AzureServiceBusTopicMessageReceivedBuilderExtensions + { + public static IActivityBuilder TopicMessageReceived(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.Then(setup, null, lineNumber, sourceFile); + + public static IActivityBuilder TopicMessageReceived(this IBuilder builder, string topicName,string subscriptionName, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.TopicMessageReceived(setup => setup.WithTopicName(topicName).WithSubscriptionName(subscriptionName).WithMessageType(), lineNumber, sourceFile); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedExtensions.cs new file mode 100644 index 000000000..d66aec526 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/AzureServiceBusMessageReceived/AzureServiceBusTopicMessageReceivedExtensions.cs @@ -0,0 +1,29 @@ +using System; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.AzureServiceBus +{ + public static class AzureServiceBusTopicMessageReceivedExtensions + { + public static ISetupActivity WithTopicName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity messageReceived, string value) => messageReceived.Set(x => x.TopicName, value!); + + public static ISetupActivity WithSubscriptionName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.SubscriptionName, value!); + public static ISetupActivity WithSubscriptionName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.SubscriptionName, value!); + public static ISetupActivity WithSubscriptionName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.SubscriptionName, value!); + public static ISetupActivity WithSubscriptionName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.SubscriptionName, value!); + public static ISetupActivity WithSubscriptionName(this ISetupActivity messageReceived, string value) => messageReceived.Set(x => x.SubscriptionName, value!); + + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!); + public static ISetupActivity WithMessageType(this ISetupActivity messageReceived) => messageReceived.WithMessageType(typeof(T)); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedBuilderExtensions.cs deleted file mode 100644 index 11e7246d8..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedBuilderExtensions.cs +++ /dev/null @@ -1,27 +0,0 @@ -using System; -using System.Runtime.CompilerServices; -using System.Threading.Tasks; -using Elsa.Builders; -using Elsa.Services.Models; - -// ReSharper disable ExplicitCallerInfoArgument -namespace Elsa.Activities.AzureServiceBus -{ - public static class SendAzureServiceBusMessageBuilderExtensions - { - public static IActivityBuilder SendMessage(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => - builder.Then(setup, null, lineNumber, sourceFile); - - public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func> message, [CallerLineNumber] int lineNumber = default, - [CallerFilePath] string? sourceFile = default) => builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); - - public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => - builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); - - public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => - builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); - - public static IActivityBuilder SendMessage(this IBuilder builder, string queueName, object message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => - builder.SendMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedExtensions.cs deleted file mode 100644 index 49e336f76..000000000 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusMessageReceivedExtensions.cs +++ /dev/null @@ -1,22 +0,0 @@ -using System; -using System.Threading.Tasks; -using Elsa.Builders; -using Elsa.Services.Models; - -// ReSharper disable once CheckNamespace -namespace Elsa.Activities.AzureServiceBus -{ - public static class SendAzureServiceBusMessageExtensions - { - public static ISetupActivity WithQueueName(this ISetupActivity activity, Func> value) => activity.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity activity, Func> value) => activity.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity activity, Func value) => activity.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity activity, Func value) => activity.Set(x => x.QueueName, value!); - public static ISetupActivity WithQueueName(this ISetupActivity activity, string value) => activity.Set(x => x.QueueName, value!); - - public static ISetupActivity WithMessage(this ISetupActivity activity, Func> value) => activity.Set(x => x.Message, value!); - public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); - public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); - public static ISetupActivity WithMessage(this ISetupActivity activity, object value) => activity.Set(x => x.Message, value!); - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusMessage.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs similarity index 81% rename from src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusMessage.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs index 5cbe36d21..c7422fdc7 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusMessage.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs @@ -11,12 +11,12 @@ using Microsoft.Azure.ServiceBus; namespace Elsa.Activities.AzureServiceBus { [Action(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 + public class SendAzureServiceBusQueueMessage : Activity { private readonly IMessageSenderFactory _messageSenderFactory; private readonly IContentSerializer _serializer; - public SendAzureServiceBusMessage(IMessageSenderFactory messageSenderFactory, IContentSerializer serializer) + public SendAzureServiceBusQueueMessage(IMessageSenderFactory messageSenderFactory, IContentSerializer serializer) { _messageSenderFactory = messageSenderFactory; _serializer = serializer; @@ -28,9 +28,8 @@ namespace Elsa.Activities.AzureServiceBus protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) { var sender = await _messageSenderFactory.GetSenderAsync(QueueName, context.CancellationToken); - var json = _serializer.Serialize(Message); - var bytes = Encoding.UTF8.GetBytes(json); - var message = new Message(bytes); + + var message = Extensions.MessageBodyExtensions.CreateMessage(_serializer, Message); if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId)) message.CorrelationId = context.WorkflowExecutionContext.CorrelationId; diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageBuilderExtensions.cs new file mode 100644 index 000000000..40393e113 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageBuilderExtensions.cs @@ -0,0 +1,27 @@ +using System; +using System.Runtime.CompilerServices; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable ExplicitCallerInfoArgument +namespace Elsa.Activities.AzureServiceBus +{ + public static class SendAzureServiceBusQueueMessageBuilderExtensions + { + public static IActivityBuilder SendQueueMessage(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.Then(setup, null, lineNumber, sourceFile); + + public static IActivityBuilder SendQueueMessage(this IBuilder builder, string queueName, Func> message, [CallerLineNumber] int lineNumber = default, + [CallerFilePath] string? sourceFile = default) => builder.SendQueueMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendQueueMessage(this IBuilder builder, string queueName, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendQueueMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendQueueMessage(this IBuilder builder, string queueName, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendQueueMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendQueueMessage(this IBuilder builder, string queueName, object message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendQueueMessage(setup => setup.WithQueueName(queueName).WithMessage(message), lineNumber, sourceFile); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageExtensions.cs new file mode 100644 index 000000000..b8d49fbcd --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessageExtensions.cs @@ -0,0 +1,22 @@ +using System; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.AzureServiceBus +{ + public static class SendAzureServiceBusQueueMessageExtensions + { + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func> value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func> value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, Func value) => activity.Set(x => x.QueueName, value!); + public static ISetupActivity WithQueueName(this ISetupActivity activity, string value) => activity.Set(x => x.QueueName, value!); + + public static ISetupActivity WithMessage(this ISetupActivity activity, Func> value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, object value) => activity.Set(x => x.Message, value!); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs new file mode 100644 index 000000000..dc46dfce2 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs @@ -0,0 +1,41 @@ +using System.Text; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Services; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Serialization; +using Elsa.Services; +using Elsa.Services.Models; +using Microsoft.Azure.ServiceBus; + +namespace Elsa.Activities.AzureServiceBus +{ + [Action(Category = "Azure Service Bus", DisplayName = "Send Service Bus Topic Message", Description = "Sends a message to the specified topic", Outcomes = new[] { OutcomeNames.Done })] + public class SendAzureServiceBusTopicMessage : Activity + { + private readonly ITopicMessageSenderFactory _messageSenderFactory; + private readonly IContentSerializer _serializer; + + public SendAzureServiceBusTopicMessage(ITopicMessageSenderFactory messageSenderFactory, IContentSerializer serializer) + { + _messageSenderFactory = messageSenderFactory; + _serializer = serializer; + } + + [ActivityProperty] public string TopicName { get; set; } = default!; + [ActivityProperty] public object Message { get; set; } = default!; + + protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) + { + var sender = await _messageSenderFactory.GetTopicSenderAsync(TopicName, context.CancellationToken); + + var message = Extensions.MessageBodyExtensions.CreateMessage(_serializer,Message); + + 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/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageBuilderExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageBuilderExtensions.cs new file mode 100644 index 000000000..86cf892f4 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageBuilderExtensions.cs @@ -0,0 +1,27 @@ +using System; +using System.Runtime.CompilerServices; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable ExplicitCallerInfoArgument +namespace Elsa.Activities.AzureServiceBus +{ + public static class SendAzureServiceBusTopicMessageBuilderExtensions + { + public static IActivityBuilder SendTopicMessage(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.Then(setup, null, lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string topicName, Func> message, [CallerLineNumber] int lineNumber = default, + [CallerFilePath] string? sourceFile = default) => builder.SendTopicMessage(setup => setup.WithTopicName(topicName).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string topicName, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithTopicName(topicName).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string topicName, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithTopicName(topicName).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string topicName, object message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithTopicName(topicName).WithMessage(message), lineNumber, sourceFile); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageExtensions.cs new file mode 100644 index 000000000..f33d52670 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessageExtensions.cs @@ -0,0 +1,22 @@ +using System; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.AzureServiceBus +{ + public static class SendAzureServiceBusTopicMessageExtensions + { + public static ISetupActivity WithTopicName(this ISetupActivity activity, Func> value) => activity.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity activity, Func> value) => activity.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity activity, Func value) => activity.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity activity, Func value) => activity.Set(x => x.TopicName, value!); + public static ISetupActivity WithTopicName(this ISetupActivity activity, string value) => activity.Set(x => x.TopicName, value!); + + public static ISetupActivity WithMessage(this ISetupActivity activity, Func> value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, Func value) => activity.Set(x => x.Message, value!); + public static ISetupActivity WithMessage(this ISetupActivity activity, object value) => activity.Set(x => x.Message, value!); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/MessageReceivedBookmark.cs b/src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/QueueMessageReceivedBookmark.cs similarity index 58% rename from src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/MessageReceivedBookmark.cs rename to src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/QueueMessageReceivedBookmark.cs index 9fdf38ef6..69e24b97b 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/MessageReceivedBookmark.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/QueueMessageReceivedBookmark.cs @@ -1,17 +1,17 @@ -using System.Collections.Generic; +using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.Bookmarks; namespace Elsa.Activities.AzureServiceBus.Bookmarks { - public class MessageReceivedBookmark : IBookmark + public class QueueMessageReceivedBookmark : IBookmark { - public MessageReceivedBookmark() + public QueueMessageReceivedBookmark() { } - public MessageReceivedBookmark(string queueName, string? correlationId = default) + public QueueMessageReceivedBookmark(string queueName, string? correlationId = default) { QueueName = queueName; CorrelationId = correlationId; @@ -21,12 +21,12 @@ namespace Elsa.Activities.AzureServiceBus.Bookmarks public string? CorrelationId { get; set; } } - public class MessageReceivedBookmarkProvider : BookmarkProvider + public class QueueMessageReceivedBookmarkProvider : BookmarkProvider { - public override async ValueTask> GetBookmarksAsync(BookmarkProviderContext context, CancellationToken cancellationToken) => + public override async ValueTask> GetBookmarksAsync(BookmarkProviderContext context, CancellationToken cancellationToken) => new[] { - new MessageReceivedBookmark + new QueueMessageReceivedBookmark { QueueName = (await context.Activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken))!, CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/TopicMessageReceivedBookmark.cs b/src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/TopicMessageReceivedBookmark.cs new file mode 100644 index 000000000..2bd88e493 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Bookmarks/TopicMessageReceivedBookmark.cs @@ -0,0 +1,39 @@ +using Elsa.Bookmarks; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Activities.AzureServiceBus.Bookmarks +{ + public class TopicMessageReceivedBookmark : IBookmark + { + public TopicMessageReceivedBookmark() + { + } + + public TopicMessageReceivedBookmark(string topicName, string subscriptionName, string? correlationId = default) + { + TopicName = topicName; + SubscriptionName = subscriptionName; + CorrelationId = correlationId; + } + + public string TopicName { get; set; } = default!; + public string SubscriptionName { get; set; } = default!; + public string? CorrelationId { get; set; } + } + + public class TopicMessageReceivedBookmarkProvider : BookmarkProvider + { + public override async ValueTask> GetBookmarksAsync(BookmarkProviderContext context, CancellationToken cancellationToken) => + new[] + { + new TopicMessageReceivedBookmark + { + TopicName = (await context.Activity.GetPropertyValueAsync(x => x.TopicName, cancellationToken))!, + SubscriptionName = (await context.Activity.GetPropertyValueAsync(x => x.SubscriptionName, cancellationToken))!, + CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId + } + }; + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/MessageBodyExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/MessageBodyExtensions.cs index e1808dbe3..e75d4b646 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/MessageBodyExtensions.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/MessageBodyExtensions.cs @@ -1,7 +1,8 @@ -using System; +using System; using System.Text; using Elsa.Activities.AzureServiceBus.Models; using Elsa.Serialization; +using Microsoft.Azure.ServiceBus; namespace Elsa.Activities.AzureServiceBus.Extensions { @@ -11,9 +12,28 @@ namespace Elsa.Activities.AzureServiceBus.Extensions public static object ReadBody(this MessageModel message, Type type, IContentSerializer serializer) { + if (type == typeof(string)) + return UTF8Encoding.UTF8.GetString(message.Body); + var bytes = message.Body; var json = Encoding.UTF8.GetString(bytes); return serializer.Deserialize(json, type)!; } + + public static Message CreateMessage(IContentSerializer serializer, object Message) + { + byte[] messageBytes; + + if (Message.GetType() == typeof(string)) + messageBytes = UTF8Encoding.UTF8.GetBytes(Message as string); + else + { + var json = serializer.Serialize(Message); + messageBytes = Encoding.UTF8.GetBytes(json); + } + + return new Message(messageBytes); + } + } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs index c2cf64438..398f25e71 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Extensions/ServiceCollectionExtensions.cs @@ -27,11 +27,21 @@ namespace Elsa.Activities.AzureServiceBus.Extensions .AddSingleton(sp => sp.GetRequiredService()) .AddSingleton(sp => sp.GetRequiredService()) .AddStartupTask() - .AddBookmarkProvider(); + .AddBookmarkProvider() + + .AddSingleton(sp => sp.GetRequiredService()) + .AddSingleton(sp => sp.GetRequiredService()) + .AddStartupTask() + .AddBookmarkProvider() + ; options - .AddActivity() - .AddActivity(); + .AddActivity() + .AddActivity() + .AddActivity() + .AddActivity() + + ; return options; } diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs new file mode 100644 index 000000000..e747140db --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageReceiverFactory.cs @@ -0,0 +1,11 @@ +using Microsoft.Azure.ServiceBus.Core; +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public interface ITopicMessageReceiverFactory + { + Task GetTopicReceiverAsync(string topicName, string subscriptionName, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs new file mode 100644 index 000000000..a0d16a411 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/ITopicMessageSenderFactory.cs @@ -0,0 +1,11 @@ +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Azure.ServiceBus.Core; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public interface ITopicMessageSenderFactory + { + Task GetTopicSenderAsync(string topicName, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs index 4a889df8d..6af468778 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/MessageBusFactory.cs @@ -1,4 +1,4 @@ -using System.Collections.Generic; +using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Microsoft.Azure.ServiceBus; @@ -7,12 +7,14 @@ using Microsoft.Azure.ServiceBus.Management; namespace Elsa.Activities.AzureServiceBus.Services { - public class MessageBusFactory : IMessageSenderFactory, IMessageReceiverFactory + public class MessageBusFactory : IMessageSenderFactory, IMessageReceiverFactory, ITopicMessageReceiverFactory, ITopicMessageSenderFactory { private readonly ServiceBusConnection _connection; private readonly ManagementClient _managementClient; private readonly IDictionary _senders = new Dictionary(); private readonly IDictionary _receivers = new Dictionary(); + + private readonly IDictionary<(string topicName, string queueName), IReceiverClient> _topicReceivers = new Dictionary<(string topicName, string queueName), IReceiverClient>(); private readonly SemaphoreSlim _semaphore = new(1); public MessageBusFactory(ServiceBusConnection connection, ManagementClient managementClient) @@ -68,5 +70,63 @@ namespace Elsa.Activities.AzureServiceBus.Services await _managementClient.CreateQueueAsync(queueName, cancellationToken); } + + public async Task GetTopicSenderAsync(string topicName, CancellationToken cancellationToken) + { + await _semaphore.WaitAsync(cancellationToken); + + try + { + if (_senders.TryGetValue(topicName, out var messageSender)) + return messageSender; + + await EnsureTopicExistsAsync(topicName, cancellationToken); + var newMessageSender = new MessageSender(_connection, topicName); + _senders.Add(topicName, newMessageSender); + return newMessageSender; + } + finally + { + _semaphore.Release(); + } + } + + public async Task GetTopicReceiverAsync(string topicName, string subscriptionName, CancellationToken cancellationToken) + { + await _semaphore.WaitAsync(cancellationToken); + + if (_topicReceivers.TryGetValue((topicName,subscriptionName), out var messageReceiver)) + return messageReceiver; + + try + { + await EnsureTopicAndSubscriptionExistsAsync(topicName, subscriptionName, cancellationToken); + + var newTopicMessageReceiver = new SubscriptionClient( + _connection, + topicPath: topicName, subscriptionName,ReceiveMode.PeekLock,RetryPolicy.Default) ; + + _topicReceivers.Add((topicName, subscriptionName), newTopicMessageReceiver); + return newTopicMessageReceiver; + } + finally + { + _semaphore.Release(); + } + } + + private async Task EnsureTopicExistsAsync(string topicName, CancellationToken cancellationToken) + { + if (!await _managementClient.TopicExistsAsync(topicName, cancellationToken)) + await _managementClient.CreateTopicAsync(topicName, cancellationToken); + } + + private async Task EnsureTopicAndSubscriptionExistsAsync(string topicName, string subscriptionName ,CancellationToken cancellationToken) + { + await EnsureTopicExistsAsync(topicName, cancellationToken); + + if(!await _managementClient.SubscriptionExistsAsync(topicName, subscriptionName, cancellationToken)) + await _managementClient.CreateSubscriptionAsync(topicName, subscriptionName, cancellationToken); + } } } \ 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 index 699e130c7..d17067792 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -87,8 +87,8 @@ namespace Elsa.Activities.AzureServiceBus.Services async Task TriggerNewWorkflowAsync() { - var bookmark = new MessageReceivedBookmark(queueName); - var triggers = await triggerFinder.FindTriggersAsync(bookmark, TenantId, cancellationToken); + var bookmark = new QueueMessageReceivedBookmark(queueName); + var triggers = await triggerFinder.FindTriggersAsync(bookmark, TenantId, cancellationToken); foreach (var trigger in triggers) { @@ -125,8 +125,8 @@ namespace Elsa.Activities.AzureServiceBus.Services { // Trigger existing workflows (if blocked on this message). _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}'. Resuming them", correlatedWorkflowInstanceCount, correlationId); - var bookmark = new MessageReceivedBookmark(queueName, correlationId); - var existingWorkflows = await bookmarkFinder.FindBookmarksAsync(bookmark, TenantId, cancellationToken).ToList(); + var bookmark = new QueueMessageReceivedBookmark(queueName, correlationId); + var existingWorkflows = await bookmarkFinder.FindBookmarksAsync(bookmark, TenantId, cancellationToken).ToList(); await workflowQueue.EnqueueWorkflowsAsync(existingWorkflows, model, model.CorrelationId, cancellationToken: cancellationToken); } else diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs new file mode 100644 index 000000000..f3620fd3e --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/TopicWorker.cs @@ -0,0 +1,166 @@ +using System; +using System.Diagnostics; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Bookmarks; +using Elsa.Activities.AzureServiceBus.Models; +using Elsa.Activities.AzureServiceBus.Options; +using Elsa.Bookmarks; +using Elsa.DistributedLock; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Persistence.Specifications; +using Elsa.Services; +using Elsa.Triggers; +using Microsoft.Azure.ServiceBus; +using Microsoft.Azure.ServiceBus.Core; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Activities.AzureServiceBus.Services +{ + public class TopicWorker : IAsyncDisposable + { + // TODO: Figure out how to start jobs across multiple tenants / how to get a list of all tenants. + private const string TenantId = default; + + private readonly IReceiverClient _messageReceiver; + private readonly IServiceScopeFactory _serviceScopeFactory; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly ILogger _logger; + + public TopicWorker( + IReceiverClient messageReceiver, + IServiceScopeFactory serviceScopeFactory, + IDistributedLockProvider distributedLockProvider, + IOptions options, + ILogger logger) + { + _messageReceiver = messageReceiver; + _serviceScopeFactory = serviceScopeFactory; + _distributedLockProvider = distributedLockProvider; + _logger = logger; + + _messageReceiver.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler) + { + AutoComplete = false, + MaxConcurrentCalls = options.Value.MaxConcurrentCalls + }); + } + + public async ValueTask DisposeAsync() => await _messageReceiver.CloseAsync(); + + private async Task OnMessageReceived(Message message, CancellationToken cancellationToken) + { + _logger.LogDebug("Message received with ID {MessageId}", message.MessageId); + await TriggerWorkflowsAsync(message, cancellationToken); + await _messageReceiver.CompleteAsync(message.SystemProperties.LockToken); + } + + private async Task TriggerWorkflowsAsync(Message message, CancellationToken cancellationToken) + { + using var scope = _serviceScopeFactory.CreateScope(); + var workflowQueue = scope.ServiceProvider.GetRequiredService(); + var topicName = _messageReceiver.Path.Split('/')[0]; + var subscriptionName = _messageReceiver.Path.Split('/')[2]; + var correlationId = message.CorrelationId; + var triggerFinder = scope.ServiceProvider.GetRequiredService(); + + var model = new MessageModel + { + Body = message.Body, + CorrelationId = message.CorrelationId, + ContentType = message.ContentType, + Label = message.Label, + To = message.To, + MessageId = message.MessageId, + PartitionKey = message.PartitionKey, + ViaPartitionKey = message.ViaPartitionKey, + ReplyTo = message.ReplyTo, + SessionId = message.SessionId, + ExpiresAtUtc = message.ExpiresAtUtc, + TimeToLive = message.TimeToLive, + ReplyToSessionId = message.ReplyToSessionId, + ScheduledEnqueueTimeUtc = message.ScheduledEnqueueTimeUtc + }; + + async Task TriggerNewWorkflowAsync() + { + var bookmark = new TopicMessageReceivedBookmark(topicName, subscriptionName); + var triggers = await triggerFinder.FindTriggersAsync(bookmark, TenantId, cancellationToken); + + foreach (var trigger in triggers) + { + var workflowBlueprint = trigger.WorkflowBlueprint; + await workflowQueue.EnqueueWorkflowDefinition(workflowBlueprint.Id, workflowBlueprint.TenantId, trigger.ActivityId, model, correlationId, null, cancellationToken); + } + } + + if (string.IsNullOrWhiteSpace(correlationId)) + { + await TriggerNewWorkflowAsync(); + return; + } + + var lockKey = $"azure-service-bus:{topicName}:{subscriptionName}:correlation-{correlationId}"; + var stopwatch = new Stopwatch(); + + _logger.LogDebug("Acquiring lock {LockKey}", lockKey); + stopwatch.Start(); + + if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken)) + { + _logger.LogDebug("Lock {LockKey} already taken", lockKey); + return; + } + + try + { + var bookmarkFinder = scope.ServiceProvider.GetRequiredService(); + var workflowInstanceStore = scope.ServiceProvider.GetRequiredService(); + var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification(model.CorrelationId), cancellationToken); + + if (correlatedWorkflowInstanceCount > 0) + { + // Trigger existing workflows (if blocked on this message). + _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}'. Resuming them", correlatedWorkflowInstanceCount, correlationId); + var bookmark = new TopicMessageReceivedBookmark(topicName, subscriptionName, correlationId); + var existingWorkflows = await bookmarkFinder.FindBookmarksAsync(bookmark, TenantId, cancellationToken).ToList(); + await workflowQueue.EnqueueWorkflowsAsync(existingWorkflows, model, model.CorrelationId, cancellationToken: cancellationToken); + } + else + { + // Trigger new workflow. + _logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId); + await TriggerNewWorkflowAsync(); + } + } + finally + { + await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken); + stopwatch.Stop(); + _logger.LogDebug("Lock held for {ElapseTime}", stopwatch.Elapsed); + } + } + + private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs e) + { + switch (e.Exception) + { + case MessageLockLostException: + _logger.LogDebug(e.Exception, "Message lock lost"); + break; + case ServiceBusCommunicationException: + _logger.LogDebug(e.Exception, "Lost service bus communication"); + break; + default: + _logger.LogError(e.Exception, "Unhandled exception"); + break; + } + + return Task.CompletedTask; + } + } +} \ 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 index bb1582ed3..d578fa73f 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusQueues.cs @@ -45,14 +45,14 @@ namespace Elsa.Activities.AzureServiceBus.StartupTasks var query = from workflow in workflows from activity in workflow.Activities - where activity.Type == nameof(AzureServiceBusMessageReceived) + where activity.Type == nameof(AzureServiceBusQueueMessageReceived) select workflow; foreach (var workflow in query) { var workflowBlueprintWrapper = await _workflowBlueprintReflector.ReflectAsync(_serviceProvider, workflow, cancellationToken); - foreach (var activity in workflowBlueprintWrapper.Filter()) + foreach (var activity in workflowBlueprintWrapper.Filter()) { var queueName = await activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken); yield return queueName!; diff --git a/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusSubscription.cs b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusSubscription.cs new file mode 100644 index 000000000..8aba8b1b2 --- /dev/null +++ b/src/activities/Elsa.Activities.AzureServiceBus/StartupTasks/StartServiceBusSubscription.cs @@ -0,0 +1,64 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Runtime.CompilerServices; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.AzureServiceBus.Services; +using Elsa.Services; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Activities.AzureServiceBus.StartupTasks +{ + public class StartServiceBusSubscription : IStartupTask + { + private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector; + private readonly ITopicMessageReceiverFactory _messageReceiverFactory; + private readonly IServiceProvider _serviceProvider; + + public StartServiceBusSubscription(IWorkflowBlueprintReflector workflowBlueprintReflector, ITopicMessageReceiverFactory messageReceiverFactory, IServiceProvider serviceProvider) + { + _workflowBlueprintReflector = workflowBlueprintReflector; + _messageReceiverFactory = messageReceiverFactory; + _serviceProvider = serviceProvider; + } + + public int Order => 2000; + + public async Task ExecuteAsync(CancellationToken stoppingToken) + { + var cancellationToken = stoppingToken; + var entities = (await GetTopicSubscriptionNamesAsync(cancellationToken).ToListAsync(cancellationToken)).Distinct(); + + foreach (var entity in entities) + { + var receiver = await _messageReceiverFactory.GetTopicReceiverAsync(entity.topicName, entity.subscriptionName, cancellationToken); + ActivatorUtilities.CreateInstance(_serviceProvider, receiver); + } + } + + private async IAsyncEnumerable<(string topicName, string subscriptionName)> GetTopicSubscriptionNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + var workflowRegistry = _serviceProvider.GetRequiredService(); + var workflows = await workflowRegistry.ListAsync(cancellationToken); + + var query = + from workflow in workflows + from activity in workflow.Activities + where activity.Type == nameof(AzureServiceBusTopicMessageReceived) + select workflow; + + foreach (var workflow in query) + { + var workflowBlueprintWrapper = await _workflowBlueprintReflector.ReflectAsync(_serviceProvider, workflow, cancellationToken); + + foreach (var activity in workflowBlueprintWrapper.Filter()) + { + var topicName = await activity.GetPropertyValueAsync(x => x.TopicName, cancellationToken); + var subscriptionName = await activity.GetPropertyValueAsync(x => x.SubscriptionName, cancellationToken); + yield return (topicName, subscriptionName)!; + } + } + } + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ConsumerWorkflow.cs b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ConsumerWorkflow.cs index 1db729fda..722d985c0 100644 --- a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ConsumerWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ConsumerWorkflow.cs @@ -10,7 +10,7 @@ namespace Elsa.Samples.AzureServiceBusWorker.Workflows public void Build(IWorkflowBuilder builder) { builder - .MessageReceived("greetings") + .MessageQueueReceived("greetings") .WriteLine(context => { var greeting = context.GetInput(); diff --git a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ProducerWorkflow.cs b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ProducerWorkflow.cs index f2ecb503c..d64b960ae 100644 --- a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ProducerWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Workflows/ProducerWorkflow.cs @@ -24,7 +24,7 @@ namespace Elsa.Samples.AzureServiceBusWorker.Workflows builder .Timer(Duration.FromSeconds(5)) .WriteLine("Sending a random greeting to the \"greetings\" queue.") - .SendMessage("greetings", GetRandomGreeting); + .SendQueueMessage("greetings", GetRandomGreeting); } private Greeting GetRandomGreeting() From 2320d216f2d5b117f32519156a6cd8a62960aa79 Mon Sep 17 00:00:00 2001 From: Craig Fowler Date: Sun, 14 Mar 2021 13:38:58 +0000 Subject: [PATCH 2/2] Resolve #683 - Tests to prove no exception (#758) As stated in the comments for WorkflowMayContainDuplicateActivitiesIntegrationTests These tests might not describe actually-desired behaviour. If they begin to "get in the way" in future, then it would probably be safe to remove them. --- ...thDuplicateActivitiesWorkflowAttributes.cs | 80 +++++++++++++++++++ .../Elsa.Core.IntegrationTests.csproj | 8 +- ...tainDuplicateActivitiesIntegrationTests.cs | 80 +++++++++++++++++++ .../Workflows/DuplicateActivitiesWorkflow.cs | 19 +++++ 4 files changed, 185 insertions(+), 2 deletions(-) create mode 100644 test/integration/Elsa.Core.IntegrationTests/Autofixture/HostBuilderWithDuplicateActivitiesWorkflowAttributes.cs create mode 100644 test/integration/Elsa.Core.IntegrationTests/Persistence/WorkflowMayContainDuplicateActivitiesIntegrationTests.cs create mode 100644 test/integration/Elsa.Core.IntegrationTests/Workflows/DuplicateActivitiesWorkflow.cs diff --git a/test/integration/Elsa.Core.IntegrationTests/Autofixture/HostBuilderWithDuplicateActivitiesWorkflowAttributes.cs b/test/integration/Elsa.Core.IntegrationTests/Autofixture/HostBuilderWithDuplicateActivitiesWorkflowAttributes.cs new file mode 100644 index 000000000..48ef3dddc --- /dev/null +++ b/test/integration/Elsa.Core.IntegrationTests/Autofixture/HostBuilderWithDuplicateActivitiesWorkflowAttributes.cs @@ -0,0 +1,80 @@ +using System.Reflection; +using AutoFixture; +using AutoFixture.Xunit2; +using Elsa.Core.IntegrationTests.Workflows; +using Elsa.Persistence.MongoDb.Extensions; +using Elsa.Testing.Shared.AutoFixture.Customizations; +using Microsoft.Extensions.DependencyInjection; +using Elsa.Persistence.EntityFramework.Core.Extensions; +using Microsoft.EntityFrameworkCore; +using Elsa.Persistence.EntityFramework.Sqlite; +using Elsa.Persistence.YesSql; +using YesSql.Provider.Sqlite; +using System.Data; + +namespace Elsa.Core.IntegrationTests.Autofixture +{ + public class HostBuilderWithDuplicateActivitiesWorkflowAttribute : CustomizeAttribute + { + public override ICustomization GetCustomization(ParameterInfo parameter) + { + return new HostBubilderUsingServicesCustomization(services => { + services + .AddElsa(elsa => { + elsa.AddWorkflow(); + }); + }, parameter); + } + } + + public class HostBuilderWithDuplicateActivitiesWorkflowAndMongoDbAttribute : CustomizeAttribute + { + public override ICustomization GetCustomization(ParameterInfo parameter) + { + return new HostBubilderUsingServicesCustomization(services => { + services + .AddElsa(elsa => { + elsa.AddWorkflow(); + elsa.UseMongoDbPersistence(opts => { + opts.ConnectionString = "mongodb://localhost:27017"; + opts.DatabaseName = "IntegrationTests"; + }); + }); + }, parameter); + } + } + + public class HostBuilderWithDuplicateActivitiesWorkflowAndEntityFrameworkAttribute : CustomizeAttribute + { + public override ICustomization GetCustomization(ParameterInfo parameter) + { + return new HostBubilderUsingServicesCustomization(services => { + services + .AddElsa(elsa => { + elsa + .AddWorkflow() + .UseEntityFrameworkPersistence(opts => { + opts.UseSqlite("Data Source=elsa.db;", db => db.MigrationsAssembly(typeof(SqliteElsaContextFactory).Assembly.GetName().Name)); + }); + }); + }, parameter); + } + } + + public class HostBuilderWithDuplicateActivitiesWorkflowAndYesSqlAttribute : CustomizeAttribute + { + public override ICustomization GetCustomization(ParameterInfo parameter) + { + return new HostBubilderUsingServicesCustomization(services => { + services + .AddElsa(elsa => { + elsa + .AddWorkflow() + .UseYesSqlPersistence(config => { + config.UseSqLite("Data Source=elsa-sqlite.db;", IsolationLevel.ReadUncommitted); + }); + }); + }, parameter); + } + } +} \ No newline at end of file diff --git a/test/integration/Elsa.Core.IntegrationTests/Elsa.Core.IntegrationTests.csproj b/test/integration/Elsa.Core.IntegrationTests/Elsa.Core.IntegrationTests.csproj index 58e23c58f..081da3096 100644 --- a/test/integration/Elsa.Core.IntegrationTests/Elsa.Core.IntegrationTests.csproj +++ b/test/integration/Elsa.Core.IntegrationTests/Elsa.Core.IntegrationTests.csproj @@ -23,11 +23,12 @@ all runtime; build; native; contentfiles; analyzers; buildtransitive - - + + + @@ -46,6 +47,9 @@ + + + diff --git a/test/integration/Elsa.Core.IntegrationTests/Persistence/WorkflowMayContainDuplicateActivitiesIntegrationTests.cs b/test/integration/Elsa.Core.IntegrationTests/Persistence/WorkflowMayContainDuplicateActivitiesIntegrationTests.cs new file mode 100644 index 000000000..498565d53 --- /dev/null +++ b/test/integration/Elsa.Core.IntegrationTests/Persistence/WorkflowMayContainDuplicateActivitiesIntegrationTests.cs @@ -0,0 +1,80 @@ +using Xunit; +using System.Threading.Tasks; +using Elsa.Core.IntegrationTests.Autofixture; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.DependencyInjection; +using System.Threading; +using Elsa.Services; +using Elsa.Core.IntegrationTests.Workflows; +using Elsa.Persistence; + +namespace Elsa.Core.IntegrationTests.Persistence +{ + public class WorkflowMayContainDuplicateActivitiesIntegrationTests + { + /* Please note that these tests might not represent _actually desired behaviour_. + * The tests do prove that issue #683 is no longer a problem, but it still does not seem logical + * that Elsa should want to allow duplicate activity IDs in a Workflow (definition or instance). + * + * It would be reasonable to remove these tests if they begin to "get in the way". + */ + + [Theory(DisplayName = "A workflow that contains duplicate activities may be run & persisted to an in-memory store"), AutoMoqData] + public async Task ADuplicateActivitiesWorkflowInstanceShouldBeRoundTrippableInMemory([HostBuilderWithDuplicateActivitiesWorkflow] IHostBuilder hostBuilder) + { + hostBuilder.ConfigureServices((ctx, services) => { + services.AddHostedService>(); + }); + var host = await hostBuilder.StartAsync(); + } + + [Theory(DisplayName = "A workflow that contains duplicate activities may be run & persisted to an EF Sqlite store"), AutoMoqData] + public async Task ADuplicateActivitiesWorkflowInstanceShouldBeRoundTrippableWithEntityFramework([HostBuilderWithDuplicateActivitiesWorkflowAndEntityFramework] IHostBuilder hostBuilder) + { + hostBuilder.ConfigureServices((ctx, services) => { + services.AddHostedService>(); + }); + var host = await hostBuilder.StartAsync(); + } + + [Theory(DisplayName = "A workflow that contains duplicate activities may be run & persisted to a MongoDb store"), AutoMoqData] + public async Task ADuplicateActivitiesWorkflowInstanceShouldBeRoundTrippableWithMongoDb([HostBuilderWithDuplicateActivitiesWorkflowAndMongoDb] IHostBuilder hostBuilder) + { + hostBuilder.ConfigureServices((ctx, services) => { + services.AddHostedService>(); + }); + var host = await hostBuilder.StartAsync(); + } + + [Theory(DisplayName = "A workflow that contains duplicate activities may be run & persisted to a YesSQL store"), AutoMoqData] + public async Task ADuplicateActivitiesWorkflowInstanceShouldBeRoundTrippableWithYesSql([HostBuilderWithDuplicateActivitiesWorkflowAndYesSql] IHostBuilder hostBuilder) + { + hostBuilder.ConfigureServices((ctx, services) => { + services.AddHostedService>(); + }); + var host = await hostBuilder.StartAsync(); + } + + class HostedWorkflowRunner : IHostedService where TWorkflow : DuplicateActivitiesWorkflow + { + readonly IWorkflowRunner workflowRunner; + readonly IWorkflowInstanceStore instanceStore; + + public async Task StartAsync(CancellationToken cancellationToken) + { + var instance = await workflowRunner.RunWorkflowAsync(); + var retrievedInstance = await instanceStore.FindByIdAsync(instance.Id); + + Assert.NotNull(retrievedInstance); + } + + public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; + + public HostedWorkflowRunner(IWorkflowRunner workflowRunner, IWorkflowInstanceStore instanceStore) + { + this.workflowRunner = workflowRunner ?? throw new System.ArgumentNullException(nameof(workflowRunner)); + this.instanceStore = instanceStore ?? throw new System.ArgumentNullException(nameof(instanceStore)); + } + } + } +} \ No newline at end of file diff --git a/test/integration/Elsa.Core.IntegrationTests/Workflows/DuplicateActivitiesWorkflow.cs b/test/integration/Elsa.Core.IntegrationTests/Workflows/DuplicateActivitiesWorkflow.cs new file mode 100644 index 000000000..00c8742c8 --- /dev/null +++ b/test/integration/Elsa.Core.IntegrationTests/Workflows/DuplicateActivitiesWorkflow.cs @@ -0,0 +1,19 @@ +using Elsa.Activities.Primitives; +using Elsa.Builders; + +namespace Elsa.Core.IntegrationTests.Workflows +{ + public class DuplicateActivitiesWorkflow : IWorkflow + { + const string duplicateId = "Duplicate"; + + public static readonly object Result = new object(); + + public virtual void Build(IWorkflowBuilder builder) + { + builder + .StartWith(a => a.Set(x => x.VariableName, "Unused").Set(x => x.Value, "Unused").Set(x => x.Id, duplicateId)) + .Then(a => a.Set(x => x.VariableName, "AlsoUnused").Set(x => x.Value, "Unused").Set(x => x.Id, duplicateId)); + } + } +} \ No newline at end of file