Merge branch 'feature/elsa-2.0' into feature/748-multi-tenant-EntityFramework

This commit is contained in:
Craig Fowler 2021-03-14 13:51:06 +00:00 committed by GitHub
commit 4bc3539c4d
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
33 changed files with 851 additions and 119 deletions

View file

@ -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<ISetupActivity<AzureServiceBusMessageReceived>> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) =>
builder.Then(setup, null, lineNumber, sourceFile);
public static IActivityBuilder MessageReceived<T>(this IBuilder builder, string queueName, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) =>
builder.MessageReceived(setup => setup.WithQueueName(queueName).WithMessageType<T>(), lineNumber, sourceFile);
}
}

View file

@ -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<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<string>> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ValueTask<string>> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<string> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, string> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, string value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<Type>> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<ActivityExecutionContext, Type> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Func<Type> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusMessageReceived> WithMessageType<T>(this ISetupActivity<AzureServiceBusMessageReceived> messageReceived) => messageReceived.WithMessageType(typeof(T));
}
}

View file

@ -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;
}

View file

@ -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<ISetupActivity<AzureServiceBusQueueMessageReceived>> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) =>
builder.Then(setup, null, lineNumber, sourceFile);
public static IActivityBuilder MessageQueueReceived<T>(this IBuilder builder, string queueName, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) =>
builder.MessageQueueReceived(setup => setup.WithQueueName(queueName).WithMessageType<T>(), lineNumber, sourceFile);
}
}

View file

@ -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<AzureServiceBusQueueMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<string>> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Func<ValueTask<string>> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Func<string> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Func<ActivityExecutionContext, string> value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithQueueName(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, string value) => messageReceived.Set(x => x.QueueName, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<Type>> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Func<ActivityExecutionContext, Type> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Func<Type> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusQueueMessageReceived> WithMessageType<T>(this ISetupActivity<AzureServiceBusQueueMessageReceived> messageReceived) => messageReceived.WithMessageType(typeof(T));
}
}

View file

@ -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);
}
}
}

View file

@ -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<ISetupActivity<AzureServiceBusTopicMessageReceived>> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) =>
builder.Then(setup, null, lineNumber, sourceFile);
public static IActivityBuilder TopicMessageReceived<T>(this IBuilder builder, string topicName,string subscriptionName, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) =>
builder.TopicMessageReceived(setup => setup.WithTopicName(topicName).WithSubscriptionName(subscriptionName).WithMessageType<T>(), lineNumber, sourceFile);
}
}

View file

@ -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<AzureServiceBusTopicMessageReceived> WithTopicName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<string>> value) => messageReceived.Set(x => x.TopicName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithTopicName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ValueTask<string>> value) => messageReceived.Set(x => x.TopicName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithTopicName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<string> value) => messageReceived.Set(x => x.TopicName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithTopicName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ActivityExecutionContext, string> value) => messageReceived.Set(x => x.TopicName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithTopicName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, string value) => messageReceived.Set(x => x.TopicName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithSubscriptionName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<string>> value) => messageReceived.Set(x => x.SubscriptionName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithSubscriptionName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ValueTask<string>> value) => messageReceived.Set(x => x.SubscriptionName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithSubscriptionName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<string> value) => messageReceived.Set(x => x.SubscriptionName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithSubscriptionName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ActivityExecutionContext, string> value) => messageReceived.Set(x => x.SubscriptionName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithSubscriptionName(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, string value) => messageReceived.Set(x => x.SubscriptionName, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ActivityExecutionContext, ValueTask<Type>> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<ActivityExecutionContext, Type> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Func<Type> value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithMessageType(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived, Type value) => messageReceived.Set(x => x.MessageType, value!);
public static ISetupActivity<AzureServiceBusTopicMessageReceived> WithMessageType<T>(this ISetupActivity<AzureServiceBusTopicMessageReceived> messageReceived) => messageReceived.WithMessageType(typeof(T));
}
}

View file

@ -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<ISetupActivity<SendAzureServiceBusMessage>> 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<ActivityExecutionContext, ValueTask<object>> 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<ActivityExecutionContext, object> 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<object> 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);
}
}

View file

@ -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<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, ValueTask<string>> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ValueTask<string>> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<string> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, string> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusMessage> activity, string value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, ValueTask<object>> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<ActivityExecutionContext, object> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, Func<object> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusMessage> WithMessage(this ISetupActivity<SendAzureServiceBusMessage> activity, object value) => activity.Set(x => x.Message, value!);
}
}

View file

@ -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<IActivityExecutionResult> 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;

View file

@ -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<ISetupActivity<SendAzureServiceBusQueueMessage>> 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<ActivityExecutionContext, ValueTask<object>> 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<ActivityExecutionContext, object> 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<object> 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);
}
}

View file

@ -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<SendAzureServiceBusQueueMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, Func<ActivityExecutionContext, ValueTask<string>> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, Func<ValueTask<string>> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, Func<string> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, Func<ActivityExecutionContext, string> value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithQueueName(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, string value) => activity.Set(x => x.QueueName, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithMessage(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, Func<ActivityExecutionContext, ValueTask<object>> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithMessage(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, Func<ActivityExecutionContext, object> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithMessage(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, Func<object> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusQueueMessage> WithMessage(this ISetupActivity<SendAzureServiceBusQueueMessage> activity, object value) => activity.Set(x => x.Message, value!);
}
}

View file

@ -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<IActivityExecutionResult> 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();
}
}
}

View file

@ -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<ISetupActivity<SendAzureServiceBusTopicMessage>> 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<ActivityExecutionContext, ValueTask<object>> 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<ActivityExecutionContext, object> 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<object> 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);
}
}

View file

@ -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<SendAzureServiceBusTopicMessage> WithTopicName(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, Func<ActivityExecutionContext, ValueTask<string>> value) => activity.Set(x => x.TopicName, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithTopicName(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, Func<ValueTask<string>> value) => activity.Set(x => x.TopicName, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithTopicName(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, Func<string> value) => activity.Set(x => x.TopicName, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithTopicName(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, Func<ActivityExecutionContext, string> value) => activity.Set(x => x.TopicName, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithTopicName(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, string value) => activity.Set(x => x.TopicName, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithMessage(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, Func<ActivityExecutionContext, ValueTask<object>> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithMessage(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, Func<ActivityExecutionContext, object> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithMessage(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, Func<object> value) => activity.Set(x => x.Message, value!);
public static ISetupActivity<SendAzureServiceBusTopicMessage> WithMessage(this ISetupActivity<SendAzureServiceBusTopicMessage> activity, object value) => activity.Set(x => x.Message, value!);
}
}

View file

@ -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<MessageReceivedBookmark, AzureServiceBusMessageReceived>
public class QueueMessageReceivedBookmarkProvider : BookmarkProvider<QueueMessageReceivedBookmark, AzureServiceBusQueueMessageReceived>
{
public override async ValueTask<IEnumerable<IBookmark>> GetBookmarksAsync(BookmarkProviderContext<AzureServiceBusMessageReceived> context, CancellationToken cancellationToken) =>
public override async ValueTask<IEnumerable<IBookmark>> GetBookmarksAsync(BookmarkProviderContext<AzureServiceBusQueueMessageReceived> context, CancellationToken cancellationToken) =>
new[]
{
new MessageReceivedBookmark
new QueueMessageReceivedBookmark
{
QueueName = (await context.Activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken))!,
CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId

View file

@ -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<TopicMessageReceivedBookmark, AzureServiceBusTopicMessageReceived>
{
public override async ValueTask<IEnumerable<IBookmark>> GetBookmarksAsync(BookmarkProviderContext<AzureServiceBusTopicMessageReceived> 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
}
};
}
}

View file

@ -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);
}
}
}

View file

@ -27,11 +27,21 @@ namespace Elsa.Activities.AzureServiceBus.Extensions
.AddSingleton<IMessageSenderFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddSingleton<IMessageReceiverFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddStartupTask<StartServiceBusQueues>()
.AddBookmarkProvider<MessageReceivedBookmarkProvider>();
.AddBookmarkProvider<QueueMessageReceivedBookmarkProvider>()
.AddSingleton<ITopicMessageSenderFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddSingleton<ITopicMessageReceiverFactory>(sp => sp.GetRequiredService<MessageBusFactory>())
.AddStartupTask<StartServiceBusSubscription>()
.AddBookmarkProvider<TopicMessageReceivedBookmarkProvider>()
;
options
.AddActivity<AzureServiceBusMessageReceived>()
.AddActivity<SendAzureServiceBusMessage>();
.AddActivity<AzureServiceBusQueueMessageReceived>()
.AddActivity<SendAzureServiceBusQueueMessage>()
.AddActivity<SendAzureServiceBusTopicMessage>()
.AddActivity<AzureServiceBusTopicMessageReceived>()
;
return options;
}

View file

@ -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<IReceiverClient> GetTopicReceiverAsync(string topicName, string subscriptionName, CancellationToken cancellationToken = default);
}
}

View file

@ -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<IMessageSender> GetTopicSenderAsync(string topicName, CancellationToken cancellationToken = default);
}
}

View file

@ -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<string, IMessageSender> _senders = new Dictionary<string, IMessageSender>();
private readonly IDictionary<string, IMessageReceiver> _receivers = new Dictionary<string, IMessageReceiver>();
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<IMessageSender> 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<IReceiverClient> 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);
}
}
}

View file

@ -87,8 +87,8 @@ namespace Elsa.Activities.AzureServiceBus.Services
async Task TriggerNewWorkflowAsync()
{
var bookmark = new MessageReceivedBookmark(queueName);
var triggers = await triggerFinder.FindTriggersAsync<AzureServiceBusMessageReceived>(bookmark, TenantId, cancellationToken);
var bookmark = new QueueMessageReceivedBookmark(queueName);
var triggers = await triggerFinder.FindTriggersAsync<AzureServiceBusQueueMessageReceived>(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<AzureServiceBusMessageReceived>(bookmark, TenantId, cancellationToken).ToList();
var bookmark = new QueueMessageReceivedBookmark(queueName, correlationId);
var existingWorkflows = await bookmarkFinder.FindBookmarksAsync<AzureServiceBusQueueMessageReceived>(bookmark, TenantId, cancellationToken).ToList();
await workflowQueue.EnqueueWorkflowsAsync(existingWorkflows, model, model.CorrelationId, cancellationToken: cancellationToken);
}
else

View file

@ -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<AzureServiceBusOptions> options,
ILogger<TopicWorker> 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<IWorkflowQueue>();
var topicName = _messageReceiver.Path.Split('/')[0];
var subscriptionName = _messageReceiver.Path.Split('/')[2];
var correlationId = message.CorrelationId;
var triggerFinder = scope.ServiceProvider.GetRequiredService<ITriggerFinder>();
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<AzureServiceBusTopicMessageReceived>(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<IBookmarkFinder>();
var workflowInstanceStore = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(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<AzureServiceBusTopicMessageReceived>(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;
}
}
}

View file

@ -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<AzureServiceBusMessageReceived>())
foreach (var activity in workflowBlueprintWrapper.Filter<AzureServiceBusQueueMessageReceived>())
{
var queueName = await activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken);
yield return queueName!;

View file

@ -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<TopicWorker>(_serviceProvider, receiver);
}
}
private async IAsyncEnumerable<(string topicName, string subscriptionName)> GetTopicSubscriptionNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken)
{
var workflowRegistry = _serviceProvider.GetRequiredService<IWorkflowRegistry>();
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<AzureServiceBusTopicMessageReceived>())
{
var topicName = await activity.GetPropertyValueAsync(x => x.TopicName, cancellationToken);
var subscriptionName = await activity.GetPropertyValueAsync(x => x.SubscriptionName, cancellationToken);
yield return (topicName, subscriptionName)!;
}
}
}
}
}

View file

