Implement activity cancellation
This commit is contained in:
parent
ae7b06059a
commit
17406c1db3
|
|
@ -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<ExecuteBackgroundActivityJob>(x => x.ExecuteAsync(scheduledBackgroundActivity, CancellationToken.None));
|
||||
return Task.FromResult(jobId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task CancelAsync(string jobId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
_backgroundJobClient.Delete(jobId);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
}
|
||||
|
|
@ -246,6 +246,21 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
await StoreBookmarksAsync(context.InstanceId, context.Diff.Added, context.CorrelationId, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
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);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<int> CountRunningWorkflowsAsync(CountRunningWorkflowsArgs args, CancellationToken cancellationToken = default)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -48,7 +48,7 @@ public class WebhookEvent : Activity<Payload>
|
|||
var eventType = EventType;
|
||||
var payload = new WebhookEventBookmarkPayload(eventType);
|
||||
|
||||
context.CreateBookmark(new CreateBookmarkOptions(payload, Resume, Type));
|
||||
context.CreateBookmark(new BookmarkOptions(payload, Resume, Type));
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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<MemoryBlock>? configure = default) =>
|
||||
public static Variable SetVariable(this ActivityExecutionContext context, string name, object? value, Type? storageDriverType = default, Action<MemoryBlock>? configure = default) =>
|
||||
context.ExpressionExecutionContext.SetVariable(name, value, storageDriverType, configure);
|
||||
|
||||
public static T? GetVariable<T>(this ActivityExecutionContext context, string id) => context.ExpressionExecutionContext.GetVariable<T?>(id);
|
||||
|
||||
/// <summary>
|
||||
|
|
@ -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<IExpressionEvaluator>();
|
||||
var stateSerializer = context.GetRequiredService<IActivityStateSerializer>();
|
||||
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
|
|||
/// </summary>
|
||||
public static async Task CancelActivityAsync(this ActivityExecutionContext context)
|
||||
{
|
||||
var publisher = context.GetRequiredService<IEventPublisher>();
|
||||
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()));
|
||||
|
|
|
|||
|
|
@ -168,34 +168,34 @@ public class ActivityExecutionContext : IExecutionContext
|
|||
public void CreateBookmarks(IEnumerable<object> payloads, ExecuteActivityDelegate? callback = default)
|
||||
{
|
||||
foreach (var payload in payloads)
|
||||
CreateBookmark(new CreateBookmarkOptions(payload, callback));
|
||||
CreateBookmark(new BookmarkOptions(payload, callback));
|
||||
}
|
||||
|
||||
public void AddBookmarks(IEnumerable<Bookmark> 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));
|
||||
|
||||
/// <summary>
|
||||
/// 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.
|
||||
/// </summary>
|
||||
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<IBookmarkHasher>();
|
||||
var identityGenerator = GetRequiredService<IIdentityGenerator>();
|
||||
var payloadSerializer = GetRequiredService<IBookmarkPayloadSerializer>();
|
||||
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,
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
@ -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);
|
||||
|
|
@ -0,0 +1,10 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
|
||||
namespace Elsa.Workflows.Core.Notifications;
|
||||
|
||||
/// <summary>
|
||||
/// A notification that is sent when an activity is cancelled.
|
||||
/// </summary>
|
||||
/// <param name="ActivityExecutionContext">The activity execution context.</param>
|
||||
public record ActivityCancelled(ActivityExecutionContext ActivityExecutionContext) : INotification;
|
||||
|
|
@ -1,28 +0,0 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using Elsa.Workflows.Runtime.Models;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Consumers;
|
||||
|
||||
/// <summary>
|
||||
/// A consumer that executes an activity in the background.
|
||||
/// </summary>
|
||||
public class ExecuteBackgroundActivityConsumer : IConsumer<ScheduledBackgroundActivity>
|
||||
{
|
||||
private readonly IBackgroundActivityInvoker _backgroundActivityInvoker;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="ExecuteBackgroundActivityConsumer"/> class.
|
||||
/// </summary>
|
||||
public ExecuteBackgroundActivityConsumer(IBackgroundActivityInvoker backgroundActivityInvoker, IWorkflowDispatcher workflowDispatcher)
|
||||
{
|
||||
_backgroundActivityInvoker = backgroundActivityInvoker;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask ConsumeAsync(ScheduledBackgroundActivity message, CancellationToken cancellationToken)
|
||||
{
|
||||
// Execute the activity.
|
||||
await _backgroundActivityInvoker.ExecuteAsync(message, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -14,4 +14,11 @@ public interface IBackgroundActivityScheduler
|
|||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
/// <returns>A handle representing the asynchronous invocation.</returns>
|
||||
Task<string> ScheduleAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Cancels the specified job.
|
||||
/// </summary>
|
||||
/// <param name="jobId">the ID of the job to cancel.</param>
|
||||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
Task CancelAsync(string jobId, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -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
|
|||
/// </summary>
|
||||
Task UpdateBookmarksAsync(UpdateBookmarksContext context, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Updates the specified bookmark.
|
||||
/// </summary>
|
||||
/// <param name="bookmark">The bookmark to update.</param>
|
||||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Counts the number of workflow instances based on the provided query args.
|
||||
/// </summary>
|
||||
|
|
|
|||
|
|
@ -4,6 +4,9 @@ using Microsoft.Extensions.DependencyInjection;
|
|||
// ReSharper disable once CheckNamespace
|
||||
namespace Elsa.Extensions;
|
||||
|
||||
/// <summary>
|
||||
/// Extension methods for <see cref="IDictionary{TKey,TValue}"/>.
|
||||
/// </summary>
|
||||
public static class WorkflowDictionaryExtensions
|
||||
{
|
||||
/// <summary>
|
||||
|
|
|
|||
|
|
@ -17,18 +17,12 @@ public static class WorkflowExecutionPipelineBuilderExtensions
|
|||
public static IWorkflowExecutionPipelineBuilder UseDefaultRuntimePipeline(this IWorkflowExecutionPipelineBuilder pipelineBuilder) =>
|
||||
pipelineBuilder
|
||||
.Reset()
|
||||
.UseBackgroundActivities()
|
||||
.UsePersistentVariables()
|
||||
.UseBookmarkPersistence()
|
||||
.UseWorkflowExecutionLogPersistence()
|
||||
.UseWorkflowStatePersistence()
|
||||
.UseDefaultActivityScheduler();
|
||||
|
||||
/// <summary>
|
||||
/// Installs middleware that schedules activities to run in the background.
|
||||
/// </summary>
|
||||
public static IWorkflowExecutionPipelineBuilder UseBackgroundActivities(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<ScheduleBackgroundActivitiesMiddleware>();
|
||||
|
||||
/// <summary>
|
||||
/// Installs middleware that persists the workflow instance before and after workflow execution.
|
||||
/// </summary>
|
||||
|
|
|
|||
|
|
@ -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<ResumeDispatchWorkflowActivityHandler, WorkflowExecuted>()
|
||||
.AddNotificationHandler<IndexWorkflowTriggersHandler, WorkflowDefinitionPublished>()
|
||||
.AddNotificationHandler<IndexWorkflowTriggersHandler, WorkflowDefinitionRetracted>()
|
||||
.AddNotificationHandler<ScheduleBackgroundActivities, WorkflowBookmarksIndexed>()
|
||||
.AddNotificationHandler<CancelBackgroundActivities, WorkflowBookmarksIndexed>()
|
||||
|
||||
// Workflow activation strategies.
|
||||
.AddSingleton<IWorkflowActivationStrategy, SingletonStrategy>()
|
||||
.AddSingleton<IWorkflowActivationStrategy, CorrelatedSingletonStrategy>()
|
||||
.AddSingleton<IWorkflowActivationStrategy, CorrelationStrategy>()
|
||||
;
|
||||
|
||||
// If the local background activity invoker is used, register the consumer too.
|
||||
Services.AddMessageChannel<ScheduledBackgroundActivity>();
|
||||
Services.AddMessageConsumer<ScheduledBackgroundActivity, ExecuteBackgroundActivityConsumer>();
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// A handler that cancels background activities.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
public class CancelBackgroundActivities : INotificationHandler<WorkflowBookmarksIndexed>
|
||||
{
|
||||
private readonly IBackgroundActivityScheduler _backgroundActivityScheduler;
|
||||
private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="ScheduledBackgroundActivity"/> class.
|
||||
/// </summary>
|
||||
public CancelBackgroundActivities(IBackgroundActivityScheduler backgroundActivityScheduler, IBookmarkPayloadSerializer bookmarkPayloadSerializer)
|
||||
{
|
||||
_backgroundActivityScheduler = backgroundActivityScheduler;
|
||||
_bookmarkPayloadSerializer = bookmarkPayloadSerializer;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
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<BackgroundActivityBookmark>(removedBookmark.Data!);
|
||||
if (payload.JobId != null)
|
||||
{
|
||||
await _backgroundActivityScheduler.CancelAsync(payload.JobId, cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// A handler that schedules background activities.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
public class ScheduleBackgroundActivities : INotificationHandler<WorkflowBookmarksIndexed>
|
||||
{
|
||||
private readonly IBackgroundActivityScheduler _backgroundActivityScheduler;
|
||||
private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer;
|
||||
private readonly IBookmarkHasher _bookmarkHasher;
|
||||
private readonly IWorkflowRuntime _workflowRuntime;
|
||||
private IWorkflowStateSerializer _workflowStateSerializer;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="ScheduledBackgroundActivity"/> class.
|
||||
/// </summary>
|
||||
public ScheduleBackgroundActivities(
|
||||
IBackgroundActivityScheduler backgroundActivityScheduler,
|
||||
IBookmarkPayloadSerializer bookmarkPayloadSerializer,
|
||||
IBookmarkHasher bookmarkHasher,
|
||||
IWorkflowRuntime workflowRuntime,
|
||||
IWorkflowStateSerializer workflowStateSerializer)
|
||||
{
|
||||
_backgroundActivityScheduler = backgroundActivityScheduler;
|
||||
_bookmarkPayloadSerializer = bookmarkPayloadSerializer;
|
||||
_bookmarkHasher = bookmarkHasher;
|
||||
_workflowRuntime = workflowRuntime;
|
||||
_workflowStateSerializer = workflowStateSerializer;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowBookmarksIndexed notification, CancellationToken cancellationToken)
|
||||
{
|
||||
var workflowExecutionContext = notification.WorkflowExecutionContext;
|
||||
|
||||
var scheduledBackgroundActivities = workflowExecutionContext
|
||||
.TransientProperties
|
||||
.GetOrAdd(BackgroundActivityInvokerMiddleware.BackgroundActivitySchedulesKey, () => new List<ScheduledBackgroundActivity>());
|
||||
|
||||
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<BackgroundActivityBookmark>(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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
|||
|
||||
/// <summary>
|
||||
/// Executes the current activity from a background job if the activity is of kind <see cref="ActivityKind.Job"/> or <see cref="ActivityKind.Task"/>.
|
||||
/// Works in tandem with <see cref="ScheduleBackgroundActivitiesMiddleware"/>.
|
||||
/// </summary>
|
||||
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";
|
||||
|
||||
/// <inheritdoc />
|
||||
public BackgroundActivityInvokerMiddleware(ActivityMiddlewareDelegate next) : base(next)
|
||||
|
|
@ -55,12 +56,13 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew
|
|||
/// </summary>
|
||||
private static void ScheduleBackgroundActivity(ActivityExecutionContext context)
|
||||
{
|
||||
var scheduledBackgroundActivities = context.WorkflowExecutionContext.TransientProperties.GetOrAdd(ScheduleBackgroundActivitiesMiddleware.BackgroundActivitySchedulesKey, () => new List<ScheduledBackgroundActivity>());
|
||||
var scheduledBackgroundActivities = context.WorkflowExecutionContext.TransientProperties.GetOrAdd(BackgroundActivitySchedulesKey, () => new List<ScheduledBackgroundActivity>());
|
||||
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));
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Schedules background activities for execution. This component works in tandem with <see cref="BackgroundActivityInvokerMiddleware"/>.
|
||||
/// </summary>
|
||||
public class ScheduleBackgroundActivitiesMiddleware : WorkflowExecutionMiddleware
|
||||
{
|
||||
private readonly IBackgroundActivityScheduler _backgroundActivityScheduler;
|
||||
internal static readonly object BackgroundActivitySchedulesKey = new();
|
||||
|
||||
/// <inheritdoc />
|
||||
public ScheduleBackgroundActivitiesMiddleware(WorkflowMiddlewareDelegate next, IBackgroundActivityScheduler backgroundActivityScheduler) : base(next)
|
||||
{
|
||||
_backgroundActivityScheduler = backgroundActivityScheduler;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
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<ScheduledBackgroundActivity>)context.TransientProperties.GetValue(BackgroundActivitySchedulesKey)!;
|
||||
|
||||
// Schedule activities.
|
||||
foreach (var scheduledBackgroundActivity in scheduledActivities)
|
||||
await _backgroundActivityScheduler.ScheduleAsync(scheduledBackgroundActivity, context.CancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,5 +1,3 @@
|
|||
using System.Text.Json.Serialization;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Models.Bookmarks;
|
||||
|
||||
/// <summary>
|
||||
|
|
@ -8,10 +6,7 @@ namespace Elsa.Workflows.Runtime.Models.Bookmarks;
|
|||
public class BackgroundActivityBookmark
|
||||
{
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="BackgroundActivityBookmark"/> class.
|
||||
/// Set retroactively after the job has been scheduled. It is used to cancel te job when the bookmark is deleted.
|
||||
/// </summary>
|
||||
[JsonConstructor]
|
||||
public BackgroundActivityBookmark()
|
||||
{
|
||||
}
|
||||
public string? JobId { get; set; }
|
||||
}
|
||||
|
|
@ -2,4 +2,15 @@ using Elsa.Workflows.Core.Models;
|
|||
|
||||
namespace Elsa.Workflows.Runtime.Models.Notifications;
|
||||
|
||||
public record IndexedWorkflowBookmarks(string InstanceId, ICollection<Bookmark> AddedBookmarks, ICollection<Bookmark> RemovedBookmarks, ICollection<Bookmark> UnchangedBookmarks);
|
||||
/// <summary>
|
||||
/// Contains the bookmarks that were added, removed, or unchanged.
|
||||
/// </summary>
|
||||
/// <param name="InstanceId">The workflow instance ID.</param>
|
||||
/// <param name="AddedBookmarks">The bookmarks that were added.</param>
|
||||
/// <param name="RemovedBookmarks">The bookmarks that were removed.</param>
|
||||
/// <param name="UnchangedBookmarks">The bookmarks that were unchanged.</param>
|
||||
public record IndexedWorkflowBookmarks(
|
||||
string InstanceId,
|
||||
ICollection<Bookmark> AddedBookmarks,
|
||||
ICollection<Bookmark> RemovedBookmarks,
|
||||
ICollection<Bookmark> UnchangedBookmarks);
|
||||
|
|
@ -4,6 +4,6 @@ namespace Elsa.Workflows.Runtime.Models;
|
|||
/// Represents a scheduled background activity
|
||||
/// </summary>
|
||||
/// <param name="WorkflowInstanceId">The ID of the workflow instance containing the activity to execute.</param>
|
||||
/// <param name="ActivityId">The ID of the activity to execute.</param>
|
||||
/// <param name="ActivityNodeId">The ID of the activity to execute.</param>
|
||||
/// <param name="BookmarkId">The ID of the bookmark to resume.</param>
|
||||
public record ScheduledBackgroundActivity(string WorkflowInstanceId, string ActivityId, string BookmarkId);
|
||||
public record ScheduledBackgroundActivity(string WorkflowInstanceId, string ActivityNodeId, string BookmarkId);
|
||||
|
|
@ -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;
|
||||
/// <summary>
|
||||
/// A notification that is sent when the bookmarks of a workflow instance have been indexed.
|
||||
/// </summary>
|
||||
/// <param name="WorkflowExecutionContext">The workflow execution context.</param>
|
||||
/// <param name="IndexedWorkflowBookmarks">The bookmarks that were added, removed, or unchanged.</param>
|
||||
public record WorkflowBookmarksIndexed(WorkflowExecutionContext WorkflowExecutionContext, IndexedWorkflowBookmarks IndexedWorkflowBookmarks) : INotification;
|
||||
|
|
@ -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
|
||||
{
|
||||
|
|
|
|||
|
|
@ -223,6 +223,12 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
await StoreBookmarksAsync(context.InstanceId, context.Diff.Added, context.CorrelationId, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await _bookmarkStore.SaveAsync(bookmark, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<int> CountRunningWorkflowsAsync(CountRunningWorkflowsArgs args, CancellationToken cancellationToken = default) => await _workflowStateStore.CountAsync(args, cancellationToken);
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
|||
/// </summary>
|
||||
public class LocalBackgroundActivityScheduler : IBackgroundActivityScheduler
|
||||
{
|
||||
private readonly Channel<ScheduledBackgroundActivity> _channel;
|
||||
private readonly IJobQueue _jobQueue;
|
||||
private readonly IBackgroundActivityInvoker _backgroundActivityInvoker;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="LocalBackgroundActivityScheduler"/> class.
|
||||
/// </summary>
|
||||
/// <param name="channel">The channel to write to.</param>
|
||||
public LocalBackgroundActivityScheduler(Channel<ScheduledBackgroundActivity> channel)
|
||||
public LocalBackgroundActivityScheduler(IJobQueue jobQueue, IBackgroundActivityInvoker backgroundActivityInvoker)
|
||||
{
|
||||
_channel = channel;
|
||||
_jobQueue = jobQueue;
|
||||
_backgroundActivityInvoker = backgroundActivityInvoker;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<string> ScheduleAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var jobId = _jobQueue.Enqueue(async ct => await InvokeBackgroundActivity(scheduledBackgroundActivity, ct));
|
||||
return Task.FromResult(jobId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task CancelAsync(string jobId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
_jobQueue.Cancel(jobId);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<string> 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);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue