WIP Azure Service Bus activities

This commit is contained in:
Sipke Schoorstra 2020-11-22 22:16:28 +01:00
parent 312fcf233a
commit 123e49c19d
46 changed files with 715 additions and 71 deletions

View file

@ -148,6 +148,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.ServiceBus.AzureServic
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.ServiceBus.RabbitMq", "src\servicebus\Elsa.ServiceBus.RabbitMq\Elsa.ServiceBus.RabbitMq.csproj", "{FC656D28-1AA3-462D-9B18-A10436E47C08}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Activities.AzureServiceBus", "src\activities\Elsa.Activities.AzureServiceBus\Elsa.Activities.AzureServiceBus.csproj", "{DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.AzureServiceBusWorker", "src\samples\Elsa.Samples.AzureServiceBusWorker\Elsa.Samples.AzureServiceBusWorker.csproj", "{1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -363,6 +367,14 @@ Global
{FC656D28-1AA3-462D-9B18-A10436E47C08}.Debug|Any CPU.Build.0 = Debug|Any CPU
{FC656D28-1AA3-462D-9B18-A10436E47C08}.Release|Any CPU.ActiveCfg = Release|Any CPU
{FC656D28-1AA3-462D-9B18-A10436E47C08}.Release|Any CPU.Build.0 = Release|Any CPU
{DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Debug|Any CPU.Build.0 = Debug|Any CPU
{DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Release|Any CPU.ActiveCfg = Release|Any CPU
{DD607BC1-F9D9-4685-BC8F-A4056C2C31FA}.Release|Any CPU.Build.0 = Release|Any CPU
{1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Debug|Any CPU.Build.0 = Debug|Any CPU
{1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Release|Any CPU.ActiveCfg = Release|Any CPU
{1D8B906A-5FCB-413C-8972-91D7A3EB8AE7}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -433,6 +445,8 @@ Global
{47EDFC14-BB99-44BF-A447-035D40D7347B} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925}
{D3440E66-A871-4AB0-A1E0-8B3864ACA16A} = {47EDFC14-BB99-44BF-A447-035D40D7347B}
{FC656D28-1AA3-462D-9B18-A10436E47C08} = {47EDFC14-BB99-44BF-A447-035D40D7347B}
{DD607BC1-F9D9-4685-BC8F-A4056C2C31FA} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180}
{1D8B906A-5FCB-413C-8972-91D7A3EB8AE7} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158}

View file

@ -0,0 +1,37 @@
using System;
using System.Text;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Serialization;
using Elsa.Services;
using Elsa.Services.Models;
using Microsoft.Azure.ServiceBus;
namespace Elsa.Activities.AzureServiceBus.Activities
{
[Trigger(Category = "Azure Service Bus", DisplayName = "Service Bus Message Received", Description = "Triggered when a message is received on the specified queue", Outcomes = new[] { OutcomeNames.Done })]
public class AzureServiceBusMessageReceived : Activity
{
private readonly IContentSerializer _serializer;
public AzureServiceBusMessageReceived(IContentSerializer serializer)
{
_serializer = serializer;
}
[ActivityProperty] public string QueueName { get; set; } = default!;
[ActivityProperty] public Type MessageType { get; set; } = default!;
protected override IActivityExecutionResult OnExecute() => Suspend();
protected override IActivityExecutionResult OnResume(ActivityExecutionContext context)
{
var message = (Message) context.Input!;
var bytes = message.Body;
var json = Encoding.UTF8.GetString(bytes);
var model = _serializer.Deserialize(json, MessageType);
return Done(model);
}
}
}

View file

@ -0,0 +1,43 @@
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Serialization;
using Elsa.Services;
using Elsa.Services.Models;
using Microsoft.Azure.ServiceBus;
using IServiceBusFactory = Elsa.Activities.AzureServiceBus.Services.IServiceBusFactory;
namespace Elsa.Activities.AzureServiceBus.Activities
{
[Trigger(Category = "Azure Service Bus", DisplayName = "Send Service Bus Message", Description = "Sends a message to the specified queue", Outcomes = new[] { OutcomeNames.Done })]
public class SendAzureServiceBusMessage : Activity
{
private readonly IServiceBusFactory _serviceBusFactory;
private readonly IContentSerializer _serializer;
public SendAzureServiceBusMessage(IServiceBusFactory serviceBusFactory, IContentSerializer serializer)
{
_serviceBusFactory = serviceBusFactory;
_serializer = serializer;
}
[ActivityProperty] public string QueueName { get; set; } = default!;
[ActivityProperty] public object Message { get; set; } = default!;
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
var sender = await _serviceBusFactory.GetSenderAsync(QueueName, cancellationToken);
var json = _serializer.Serialize(Message);
var bytes = Encoding.UTF8.GetBytes(json);
var message = new Message(bytes);
if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId))
message.CorrelationId = context.WorkflowExecutionContext.CorrelationId;
await sender.SendAsync(message);
return Done();
}
}
}