@ -10,7 +10,7 @@ namespace Elsa.Samples.AzureServiceBusWorker.Workflows
public void Build(IWorkflowBuilder builder)
{
builder
.MessageReceived<Greeting>("greetings")
.MessageQueueReceived<Greeting>("greetings")
.WriteLine(context =>
{
var greeting = context.GetInput<Greeting>();

View file

@ -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()

View file

@ -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<DuplicateActivitiesWorkflow>();
});
}, parameter);
}
}
public class HostBuilderWithDuplicateActivitiesWorkflowAndMongoDbAttribute : CustomizeAttribute
{
public override ICustomization GetCustomization(ParameterInfo parameter)
{
return new HostBubilderUsingServicesCustomization(services => {
services
.AddElsa(elsa => {
elsa.AddWorkflow<DuplicateActivitiesWorkflow>();
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<DuplicateActivitiesWorkflow>()
.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<DuplicateActivitiesWorkflow>()
.UseYesSqlPersistence(config => {
config.UseSqLite("Data Source=elsa-sqlite.db;", IsolationLevel.ReadUncommitted);
});
});
}, parameter);
}
}
}

View file

@ -23,8 +23,8 @@
<PrivateAssets>all</PrivateAssets>
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
</PackageReference>
<PackageReference Include="YesSql.Core" Version="2.0.0-beta-1620" />
<PackageReference Include="YesSql.Provider.Sqlite" Version="1.0.0-beta-1620" />
<PackageReference Include="YesSql.Core" Version="2.0.0-beta-1637" />
<PackageReference Include="YesSql.Provider.Sqlite" Version="1.0.0-beta-1637" />
<PackageReference Include="Microsoft.Extensions.Hosting" Version="5.0.0" />
<PackageReference Include="Hangfire.AspNetCore" Version="1.7.19" />
<PackageReference Include="Hangfire.InMemory" Version="0.3.4" />
@ -47,6 +47,7 @@
<ProjectReference Include="..\..\..\src\activities\Elsa.Activities.Temporal.Hangfire\Elsa.Activities.Temporal.Hangfire.csproj" />
<ProjectReference Include="..\..\..\src\activities\Elsa.Activities.Temporal.Quartz\Elsa.Activities.Temporal.Quartz.csproj" />
<ProjectReference Include="..\..\..\src\persistence\Elsa.Persistence.MongoDb\Elsa.Persistence.MongoDb.csproj" />
<ProjectReference Include="..\..\..\src\persistence\Elsa.Persistence.YesSql\Elsa.Persistence.YesSql.csproj" />
<ProjectReference Include="..\..\..\src\persistence\Elsa.Persistence.EntityFramework\Elsa.Persistence.EntityFramework.Core\Elsa.Persistence.EntityFramework.Core.csproj" />
<ProjectReference Include="..\..\..\src\persistence\Elsa.Persistence.EntityFramework\Elsa.Persistence.EntityFramework.Sqlite\Elsa.Persistence.EntityFramework.Sqlite.csproj" />
</ItemGroup>

View file

@ -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<HostedWorkflowRunner<DuplicateActivitiesWorkflow>>();
});
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<HostedWorkflowRunner<DuplicateActivitiesWorkflow>>();
});
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<HostedWorkflowRunner<DuplicateActivitiesWorkflow>>();
});
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<HostedWorkflowRunner<DuplicateActivitiesWorkflow>>();
});
var host = await hostBuilder.StartAsync();
}
class HostedWorkflowRunner<TWorkflow> : IHostedService where TWorkflow : DuplicateActivitiesWorkflow
{
readonly IWorkflowRunner workflowRunner;
readonly IWorkflowInstanceStore instanceStore;
public async Task StartAsync(CancellationToken cancellationToken)
{
var instance = await workflowRunner.RunWorkflowAsync<TWorkflow>();
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));
}
}
}
}

View file

@ -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<SetVariable>(a => a.Set(x => x.VariableName, "Unused").Set(x => x.Value, "Unused").Set(x => x.Id, duplicateId))
.Then<SetVariable>(a => a.Set(x => x.VariableName, "AlsoUnused").Set(x => x.Value, "Unused").Set(x => x.Id, duplicateId));
}
}
}