diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs index 1ccdfd811..98620e04a 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Services/QueueWorker.cs @@ -33,7 +33,7 @@ namespace Elsa.Activities.AzureServiceBus.Services private async Task OnMessageReceived(Message message, CancellationToken cancellationToken) { - await _backgroundWorker.ScheduleTask(async () => await TriggerWorkflowsAsync(message, cancellationToken), cancellationToken); + await _backgroundWorker.ScheduleTask(GetType().FullName, async () => await TriggerWorkflowsAsync(message, cancellationToken), cancellationToken); await _messageReceiver.CompleteAsync(message.SystemProperties.LockToken); } diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Extensions/TimersOptionsExtensions.cs b/src/activities/Elsa.Activities.Timers.Quartz/Extensions/TimersOptionsExtensions.cs index 19f96f992..d066136b1 100644 --- a/src/activities/Elsa.Activities.Timers.Quartz/Extensions/TimersOptionsExtensions.cs +++ b/src/activities/Elsa.Activities.Timers.Quartz/Extensions/TimersOptionsExtensions.cs @@ -3,6 +3,7 @@ using Elsa.Activities.Timers.Options; using Elsa.Activities.Timers.Quartz.Jobs; using Elsa.Activities.Timers.Quartz.Services; using Elsa.Activities.Timers.Services; +using Elsa.Services; using Microsoft.Extensions.DependencyInjection; using Quartz; @@ -22,7 +23,6 @@ namespace Elsa .AddQuartzHostedService(ConfigureQuartzHostedService) .AddSingleton() .AddSingleton() - .AddSingleton() .AddTransient(); } diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs b/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs index 1a15fe546..b8eebc64d 100644 --- a/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs +++ b/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs @@ -13,14 +13,14 @@ namespace Elsa.Activities.Timers.Quartz.Jobs private readonly IWorkflowRunner _workflowRunner; private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowInstanceStore _workflowInstanceStore; - private readonly WorkflowRunnerQueue _workflowRunnerQueue; + private readonly IWorkflowQueue _workflowQueue; - public RunQuartzWorkflowJob(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, WorkflowRunnerQueue workflowRunnerQueue) + public RunQuartzWorkflowJob(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, IWorkflowQueue workflowQueue) { _workflowRunner = workflowRunner; _workflowRegistry = workflowRegistry; _workflowInstanceStore = workflowInstanceStore; - _workflowRunnerQueue = workflowRunnerQueue; + _workflowQueue = workflowQueue; } public async Task Execute(IJobExecutionContext context) @@ -47,7 +47,7 @@ namespace Elsa.Activities.Timers.Quartz.Jobs } else { - await _workflowRunnerQueue.Enqueue(workflowInstanceId, activityId, cancellationToken); + await _workflowQueue.Enqueue(workflowInstanceId, activityId, cancellationToken); } } } diff --git a/src/core/Elsa.Abstractions/Services/IBackgroundWorker.cs b/src/core/Elsa.Abstractions/Services/IBackgroundWorker.cs index d7a7ebb44..a157ee745 100644 --- a/src/core/Elsa.Abstractions/Services/IBackgroundWorker.cs +++ b/src/core/Elsa.Abstractions/Services/IBackgroundWorker.cs @@ -1,13 +1,13 @@ using System; using System.Threading; using System.Threading.Tasks; +using NodaTime; namespace Elsa.Services { public interface IBackgroundWorker { - Task ScheduleTask(Func task, CancellationToken cancellationToken = default); - Task ScheduleTask(Action task, CancellationToken cancellationToken = default); - Task RunAsync(CancellationToken cancellationToken = default); + Task ScheduleTask(string channelName, Func task, CancellationToken cancellationToken = default, Duration? channelTtl = default); + Task ScheduleTask(string channelName, Action task, CancellationToken cancellationToken = default, Duration? channelTtl = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs b/src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs new file mode 100644 index 000000000..f5b97ffe6 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/IWorkflowQueue.cs @@ -0,0 +1,13 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Services +{ + public interface IWorkflowQueue + { + /// + /// Enqueues the specified workflow instance and activity for execution in the background. + /// + Task Enqueue(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/Activities/ReceiveSignal/SignalReceived.cs b/src/core/Elsa.Core/Activities/Signaling/Activities/ReceiveSignal/SignalReceived.cs index 8c9715553..8c38a78dd 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Activities/ReceiveSignal/SignalReceived.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Activities/ReceiveSignal/SignalReceived.cs @@ -33,7 +33,7 @@ namespace Elsa.Activities.Signaling protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) { - var triggeredSignal = context.GetInput(); + var triggeredSignal = context.GetInput()!; return Done(triggeredSignal.Input); } } diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index 8d92a0c41..7e2e53ba8 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -36,7 +36,6 @@ namespace Microsoft.Extensions.DependencyInjection this IServiceCollection services, Action? configure = default) { - services.AddHostedService(); services.AddStartupRunner(); var options = new ElsaOptions(services); @@ -97,6 +96,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/core/Elsa.Core/HostedServices/StartBackgroundWorker.cs b/src/core/Elsa.Core/HostedServices/StartBackgroundWorker.cs deleted file mode 100644 index 61578f263..000000000 --- a/src/core/Elsa.Core/HostedServices/StartBackgroundWorker.cs +++ /dev/null @@ -1,14 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.Services; -using Microsoft.Extensions.Hosting; - -namespace Elsa.HostedServices -{ - public class StartBackgroundWorker : BackgroundService - { - private readonly IBackgroundWorker _backgroundWorker; - public StartBackgroundWorker(IBackgroundWorker backgroundWorker) => _backgroundWorker = backgroundWorker; - protected override async Task ExecuteAsync(CancellationToken stoppingToken) => await _backgroundWorker.RunAsync(stoppingToken); - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/BackgroundWorker.cs b/src/core/Elsa.Core/Services/BackgroundWorker.cs index 7fc54b296..db16fadb6 100644 --- a/src/core/Elsa.Core/Services/BackgroundWorker.cs +++ b/src/core/Elsa.Core/Services/BackgroundWorker.cs @@ -1,29 +1,88 @@ using System; +using System.Collections.Generic; using System.Threading; using System.Threading.Channels; using System.Threading.Tasks; +using NodaTime; namespace Elsa.Services { public class BackgroundWorker : IBackgroundWorker { - private readonly Channel> _channel; - public BackgroundWorker() => _channel = Channel.CreateBounded>(10); + private readonly IDictionary>> _channels = new Dictionary>>(); - public async Task ScheduleTask(Func task, CancellationToken cancellationToken) => await _channel.Writer.WriteAsync(task, cancellationToken); + public async Task ScheduleTask(string channelName, Func task, CancellationToken cancellationToken, Duration? channelTtl) + { + var channel = GetOrCreateChannel(channelName, channelTtl, cancellationToken); + await channel.Writer.WriteAsync(task, cancellationToken); + } - public async Task ScheduleTask(Action task, CancellationToken cancellationToken) => - await _channel.Writer.WriteAsync(() => + public async Task ScheduleTask(string channelName, Action task, CancellationToken cancellationToken, Duration? channelTtl) + { + await ScheduleTask(channelName, () => { task(); return new ValueTask(); - }, cancellationToken); + }, cancellationToken, channelTtl); + } - public async Task RunAsync(CancellationToken cancellationToken) + private Channel> GetOrCreateChannel(string channelName, Duration? channelTtl, CancellationToken cancellationToken) { - while (await _channel.Reader.WaitToReadAsync(cancellationToken)) - while (_channel.Reader.TryRead(out var task)) - await task(); + var channel = _channels.ContainsKey(channelName) ? _channels[channelName] : default; + + if (channel != null) + return channel; + + channel = Channel.CreateBounded>(10); + _channels[channelName] = channel; + var reader = new BackgroundChannelReader(channel, () => CompleteChannel(channelName)); + reader.Start(channelTtl, cancellationToken); + + return channel; + } + + private void CompleteChannel(string name) + { + if (_channels.ContainsKey(name)) + _channels.Remove(name); + } + + private class BackgroundChannelReader + { + private readonly Channel> _channel; + private readonly Action _onComplete; + + public BackgroundChannelReader(Channel> channel, Action onComplete) + { + _channel = channel; + _onComplete = onComplete; + } + + public void Start(Duration? channelTtl, CancellationToken cancellationToken) => Task.Factory.StartNew(() => ProcessAsync(channelTtl, cancellationToken), cancellationToken); + + private async Task ProcessAsync(Duration? channelTtl, CancellationToken cancellationToken) + { + var timeOutTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + + if (channelTtl != null) + timeOutTokenSource.CancelAfter(channelTtl.Value.ToTimeSpan()); + + try + { + while (await _channel.Reader.WaitToReadAsync(timeOutTokenSource.Token)) + { + var task = await _channel.Reader.ReadAsync(CancellationToken.None); + await task(); + + if (channelTtl != null) + timeOutTokenSource.CancelAfter(channelTtl.Value.ToTimeSpan()); // Reset timeout. + } + } + catch(TaskCanceledException) + { + _onComplete(); + } + } } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Services/WorkflowRunnerQueue.cs b/src/core/Elsa.Core/Services/WorkflowQueue.cs similarity index 72% rename from src/activities/Elsa.Activities.Timers.Quartz/Services/WorkflowRunnerQueue.cs rename to src/core/Elsa.Core/Services/WorkflowQueue.cs index d71e8f959..7302f6602 100644 --- a/src/activities/Elsa.Activities.Timers.Quartz/Services/WorkflowRunnerQueue.cs +++ b/src/core/Elsa.Core/Services/WorkflowQueue.cs @@ -3,32 +3,29 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Models; using Elsa.Persistence; -using Elsa.Services; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; +using NodaTime; -namespace Elsa.Activities.Timers.Quartz.Services +namespace Elsa.Services { - // TODO: Consider turning this into a global service to allow background, sequential execution of a given workflow instance (but allow for parallel execution of workflow instances with different definitions). Executing the same workflow instances sequentially prevents loss of data during update concurrency conflicts. - public class WorkflowRunnerQueue + public class WorkflowQueue : IWorkflowQueue { private readonly IBackgroundWorker _backgroundWorker; private readonly IServiceProvider _serviceProvider; private readonly ILogger _logger; - public WorkflowRunnerQueue(IBackgroundWorker backgroundWorker, IServiceProvider serviceProvider, ILogger logger) + public WorkflowQueue(IBackgroundWorker backgroundWorker, IServiceProvider serviceProvider, ILogger logger) { _backgroundWorker = backgroundWorker; _serviceProvider = serviceProvider; _logger = logger; } - public async Task Enqueue(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default) - { - await _backgroundWorker.ScheduleTask(async () => await RunWorkflowAsync(workflowInstanceId, activityId, cancellationToken), cancellationToken); - } + public async Task Enqueue(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default) => + await _backgroundWorker.ScheduleTask(workflowInstanceId, async () => await RunWorkflowAsync(workflowInstanceId, activityId, cancellationToken), cancellationToken, Duration.FromMinutes(1)); - private async Task RunWorkflowAsync(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default) + private async ValueTask RunWorkflowAsync(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default) { using var scope = _serviceProvider.CreateScope(); var store = scope.ServiceProvider.GetRequiredService(); diff --git a/src/samples/aspnet/Elsa.Samples.ContextualWorkflowHttp/Workflows/DocumentApprovalWorkflow.cs b/src/samples/aspnet/Elsa.Samples.ContextualWorkflowHttp/Workflows/DocumentApprovalWorkflow.cs index f81be3bfb..a5a41d1c5 100644 --- a/src/samples/aspnet/Elsa.Samples.ContextualWorkflowHttp/Workflows/DocumentApprovalWorkflow.cs +++ b/src/samples/aspnet/Elsa.Samples.ContextualWorkflowHttp/Workflows/DocumentApprovalWorkflow.cs @@ -19,7 +19,7 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows // Demonstrating that we can create activities and connect to them later on by using the activity builder reference. var join = workflow.Add(x => x.WithMode(Join.JoinMode.WaitAny)); join.Finish(); - + workflow // The workflow context type of this workflow. .WithContextType() @@ -28,7 +28,7 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows .ReceiveHttpPostRequest("/documents") // Store the document as the workflow context. It will be saved automatically when the workflow gets suspended. - .Then(context => context.SetWorkflowContext(context.GetInput().GetBody())).LoadWorkflowContext() + .Then(context => context.SetWorkflowContext(context.GetInput()!.GetBody())).LoadWorkflowContext() // Write an HTTP response. .WriteHttpResponse( @@ -62,8 +62,8 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows private static void StoreComment(ActivityExecutionContext context) { - var document = (Document)context.WorkflowExecutionContext.WorkflowContext!; - var comment = (Comment)((HttpRequestModel)context.Input!).Body!; + var document = (Document) context.WorkflowExecutionContext.WorkflowContext!; + var comment = (Comment) ((HttpRequestModel) context.Input!).Body!; document.Comments.Add(comment); } diff --git a/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/ConditionWorkflow.cs b/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/ConditionWorkflow.cs index b7d65ef21..c03ed6f61 100644 --- a/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/ConditionWorkflow.cs +++ b/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/ConditionWorkflow.cs @@ -17,7 +17,7 @@ namespace ElsaDashboard.Samples.Monolith.Workflows .Then(() => Console.WriteLine("What is your age?")).WithDisplayName("Write").WithDescription("What is your age?") .ReadLine() .Timer(Duration.FromMinutes(5)) - .SetVariable("Age", context => int.Parse(context.GetInput())) + .SetVariable("Age", context => int.Parse(context.GetInput()!)) .IfElse( context => context.GetVariable("Age") < 18, whenTrue => whenTrue.WriteLine("You are not allowed to drink beer."), diff --git a/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/NamingWorkflow.cs b/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/NamingWorkflow.cs index 96e05d93f..141b2af26 100644 --- a/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/NamingWorkflow.cs +++ b/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Workflows/NamingWorkflow.cs @@ -19,7 +19,7 @@ namespace ElsaDashboard.Samples.Monolith.Workflows .Correlate(() => Guid.NewGuid().ToString("N")) .WriteHttpResponse(x => x.WithStatusCode(HttpStatusCode.OK).WithContent(context => $"Tell me your name please. Use correlation ID {context.WorkflowExecutionContext.CorrelationId}").WithContentType("text/plain")) .HttpRequestReceived(x => x.WithPath("/signup").WithMethod(HttpMethod.Post.ToString()).WithReadContent()) - .SetVariable("Name", context => (string?) context.GetInput().Body) + .SetVariable("Name", context => (string?) context.GetInput()!.Body) .SetName(context => context.GetVariable("Name")) .WriteHttpResponse(x => x diff --git a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/GenerateOrdersWorkflow.cs b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/GenerateOrdersWorkflow.cs index 3095a7216..f165f0ca4 100644 --- a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/GenerateOrdersWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/GenerateOrdersWorkflow.cs @@ -23,7 +23,7 @@ namespace Elsa.Samples.CustomAttributesChildWorker.Workflows .Timer(Duration.FromSeconds(5)) .SetVariable("CustomerId", SelectRandomCustomerId) .WriteLine(context => $"Creating a new order for customer {context.GetVariable("CustomerId")}.") - .Then(message => message.Set(x => x.Message, context => new OrderReceived { CustomerId = context.GetVariable("CustomerId") })); + .Then(message => message.Set(x => x.Message, context => new OrderReceived { CustomerId = context.GetVariable("CustomerId")! })); } private string SelectRandomCustomerId() diff --git a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/OrderReceivedWorkflow.cs b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/OrderReceivedWorkflow.cs index 943ea511c..af9b3ee76 100644 --- a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/OrderReceivedWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Workflows/OrderReceivedWorkflow.cs @@ -18,8 +18,8 @@ namespace Elsa.Samples.CustomAttributesChildWorker.Workflows workflow .StartWith(activity => activity.Set(x => x.MessageType, typeof(OrderReceived))) .SetVariable(context => context.GetInput()) - .WriteLine(context => $"Received a new order for {context.GetVariable().CustomerId}.") - .RunWorkflow(activity => activity.WithCustomAttributes(context => new Variables().Set("Customer", context.GetVariable().CustomerId))) + .WriteLine(context => $"Received a new order for {context.GetVariable()!.CustomerId}.") + .RunWorkflow(activity => activity.WithCustomAttributes(context => new Variables().Set("Customer", context.GetVariable()!.CustomerId))) .WriteLine("Returned back from child workflow."); } } diff --git a/src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs index f7b8eb626..a43ae6d78 100644 --- a/src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs @@ -1,19 +1,11 @@ using Elsa.Activities.Console; using Elsa.Builders; using Elsa.Samples.Timers.Activities; -using NodaTime; namespace Elsa.Samples.Timers.Workflows { public class RecurringTaskWorkflow : IWorkflow { - private readonly IClock _clock; - - public RecurringTaskWorkflow(IClock clock) - { - _clock = clock; - } - public void Build(IWorkflowBuilder workflow) { workflow