View file

@ -0,0 +1,18 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>netstandard2.0</TargetFramework>
<LangVersion>latest</LangVersion>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<!-- <PackageReference Include="Microsoft.Azure.Management.ServiceBus.Fluent" Version="1.35.0" />-->
<PackageReference Include="Microsoft.Azure.ServiceBus" Version="5.1.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,46 @@
using System;
using Elsa.Activities.AzureServiceBus.Activities;
using Elsa.Activities.AzureServiceBus.Options;
using Elsa.Activities.AzureServiceBus.Services;
using Elsa.Activities.AzureServiceBus.StartupTasks;
using Elsa.Runtime;
using Microsoft.Azure.ServiceBus;
using Microsoft.Azure.ServiceBus.Management;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
namespace Elsa.Activities.AzureServiceBus.Extensions
{
public static class ServiceCollectionExtensions
{
public static IServiceCollection AddAzureServiceBusActivities(this IServiceCollection services, Action<AzureServiceBusOptions>? configure)
{
if (configure != null)
services.Configure(configure);
else
services.AddOptions<AzureServiceBusOptions>();
return services
.AddSingleton(CreateServiceBusConnection)
.AddSingleton(CreateServiceBusManagementClient)
.AddSingleton<IServiceBusFactory, ServiceBusFactory>()
.AddStartupTask<StartServiceBusQueues>()
.AddActivity<AzureServiceBusMessageReceived>()
.AddActivity<SendAzureServiceBusMessage>();
}
private static ServiceBusConnection CreateServiceBusConnection(IServiceProvider serviceProvider)
{
var options = serviceProvider.GetRequiredService<IOptions<AzureServiceBusOptions>>().Value;
var connectionString = options.ConnectionString;
return new ServiceBusConnection(connectionString, RetryPolicy.Default);
}
private static ManagementClient CreateServiceBusManagementClient(IServiceProvider serviceProvider)
{
var options = serviceProvider.GetRequiredService<IOptions<AzureServiceBusOptions>>().Value;
var connectionString = options.ConnectionString;
return new ManagementClient(connectionString);
}
}
}

View file

@ -0,0 +1,7 @@
namespace Elsa.Activities.AzureServiceBus.Options
{
public class AzureServiceBusOptions
{
public string ConnectionString { get; set; }
}
}

View file

@ -0,0 +1,12 @@
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Azure.ServiceBus.Core;
namespace Elsa.Activities.AzureServiceBus.Services
{
public interface IServiceBusFactory
{
Task<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken = default);
Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken = default);
}
}

View file

@ -0,0 +1,40 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Triggers;
using Elsa.Services;
using Microsoft.Azure.ServiceBus.Core;
namespace Elsa.Activities.AzureServiceBus.Services
{
public class QueueWorker
{
private readonly IMessageReceiver _messageReceiver;
private readonly IWorkflowScheduler _workflowScheduler;
public QueueWorker(IMessageReceiver messageReceiver, IWorkflowScheduler workflowScheduler)
{
_messageReceiver = messageReceiver;
_workflowScheduler = workflowScheduler;
}
public async Task StartAsync(CancellationToken cancellationToken)
{
var cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
await Task.Factory.StartNew(() => ReadQueueAsync(cancellationTokenSource.Token), cancellationToken);
}
private async Task ReadQueueAsync(CancellationToken cancellationToken)
{
while (!cancellationToken.IsCancellationRequested)
{
var message = await _messageReceiver.ReceiveAsync();
if(message == null)
continue;
await _workflowScheduler.TriggerWorkflowsAsync<MessageReceivedTrigger>(x => x.QueueName == _messageReceiver.Path && (string.IsNullOrWhiteSpace(x.CorrelationId) || x.CorrelationId == message.CorrelationId), message,
message.CorrelationId, cancellationToken: cancellationToken);
}
}
}
}

View file

