diff --git a/src/modules/Elsa.Http/Features/HttpFeature.cs b/src/modules/Elsa.Http/Features/HttpFeature.cs index fbe19748a..9dca6a73a 100644 --- a/src/modules/Elsa.Http/Features/HttpFeature.cs +++ b/src/modules/Elsa.Http/Features/HttpFeature.cs @@ -27,8 +27,8 @@ public class HttpFeature : FeatureBase /// public Action? ConfigureHttpOptions { get; set; } - public Func HttpEndpointAuthorizationHandlerFactory { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance; - public Func HttpEndpointWorkflowFaultHandlerFactory { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance; + public Func HttpEndpointAuthorizationHandler { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance; + public Func HttpEndpointWorkflowFaultHandler { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance; /// /// A delegate to configure the used when by the activity. @@ -74,6 +74,7 @@ public class HttpFeature : FeatureBase .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddNotificationHandlersFrom() .AddHttpContextAccessor() diff --git a/src/modules/Elsa.Http/Implementations/HttpBookmarkProcessor.cs b/src/modules/Elsa.Http/Implementations/HttpBookmarkProcessor.cs new file mode 100644 index 000000000..037c0aca5 --- /dev/null +++ b/src/modules/Elsa.Http/Implementations/HttpBookmarkProcessor.cs @@ -0,0 +1,94 @@ +using Elsa.Common.Models; +using Elsa.Http.Services; +using Elsa.Workflows.Core.Helpers; +using Elsa.Workflows.Runtime.Services; +using Microsoft.AspNetCore.Http; + +namespace Elsa.Http.Implementations; + +/// +public class HttpBookmarkProcessor : IHttpBookmarkProcessor +{ + private readonly IWorkflowRuntime _workflowRuntime; + private readonly IWorkflowDefinitionService _workflowDefinitionService; + private readonly IWorkflowHostFactory _workflowHostFactory; + private readonly IHttpContextAccessor _httpContextAccessor; + + /// + /// Constructor. + /// + public HttpBookmarkProcessor( + IWorkflowRuntime workflowRuntime, + IWorkflowDefinitionService workflowDefinitionService, + IWorkflowHostFactory workflowHostFactory, + IHttpContextAccessor httpContextAccessor) + { + _workflowRuntime = workflowRuntime; + _workflowDefinitionService = workflowDefinitionService; + _workflowHostFactory = workflowHostFactory; + _httpContextAccessor = httpContextAccessor; + } + + /// + public async Task ProcessBookmarks( + IEnumerable executionResults, + string? correlationId, + IDictionary? input, + CancellationToken cancellationToken = default) + { + var httpContext = _httpContextAccessor.HttpContext; + + if (httpContext == null) + throw new Exception("Invalid use of this method, because there is no HTTP context"); + + // We must assume that the workflow executed in a different process (when e.g. using Proto.Actor) + // and check if we received any `HttpEndpoint` or `WriteHttpResponse` activity bookmarks. + // If we did, acquire a lock on the workflow instance and resume it from here within an actual HTTP context so that the activity can complete its HTTP response. + var httpEndpointTypeName = ActivityTypeNameHelper.GenerateTypeName(); + var writeHttpResponseTypeName = ActivityTypeNameHelper.GenerateTypeName(); + + var query = + from executionResult in executionResults + from bookmark in executionResult.Bookmarks + where bookmark.Name == writeHttpResponseTypeName || bookmark.Name == httpEndpointTypeName + select (executionResult.InstanceId, bookmark.Id); + + var workflowExecutionResults = new Stack<(string InstanceId, string BookmarkId)>(query); + + while (workflowExecutionResults.TryPop(out var result)) + { + // Resume the workflow "in-process". + var workflowState = await _workflowRuntime.ExportWorkflowStateAsync( + result.InstanceId, + cancellationToken); + + if (workflowState == null) + { + // TODO: log this, shouldn't normally happen. + continue; + } + + var workflowDefinition = await _workflowDefinitionService.FindAsync( + workflowState.DefinitionId, + VersionOptions.SpecificVersion(workflowState.DefinitionVersion), + cancellationToken); + + if (workflowDefinition == null) + { + // TODO: Log this, shouldn't normally happen. + continue; + } + + var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync( + workflowDefinition, + cancellationToken); + + var workflowHost = await _workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken); + var options = new ResumeWorkflowHostOptions(correlationId, result.BookmarkId, Input: input); + await workflowHost.ResumeWorkflowAsync(options, cancellationToken); + + // Import the updated workflow state into the runtime. + await _workflowRuntime.ImportWorkflowStateAsync(workflowState, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs b/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs index d07873e92..5cd6ccb02 100644 --- a/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs +++ b/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs @@ -22,6 +22,7 @@ public class WorkflowsMiddleware private readonly IWorkflowRuntime _workflowRuntime; private readonly IWorkflowHostFactory _workflowHostFactory; private readonly IWorkflowDefinitionService _workflowDefinitionService; + private readonly IHttpBookmarkProcessor _httpBookmarkProcessor; private readonly HttpActivityOptions _options; private readonly string _activityTypeName = ActivityTypeNameHelper.GenerateTypeName(); @@ -33,12 +34,14 @@ public class WorkflowsMiddleware IWorkflowRuntime workflowRuntime, IWorkflowHostFactory workflowHostFactory, IWorkflowDefinitionService workflowDefinitionService, + IHttpBookmarkProcessor httpBookmarkProcessor, IOptions options) { _next = next; _workflowRuntime = workflowRuntime; _workflowHostFactory = workflowHostFactory; _workflowDefinitionService = workflowDefinitionService; + _httpBookmarkProcessor = httpBookmarkProcessor; _options = options.Value; } @@ -79,61 +82,10 @@ public class WorkflowsMiddleware var cancellationToken = httpContext.RequestAborted; // Trigger the workflow. - var triggerResult = await _workflowRuntime.TriggerWorkflowsAsync( - _activityTypeName, - bookmarkPayload, - triggerOptions, - cancellationToken); - - // We must assume that the workflow executed in a different process (when e.g. using Proto.Actor) - // and check if we received any `HttpEndpoint` or `WriteHttpResponse` activity bookmarks. - // If we did, acquire a lock on the workflow instance and resume it from here within an actual HTTP context so that the activity can complete its HTTP response. - var httpEndpointTypeName = ActivityTypeNameHelper.GenerateTypeName(); - var writeHttpResponseTypeName = ActivityTypeNameHelper.GenerateTypeName(); - - var query = - from triggeredWorkflow in triggerResult.TriggeredWorkflows - from bookmark in triggeredWorkflow.Bookmarks - where bookmark.Name == writeHttpResponseTypeName || bookmark.Name == httpEndpointTypeName - select (triggeredWorkflow.InstanceId, bookmark.Id); - - var workflowExecutionResults = new Stack<(string InstanceId, string BookmarkId)>(query); - - while (workflowExecutionResults.TryPop(out var result)) - { - // Resume the workflow "in-process". - var workflowState = await _workflowRuntime.ExportWorkflowStateAsync( - result.InstanceId, - cancellationToken); - - if (workflowState == null) - { - // TODO: log this, shouldn't normally happen. - continue; - } - - var workflowDefinition = await _workflowDefinitionService.FindAsync( - workflowState.DefinitionId, - VersionOptions.SpecificVersion(workflowState.DefinitionVersion), - cancellationToken); - - if (workflowDefinition == null) - { - // TODO: Log this, shouldn't normally happen. - continue; - } - - var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync( - workflowDefinition, - cancellationToken); - - var workflowHost = await _workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken); - var options = new ResumeWorkflowHostOptions(correlationId, result.BookmarkId, Input: input); - await workflowHost.ResumeWorkflowAsync(options, cancellationToken); - - // Import the updated workflow state into the runtime. - await _workflowRuntime.ImportWorkflowStateAsync(workflowState, cancellationToken); - } + var triggerResult = await _workflowRuntime.TriggerWorkflowsAsync(_activityTypeName, bookmarkPayload, triggerOptions, cancellationToken); + + // Process the trigger result by resuming each HTTP bookmark, if any. + await _httpBookmarkProcessor.ProcessBookmarks(triggerResult.TriggeredWorkflows, correlationId, input, cancellationToken); } private static async Task WriteResponseAsync(HttpContext httpContext, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Http/Services/IHttpBookmarkProcessor.cs b/src/modules/Elsa.Http/Services/IHttpBookmarkProcessor.cs new file mode 100644 index 000000000..09ff66197 --- /dev/null +++ b/src/modules/Elsa.Http/Services/IHttpBookmarkProcessor.cs @@ -0,0 +1,18 @@ +using Elsa.Workflows.Runtime.Services; + +namespace Elsa.Http.Services; + +/// +/// A helper service that can process s within the current HTTP context. +/// +public interface IHttpBookmarkProcessor +{ + /// + /// Processes the specified by resuming each HTTP bookmark while we are in an HTTP context. + /// + Task ProcessBookmarks( + IEnumerable executionResults, + string? correlationId, + IDictionary? input, + CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs index e9ff4d618..933a42e96 100644 --- a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs @@ -66,7 +66,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime } /// - public async Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) + public async Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) { var versionOptions = options.VersionOptions; var correlationId = options.CorrelationId; @@ -87,7 +87,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime var response = await client.Start(request, cancellationToken); var bookmarks = Map(response!.Bookmarks).ToList(); - return new StartWorkflowResult(workflowInstanceId, bookmarks); + return new WorkflowExecutionResult(workflowInstanceId, bookmarks); } /// @@ -110,7 +110,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime } /// - public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) + public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) { var hash = _hasher.Hash(activityTypeName, bookmarkPayload); var client = _cluster.GetNamedBookmarkGrain(hash); @@ -129,7 +129,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime /// public async Task TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) { - var triggeredWorkflows = new List(); + var triggeredWorkflows = new List(); var hash = _hasher.Hash(activityTypeName, bookmarkPayload); // Start new workflows. @@ -147,7 +147,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken); - triggeredWorkflows.Add(new TriggeredWorkflow(startResult.InstanceId, startResult.Bookmarks)); + triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks)); } // Resume existing workflow instances. @@ -171,7 +171,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime new ResumeWorkflowRuntimeOptions(options.CorrelationId, bookmark.BookmarkId, null, options.Input), cancellationToken); - triggeredWorkflows.Add(new TriggeredWorkflow(workflowInstanceId, resumeResult.Bookmarks)); + triggeredWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks)); } return new TriggerWorkflowsResult(triggeredWorkflows); @@ -229,9 +229,9 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime return response!.Count; } - private async Task> ResumeWorkflowsAsync(IEnumerable bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default) + private async Task> ResumeWorkflowsAsync(IEnumerable bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default) { - var resumedWorkflows = new List(); + var resumedWorkflows = new List(); foreach (var bookmark in bookmarks) { @@ -242,7 +242,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime runtimeOptions with { BookmarkId = bookmark.BookmarkId }, cancellationToken); - resumedWorkflows.Add(new ResumedWorkflow(workflowInstanceId, resumeResult.Bookmarks)); + resumedWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks)); } return resumedWorkflows; diff --git a/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj b/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj index bfe55e6a1..e3c2bb372 100644 --- a/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj +++ b/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj @@ -14,13 +14,10 @@ + - - - - diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Execute/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Execute/Endpoint.cs index 42baf6556..791a6eada 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Execute/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Execute/Endpoint.cs @@ -1,5 +1,6 @@ using Elsa.Abstractions; using Elsa.Common.Models; +using Elsa.Http.Services; using Elsa.Workflows.Management.Services; using Elsa.Workflows.Runtime.Services; using JetBrains.Annotations; @@ -14,12 +15,14 @@ public class Execute : ElsaEndpoint { private readonly IWorkflowDefinitionStore _store; private readonly IWorkflowRuntime _workflowRuntime; + private readonly IHttpBookmarkProcessor _httpBookmarkProcessor; /// - public Execute(IWorkflowDefinitionStore store, IWorkflowRuntime workflowRuntime) + public Execute(IWorkflowDefinitionStore store, IWorkflowRuntime workflowRuntime, IHttpBookmarkProcessor httpBookmarkProcessor) { _store = store; _workflowRuntime = workflowRuntime; + _httpBookmarkProcessor = httpBookmarkProcessor; } /// @@ -45,6 +48,9 @@ public class Execute : ElsaEndpoint var startWorkflowOptions = new StartWorkflowRuntimeOptions(correlationId, VersionOptions: VersionOptions.Published); var result = await _workflowRuntime.StartWorkflowAsync(definitionId, startWorkflowOptions, cancellationToken); + // Resume any HTTP bookmarks. + await _httpBookmarkProcessor.ProcessBookmarks(new[] { result }, correlationId, default, cancellationToken); + if (!HttpContext.Response.HasStarted) await SendOkAsync(new Response(result.InstanceId), cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Import/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Import/Endpoint.cs index 9875ab927..d65b1ddf0 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Import/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Import/Endpoint.cs @@ -38,7 +38,7 @@ internal class Import : ElsaEndpoint - public async Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) + public async Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) { var input = options.Input; var correlationId = options.CorrelationId; @@ -65,7 +65,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime await SaveWorkflowStateAsync(workflowState, cancellationToken); - return new StartWorkflowResult(workflowState.Id, workflowState.Bookmarks); + return new WorkflowExecutionResult(workflowState.Id, workflowState.Bookmarks); } /// @@ -104,7 +104,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime } /// - public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) + public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default) { var hash = _hasher.Hash(activityTypeName, bookmarkPayload); var correlationId = options.CorrelationId; @@ -119,7 +119,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) { - var triggeredWorkflows = new List(); + var triggeredWorkflows = new List(); var hash = _hasher.Hash(activityTypeName, bookmarkPayload); // Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances being created. @@ -139,7 +139,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime continue; var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken); - triggeredWorkflows.Add(new TriggeredWorkflow(startResult.InstanceId, startResult.Bookmarks)); + triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks)); } } @@ -148,7 +148,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime var bookmarks = (string.IsNullOrEmpty(correlationId) ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken)).ToList(); var resumedWorkflows = await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken); - triggeredWorkflows.AddRange(resumedWorkflows.Select(x => new TriggeredWorkflow(x.InstanceId, x.Bookmarks))); + triggeredWorkflows.AddRange(resumedWorkflows.Select(x => new WorkflowExecutionResult(x.InstanceId, x.Bookmarks))); return new TriggerWorkflowsResult(triggeredWorkflows); } @@ -180,9 +180,9 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime return await _workflowHostFactory.CreateAsync(workflow, cancellationToken); } - private async Task> ResumeWorkflowsAsync(IEnumerable bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default) + private async Task> ResumeWorkflowsAsync(IEnumerable bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default) { - var resumedWorkflows = new List(); + var resumedWorkflows = new List(); foreach (var bookmark in bookmarks) { @@ -190,7 +190,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime var resumeOptions = new ResumeWorkflowRuntimeOptions(runtimeOptions.CorrelationId, bookmark.BookmarkId, Input: runtimeOptions.Input); var resumeResult = await ResumeWorkflowAsync(workflowInstanceId, resumeOptions, cancellationToken); - resumedWorkflows.Add(new ResumedWorkflow(workflowInstanceId, resumeResult.Bookmarks)); + resumedWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks)); } return resumedWorkflows; diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs index fc898c052..d38e7e462 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs @@ -21,7 +21,7 @@ public interface IWorkflowRuntime /// The workflow definition ID to run. /// /// - Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default); + Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default); /// /// Resumes an existing workflow instance. @@ -34,7 +34,7 @@ public interface IWorkflowRuntime /// /// Resumes all workflows that are bookmarked on the specified activity type. /// - Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default); + Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default); /// /// Starts all workflows and resumes existing workflow instances based on the specified activity type and bookmark payload. @@ -65,12 +65,10 @@ public interface IWorkflowRuntime public record StartWorkflowRuntimeOptions(string? CorrelationId = default, IDictionary? Input = default, VersionOptions VersionOptions = default, string? TriggerActivityId = default); public record ResumeWorkflowRuntimeOptions(string? CorrelationId = default, string? BookmarkId = default, string? ActivityId = default, IDictionary? Input = default); public record CanStartWorkflowResult(string? InstanceId, bool CanStart); -public record StartWorkflowResult(string InstanceId, ICollection Bookmarks); public record ResumeWorkflowResult(ICollection Bookmarks); public record TriggerWorkflowsRuntimeOptions(string? CorrelationId = default, IDictionary? Input = default); -public record TriggerWorkflowsResult(ICollection TriggeredWorkflows); -public record ResumedWorkflow(string InstanceId, ICollection Bookmarks); -public record TriggeredWorkflow(string InstanceId, ICollection Bookmarks); +public record TriggerWorkflowsResult(ICollection TriggeredWorkflows); +public record WorkflowExecutionResult(string InstanceId, ICollection Bookmarks); public record UpdateBookmarksContext(string InstanceId, Diff Diff, string? CorrelationId); ///