From 17406c1db397d01fc01cb2d6e5015ad9e9fa0b8a Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 13 Apr 2023 20:05:14 +0200 Subject: [PATCH] Implement activity cancellation --- .../HangfireBackgroundActivityScheduler.cs | 8 ++ .../Services/ProtoActorWorkflowRuntime.cs | 15 +++ .../Elsa.Telnyx/Activities/WebhookEvent.cs | 2 +- .../Flowchart/Activities/FlowJoin.cs | 12 ++- .../ActivityExecutionContextExtensions.cs | 16 +++- .../Models/ActivityExecutionContext.cs | 16 ++-- .../Models/BookmarkOptions.cs | 5 + .../Models/CreateBookmarkOptions.cs | 5 - .../Notifications/ActivityCancelled.cs | 10 ++ .../ExecuteBackgroundActivityConsumer.cs | 28 ------ .../Contracts/IBackgroundActivityScheduler.cs | 7 ++ .../Contracts/IWorkflowRuntime.cs | 8 ++ .../WorkflowDictionaryExtensions.cs | 3 + ...kflowExecutionPipelineBuilderExtensions.cs | 6 -- .../Features/WorkflowRuntimeFeature.cs | 8 +- .../Handlers/CancelBackgroundActivities.cs | 44 +++++++++ .../Handlers/ScheduleBackgroundActivities.cs | 93 +++++++++++++++++++ .../BackgroundActivityInvokerMiddleware.cs | 14 +-- .../Workflows/PersistBookmarkMiddleware.cs | 2 +- .../ScheduleBackgroundActivitiesMiddleware.cs | 40 -------- .../Bookmarks/BackgroundActivityBookmark.cs | 9 +- .../Notifications/IndexedWorkflowBookmarks.cs | 13 ++- .../Models/ScheduledBackgroundActivity.cs | 4 +- .../Notifications/WorkflowBookmarksIndexed.cs | 8 +- .../DefaultBackgroundActivityInvoker.cs | 9 +- .../Services/DefaultWorkflowRuntime.cs | 6 ++ .../LocalBackgroundActivityScheduler.cs | 31 +++++-- 27 files changed, 287 insertions(+), 135 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Core/Models/BookmarkOptions.cs delete mode 100644 src/modules/Elsa.Workflows.Core/Models/CreateBookmarkOptions.cs create mode 100644 src/modules/Elsa.Workflows.Core/Notifications/ActivityCancelled.cs delete mode 100644 src/modules/Elsa.Workflows.Runtime/Consumers/ExecuteBackgroundActivityConsumer.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Handlers/CancelBackgroundActivities.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Handlers/ScheduleBackgroundActivities.cs delete mode 100644 src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs diff --git a/src/modules/Elsa.Hangfire/Services/HangfireBackgroundActivityScheduler.cs b/src/modules/Elsa.Hangfire/Services/HangfireBackgroundActivityScheduler.cs index fa921d7d7..2fedb6f30 100644 --- a/src/modules/Elsa.Hangfire/Services/HangfireBackgroundActivityScheduler.cs +++ b/src/modules/Elsa.Hangfire/Services/HangfireBackgroundActivityScheduler.cs @@ -2,6 +2,7 @@ using Elsa.Hangfire.Jobs; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Models; using Hangfire; +using Hangfire.States; namespace Elsa.Hangfire.Services; @@ -26,4 +27,11 @@ public class HangfireBackgroundActivityScheduler : IBackgroundActivityScheduler var jobId = _backgroundJobClient.Enqueue(x => x.ExecuteAsync(scheduledBackgroundActivity, CancellationToken.None)); return Task.FromResult(jobId); } + + /// + public Task CancelAsync(string jobId, CancellationToken cancellationToken = default) + { + _backgroundJobClient.Delete(jobId); + return Task.CompletedTask; + } } \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs index 379bffb6e..2ee043040 100644 --- a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs @@ -246,6 +246,21 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime await StoreBookmarksAsync(context.InstanceId, context.Diff.Added, context.CorrelationId, cancellationToken); } + /// + public async Task UpdateBookmarkAsync(Workflows.Runtime.Models.StoredBookmark bookmark, CancellationToken cancellationToken = default) + { + var bookmarkClient = _cluster.GetNamedBookmarkGrain(bookmark.Hash); + + var storeBookmarkRequest = new StoreBookmarksRequest + { + WorkflowInstanceId = bookmark.WorkflowInstanceId, + CorrelationId = bookmark.CorrelationId.EmptyIfNull() + }; + + storeBookmarkRequest.BookmarkIds.Add(bookmark.BookmarkId); + await bookmarkClient.Store(storeBookmarkRequest, cancellationToken); + } + /// public async Task CountRunningWorkflowsAsync(CountRunningWorkflowsArgs args, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs b/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs index e740e004a..2cc2c9dad 100644 --- a/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs +++ b/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs @@ -48,7 +48,7 @@ public class WebhookEvent : Activity var eventType = EventType; var payload = new WebhookEventBookmarkPayload(eventType); - context.CreateBookmark(new CreateBookmarkOptions(payload, Resume, Type)); + context.CreateBookmark(new BookmarkOptions(payload, Resume, Type)); } } diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs index abe79bb60..3e9707a61 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/FlowJoin.cs @@ -47,18 +47,22 @@ public class FlowJoin : Activity, IJoinNode if (!alreadyExecuted) { await context.CompleteActivityAsync(); - ClearBookmarks(flowchart, context); + await ClearBookmarksAsync(flowchart, context); } break; } } - private void ClearBookmarks(Flowchart flowchart, ActivityExecutionContext context) + private async Task ClearBookmarksAsync(Flowchart flowchart, ActivityExecutionContext context) { // Clear any bookmarks created between this join and its most recent fork. var connections = flowchart.Connections; var workflowExecutionContext = context.WorkflowExecutionContext; - var inboundActivities = connections.LeftAncestorActivities(this).Select(x => workflowExecutionContext.FindNodeByActivity(x)).Select(x => x.NodeId).ToList(); - context.WorkflowExecutionContext.Bookmarks.RemoveWhere(x => inboundActivities.Contains(x.ActivityNodeId)); + var inboundActivities = connections.LeftAncestorActivities(this).Select(x => workflowExecutionContext.FindNodeByActivity(x)).Select(x => x.Activity).ToList(); + var inboundActivityExecutionContexts = workflowExecutionContext.ActivityExecutionContexts.Where(x => inboundActivities.Contains(x.Activity)).ToList(); + + // Cancel each inbound activity. + foreach (var activityExecutionContext in inboundActivityExecutionContexts) + await activityExecutionContext.CancelActivityAsync(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs index 08a827f41..8d0d42535 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs @@ -6,10 +6,12 @@ using Elsa.Common.Contracts; using Elsa.Expressions.Contracts; using Elsa.Expressions.Helpers; using Elsa.Expressions.Models; +using Elsa.Mediator.Contracts; using Elsa.Workflows.Core.Activities.Flowchart.Models; using Elsa.Workflows.Core.Attributes; using Elsa.Workflows.Core.Contracts; using Elsa.Workflows.Core.Models; +using Elsa.Workflows.Core.Notifications; using Elsa.Workflows.Core.Signals; using Microsoft.Extensions.Logging; @@ -97,8 +99,9 @@ public static class ActivityExecutionContextExtensions return logEntry; } - public static Variable SetVariable(this ActivityExecutionContext context, string name, object? value, Type? storageDriverType = default, Action? configure = default) => + public static Variable SetVariable(this ActivityExecutionContext context, string name, object? value, Type? storageDriverType = default, Action? configure = default) => context.ExpressionExecutionContext.SetVariable(name, value, storageDriverType, configure); + public static T? GetVariable(this ActivityExecutionContext context, string id) => context.ExpressionExecutionContext.GetVariable(id); /// @@ -119,7 +122,7 @@ public static class ActivityExecutionContextExtensions .GetWrappedInputProperties(activity) .Where(x => x.Value is { MemoryBlockReference: { } }) .ToDictionary(x => x.Key, x => x.Value); - + var evaluator = context.GetRequiredService(); var stateSerializer = context.GetRequiredService(); var expressionExecutionContext = context.ExpressionExecutionContext; @@ -132,11 +135,11 @@ public static class ActivityExecutionContextExtensions // Store the evaluated input value in the activity state. var serializedValue = await stateSerializer.SerializeAsync(value); - - if(serializedValue.ValueKind != JsonValueKind.Undefined) + + if (serializedValue.ValueKind != JsonValueKind.Undefined) context.ActivityState[input.Key] = serializedValue; } - + context.SetHasEvaluatedProperties(); } @@ -293,8 +296,11 @@ public static class ActivityExecutionContextExtensions /// public static async Task CancelActivityAsync(this ActivityExecutionContext context) { + var publisher = context.GetRequiredService(); context.ClearBookmarks(); + context.WorkflowExecutionContext.Bookmarks.RemoveWhere(x => x.ActivityNodeId == context.NodeId); await context.SendSignalAsync(new CancelSignal()); + await publisher.PublishAsync(new ActivityCancelled(context)); } public static ILogger GetLogger(this ActivityExecutionContext context) => (ILogger)context.GetRequiredService(typeof(ILogger<>).MakeGenericType(context.Activity.GetType())); diff --git a/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs index da2943b29..b6deebc8a 100644 --- a/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs @@ -168,34 +168,34 @@ public class ActivityExecutionContext : IExecutionContext public void CreateBookmarks(IEnumerable payloads, ExecuteActivityDelegate? callback = default) { foreach (var payload in payloads) - CreateBookmark(new CreateBookmarkOptions(payload, callback)); + CreateBookmark(new BookmarkOptions(payload, callback)); } public void AddBookmarks(IEnumerable bookmarks) => _bookmarks.AddRange(bookmarks); public void AddBookmark(Bookmark bookmark) => _bookmarks.Add(bookmark); - public Bookmark CreateBookmark(ExecuteActivityDelegate callback) => CreateBookmark(new CreateBookmarkOptions(default, callback)); - public Bookmark CreateBookmark(object payload, ExecuteActivityDelegate callback) => CreateBookmark(new CreateBookmarkOptions(payload, callback)); - public Bookmark CreateBookmark(object payload) => CreateBookmark(new CreateBookmarkOptions(payload)); + public Bookmark CreateBookmark(ExecuteActivityDelegate callback) => CreateBookmark(new BookmarkOptions(default, callback)); + public Bookmark CreateBookmark(object payload, ExecuteActivityDelegate callback) => CreateBookmark(new BookmarkOptions(payload, callback)); + public Bookmark CreateBookmark(object payload) => CreateBookmark(new BookmarkOptions(payload)); /// /// Creates a bookmark so that this activity can be resumed at a later time. /// Creating a bookmark will automatically suspend the workflow after all pending activities have executed. /// - public Bookmark CreateBookmark(CreateBookmarkOptions? options = default) + public Bookmark CreateBookmark(BookmarkOptions? options = default) { var payload = options?.Payload; var callback = options?.Callback; - var activityTypeName = options?.ActivityTypeName ?? Activity.Type; + var bookmarkName = options?.BookmarkName ?? Activity.Type; var bookmarkHasher = GetRequiredService(); var identityGenerator = GetRequiredService(); var payloadSerializer = GetRequiredService(); var payloadJson = payload != null ? payloadSerializer.Serialize(payload) : default; - var hash = bookmarkHasher.Hash(activityTypeName, payload); + var hash = bookmarkHasher.Hash(bookmarkName, payload); var bookmark = new Bookmark( identityGenerator.GenerateId(), - activityTypeName, + bookmarkName, hash, payloadJson, ActivityNode.NodeId, diff --git a/src/modules/Elsa.Workflows.Core/Models/BookmarkOptions.cs b/src/modules/Elsa.Workflows.Core/Models/BookmarkOptions.cs new file mode 100644 index 000000000..18d16e7ad --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Models/BookmarkOptions.cs @@ -0,0 +1,5 @@ +using Elsa.Workflows.Core.Services; + +namespace Elsa.Workflows.Core.Models; + +public record BookmarkOptions(object? Payload = default, ExecuteActivityDelegate? Callback = default, string? BookmarkName = default, bool AutoBurn = true); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Models/CreateBookmarkOptions.cs b/src/modules/Elsa.Workflows.Core/Models/CreateBookmarkOptions.cs deleted file mode 100644 index 83a74fe9d..000000000 --- a/src/modules/Elsa.Workflows.Core/Models/CreateBookmarkOptions.cs +++ /dev/null @@ -1,5 +0,0 @@ -using Elsa.Workflows.Core.Services; - -namespace Elsa.Workflows.Core.Models; - -public record CreateBookmarkOptions(object? Payload = default, ExecuteActivityDelegate? Callback = default, string? ActivityTypeName = default, bool AutoBurn = true); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Notifications/ActivityCancelled.cs b/src/modules/Elsa.Workflows.Core/Notifications/ActivityCancelled.cs new file mode 100644 index 000000000..0bec4f928 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Notifications/ActivityCancelled.cs @@ -0,0 +1,10 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Core.Models; + +namespace Elsa.Workflows.Core.Notifications; + +/// +/// A notification that is sent when an activity is cancelled. +/// +/// The activity execution context. +public record ActivityCancelled(ActivityExecutionContext ActivityExecutionContext) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Consumers/ExecuteBackgroundActivityConsumer.cs b/src/modules/Elsa.Workflows.Runtime/Consumers/ExecuteBackgroundActivityConsumer.cs deleted file mode 100644 index 76f5ef725..000000000 --- a/src/modules/Elsa.Workflows.Runtime/Consumers/ExecuteBackgroundActivityConsumer.cs +++ /dev/null @@ -1,28 +0,0 @@ -using Elsa.Mediator.Contracts; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Models; - -namespace Elsa.Workflows.Runtime.Consumers; - -/// -/// A consumer that executes an activity in the background. -/// -public class ExecuteBackgroundActivityConsumer : IConsumer -{ - private readonly IBackgroundActivityInvoker _backgroundActivityInvoker; - - /// - /// Initializes a new instance of the class. - /// - public ExecuteBackgroundActivityConsumer(IBackgroundActivityInvoker backgroundActivityInvoker, IWorkflowDispatcher workflowDispatcher) - { - _backgroundActivityInvoker = backgroundActivityInvoker; - } - - /// - public async ValueTask ConsumeAsync(ScheduledBackgroundActivity message, CancellationToken cancellationToken) - { - // Execute the activity. - await _backgroundActivityInvoker.ExecuteAsync(message, cancellationToken); - } -} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBackgroundActivityScheduler.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBackgroundActivityScheduler.cs index 599e57272..eb89c8523 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IBackgroundActivityScheduler.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBackgroundActivityScheduler.cs @@ -14,4 +14,11 @@ public interface IBackgroundActivityScheduler /// The cancellation token. /// A handle representing the asynchronous invocation. Task ScheduleAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default); + + /// + /// Cancels the specified job. + /// + /// the ID of the job to cancel. + /// The cancellation token. + Task CancelAsync(string jobId, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs index bcb406810..c7a4e50ef 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs @@ -3,6 +3,7 @@ using Elsa.Workflows.Core.Helpers; using Elsa.Workflows.Core.Models; using Elsa.Workflows.Core.State; using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Runtime.Models; namespace Elsa.Workflows.Runtime.Contracts; @@ -89,6 +90,13 @@ public interface IWorkflowRuntime /// Task UpdateBookmarksAsync(UpdateBookmarksContext context, CancellationToken cancellationToken = default); + /// + /// Updates the specified bookmark. + /// + /// The bookmark to update. + /// The cancellation token. + Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default); + /// /// Counts the number of workflow instances based on the provided query args. /// diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowDictionaryExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowDictionaryExtensions.cs index e22aba356..3fae58787 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowDictionaryExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowDictionaryExtensions.cs @@ -4,6 +4,9 @@ using Microsoft.Extensions.DependencyInjection; // ReSharper disable once CheckNamespace namespace Elsa.Extensions; +/// +/// Extension methods for . +/// public static class WorkflowDictionaryExtensions { /// diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs index 678f95baa..d169a8670 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs @@ -17,18 +17,12 @@ public static class WorkflowExecutionPipelineBuilderExtensions public static IWorkflowExecutionPipelineBuilder UseDefaultRuntimePipeline(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder .Reset() - .UseBackgroundActivities() .UsePersistentVariables() .UseBookmarkPersistence() .UseWorkflowExecutionLogPersistence() .UseWorkflowStatePersistence() .UseDefaultActivityScheduler(); - /// - /// Installs middleware that schedules activities to run in the background. - /// - public static IWorkflowExecutionPipelineBuilder UseBackgroundActivities(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); - /// /// Installs middleware that persists the workflow instance before and after workflow execution. /// diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index 3bf63e7df..938e68e68 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -9,12 +9,12 @@ using Elsa.Workflows.Core.State; using Elsa.Workflows.Management.Notifications; using Elsa.Workflows.Runtime.ActivationValidators; using Elsa.Workflows.Runtime.Commands; -using Elsa.Workflows.Runtime.Consumers; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Handlers; using Elsa.Workflows.Runtime.HostedServices; using Elsa.Workflows.Runtime.Models; +using Elsa.Workflows.Runtime.Notifications; using Elsa.Workflows.Runtime.Options; using Elsa.Workflows.Runtime.Services; using Medallion.Threading; @@ -169,15 +169,13 @@ public class WorkflowRuntimeFeature : FeatureBase .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() + .AddNotificationHandler() // Workflow activation strategies. .AddSingleton() .AddSingleton() .AddSingleton() ; - - // If the local background activity invoker is used, register the consumer too. - Services.AddMessageChannel(); - Services.AddMessageConsumer(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/CancelBackgroundActivities.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/CancelBackgroundActivities.cs new file mode 100644 index 000000000..41975a1e5 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/CancelBackgroundActivities.cs @@ -0,0 +1,44 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Core.Contracts; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Middleware.Activities; +using Elsa.Workflows.Runtime.Models; +using Elsa.Workflows.Runtime.Models.Bookmarks; +using Elsa.Workflows.Runtime.Notifications; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Handlers; + +/// +/// A handler that cancels background activities. +/// +[PublicAPI] +public class CancelBackgroundActivities : INotificationHandler +{ + private readonly IBackgroundActivityScheduler _backgroundActivityScheduler; + private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer; + + /// + /// Initializes a new instance of the class. + /// + public CancelBackgroundActivities(IBackgroundActivityScheduler backgroundActivityScheduler, IBookmarkPayloadSerializer bookmarkPayloadSerializer) + { + _backgroundActivityScheduler = backgroundActivityScheduler; + _bookmarkPayloadSerializer = bookmarkPayloadSerializer; + } + + /// + public async Task HandleAsync(WorkflowBookmarksIndexed notification, CancellationToken cancellationToken) + { + var removedBookmarks = notification.IndexedWorkflowBookmarks.RemovedBookmarks.Where(x => x.Name == BackgroundActivityInvokerMiddleware.BackgroundActivityBookmarkName); + + foreach (var removedBookmark in removedBookmarks) + { + var payload = _bookmarkPayloadSerializer.Deserialize(removedBookmark.Data!); + if (payload.JobId != null) + { + await _backgroundActivityScheduler.CancelAsync(payload.JobId, cancellationToken); + } + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/ScheduleBackgroundActivities.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/ScheduleBackgroundActivities.cs new file mode 100644 index 000000000..8f6949124 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/ScheduleBackgroundActivities.cs @@ -0,0 +1,93 @@ +using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Core.Contracts; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Middleware.Activities; +using Elsa.Workflows.Runtime.Middleware.Workflows; +using Elsa.Workflows.Runtime.Models; +using Elsa.Workflows.Runtime.Models.Bookmarks; +using Elsa.Workflows.Runtime.Notifications; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Handlers; + +/// +/// A handler that schedules background activities. +/// +[PublicAPI] +public class ScheduleBackgroundActivities : INotificationHandler +{ + private readonly IBackgroundActivityScheduler _backgroundActivityScheduler; + private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer; + private readonly IBookmarkHasher _bookmarkHasher; + private readonly IWorkflowRuntime _workflowRuntime; + private IWorkflowStateSerializer _workflowStateSerializer; + + /// + /// Initializes a new instance of the class. + /// + public ScheduleBackgroundActivities( + IBackgroundActivityScheduler backgroundActivityScheduler, + IBookmarkPayloadSerializer bookmarkPayloadSerializer, + IBookmarkHasher bookmarkHasher, + IWorkflowRuntime workflowRuntime, + IWorkflowStateSerializer workflowStateSerializer) + { + _backgroundActivityScheduler = backgroundActivityScheduler; + _bookmarkPayloadSerializer = bookmarkPayloadSerializer; + _bookmarkHasher = bookmarkHasher; + _workflowRuntime = workflowRuntime; + _workflowStateSerializer = workflowStateSerializer; + } + + /// + public async Task HandleAsync(WorkflowBookmarksIndexed notification, CancellationToken cancellationToken) + { + var workflowExecutionContext = notification.WorkflowExecutionContext; + + var scheduledBackgroundActivities = workflowExecutionContext + .TransientProperties + .GetOrAdd(BackgroundActivityInvokerMiddleware.BackgroundActivitySchedulesKey, () => new List()); + + var bookmarks = notification.IndexedWorkflowBookmarks.AddedBookmarks; + + foreach (var scheduledBackgroundActivity in scheduledBackgroundActivities) + { + // Schedule the background activity. + var jobId = await _backgroundActivityScheduler.ScheduleAsync(scheduledBackgroundActivity, cancellationToken); + + // Select the bookmark associated with the background activity. + var bookmark = workflowExecutionContext.Bookmarks.First(x => x.Id == scheduledBackgroundActivity.BookmarkId); + var payload = _bookmarkPayloadSerializer.Deserialize(bookmark.Data!); + + // Store the created job ID. + workflowExecutionContext.Bookmarks.Remove(bookmark); + payload.JobId = jobId; + bookmark = bookmark with + { + Data = _bookmarkPayloadSerializer.Serialize(payload), + Hash = _bookmarkHasher.Hash(bookmark.Name, payload) + }; + workflowExecutionContext.Bookmarks.Add(bookmark); + + // Update the bookmark. + var storedBookmark = new StoredBookmark( + bookmark.Name, + bookmark.Hash, + workflowExecutionContext.Id, + bookmark.Id, + workflowExecutionContext.CorrelationId, + bookmark.Data + ); + + await _workflowRuntime.UpdateBookmarkAsync(storedBookmark, cancellationToken); + } + + if (scheduledBackgroundActivities.Any()) + { + // Bookmarks got updated, so we need to update the workflow state. + var workflowState = _workflowStateSerializer.SerializeState(workflowExecutionContext); + await _workflowRuntime.ImportWorkflowStateAsync(workflowState, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs index 306098805..4da2a07c2 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs @@ -1,8 +1,8 @@ +using System.Xml; using Elsa.Extensions; using Elsa.Workflows.Core.Middleware.Activities; using Elsa.Workflows.Core.Models; using Elsa.Workflows.Core.Pipelines.ActivityExecution; -using Elsa.Workflows.Runtime.Middleware.Workflows; using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Models.Bookmarks; @@ -10,12 +10,13 @@ namespace Elsa.Workflows.Runtime.Middleware.Activities; /// /// Executes the current activity from a background job if the activity is of kind or . -/// Works in tandem with . /// public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddleware { internal static readonly object IsBackgroundExecution = new(); internal static string GetBackgroundActivityOutputKey(string activityId) => $"__BackgroundActivityOutput:{activityId}"; + internal static readonly object BackgroundActivitySchedulesKey = new(); + internal const string BackgroundActivityBookmarkName = "BackgroundActivity"; /// public BackgroundActivityInvokerMiddleware(ActivityMiddlewareDelegate next) : base(next) @@ -55,12 +56,13 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew /// private static void ScheduleBackgroundActivity(ActivityExecutionContext context) { - var scheduledBackgroundActivities = context.WorkflowExecutionContext.TransientProperties.GetOrAdd(ScheduleBackgroundActivitiesMiddleware.BackgroundActivitySchedulesKey, () => new List()); + var scheduledBackgroundActivities = context.WorkflowExecutionContext.TransientProperties.GetOrAdd(BackgroundActivitySchedulesKey, () => new List()); var workflowInstanceId = context.WorkflowExecutionContext.Id; - var activityId = context.Activity.Id; + var activityNodeId = context.NodeId; var bookmarkPayload = new BackgroundActivityBookmark(); - var bookmark = context.CreateBookmark(bookmarkPayload); - scheduledBackgroundActivities.Add(new ScheduledBackgroundActivity(workflowInstanceId, activityId, bookmark.Id)); + var bookmarkOptions = new BookmarkOptions { BookmarkName = BackgroundActivityBookmarkName, Payload = bookmarkPayload }; + var bookmark = context.CreateBookmark(bookmarkOptions); + scheduledBackgroundActivities.Add(new ScheduledBackgroundActivity(workflowInstanceId, activityNodeId, bookmark.Id)); } /// diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs index c72b69e89..768a7a56c 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs @@ -45,7 +45,7 @@ public class PersistBookmarkMiddleware : WorkflowExecutionMiddleware await _workflowRuntime.UpdateBookmarksAsync(updateBookmarksContext, cancellationToken); // Publish domain event. - await _eventPublisher.PublishAsync(new WorkflowBookmarksIndexed(new IndexedWorkflowBookmarks(context.Id, diff.Added, diff.Removed, diff.Unchanged)), cancellationToken); + await _eventPublisher.PublishAsync(new WorkflowBookmarksIndexed(context, new IndexedWorkflowBookmarks(context.Id, diff.Added, diff.Removed, diff.Unchanged)), cancellationToken); // Notify all interested activities that the bookmarks have been persisted. var activityExecutionContexts = context.ActivityExecutionContexts.Where(x => x.Activity is IBookmarksPersistedHandler && x.Bookmarks.Any()).ToList(); diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs deleted file mode 100644 index a967b1228..000000000 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs +++ /dev/null @@ -1,40 +0,0 @@ -using Elsa.Extensions; -using Elsa.Workflows.Core.Models; -using Elsa.Workflows.Core.Pipelines.WorkflowExecution; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Middleware.Activities; -using Elsa.Workflows.Runtime.Models; - -namespace Elsa.Workflows.Runtime.Middleware.Workflows; - -/// -/// Schedules background activities for execution. This component works in tandem with . -/// -public class ScheduleBackgroundActivitiesMiddleware : WorkflowExecutionMiddleware -{ - private readonly IBackgroundActivityScheduler _backgroundActivityScheduler; - internal static readonly object BackgroundActivitySchedulesKey = new(); - - /// - public ScheduleBackgroundActivitiesMiddleware(WorkflowMiddlewareDelegate next, IBackgroundActivityScheduler backgroundActivityScheduler) : base(next) - { - _backgroundActivityScheduler = backgroundActivityScheduler; - } - - /// - public override async ValueTask InvokeAsync(WorkflowExecutionContext context) - { - // Invoke next middleware. - await Next(context); - - // Get activities to schedule. - if (context.TransientProperties.ContainsKey(BackgroundActivitySchedulesKey)) - { - var scheduledActivities = (ICollection)context.TransientProperties.GetValue(BackgroundActivitySchedulesKey)!; - - // Schedule activities. - foreach (var scheduledBackgroundActivity in scheduledActivities) - await _backgroundActivityScheduler.ScheduleAsync(scheduledBackgroundActivity, context.CancellationToken); - } - } -} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/Bookmarks/BackgroundActivityBookmark.cs b/src/modules/Elsa.Workflows.Runtime/Models/Bookmarks/BackgroundActivityBookmark.cs index d8268f9a5..44fa4c6f5 100644 --- a/src/modules/Elsa.Workflows.Runtime/Models/Bookmarks/BackgroundActivityBookmark.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/Bookmarks/BackgroundActivityBookmark.cs @@ -1,5 +1,3 @@ -using System.Text.Json.Serialization; - namespace Elsa.Workflows.Runtime.Models.Bookmarks; /// @@ -8,10 +6,7 @@ namespace Elsa.Workflows.Runtime.Models.Bookmarks; public class BackgroundActivityBookmark { /// - /// Initializes a new instance of the class. + /// Set retroactively after the job has been scheduled. It is used to cancel te job when the bookmark is deleted. /// - [JsonConstructor] - public BackgroundActivityBookmark() - { - } + public string? JobId { get; set; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/Notifications/IndexedWorkflowBookmarks.cs b/src/modules/Elsa.Workflows.Runtime/Models/Notifications/IndexedWorkflowBookmarks.cs index e8cf57208..2665cde5d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Models/Notifications/IndexedWorkflowBookmarks.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/Notifications/IndexedWorkflowBookmarks.cs @@ -2,4 +2,15 @@ using Elsa.Workflows.Core.Models; namespace Elsa.Workflows.Runtime.Models.Notifications; -public record IndexedWorkflowBookmarks(string InstanceId, ICollection AddedBookmarks, ICollection RemovedBookmarks, ICollection UnchangedBookmarks); \ No newline at end of file +/// +/// Contains the bookmarks that were added, removed, or unchanged. +/// +/// The workflow instance ID. +/// The bookmarks that were added. +/// The bookmarks that were removed. +/// The bookmarks that were unchanged. +public record IndexedWorkflowBookmarks( + string InstanceId, + ICollection AddedBookmarks, + ICollection RemovedBookmarks, + ICollection UnchangedBookmarks); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/ScheduledBackgroundActivity.cs b/src/modules/Elsa.Workflows.Runtime/Models/ScheduledBackgroundActivity.cs index d370bd5ae..2dc09e099 100644 --- a/src/modules/Elsa.Workflows.Runtime/Models/ScheduledBackgroundActivity.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/ScheduledBackgroundActivity.cs @@ -4,6 +4,6 @@ namespace Elsa.Workflows.Runtime.Models; /// Represents a scheduled background activity /// /// The ID of the workflow instance containing the activity to execute. -/// The ID of the activity to execute. +/// The ID of the activity to execute. /// The ID of the bookmark to resume. -public record ScheduledBackgroundActivity(string WorkflowInstanceId, string ActivityId, string BookmarkId); \ No newline at end of file +public record ScheduledBackgroundActivity(string WorkflowInstanceId, string ActivityNodeId, string BookmarkId); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowBookmarksIndexed.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowBookmarksIndexed.cs index 61df8082a..96c7a4a1b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowBookmarksIndexed.cs +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowBookmarksIndexed.cs @@ -1,7 +1,13 @@ using Elsa.Mediator.Contracts; +using Elsa.Workflows.Core.Models; using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Models.Notifications; namespace Elsa.Workflows.Runtime.Notifications; -public record WorkflowBookmarksIndexed(IndexedWorkflowBookmarks IndexedWorkflowBookmarks) : INotification; \ No newline at end of file +/// +/// A notification that is sent when the bookmarks of a workflow instance have been indexed. +/// +/// The workflow execution context. +/// The bookmarks that were added, removed, or unchanged. +public record WorkflowBookmarksIndexed(WorkflowExecutionContext WorkflowExecutionContext, IndexedWorkflowBookmarks IndexedWorkflowBookmarks) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs index 12f34a11a..0c69ec01a 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs @@ -18,7 +18,6 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker private readonly IWorkflowDispatcher _workflowDispatcher; private readonly IWorkflowDefinitionService _workflowDefinitionService; private readonly IWorkflowExecutionContextFactory _workflowExecutionContextFactory; - private readonly IWorkflowStateSerializer _workflowStateSerializer; private readonly IVariablePersistenceManager _variablePersistenceManager; private readonly IActivityInvoker _activityInvoker; private readonly IServiceProvider _serviceProvider; @@ -31,7 +30,6 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker IWorkflowDispatcher workflowDispatcher, IWorkflowDefinitionService workflowDefinitionService, IWorkflowExecutionContextFactory workflowExecutionContextFactory, - IWorkflowStateSerializer workflowStateSerializer, IVariablePersistenceManager variablePersistenceManager, IActivityInvoker activityInvoker, IServiceProvider serviceProvider) @@ -40,7 +38,6 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker _workflowDispatcher = workflowDispatcher; _workflowDefinitionService = workflowDefinitionService; _workflowExecutionContextFactory = workflowExecutionContextFactory; - _workflowStateSerializer = workflowStateSerializer; _variablePersistenceManager = variablePersistenceManager; _activityInvoker = activityInvoker; _serviceProvider = serviceProvider; @@ -62,8 +59,8 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken); var workflowExecutionContext = await _workflowExecutionContextFactory.CreateAsync(_serviceProvider, workflow, workflowState.Id, workflowState, cancellationToken: cancellationToken); - var activityId = scheduledBackgroundActivity.ActivityId; - var activityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.Activity.Id == activityId); + var activityNodeId = scheduledBackgroundActivity.ActivityNodeId; + var activityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.NodeId == activityNodeId); // Load persistent variables for the activity to use. await _variablePersistenceManager.LoadVariablesAsync(workflowExecutionContext); @@ -109,7 +106,7 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker // Resume the workflow, passing along the activity output. // TODO: This approach will fail if the output is non-serializable. We need to find a way to pass the output to the workflow without serializing it. var bookmarkId = scheduledBackgroundActivity.BookmarkId; - var inputKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityOutputKey(activityId); + var inputKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityOutputKey(activityNodeId); var dispatchRequest = new DispatchWorkflowInstanceRequest { diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs index 1cfbe1947..c0ddf9e73 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs @@ -223,6 +223,12 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime await StoreBookmarksAsync(context.InstanceId, context.Diff.Added, context.CorrelationId, cancellationToken); } + /// + public async Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default) + { + await _bookmarkStore.SaveAsync(bookmark, cancellationToken); + } + /// public async Task CountRunningWorkflowsAsync(CountRunningWorkflowsArgs args, CancellationToken cancellationToken = default) => await _workflowStateStore.CountAsync(args, cancellationToken); diff --git a/src/modules/Elsa.Workflows.Runtime/Services/LocalBackgroundActivityScheduler.cs b/src/modules/Elsa.Workflows.Runtime/Services/LocalBackgroundActivityScheduler.cs index 073078b62..c00dc9fac 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/LocalBackgroundActivityScheduler.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/LocalBackgroundActivityScheduler.cs @@ -1,4 +1,4 @@ -using System.Threading.Channels; +using Elsa.Mediator.Contracts; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Models; @@ -9,21 +9,34 @@ namespace Elsa.Workflows.Runtime.Services; /// public class LocalBackgroundActivityScheduler : IBackgroundActivityScheduler { - private readonly Channel _channel; + private readonly IJobQueue _jobQueue; + private readonly IBackgroundActivityInvoker _backgroundActivityInvoker; /// /// Initializes a new instance of the class. /// - /// The channel to write to. - public LocalBackgroundActivityScheduler(Channel channel) + public LocalBackgroundActivityScheduler(IJobQueue jobQueue, IBackgroundActivityInvoker backgroundActivityInvoker) { - _channel = channel; + _jobQueue = jobQueue; + _backgroundActivityInvoker = backgroundActivityInvoker; + } + + /// + public Task ScheduleAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default) + { + var jobId = _jobQueue.Enqueue(async ct => await InvokeBackgroundActivity(scheduledBackgroundActivity, ct)); + return Task.FromResult(jobId); + } + + /// + public Task CancelAsync(string jobId, CancellationToken cancellationToken = default) + { + _jobQueue.Cancel(jobId); + return Task.CompletedTask; } - /// - public async Task ScheduleAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default) + private async Task InvokeBackgroundActivity(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken) { - await _channel.Writer.WriteAsync(scheduledBackgroundActivity, cancellationToken); - return ""; + await _backgroundActivityInvoker.ExecuteAsync(scheduledBackgroundActivity, cancellationToken); } } \ No newline at end of file