@ -0,0 +1,54 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Azure.ServiceBus;
using Microsoft.Azure.ServiceBus.Core;
using Microsoft.Azure.ServiceBus.Management;
namespace Elsa.Activities.AzureServiceBus.Services
{
public class ServiceBusFactory : IServiceBusFactory
{
private readonly ServiceBusConnection _connection;
private readonly ManagementClient _managementClient;
private readonly IDictionary<string, IMessageSender> _senders = new Dictionary<string, IMessageSender>();
private readonly IDictionary<string, IMessageReceiver> _receivers = new Dictionary<string, IMessageReceiver>();
public ServiceBusFactory(ServiceBusConnection connection, ManagementClient managementClient)
{
_connection = connection;
_managementClient = managementClient;
}
public async Task<IMessageSender> GetSenderAsync(string queueName, CancellationToken cancellationToken)
{
if (_senders.TryGetValue(queueName, out var messageSender))
return messageSender;
await EnsureQueueExistsAsync(queueName, cancellationToken);
var newMessageSender = new MessageSender(_connection, queueName);
_senders.Add(queueName, newMessageSender);
return newMessageSender;
}
public async Task<IMessageReceiver> GetReceiverAsync(string queueName, CancellationToken cancellationToken)
{
if (_receivers.TryGetValue(queueName, out var messageReceiver))
return messageReceiver;
await EnsureQueueExistsAsync(queueName, cancellationToken);
var newMessageReceiver = new MessageReceiver(_connection, queueName);
_receivers.Add(queueName, newMessageReceiver);
return newMessageReceiver;
}
private async Task EnsureQueueExistsAsync(string queueName, CancellationToken cancellationToken)
{
if (await _managementClient.QueueExistsAsync(queueName, cancellationToken))
return;
await _managementClient.CreateQueueAsync(queueName, cancellationToken);
}
}
}

View file

@ -0,0 +1,65 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Runtime.CompilerServices;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Activities;
using Elsa.Activities.AzureServiceBus.Services;
using Elsa.Services;
using Microsoft.Extensions.DependencyInjection;
using IServiceBusFactory = Elsa.Activities.AzureServiceBus.Services.IServiceBusFactory;
namespace Elsa.Activities.AzureServiceBus.StartupTasks
{
public class StartServiceBusQueues : IStartupTask
{
private readonly IWorkflowRegistry _workflowRegistry;
private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector;
private readonly IServiceBusFactory _serviceBusFactory;
private readonly IServiceProvider _serviceProvider;
public StartServiceBusQueues(IWorkflowRegistry workflowRegistry, IWorkflowBlueprintReflector workflowBlueprintReflector, IServiceBusFactory serviceBusFactory, IServiceProvider serviceProvider)
{
_workflowRegistry = workflowRegistry;
_workflowBlueprintReflector = workflowBlueprintReflector;
_serviceBusFactory = serviceBusFactory;
_serviceProvider = serviceProvider;
}
public async Task ExecuteAsync(CancellationToken cancellationToken = default)
{
var queueNames = await GetQueueNamesAsync(cancellationToken).ToListAsync(cancellationToken);
foreach (var queueName in queueNames)
{
var receiver = await _serviceBusFactory.GetReceiverAsync(queueName, cancellationToken);
var worker = ActivatorUtilities.CreateInstance<QueueWorker>(_serviceProvider, receiver);
await worker.StartAsync(cancellationToken);
}
}
private async IAsyncEnumerable<string> GetQueueNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken)
{
var workflows = await _workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken);
var query =
from workflow in workflows
from activity in workflow.Activities
where activity.Type == nameof(AzureServiceBusMessageReceived)
select workflow;
foreach (var workflow in query)
{
var workflowBlueprintWrapper = await _workflowBlueprintReflector.ReflectAsync(workflow, cancellationToken);
foreach (var activity in workflowBlueprintWrapper.Filter<AzureServiceBusMessageReceived>())
{
var queueName = await activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken);
yield return queueName;
}
}
}
}
}

View file

@ -0,0 +1,23 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Activities;
using Elsa.Triggers;
namespace Elsa.Activities.AzureServiceBus.Triggers
{
public class MessageReceivedTrigger : Trigger
{
public string QueueName { get; set; } = default!;
public string? CorrelationId { get; set; }
}
public class MessageReceivedTriggerProvider : TriggerProvider<MessageReceivedTrigger, AzureServiceBusMessageReceived>
{
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<AzureServiceBusMessageReceived> context, CancellationToken cancellationToken) =>
new MessageReceivedTrigger
{
QueueName = (await context.Activity.GetPropertyValueAsync(x => x.QueueName, cancellationToken)),
CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId
};
}
}

View file

