Add scheduling function for background activities
This commit achieves two main goals. Firstly, it introduces two new classes called ScheduledActivity and ScheduledActivityOptions to store scheduled activities' information. Secondly, it modifies how activities are executed in the background by capturing the scheduling information as a serializable format and storing it in the workflow execution context properties dictionary. This change allows the workflow execution context to resume the activity execution context.
This commit is contained in:
parent
2b6294bcc1
commit
a950a6de85
|
|
@ -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;
|
||||
|
||||
|
|
|
|||
|
|
@ -1,3 +1,5 @@
|
|||
using Elsa.Workflows.Models;
|
||||
|
||||
namespace Elsa.Workflows;
|
||||
|
||||
/// <summary>
|
||||
|
|
@ -42,4 +44,22 @@ public static class BackgroundActivityExecutionContextExtensions
|
|||
{
|
||||
return activityExecutionContext.GetProperty<IEnumerable<string>>("BackgroundOutcomes") ?? Enumerable.Empty<string>();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Sets the background scheduled activities.
|
||||
/// </summary>
|
||||
public static void SetBackgroundScheduledActivities(this ActivityExecutionContext activityExecutionContext, IEnumerable<ScheduledActivity> scheduledActivities)
|
||||
{
|
||||
var scheduledActivitiesList = scheduledActivities.ToList();
|
||||
activityExecutionContext.SetProperty("BackgroundScheduledActivities", scheduledActivitiesList);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the background scheduled activities.
|
||||
/// </summary>
|
||||
/// <param name="activityExecutionContext"></param>
|
||||
public static IEnumerable<ScheduledActivity> GetBackgroundScheduledActivities(this ActivityExecutionContext activityExecutionContext)
|
||||
{
|
||||
return activityExecutionContext.GetProperty<IEnumerable<ScheduledActivity>>("BackgroundScheduledActivities") ?? Enumerable.Empty<ScheduledActivity>();
|
||||
}
|
||||
}
|
||||
|
|
@ -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; }
|
||||
}
|
||||
|
|
@ -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<Variable>? Variables { get; set; }
|
||||
public string? ExistingActivityInstanceId { get; set; }
|
||||
public bool PreventDuplicateScheduling { get; set; }
|
||||
public IDictionary<string,object>? Input { get; set; }
|
||||
}
|
||||
|
|
@ -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<IDictionary<string, object>>(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<ICollection<string>>(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<string>(scheduledActivitiesKey);
|
||||
var scheduledActivities = scheduledActivitiesJson != null ? JsonSerializer.Deserialize<ICollection<ScheduledActivity>>(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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<string, object>();
|
||||
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<string, object>
|
||||
{
|
||||
[outcomesKey] = outcomes,
|
||||
[scheduledActivitiesKey] = JsonSerializer.Serialize(scheduledActivities),
|
||||
[inputKey] = outputValues,
|
||||
[journalDataKey] = activityExecutionContext.JournalData
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue