diff --git a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs index 1aa60ee4a..b53b7958a 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs @@ -232,10 +232,27 @@ public class ActivityExecutionContext : IExecutionContext { if (this.GetIsBackgroundExecution()) { - // TODO: Capture the information in a serializable format and store it in the workflow execution context Properties dictionary. - // The information should be stored in a way that allows the workflow execution context to resume the activity execution context. + var scheduledActivity = new ScheduledActivity + { + ActivityNodeId = activityNode?.NodeId, + OwnerActivityInstanceId = owner?.Id, + Options = options != null ? new ScheduledActivityOptions + { + CompletionCallback = options?.CompletionCallback?.Method.Name, + Tag = options?.Tag, + ExistingActivityInstanceId = options?.ExistingActivityExecutionContext?.Id, + PreventDuplicateScheduling = options?.PreventDuplicateScheduling ?? false, + Variables = options?.Variables?.ToList(), + Input = options?.Input + } : default + }; + + var scheduledActivities = this.GetBackgroundScheduledActivities().ToList(); + scheduledActivities.Add(scheduledActivity); + this.SetBackgroundScheduledActivities(scheduledActivities); + return; } - + var completionCallback = options?.CompletionCallback; owner ??= this; diff --git a/src/modules/Elsa.Workflows.Core/Extensions/BackgroundActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/BackgroundActivityExecutionContextExtensions.cs index 72d84e745..b427d6751 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/BackgroundActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/BackgroundActivityExecutionContextExtensions.cs @@ -1,3 +1,5 @@ +using Elsa.Workflows.Models; + namespace Elsa.Workflows; /// @@ -42,4 +44,22 @@ public static class BackgroundActivityExecutionContextExtensions { return activityExecutionContext.GetProperty>("BackgroundOutcomes") ?? Enumerable.Empty(); } + + /// + /// Sets the background scheduled activities. + /// + public static void SetBackgroundScheduledActivities(this ActivityExecutionContext activityExecutionContext, IEnumerable scheduledActivities) + { + var scheduledActivitiesList = scheduledActivities.ToList(); + activityExecutionContext.SetProperty("BackgroundScheduledActivities", scheduledActivitiesList); + } + + /// + /// Gets the background scheduled activities. + /// + /// + public static IEnumerable GetBackgroundScheduledActivities(this ActivityExecutionContext activityExecutionContext) + { + return activityExecutionContext.GetProperty>("BackgroundScheduledActivities") ?? Enumerable.Empty(); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Models/ScheduledActivity.cs b/src/modules/Elsa.Workflows.Core/Models/ScheduledActivity.cs new file mode 100644 index 000000000..768bbb112 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Models/ScheduledActivity.cs @@ -0,0 +1,8 @@ +namespace Elsa.Workflows.Models; + +public class ScheduledActivity +{ + public string? ActivityNodeId { get; set; } + public string? OwnerActivityInstanceId { get; set; } + public ScheduledActivityOptions? Options { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Models/ScheduledActivityOptions.cs b/src/modules/Elsa.Workflows.Core/Models/ScheduledActivityOptions.cs new file mode 100644 index 000000000..403c91e65 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Models/ScheduledActivityOptions.cs @@ -0,0 +1,13 @@ +using Elsa.Workflows.Memory; + +namespace Elsa.Workflows.Models; + +public class ScheduledActivityOptions +{ + public string? CompletionCallback { get; set; } + public object? Tag { get; set; } + public ICollection? Variables { get; set; } + public string? ExistingActivityInstanceId { get; set; } + public bool PreventDuplicateScheduling { get; set; } + public IDictionary? Input { get; set; } +} \ 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 6f440c8f5..defca4956 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs @@ -1,6 +1,8 @@ +using System.Text.Json; using Elsa.Extensions; using Elsa.Workflows.Middleware.Activities; using Elsa.Workflows.Models; +using Elsa.Workflows.Options; using Elsa.Workflows.Pipelines.ActivityExecution; using Elsa.Workflows.Runtime.Bookmarks; using Elsa.Workflows.Runtime.Middleware.Workflows; @@ -17,6 +19,7 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew internal static string GetBackgroundActivityOutputKey(string activityNodeId) => $"__BackgroundActivityOutput:{activityNodeId}"; internal static string GetBackgroundActivityOutcomesKey(string activityNodeId) => $"__BackgroundActivityOutcomes:{activityNodeId}"; internal static string GetBackgroundActivityJournalDataKey(string activityNodeId) => $"__BackgroundActivityJournalData:{activityNodeId}"; + internal static string GetBackgroundActivityScheduledActivitiesKey(string activityNodeId) => $"__BackgroundActivityScheduledActivities:{activityNodeId}"; internal static readonly object BackgroundActivitySchedulesKey = new(); internal const string BackgroundActivityBookmarkName = "BackgroundActivity"; @@ -42,7 +45,8 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew { CaptureOutputIfAny(context); CaptureJournalData(context); - await CompleteBackgroundActivityAsync(context); + await CompleteBackgroundActivityOutcomesAsync(context); + await CompleteBackgroundActivityScheduledActivitiesAsync(context); } } } @@ -86,10 +90,10 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew var activity = context.Activity; var inputKey = GetBackgroundActivityOutputKey(activity.NodeId); var capturedOutput = context.WorkflowExecutionContext.GetProperty>(inputKey); - - if(capturedOutput == null) + + if (capturedOutput == null) return; - + foreach (var outputEntry in capturedOutput) { var outputDescriptor = context.ActivityDescriptor.Outputs.FirstOrDefault(x => x.Name == outputEntry.Key); @@ -101,7 +105,7 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew context.Set(output, outputEntry.Value); } } - + private void CaptureJournalData(ActivityExecutionContext context) { var activity = context.Activity; @@ -114,8 +118,8 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew foreach (var journalEntry in journalData) context.JournalData[journalEntry.Key] = journalEntry.Value; } - - private async Task CompleteBackgroundActivityAsync(ActivityExecutionContext context) + + private async Task CompleteBackgroundActivityOutcomesAsync(ActivityExecutionContext context) { var outcomesKey = GetBackgroundActivityOutcomesKey(context.NodeId); var outcomes = context.WorkflowExecutionContext.GetProperty>(outcomesKey); @@ -123,9 +127,40 @@ public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddlew if (outcomes != null) { await context.CompleteActivityWithOutcomesAsync(outcomes.ToArray()); - + // Remove the outcomes from the workflow execution context. context.WorkflowExecutionContext.Properties.Remove(outcomesKey); } } + + private async Task CompleteBackgroundActivityScheduledActivitiesAsync(ActivityExecutionContext context) + { + var scheduledActivitiesKey = GetBackgroundActivityScheduledActivitiesKey(context.NodeId); + var scheduledActivitiesJson = context.WorkflowExecutionContext.GetProperty(scheduledActivitiesKey); + var scheduledActivities = scheduledActivitiesJson != null ? JsonSerializer.Deserialize>(scheduledActivitiesJson) : null; + + if (scheduledActivities != null) + { + foreach (var scheduledActivity in scheduledActivities) + { + var activityNode = scheduledActivity.ActivityNodeId != null ? context.WorkflowExecutionContext.FindActivityByNodeId(scheduledActivity.ActivityNodeId) : null; + var owner = scheduledActivity.OwnerActivityInstanceId != null ? context.WorkflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Id == scheduledActivity.OwnerActivityInstanceId) : null; + var options = scheduledActivity.Options != null + ? new ScheduleWorkOptions + { + ExistingActivityExecutionContext = scheduledActivity.Options.ExistingActivityInstanceId != null ? context.WorkflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Id == scheduledActivity.Options.ExistingActivityInstanceId) : null, + Variables = scheduledActivity.Options?.Variables, + CompletionCallback = !string.IsNullOrEmpty(scheduledActivity.Options?.CompletionCallback) && owner != null ? owner.Activity.GetActivityCompletionCallback(scheduledActivity.Options.CompletionCallback) : default, + PreventDuplicateScheduling = scheduledActivity.Options?.PreventDuplicateScheduling ?? false, + Input = scheduledActivity.Options?.Input, + Tag = scheduledActivity.Options?.Tag + } + : default; + await context.ScheduleActivityAsync(activityNode, owner, options); + } + + // Remove the scheduled activities from the workflow execution context. + context.WorkflowExecutionContext.Properties.Remove(scheduledActivitiesKey); + } + } } \ 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 90b7557a2..372271f22 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs @@ -1,3 +1,4 @@ +using System.Text.Json; using Elsa.Common.Models; using Elsa.Workflows.Contracts; using Elsa.Workflows.Helpers; @@ -57,7 +58,7 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker public async Task ExecuteAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default) { var workflowInstanceId = scheduledBackgroundActivity.WorkflowInstanceId; - + var workflowState = await _workflowRuntime.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken); if (workflowState == null) @@ -85,7 +86,6 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker // Capture any activity output produced by the activity (but only if the associated memory block is stored in the workflow itself). var outputDescriptors = activityExecutionContext.ActivityDescriptor.Outputs; var outputValues = new Dictionary(); - var outcomes = activityExecutionContext.GetBackgroundOutcomes().ToList(); foreach (var outputDescriptor in outputDescriptors) { @@ -112,19 +112,23 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker outputValues[outputDescriptor.Name] = outputValue; } - // Resume the workflow, passing along the activity output. + // Resume the workflow, passing along activity output, outcomes and scheduled activities. var bookmarkId = scheduledBackgroundActivity.BookmarkId; var inputKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityOutputKey(activityNodeId); var outcomesKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityOutcomesKey(activityNodeId); var journalDataKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityJournalDataKey(activityNodeId); + var scheduledActivitiesKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityScheduledActivitiesKey(activityNodeId); + var outcomes = activityExecutionContext.GetBackgroundOutcomes().ToList(); + var scheduledActivities = activityExecutionContext.GetBackgroundScheduledActivities().ToList(); var dispatchRequest = new DispatchWorkflowInstanceRequest { - InstanceId = workflowInstanceId, - BookmarkId = bookmarkId, + InstanceId = workflowInstanceId, + BookmarkId = bookmarkId, Properties = new Dictionary { [outcomesKey] = outcomes, + [scheduledActivitiesKey] = JsonSerializer.Serialize(scheduledActivities), [inputKey] = outputValues, [journalDataKey] = activityExecutionContext.JournalData }