@ -9,11 +9,11 @@ using Elsa.Services.Models;
namespace Elsa.Activities.Rebus
{
[Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })]
public class PublishMessage : Activity
public class PublishRebusMessage : Activity
{
private readonly IEventPublisher _eventPublisher;
public PublishMessage(IEventPublisher eventPublisher)
public PublishRebusMessage(IEventPublisher eventPublisher)
{
_eventPublisher = eventPublisher;
}

View file

@ -7,7 +7,7 @@ using Elsa.Services.Models;
namespace Elsa.Activities.Rebus
{
[Trigger(Category = "Rebus", Description = "Triggered when a message is received.", Outcomes = new[] { OutcomeNames.Done })]
public class MessageReceived : Activity
public class RebusMessageReceived : Activity
{
[ActivityProperty(Hint = "The type of message to receive.")]
public Type MessageType { get; set; } = default!;

View file

@ -9,11 +9,11 @@ using Elsa.Services.Models;
namespace Elsa.Activities.Rebus
{
[Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })]
public class SendMessage : Activity
public class SendRebusMessage : Activity
{
private readonly ICommandSender _bus;
public SendMessage(ICommandSender bus)
public SendRebusMessage(ICommandSender bus)
{
_bus = bus;
}

View file

@ -19,13 +19,13 @@
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj"/>
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
</ItemGroup>
<ItemGroup>
<None Include="icon.png">
<Pack>True</Pack>
<PackagePath/>
<PackagePath />
</None>
</ItemGroup>

View file

@ -22,9 +22,9 @@ namespace Elsa.Activities.Rebus.Extensions
return services
.AddTriggerProvider<MessageReceivedTriggerProvider>()
.AddStartupTask(sp => ActivatorUtilities.CreateInstance<CreateSubscriptions>(sp, (object)messageTypes))
.AddActivity<PublishMessage>()
.AddActivity<SendMessage>()
.AddActivity<MessageReceived>();
.AddActivity<PublishRebusMessage>()
.AddActivity<SendRebusMessage>()
.AddActivity<RebusMessageReceived>();
}
public static IServiceCollection AddRebusActivities<T>(this IServiceCollection services) => services.AddRebusActivities(typeof(T));

View file

@ -10,9 +10,9 @@ namespace Elsa.Activities.Rebus.Triggers
public string? CorrelationId { get; set; }
}
public class MessageReceivedTriggerProvider : TriggerProvider<MessageReceivedTrigger, MessageReceived>
public class MessageReceivedTriggerProvider : TriggerProvider<MessageReceivedTrigger, RebusMessageReceived>
{
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<MessageReceived> context, CancellationToken cancellationToken) =>
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<RebusMessageReceived> context, CancellationToken cancellationToken) =>
new MessageReceivedTrigger
{
MessageType = (await context.Activity.GetPropertyValueAsync(x => x.MessageType, cancellationToken)).Name,

View file

@ -26,16 +26,17 @@ namespace Elsa.Activities.Timers
[ActivityProperty(Hint = "An instant in the future at which this activity should execute.")]
public Instant Instant { get; set; }
public Instant ExecuteAt
public Instant? ExecuteAt
{
get => GetState<Instant>();
get => GetState<Instant?>();
set => SetState(value);
}
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context)
{
ExecuteAt = Instant;
return Suspend();
var now = _clock.GetCurrentInstant();
return ExecuteAt <= now ? Done() : Suspend();
}
protected override IActivityExecutionResult OnResume() => Done();

View file

@ -1,6 +1,7 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Extensions;
using Elsa.Models;
using Elsa.Services;
using Elsa.Triggers;
using NodaTime;
@ -32,18 +33,25 @@ namespace Elsa.Activities.Timers.Triggers
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<InstantEvent> context, CancellationToken cancellationToken)
{
// Only provide a trigger if the workflow hasn't executed already sometime in the past.
var workflowDefinitionId = context.ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint.Id;
var instanceCount = await _workflowInstanceManager.ListByDefinitionAsync(workflowDefinitionId, cancellationToken).Count();
var activity = context.GetActivity<InstantEvent>();
var executeAt = context.Activity.GetState(x => x.ExecuteAt);
var workflowDefinitionId = context.ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint.Id;
var instanceCount = await _workflowInstanceManager.ListByDefinitionAsync(workflowDefinitionId, cancellationToken).Where(x => x.Status == WorkflowStatus.Finished).Count();
var now = _clock.GetCurrentInstant();
if (executeAt == null)
{
var futureInstant = await activity.GetPropertyValueAsync(x => x.Instant, cancellationToken);
executeAt = futureInstant;
}
// If the configured instant lies in the past, and the workflow was already executed once, we don't trigger again.
if (executeAt <= now && instanceCount > 0)
return NullTrigger.Instance;
return new InstantEventTrigger
{
ExecuteAt = executeAt
ExecuteAt = executeAt.Value
};
}
}

View file

@ -0,0 +1,14 @@
using System;
using System.Collections.Generic;
using System.Linq;
using Elsa.Services;
using Elsa.Services.Models;
namespace Elsa
{
public static class WorkflowBlueprintWrapperExtensions
{
public static IEnumerable<IActivityBlueprintWrapper<TActivity>> Filter<TActivity>(this IWorkflowBlueprintWrapper workflowBlueprintWrapper, Func<TActivity, bool>? predicate = default) where TActivity : IActivity =>
workflowBlueprintWrapper.Activities.Where(x => x.ActivityBlueprint.Type == typeof(TActivity).Name).Select(x => x.As<TActivity>());
}
}

