Add background execution to activities and update HTTP requests
Significantly enhanced the capabilities of background execution of activities. Included a change in activity type of "SendHttpRequest" from 'Task' to 'Action'. Introduced new classes for handling outcomes of context in background execution. Made some necessary adjustments to HTTP Request Task to handle sending HTTP requests from a background task. Updated several middleware classes to align with these modifications.
This commit is contained in:
parent
8002f20e50
commit
853355bcbf
|
|
@ -12,7 +12,7 @@ namespace Elsa.Http;
|
|||
/// <summary>
|
||||
/// Send an HTTP request.
|
||||
/// </summary>
|
||||
[Activity("Elsa", "HTTP", "Send an HTTP request.", DisplayName = "HTTP Request (flow)", Kind = ActivityKind.Task)]
|
||||
[Activity("Elsa", "HTTP", "Send an HTTP request.", DisplayName = "HTTP Request (flow)", Kind = ActivityKind.Action)]
|
||||
public class FlowSendHttpRequest : SendHttpRequestBase, IActivityPropertyDefaultValueProvider
|
||||
{
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ namespace Elsa.Http;
|
|||
/// <summary>
|
||||
/// Send an HTTP request.
|
||||
/// </summary>
|
||||
[Activity("Elsa", "HTTP", "Send an HTTP request.", DisplayName = "HTTP Request", Kind = ActivityKind.Task)]
|
||||
[Activity("Elsa", "HTTP", "Send an HTTP request.", DisplayName = "HTTP Request", Kind = ActivityKind.Action)]
|
||||
public class SendHttpRequest : SendHttpRequestBase
|
||||
{
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
220
src/modules/Elsa.Http/Activities/SendHttpRequestTask.cs
Normal file
220
src/modules/Elsa.Http/Activities/SendHttpRequestTask.cs
Normal file
|
|
@ -0,0 +1,220 @@
|
|||
using System.Net;
|
||||
using System.Net.Http.Headers;
|
||||
using System.Runtime.CompilerServices;
|
||||
using Elsa.Extensions;
|
||||
using Elsa.Http.ContentWriters;
|
||||
using Elsa.Http.UIHints;
|
||||
using Elsa.Workflows;
|
||||
using Elsa.Workflows.Attributes;
|
||||
using Elsa.Workflows.UIHints;
|
||||
using Elsa.Workflows.Models;
|
||||
using HttpHeaders = Elsa.Http.Models.HttpHeaders;
|
||||
|
||||
namespace Elsa.Http;
|
||||
|
||||
/// <summary>
|
||||
/// Sends HTTP requests from a background task.
|
||||
/// </summary>
|
||||
[Activity("Elsa", "HTTP", "Send an HTTP request from a background task.", DisplayName = "HTTP Request Task", Kind = ActivityKind.Job)]
|
||||
public class SendHttpRequestTask : Activity<HttpStatusCode>
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public SendHttpRequestTask([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
|
||||
{
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The URL to send the request to.
|
||||
/// </summary>
|
||||
[Input(Description = "The URL to send the request to.")]
|
||||
public Input<Uri?> Url { get; set; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// The HTTP method to use when sending the request.
|
||||
/// </summary>
|
||||
[Input(
|
||||
Description = "The HTTP method to use when sending the request.",
|
||||
Options = new[] { "GET", "POST", "PUT", "DELETE", "PATCH", "OPTIONS", "HEAD" },
|
||||
DefaultValue = "GET",
|
||||
UIHint = InputUIHints.DropDown
|
||||
)]
|
||||
public Input<string> Method { get; set; } = new("GET");
|
||||
|
||||
/// <summary>
|
||||
/// The content to send with the request. Can be a string, an object, a byte array or a stream.
|
||||
/// </summary>
|
||||
[Input(Description = "The content to send with the request. Can be a string, an object, a byte array or a stream.")]
|
||||
public Input<object?> Content { get; set; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// The content type to use when sending the request.
|
||||
/// </summary>
|
||||
[Input(
|
||||
Description = "The content type to use when sending the request.",
|
||||
UIHandler = typeof(HttpContentTypeOptionsProvider),
|
||||
UIHint = InputUIHints.DropDown
|
||||
)]
|
||||
public Input<string?> ContentType { get; set; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// The Authorization header value to send with the request.
|
||||
/// </summary>
|
||||
/// <example>Bearer {some-access-token}</example>
|
||||
[Input(Description = "The Authorization header value to send with the request. For example: Bearer {some-access-token}", Category = "Security")]
|
||||
public Input<string?> Authorization { get; set; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// A list of expected status codes to handle.
|
||||
/// </summary>
|
||||
[Input(
|
||||
Description = "A list of expected status codes to handle.",
|
||||
UIHint = InputUIHints.MultiText,
|
||||
DefaultValueProvider = typeof(FlowSendHttpRequest)
|
||||
)]
|
||||
public Input<ICollection<int>> ExpectedStatusCodes { get; set; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// A value that allows to add the Authorization header without validation.
|
||||
/// </summary>
|
||||
[Input(Description = "A value that allows to add the Authorization header without validation.", Category = "Security")]
|
||||
public Input<bool> DisableAuthorizationHeaderValidation { get; set; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// The headers to send along with the request.
|
||||
/// </summary>
|
||||
[Input(
|
||||
Description = "The headers to send along with the request.",
|
||||
UIHint = InputUIHints.JsonEditor,
|
||||
Category = "Advanced"
|
||||
)]
|
||||
public Input<HttpHeaders?> RequestHeaders { get; set; } = new(new HttpHeaders());
|
||||
|
||||
/// <summary>
|
||||
/// The parsed content, if any.
|
||||
/// </summary>
|
||||
[Output(Description = "The parsed content, if any.")]
|
||||
public Output<object?> ParsedContent { get; set; } = default!;
|
||||
|
||||
/// <inheritdoc />
|
||||
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
|
||||
{
|
||||
await TrySendAsync(context);
|
||||
}
|
||||
|
||||
private async Task TrySendAsync(ActivityExecutionContext context)
|
||||
{
|
||||
var request = PrepareRequest(context);
|
||||
var httpClientFactory = context.GetRequiredService<IHttpClientFactory>();
|
||||
var httpClient = httpClientFactory.CreateClient(nameof(SendHttpRequestBase));
|
||||
var cancellationToken = context.CancellationToken;
|
||||
|
||||
try
|
||||
{
|
||||
var response = await httpClient.SendAsync(request, cancellationToken);
|
||||
var parsedContent = await ParseContentAsync(context, response.Content);
|
||||
context.SetResult(response.StatusCode);
|
||||
context.Set(ParsedContent, parsedContent);
|
||||
|
||||
HandleResponse(context, response);
|
||||
}
|
||||
catch (HttpRequestException e)
|
||||
{
|
||||
context.AddExecutionLogEntry("Error", e.Message, payload: new { StackTrace = e.StackTrace });
|
||||
context.JournalData.Add("Error", e.Message);
|
||||
HandleRequestException(context, e);
|
||||
}
|
||||
catch (TaskCanceledException e)
|
||||
{
|
||||
context.AddExecutionLogEntry("Error", e.Message, payload: new { StackTrace = e.StackTrace });
|
||||
context.JournalData.Add("Cancelled", true);
|
||||
HandleTaskCanceledException(context, e);
|
||||
}
|
||||
}
|
||||
|
||||
private void HandleResponse(ActivityExecutionContext context, HttpResponseMessage response)
|
||||
{
|
||||
var expectedStatusCodes = ExpectedStatusCodes.GetOrDefault(context) ?? new List<int>(0);
|
||||
var statusCode = (int)response.StatusCode;
|
||||
var hasMatchingStatusCode = expectedStatusCodes.Contains(statusCode);
|
||||
var outcome = expectedStatusCodes.Any() ? hasMatchingStatusCode ? statusCode.ToString() : "Unmatched status code" : default;
|
||||
var outcomes = new List<string>();
|
||||
|
||||
if (outcome != null)
|
||||
outcomes.Add(outcome);
|
||||
|
||||
outcomes.Add("Done");
|
||||
context.JournalData["StatusCode"] = statusCode;
|
||||
context.SetBackgroundOutcomes(outcomes);
|
||||
}
|
||||
|
||||
private void HandleRequestException(ActivityExecutionContext context, HttpRequestException exception)
|
||||
{
|
||||
context.SetBackgroundOutcomes(new[] { "Failed to connect" });
|
||||
}
|
||||
|
||||
private void HandleTaskCanceledException(ActivityExecutionContext context, TaskCanceledException exception)
|
||||
{
|
||||
context.SetBackgroundOutcomes(new[] { "Timeout" });
|
||||
}
|
||||
|
||||
private async Task<object?> ParseContentAsync(ActivityExecutionContext context, HttpContent httpContent)
|
||||
{
|
||||
if (!HasContent(httpContent))
|
||||
return null;
|
||||
|
||||
var cancellationToken = context.CancellationToken;
|
||||
var targetType = ParsedContent.GetTargetType(context);
|
||||
var contentStream = await httpContent.ReadAsStreamAsync(cancellationToken);
|
||||
var contentType = httpContent.Headers.ContentType?.MediaType!;
|
||||
|
||||
targetType ??= contentType switch
|
||||
{
|
||||
"application/json" => typeof(object),
|
||||
_ => typeof(string)
|
||||
};
|
||||
|
||||
return await context.ParseContentAsync(contentStream, contentType, targetType, cancellationToken);
|
||||
}
|
||||
|
||||
private static bool HasContent(HttpContent httpContent) => httpContent.Headers.ContentLength > 0;
|
||||
|
||||
private HttpRequestMessage PrepareRequest(ActivityExecutionContext context)
|
||||
{
|
||||
var method = Method.GetOrDefault(context) ?? "GET";
|
||||
var url = Url.Get(context);
|
||||
var request = new HttpRequestMessage(new HttpMethod(method), url);
|
||||
var headers = context.GetHeaders(RequestHeaders);
|
||||
var authorization = Authorization.GetOrDefault(context);
|
||||
var addAuthorizationWithoutValidation = DisableAuthorizationHeaderValidation.GetOrDefault(context);
|
||||
|
||||
if (!string.IsNullOrWhiteSpace(authorization))
|
||||
if (addAuthorizationWithoutValidation)
|
||||
request.Headers.TryAddWithoutValidation("Authorization", authorization);
|
||||
else
|
||||
request.Headers.Authorization = AuthenticationHeaderValue.Parse(authorization);
|
||||
|
||||
foreach (var header in headers)
|
||||
request.Headers.Add(header.Key, header.Value.AsEnumerable());
|
||||
|
||||
var contentType = ContentType.GetOrDefault(context);
|
||||
var content = Content.GetOrDefault(context);
|
||||
|
||||
if (contentType != null && content != null)
|
||||
{
|
||||
var factories = context.GetServices<IHttpContentFactory>();
|
||||
var factory = SelectContentWriter(contentType, factories);
|
||||
request.Content = factory.CreateHttpContent(content, contentType);
|
||||
}
|
||||
|
||||
return request;
|
||||
}
|
||||
|
||||
private IHttpContentFactory SelectContentWriter(string? contentType, IEnumerable<IHttpContentFactory> factories)
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(contentType))
|
||||
return new JsonContentFactory();
|
||||
|
||||
var parsedContentType = new System.Net.Mime.ContentType(contentType);
|
||||
return factories.FirstOrDefault(httpContentFactory => httpContentFactory.SupportedContentTypes.Any(c => c == parsedContentType.MediaType)) ?? new JsonContentFactory();
|
||||
}
|
||||
}
|
||||
|
|
@ -46,6 +46,7 @@ public class MassTransitWorkflowDispatcher : IWorkflowDispatcher
|
|||
ActivityInstanceId = request.ActivityInstanceId,
|
||||
ActivityHash = request.ActivityHash,
|
||||
Input = request.Input,
|
||||
Properties = request.Properties,
|
||||
CorrelationId = request.CorrelationId
|
||||
}, cancellationToken);
|
||||
return new();
|
||||
|
|
|
|||
|
|
@ -36,7 +36,7 @@ public class WorkflowContextActivityExecutionMiddleware : IActivityExecutionMidd
|
|||
}
|
||||
|
||||
// Check if this is a background execution.
|
||||
var isBackgroundExecution = context.TransientProperties.GetValueOrDefault<object, bool>(BackgroundActivityCollectorMiddleware.IsBackgroundExecution);
|
||||
var isBackgroundExecution = context.TransientProperties.GetValueOrDefault<object, bool>(BackgroundActivityInvokerMiddleware.IsBackgroundExecution);
|
||||
|
||||
// Is the activity configured to load the context?
|
||||
foreach (var providerType in providerTypes)
|
||||
|
|
|
|||
|
|
@ -0,0 +1,24 @@
|
|||
namespace Elsa.Workflows;
|
||||
|
||||
/// <summary>
|
||||
/// Adds extension methods to <see cref="ActivityExecutionContext"/>.
|
||||
/// </summary>
|
||||
public static class BackgroundActivityExecutionContextExtensions
|
||||
{
|
||||
/// <summary>
|
||||
/// Sets the background outcomes.
|
||||
/// </summary>
|
||||
public static void SetBackgroundOutcomes(this ActivityExecutionContext activityExecutionContext, IEnumerable<string> outcomes)
|
||||
{
|
||||
var outcomesList = outcomes.ToList();
|
||||
activityExecutionContext.SetProperty("BackgroundOutcomes", outcomesList);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the background outcomes.
|
||||
/// </summary>
|
||||
public static IEnumerable<string> GetBackgroundOutcomes(this ActivityExecutionContext activityExecutionContext)
|
||||
{
|
||||
return activityExecutionContext.GetProperty<IEnumerable<string>>("BackgroundOutcomes") ?? Enumerable.Empty<string>();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,3 @@
|
|||
namespace Elsa.Workflows.Models;
|
||||
|
||||
public record BackgroundExecutionOutcome(string Name, object? Payload);
|
||||
|
|
@ -0,0 +1,8 @@
|
|||
namespace Elsa.Workflows.Models;
|
||||
|
||||
public class BackgroundExecutionResult
|
||||
{
|
||||
public ICollection<BackgroundExecutionOutcome> Outcomes { get; set; } = new List<BackgroundExecutionOutcome>();
|
||||
public ICollection<WorkflowExecutionLogEntry> ExecutionLog { get; set; } = new List<WorkflowExecutionLogEntry>();
|
||||
public IDictionary<string, object?> JournalData { get; } = new Dictionary<string, object?>();
|
||||
}
|
||||
|
|
@ -77,7 +77,9 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
|
|||
|
||||
private void ApplyProperties(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
|
||||
{
|
||||
workflowExecutionContext.Properties = state.Properties;
|
||||
// Merge properties.
|
||||
foreach (var property in state.Properties)
|
||||
workflowExecutionContext.Properties[property.Key] = property.Value;
|
||||
}
|
||||
|
||||
private static void ApplyActivityExecutionContexts(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
|
||||
|
|
@ -248,7 +250,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
|
|||
|
||||
// // If there are any faulted contexts, keep everything so that the user can fix the issue and potentially reschedule existing instances.
|
||||
// if (contexts.Any(x => x.Status == ActivityStatus.Faulted))
|
||||
return contexts;
|
||||
return contexts;
|
||||
|
||||
// return contexts
|
||||
// .Where(x => !x.IsCompleted)
|
||||
|
|
|
|||
|
|
@ -88,9 +88,4 @@ public interface IWorkflowRuntime
|
|||
/// Counts the number of workflow instances based on the provided query args.
|
||||
/// </summary>
|
||||
Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Merges the specified workflow state into the workflow runtime.
|
||||
/// </summary>
|
||||
Task MergeWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -11,7 +11,7 @@ namespace Elsa.Extensions;
|
|||
public static class ActivityExecutionPipelineBuilderExtensions
|
||||
{
|
||||
/// <summary>
|
||||
/// Installs the <see cref="BackgroundActivityCollectorMiddleware"/>.
|
||||
/// Installs the <see cref="BackgroundActivityInvokerMiddleware"/>.
|
||||
/// </summary>
|
||||
public static IActivityExecutionPipelineBuilder UseBackgroundActivityInvoker(this IActivityExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<BackgroundActivityCollectorMiddleware>();
|
||||
public static IActivityExecutionPipelineBuilder UseBackgroundActivityInvoker(this IActivityExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<BackgroundActivityInvokerMiddleware>();
|
||||
}
|
||||
|
|
@ -31,7 +31,7 @@ public class CancelBackgroundActivities : INotificationHandler<WorkflowBookmarks
|
|||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowBookmarksIndexed notification, CancellationToken cancellationToken)
|
||||
{
|
||||
var removedBookmarks = notification.IndexedWorkflowBookmarks.RemovedBookmarks.Where(x => x.Name == BackgroundActivityCollectorMiddleware.BackgroundActivityBookmarkName);
|
||||
var removedBookmarks = notification.IndexedWorkflowBookmarks.RemovedBookmarks.Where(x => x.Name == BackgroundActivityInvokerMiddleware.BackgroundActivityBookmarkName);
|
||||
|
||||
foreach (var removedBookmark in removedBookmarks)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -12,19 +12,21 @@ namespace Elsa.Workflows.Runtime.Middleware.Activities;
|
|||
/// Collects the current activity for scheduling for execution from a background job if the activity is of kind <see cref="ActivityKind.Job"/> or <see cref="Task"/>.
|
||||
/// The actual scheduling of the activity happens in <see cref="ScheduleBackgroundActivitiesMiddleware"/>.
|
||||
/// </summary>
|
||||
public class BackgroundActivityCollectorMiddleware : DefaultActivityInvokerMiddleware
|
||||
public class BackgroundActivityInvokerMiddleware : DefaultActivityInvokerMiddleware
|
||||
{
|
||||
/// <summary>
|
||||
/// A key into the activity execution context's transient properties that indicates whether the current activity is being executed in the background.
|
||||
/// </summary>
|
||||
public static readonly object IsBackgroundExecution = new();
|
||||
|
||||
internal static string GetBackgroundActivityOutputKey(string activityId) => $"__BackgroundActivityOutput:{activityId}";
|
||||
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 readonly object BackgroundActivitySchedulesKey = new();
|
||||
internal const string BackgroundActivityBookmarkName = "BackgroundActivity";
|
||||
|
||||
/// <inheritdoc />
|
||||
public BackgroundActivityCollectorMiddleware(ActivityMiddlewareDelegate next) : base(next)
|
||||
public BackgroundActivityInvokerMiddleware(ActivityMiddlewareDelegate next) : base(next)
|
||||
{
|
||||
}
|
||||
|
||||
|
|
@ -37,8 +39,16 @@ public class BackgroundActivityCollectorMiddleware : DefaultActivityInvokerMiddl
|
|||
ScheduleBackgroundActivity(context);
|
||||
else
|
||||
{
|
||||
CaptureOutputIfAny(context);
|
||||
await base.ExecuteActivityAsync(context);
|
||||
|
||||
// This part is either executed from the background, or in the foreground when the activity is resumed.
|
||||
var isResuming = !GetIsBackgroundExecution(context) && context.ActivityDescriptor.Kind is ActivityKind.Task or ActivityKind.Job;
|
||||
if (isResuming)
|
||||
{
|
||||
CaptureOutputIfAny(context);
|
||||
CaptureJournalData(context);
|
||||
await CompleteBackgroundActivityAsync(context);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -51,11 +61,13 @@ public class BackgroundActivityCollectorMiddleware : DefaultActivityInvokerMiddl
|
|||
var activityDescriptor = context.ActivityDescriptor;
|
||||
var kind = activityDescriptor.Kind;
|
||||
|
||||
return !context.TransientProperties.ContainsKey(IsBackgroundExecution)
|
||||
return !GetIsBackgroundExecution(context)
|
||||
&& context.WorkflowExecutionContext.ExecuteDelegate == null
|
||||
&& (kind is ActivityKind.Job || (kind == ActivityKind.Task && activity.GetRunAsynchronously()));
|
||||
}
|
||||
|
||||
private static bool GetIsBackgroundExecution(ActivityExecutionContext context) => context.TransientProperties.ContainsKey(IsBackgroundExecution);
|
||||
|
||||
/// <summary>
|
||||
/// Schedules the current activity for execution in the background.
|
||||
/// </summary>
|
||||
|
|
@ -77,21 +89,48 @@ public class BackgroundActivityCollectorMiddleware : DefaultActivityInvokerMiddl
|
|||
private static void CaptureOutputIfAny(ActivityExecutionContext context)
|
||||
{
|
||||
var activity = context.Activity;
|
||||
var inputKey = GetBackgroundActivityOutputKey(activity.Id);
|
||||
|
||||
if (!context.WorkflowInput.TryGetValue(inputKey, out var capturedOutput))
|
||||
var inputKey = GetBackgroundActivityOutputKey(activity.NodeId);
|
||||
var capturedOutput = context.WorkflowExecutionContext.GetProperty<IDictionary<string, object>>(inputKey);
|
||||
|
||||
if(capturedOutput == null)
|
||||
return;
|
||||
|
||||
var input = (IDictionary<string, object>)capturedOutput;
|
||||
foreach (var inputEntry in input)
|
||||
|
||||
foreach (var outputEntry in capturedOutput)
|
||||
{
|
||||
var outputDescriptor = context.ActivityDescriptor.Outputs.FirstOrDefault(x => x.Name == inputEntry.Key);
|
||||
var outputDescriptor = context.ActivityDescriptor.Outputs.FirstOrDefault(x => x.Name == outputEntry.Key);
|
||||
|
||||
if (outputDescriptor == null)
|
||||
continue;
|
||||
|
||||
var output = (Output?)outputDescriptor.ValueGetter(activity);
|
||||
context.Set(output, inputEntry.Value);
|
||||
context.Set(output, outputEntry.Value);
|
||||
}
|
||||
}
|
||||
|
||||
private void CaptureJournalData(ActivityExecutionContext context)
|
||||
{
|
||||
var activity = context.Activity;
|
||||
var journalDataKey = GetBackgroundActivityJournalDataKey(activity.NodeId);
|
||||
var journalData = context.WorkflowExecutionContext.GetProperty<IDictionary<string, object>>(journalDataKey);
|
||||
|
||||
if (journalData == null)
|
||||
return;
|
||||
|
||||
foreach (var journalEntry in journalData)
|
||||
context.JournalData[journalEntry.Key] = journalEntry.Value;
|
||||
}
|
||||
|
||||
private async Task CompleteBackgroundActivityAsync(ActivityExecutionContext context)
|
||||
{
|
||||
var outcomesKey = GetBackgroundActivityOutcomesKey(context.NodeId);
|
||||
var outcomes = context.WorkflowExecutionContext.GetProperty<ICollection<string>>(outcomesKey);
|
||||
|
||||
if (outcomes != null)
|
||||
{
|
||||
await context.CompleteActivityWithOutcomesAsync(outcomes.ToArray());
|
||||
|
||||
// Remove the outcomes from the workflow execution context.
|
||||
context.WorkflowExecutionContext.Properties.Remove(outcomesKey);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -43,7 +43,7 @@ public class ScheduleBackgroundActivitiesMiddleware : WorkflowExecutionMiddlewar
|
|||
|
||||
var scheduledBackgroundActivities = workflowExecutionContext
|
||||
.TransientProperties
|
||||
.GetOrAdd(BackgroundActivityCollectorMiddleware.BackgroundActivitySchedulesKey, () => new List<ScheduledBackgroundActivity>());
|
||||
.GetOrAdd(BackgroundActivityInvokerMiddleware.BackgroundActivitySchedulesKey, () => new List<ScheduledBackgroundActivity>());
|
||||
|
||||
if (scheduledBackgroundActivities.Any())
|
||||
{
|
||||
|
|
|
|||
|
|
@ -70,7 +70,6 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker
|
|||
|
||||
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
|
||||
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflow, workflowState, cancellationTokens: cancellationToken);
|
||||
var originalBookmarks = workflowExecutionContext.Bookmarks.ToList();
|
||||
var activityNodeId = scheduledBackgroundActivity.ActivityNodeId;
|
||||
var activityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.NodeId == activityNodeId);
|
||||
|
||||
|
|
@ -78,7 +77,7 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker
|
|||
await _variablePersistenceManager.LoadVariablesAsync(workflowExecutionContext);
|
||||
|
||||
// Mark the activity as being invoked from a background worker.
|
||||
activityExecutionContext.TransientProperties[BackgroundActivityCollectorMiddleware.IsBackgroundExecution] = true;
|
||||
activityExecutionContext.TransientProperties[BackgroundActivityInvokerMiddleware.IsBackgroundExecution] = true;
|
||||
|
||||
// Invoke the activity.
|
||||
await _activityInvoker.InvokeAsync(activityExecutionContext);
|
||||
|
|
@ -86,6 +85,7 @@ 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,33 +112,21 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker
|
|||
outputValues[outputDescriptor.Name] = outputValue;
|
||||
}
|
||||
|
||||
// TODO: Instead of importing the entire workflow state, we should only import the following:
|
||||
// - Variables
|
||||
// - Activity state
|
||||
// - Activity output
|
||||
// - Bookmarks
|
||||
workflowState = _workflowStateExtractor.Extract(workflowExecutionContext);
|
||||
await _variablePersistenceManager.SaveVariablesAsync(workflowExecutionContext);
|
||||
await _workflowRuntime.MergeWorkflowStateAsync(workflowState, cancellationToken);
|
||||
//await _workflowRuntime.ImportWorkflowStateAsync(workflowState, cancellationToken);
|
||||
|
||||
// Process bookmarks.
|
||||
var newBookmarks = workflowExecutionContext.Bookmarks.ToList();
|
||||
var diff = Diff.For(originalBookmarks, newBookmarks);
|
||||
await _bookmarksPersister.PersistBookmarksAsync(workflowExecutionContext, diff);
|
||||
|
||||
// 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 = BackgroundActivityCollectorMiddleware.GetBackgroundActivityOutputKey(activityNodeId);
|
||||
var inputKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityOutputKey(activityNodeId);
|
||||
var outcomesKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityOutcomesKey(activityNodeId);
|
||||
var journalDataKey = BackgroundActivityInvokerMiddleware.GetBackgroundActivityJournalDataKey(activityNodeId);
|
||||
|
||||
var dispatchRequest = new DispatchWorkflowInstanceRequest
|
||||
{
|
||||
InstanceId = workflowInstanceId,
|
||||
BookmarkId = bookmarkId,
|
||||
Input = new Dictionary<string, object>
|
||||
InstanceId = workflowInstanceId,
|
||||
BookmarkId = bookmarkId,
|
||||
Properties = new Dictionary<string, object>
|
||||
{
|
||||
[inputKey] = outputValues
|
||||
[outcomesKey] = outcomes,
|
||||
[inputKey] = outputValues,
|
||||
[journalDataKey] = activityExecutionContext.JournalData
|
||||
}
|
||||
};
|
||||
|
||||
|
|
|
|||
|
|
@ -292,19 +292,6 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
return await _workflowInstanceStore.CountAsync(filter, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task MergeWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var existingWorkflowInstance = (await _workflowInstanceStore.FindAsync(workflowState.Id, cancellationToken))!;
|
||||
var workflowInstance = _workflowStateMapper.Map(workflowState)!;
|
||||
|
||||
foreach (var bookmark in workflowState.Bookmarks)
|
||||
{
|
||||
existingWorkflowInstance.WorkflowState.Bookmarks.RemoveWhere(x => x.Id == bookmark.Id);
|
||||
existingWorkflowInstance.WorkflowState.Bookmarks.Add(bookmark);
|
||||
}
|
||||
await _workflowInstanceManager.SaveAsync(workflowInstance, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<WorkflowExecutionResult> StartWorkflowAsync(IWorkflowHost workflowHost, StartWorkflowRuntimeOptions options)
|
||||
{
|
||||
var workflowInstanceId = string.IsNullOrEmpty(options.InstanceId) ? _identityGenerator.GenerateId() : options.InstanceId;
|
||||
|
|
|
|||
Loading…
Reference in a new issue