View file

@ -0,0 +1,11 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Services.Models;
namespace Elsa.Services
{
public interface IWorkflowBlueprintReflector
{
public Task<IWorkflowBlueprintWrapper> ReflectAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default);
}
}

View file

@ -0,0 +1,41 @@
using System;
using System.Linq.Expressions;
using System.Threading;
using System.Threading.Tasks;
namespace Elsa.Services.Models
{
public class ActivityBlueprintWrapper : IActivityBlueprintWrapper
{
protected ActivityExecutionContext ActivityExecutionContext { get; }
public ActivityBlueprintWrapper(ActivityExecutionContext activityExecutionContext)
{
ActivityExecutionContext = activityExecutionContext;
}
public IActivityBlueprint ActivityBlueprint => ActivityExecutionContext.ActivityBlueprint;
public IActivityBlueprintWrapper<TActivity> As<TActivity>() where TActivity : IActivity => new ActivityBlueprintWrapper<TActivity>(ActivityExecutionContext);
}
public class ActivityBlueprintWrapper<TActivity> : ActivityBlueprintWrapper, IActivityBlueprintWrapper<TActivity> where TActivity : IActivity
{
public ActivityBlueprintWrapper(ActivityExecutionContext activityExecutionContext) : base(activityExecutionContext)
{
}
public async ValueTask<T> GetPropertyValueAsync<T>(Expression<Func<TActivity, T>> propertyExpression, CancellationToken cancellationToken = default)
{
var workflowBlueprint = ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint;
var activityId = ActivityExecutionContext.ActivityBlueprint.Id;
return await workflowBlueprint.GetActivityPropertyValue(activityId, propertyExpression, ActivityExecutionContext, cancellationToken);
}
public T? GetState<T>(Expression<Func<TActivity, T>> propertyExpression)
{
var workflowBlueprint = ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint;
return workflowBlueprint.GetActivityState(propertyExpression, ActivityExecutionContext);
}
}
}

View file

@ -0,0 +1,19 @@
using System;
using System.Linq.Expressions;
using System.Threading;
using System.Threading.Tasks;
namespace Elsa.Services.Models
{
public interface IActivityBlueprintWrapper
{
IActivityBlueprint ActivityBlueprint { get; }
IActivityBlueprintWrapper<TActivity> As<TActivity>() where TActivity : IActivity;
}
public interface IActivityBlueprintWrapper<TActivity> : IActivityBlueprintWrapper where TActivity:IActivity
{
ValueTask<T> GetPropertyValueAsync<T>(Expression<Func<TActivity, T>> propertyExpression, CancellationToken cancellationToken = default);
T? GetState<T>(Expression<Func<TActivity, T>> propertyExpression);
}
}

View file

@ -0,0 +1,10 @@
using System.Collections.Generic;
namespace Elsa.Services.Models
{
public interface IWorkflowBlueprintWrapper
{
IWorkflowBlueprint WorkflowBlueprint { get; }
IEnumerable<IActivityBlueprintWrapper> Activities { get; }
}
}

View file

@ -0,0 +1,31 @@
using System.Collections.Generic;
namespace Elsa.Services.Models
{
public class WorkflowBlueprintWrapper : IWorkflowBlueprintWrapper
{
private readonly WorkflowExecutionContext _workflowExecutionContext;
public WorkflowBlueprintWrapper(IWorkflowBlueprint workflowBlueprint, WorkflowExecutionContext workflowExecutionContext)
{
_workflowExecutionContext = workflowExecutionContext;
WorkflowBlueprint = workflowBlueprint;
Activities = GetActivityBlueprintWrappers();
}
public IWorkflowBlueprint WorkflowBlueprint { get; }
public IEnumerable<IActivityBlueprintWrapper> Activities { get; }
private IEnumerable<IActivityBlueprintWrapper> GetActivityBlueprintWrappers()
{
var activities = WorkflowBlueprint.Activities;
foreach (var activity in activities)
{
var activityExecutionContext = new ActivityExecutionContext(_workflowExecutionContext, _workflowExecutionContext.ServiceProvider, activity);
yield return new ActivityBlueprintWrapper(activityExecutionContext);
}
}
}
}

View file

@ -1,32 +0,0 @@
using System;
using System.Linq.Expressions;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Services;
using Elsa.Services.Models;
namespace Elsa.Triggers
{
public class ActivityBlueprintWrapper<TActivity> where TActivity : IActivity
{
private readonly ActivityExecutionContext _activityExecutionContext;
public ActivityBlueprintWrapper(ActivityExecutionContext activityExecutionContext)
{
_activityExecutionContext = activityExecutionContext;
}
public async ValueTask<T> GetPropertyValueAsync<T>(Expression<Func<TActivity, T>> propertyExpression, CancellationToken cancellationToken = default)
{
var workflowBlueprint = _activityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint;
var activityId = _activityExecutionContext.ActivityBlueprint.Id;
return await workflowBlueprint.GetActivityPropertyValue(activityId, propertyExpression, _activityExecutionContext, cancellationToken);
}
public T? GetState<T>(Expression<Func<TActivity, T>> propertyExpression)
{
var workflowBlueprint = _activityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint;
return workflowBlueprint.GetActivityState(propertyExpression, _activityExecutionContext);
}
}
}

View file

@ -9,7 +9,7 @@ namespace Elsa.Triggers
{
public Type ForType() => typeof(T);
public Type ForActivityType() => typeof(TActivity);
public virtual ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<TActivity> context, CancellationToken cancellationToken) => new ValueTask<ITrigger>(GetTrigger(context));
public virtual ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<TActivity> context, CancellationToken cancellationToken) => new(GetTrigger(context));
public virtual ITrigger GetTrigger(TriggerProviderContext<TActivity> context) => NullTrigger.Instance;
async ValueTask<ITrigger> ITriggerProvider.GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken)

View file

@ -11,7 +11,7 @@ namespace Elsa.Triggers
}
public ActivityExecutionContext ActivityExecutionContext { get; }
public ActivityBlueprintWrapper<TActivity> GetActivity<TActivity>() where TActivity : IActivity => new(ActivityExecutionContext);
public IActivityBlueprintWrapper<TActivity> GetActivity<TActivity>() where TActivity : IActivity => new ActivityBlueprintWrapper<TActivity>(ActivityExecutionContext);
}
public class TriggerProviderContext<T> : TriggerProviderContext where T:IActivity
@ -20,6 +20,6 @@ namespace Elsa.Triggers
{
}
public ActivityBlueprintWrapper<T> Activity => GetActivity<T>();
public IActivityBlueprintWrapper<T> Activity => GetActivity<T>();
}
}

View file

@ -100,6 +100,7 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton<IWorkflowFactory, WorkflowFactory>()
.AddSingleton<IActivityFactory, ActivityFactory>()
.AddSingleton<IWorkflowBlueprintMaterializer, WorkflowBlueprintMaterializer>()
.AddSingleton<IWorkflowBlueprintReflector, WorkflowBlueprintReflector>()
.AddScoped<IWorkflowSelector, WorkflowSelector>()
.AddScoped<IWorkflowDefinitionManager, WorkflowDefinitionManager>()
.AddScoped<IWorkflowInstanceManager, WorkflowInstanceManager>()

View file

@ -32,7 +32,7 @@ namespace Elsa.Services
var configureContext = new ServiceBusEndpointConfigurationContext(configurer, queueName, map, _serviceProvider);
_elsaOptions.ConfigureServiceBusEndpoint(configureContext);
var newBus = configurer.Start();
_serviceBuses.Add(queueName, newBus);

View file

@ -0,0 +1,26 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Services.Models;
namespace Elsa.Services
{
public class WorkflowBlueprintReflector : IWorkflowBlueprintReflector
{
private readonly IWorkflowFactory _workflowFactory;
private readonly IServiceProvider _serviceProvider;
public WorkflowBlueprintReflector(IWorkflowFactory workflowFactory, IServiceProvider serviceProvider)
{
_workflowFactory = workflowFactory;
_serviceProvider = serviceProvider;
}
public async Task<IWorkflowBlueprintWrapper> ReflectAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default)
{
var workflowInstance = await _workflowFactory.InstantiateAsync(workflowBlueprint, cancellationToken: cancellationToken);
var workflowExecutionContext = new WorkflowExecutionContext(_serviceProvider, workflowBlueprint, workflowInstance);
return new WorkflowBlueprintWrapper(workflowBlueprint, workflowExecutionContext);
}
}
}

View file

@ -19,6 +19,7 @@ namespace Elsa.Triggers
private readonly IWorkflowFactory _workflowFactory;
private readonly IWorkflowInstanceManager _workflowInstanceManager;
private readonly IWorkflowContextManager _workflowContextManager;
private readonly IWorkflowBlueprintReflector _workflowBlueprintReflector;
private readonly IEnumerable<ITriggerProvider> _triggerProviders;
private readonly IMemoryCache _memoryCache;
private readonly IServiceProvider _serviceProvider;
@ -29,6 +30,7 @@ namespace Elsa.Triggers
IWorkflowFactory workflowFactory,
IWorkflowInstanceManager workflowInstanceManager,
IWorkflowContextManager workflowContextManager,
IWorkflowBlueprintReflector workflowBlueprintReflector,
IEnumerable<ITriggerProvider> triggerProviders,
IMemoryCache memoryCache,
IServiceProvider serviceProvider)
@ -37,6 +39,7 @@ namespace Elsa.Triggers
_workflowFactory = workflowFactory;
_workflowInstanceManager = workflowInstanceManager;
_workflowContextManager = workflowContextManager;
_workflowBlueprintReflector = workflowBlueprintReflector;
_triggerProviders = triggerProviders;
_memoryCache = memoryCache;
_serviceProvider = serviceProvider;

View file

@ -1,7 +1,6 @@
using System;
using ElsaDashboard.Application.Models;
namespace ElsaDashboard.Application.Shared
namespace ElsaDashboard.Application.Models
{
public class WorkflowModelChangedEventArgs : EventArgs
{

View file

@ -0,0 +1,15 @@
<Project Sdk="Microsoft.NET.Sdk.Worker">
<PropertyGroup>
<TargetFramework>net5.0</TargetFramework>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.Hosting" Version="5.0.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\activities\Elsa.Activities.AzureServiceBus\Elsa.Activities.AzureServiceBus.csproj" />
<ProjectReference Include="..\..\core\Elsa\Elsa.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,10 @@
namespace Elsa.Samples.AzureServiceBusWorker.Messages
{
public class Greeting
{
public string From { get; set; }
public string To { get; set; }
public string Message { get; set; }
public override string ToString() => $"{From} says \"{Message}\" to {To}.";
}
}

View file

@ -0,0 +1,30 @@
using Elsa.Activities.AzureServiceBus.Extensions;
using Elsa.Samples.AzureServiceBusWorker.Workflows;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NodaTime;
using YesSql.Provider.Sqlite;
namespace Elsa.Samples.AzureServiceBusWorker
{
public class Program
{
public static void Main(string[] args)
{
CreateHostBuilder(args).Build().Run();
}
public static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureServices((hostContext, services) =>
{
services
.AddElsa()
.AddConsoleActivities()
.AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1))
.AddAzureServiceBusActivities(options => options.ConnectionString = "Endpoint=sb://elsa-workflows-2.servicebus.windows.net/;SharedAccessKeyName=Elsa;SharedAccessKey=hAIa+fFuUbHi94y1Z0uO/2UTccjN/y4W0xvpaUd/cr4=")
.AddWorkflow<ProducerWorkflow>()
.AddWorkflow<ConsumerWorkflow>();
});
}
}

View file

@ -0,0 +1,11 @@
{
"profiles": {
"Elsa.Samples.Rebus.AzureServiceBusWorker": {
"commandName": "Project",
"dotnetRunMessages": "true",
"environmentVariables": {
"DOTNET_ENVIRONMENT": "Development"
}
}
}
}

View file

@ -0,0 +1,21 @@
using Elsa.Activities.AzureServiceBus.Activities;
using Elsa.Activities.Console;
using Elsa.Builders;
using Elsa.Samples.AzureServiceBusWorker.Messages;
namespace Elsa.Samples.AzureServiceBusWorker.Workflows
{
public class ConsumerWorkflow : IWorkflow
{
public void Build(IWorkflowBuilder workflow)
{
workflow
.StartWith<AzureServiceBusMessageReceived>(messageReceived => messageReceived.Set(x => x.QueueName, "greetings"))
.WriteLine(context =>
{
var greeting = context.GetInput<Greeting>();
return $"Received a greeting from {greeting.From}, saying \"{greeting.Message}\" to {greeting.To}!";
});
}
}
}

View file

@ -0,0 +1,48 @@
using System;
using Elsa.Activities.AzureServiceBus.Activities;
using Elsa.Activities.Console;
using Elsa.Activities.Timers;
using Elsa.Builders;
using Elsa.Samples.AzureServiceBusWorker.Messages;
using NodaTime;
namespace Elsa.Samples.AzureServiceBusWorker.Workflows
{
public class ProducerWorkflow : IWorkflow
{
private readonly IClock _clock;
private readonly Random _random;
public ProducerWorkflow(IClock clock)
{
_clock = clock;
_random = new Random();
}
public void Build(IWorkflowBuilder workflow)
{
workflow
.InstantEvent(_clock.GetCurrentInstant().Plus(Duration.FromSeconds(5)))
.WriteLine("Sending a random greeting to the \"greetings\" queue.")
.Then<SendAzureServiceBusMessage>(sendMessage => sendMessage
.Set(x => x.Message, GetRandomGreeting)
.Set(x => x.QueueName, "greetings"));
}
private Greeting GetRandomGreeting()
{
var names = new[] { "John", "Jill", "Julia", "Miriam", "Jack", "Bob" };
var messages = new[] { "Hello!", "How do you do?", "Happy Monday!" };
var from = _random.Next(0, names.Length);
var to = _random.Next(0, names.Length);
var message = _random.Next(0, messages.Length);
return new Greeting
{
From = names[from],
To = names[to],
Message = messages[message]
};
}
}
}

View file

@ -0,0 +1,9 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
}
}

View file

@ -0,0 +1,9 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
}
}

View file

@ -7,7 +7,7 @@
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.Hosting" Version="5.0.0" />
<PackageReference Include="Rebus.AzureServiceBus" Version="7.1.6" />
<PackageReference Include="Rebus.AzureServiceBus" Version="8.0.0-a1" />
</ItemGroup>
<ItemGroup>

View file

@ -2,7 +2,6 @@ using Elsa.Activities.Rebus.Extensions;
using Elsa.Samples.RebusWorker.Messages;
using Elsa.Samples.RebusWorker.Workflows;
using Elsa.ServiceBus.AzureServiceBus;
using Elsa.ServiceBus.RabbitMq.Extensions;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NodaTime;
@ -23,15 +22,14 @@ namespace Elsa.Samples.RebusWorker
.ConfigureServices((hostContext, services) =>
{
services
.AddElsa(option => option
.UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared"))
.UseAzureServiceBus("Endpoint=sb://elsa-workflows.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=n4NBTw9eSX12AG5BdIkyxCRroJGvh+EMOOM8ypWxWrQ=", LogLevel.Debug))
//.UseRabbitMq("amqp://localhost"))
.AddElsa(option => option.UseAzureServiceBus("Endpoint=sb://elsa-workflows-2.servicebus.windows.net/;SharedAccessKeyName=Elsa;SharedAccessKey=hAIa+fFuUbHi94y1Z0uO/2UTccjN/y4W0xvpaUd/cr4=", LogLevel.Debug))
//.UseRabbitMq("amqp://localhost"))
.AddConsoleActivities()
.AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1))
.AddRebusActivities<Greeting>()
.AddWorkflow<ProducerWorkflow>()
.AddWorkflow<ConsumerWorkflow>();
//.AddHostedService<Sender>()
.AddWorkflow<ProducerWorkflow>()
.AddWorkflow<ConsumerWorkflow>();
});
}
}

View file

@ -10,7 +10,7 @@ namespace Elsa.Samples.RebusWorker.Workflows
public void Build(IWorkflowBuilder workflow)
{
workflow
.StartWith<MessageReceived>(messageReceived => messageReceived.Set(x => x.MessageType, typeof(Greeting)))
.StartWith<RebusMessageReceived>(messageReceived => messageReceived.Set(x => x.MessageType, typeof(Greeting)))
.WriteLine(context =>
{
var greeting = context.GetInput<Greeting>();

View file

@ -22,9 +22,9 @@ namespace Elsa.Samples.RebusWorker.Workflows
public void Build(IWorkflowBuilder workflow)
{
workflow
.TimerEvent(Duration.FromSeconds(5))
.InstantEvent(_clock.GetCurrentInstant().Plus(Duration.FromSeconds(5)))
.WriteLine("Sending a random greeting to the \"greetings\" queue.")
.Then<SendMessage>(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting));
.Then<SendRebusMessage>(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting));
}
private Greeting GetRandomGreeting()

View file

@ -1,6 +1,8 @@
using Elsa.Extensions;
using System.Threading.Tasks;
using Elsa.Extensions;
using Elsa.Services;
using Microsoft.Azure.ServiceBus.Primitives;
using Microsoft.Azure.Services.AppAuthentication;
using Rebus.Config;
using Rebus.Logging;
using Rebus.Routing.TypeBased;
@ -9,12 +11,12 @@ namespace Elsa.ServiceBus.AzureServiceBus
{
public static class ElsaOptionsExtensions
{
public static ElsaOptions UseAzureServiceBus(this ElsaOptions elsaOptions, string connectionString, LogLevel logLevel = LogLevel.Info, ITokenProvider tokenProvider = default)
public static ElsaOptions UseAzureServiceBus(this ElsaOptions elsaOptions, string connectionString, LogLevel logLevel = LogLevel.Info, ITokenProvider? tokenProvider = default)
{
return elsaOptions.UseServiceBus(context => ConfigureAzureServiceBusEndpoint(context, connectionString, logLevel, tokenProvider));
}
private static void ConfigureAzureServiceBusEndpoint(ServiceBusEndpointConfigurationContext context, string connectionString, LogLevel logLevel, ITokenProvider tokenProvider)
private static void ConfigureAzureServiceBusEndpoint(ServiceBusEndpointConfigurationContext context, string connectionString, LogLevel logLevel, ITokenProvider? tokenProvider)
{
var queueName = context.QueueName;