diff --git a/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs b/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs index f1645737e..23393f3df 100644 --- a/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs +++ b/src/modules/Elsa.Http/Middleware/WorkflowsMiddleware.cs @@ -98,39 +98,37 @@ public class WorkflowsMiddleware [HttpEndpoint.RequestPathInputKey] = path }; - // TODO: Get correlation ID from query string or header etc. + // TODO: Get correlation ID from query string or header. var correlationId = default(string); + var request = httpContext.Request; var method = request.Method!.ToLowerInvariant(); - var bookmarkPayload = new HttpEndpointBookmarkPayload(matchingPath, method); var triggerOptions = new TriggerWorkflowsRuntimeOptions(correlationId, input); var cancellationToken = httpContext.RequestAborted; + var workflowsFilter = new WorkflowsFilter(_activityTypeName, bookmarkPayload, triggerOptions); + var workflowMatches = (await _workflowRuntime.FindWorkflowsAsync(workflowsFilter, cancellationToken)).ToList(); - var workflowsQuery = new WorkflowsQuery(_activityTypeName, bookmarkPayload, triggerOptions); - var pendingWorkflows = await _workflowRuntime.FindWorkflowsAsync(workflowsQuery, cancellationToken); - - if (await HandleNoWorkflowsFoundAsync(httpContext, pendingWorkflows, basePath)) + if (await HandleNoWorkflowsFoundAsync(httpContext, workflowMatches, basePath)) return; - if (await HandleMultipleWorkflowsFoundAsync(httpContext, pendingWorkflows, cancellationToken)) + if (await HandleMultipleWorkflowsFoundAsync(httpContext, workflowMatches, cancellationToken)) return; - if (await HandleWorkflowFaultAsync(httpContext, pendingWorkflows.Single(), cancellationToken)) + if (await HandleWorkflowFaultAsync(httpContext, workflowMatches.Single(), cancellationToken)) return; - if (await AuthorizeAsync(httpContext, pendingWorkflows.Single(), bookmarkPayload, cancellationToken)) + if (await AuthorizeAsync(httpContext, workflowMatches.Single(), bookmarkPayload, cancellationToken)) return; - var executionResult = await _workflowRuntime.ExecutePendingWorkflowAsync(pendingWorkflows.Single(), input, cancellationToken); + var executionResult = await _workflowRuntime.ExecuteWorkflowAsync(workflowMatches.Single(), input, cancellationToken); // Process the trigger result by resuming each HTTP bookmark, if any. await _httpBookmarkProcessor.ProcessBookmarks(new List { executionResult }, correlationId, input, cancellationToken); } - private string? GetMatchingRoute(string? path) + private string GetMatchingRoute(string path) { - var matchingRouteQuery = from route in _routeTable let routeValues = _routeMatcher.Match(route, path) @@ -164,9 +162,9 @@ public class WorkflowsMiddleware private string GetPath(HttpContext httpContext) => httpContext.Request.Path.Value.ToLowerInvariant(); - private async Task HandleNoWorkflowsFoundAsync(HttpContext httpContext, IEnumerable pendingWorkflows, PathString? basePath) + private async Task HandleNoWorkflowsFoundAsync(HttpContext httpContext, ICollection workflowMatches, PathString? basePath) { - if (pendingWorkflows.Any()) + if (workflowMatches.Any()) return false; // If a base path was configured, we are sure the requester tried to execute a workflow that doesn't exist. @@ -183,9 +181,9 @@ public class WorkflowsMiddleware return true; } - private async Task HandleMultipleWorkflowsFoundAsync(HttpContext httpContext, IEnumerable pendingWorkflows, CancellationToken cancellationToken) + private async Task HandleMultipleWorkflowsFoundAsync(HttpContext httpContext, ICollection workflowMatches, CancellationToken cancellationToken) { - if (pendingWorkflows.ToList().Count <= 1) + if (workflowMatches.Count <= 1) return false; httpContext.Response.ContentType = "application/json"; @@ -194,16 +192,16 @@ public class WorkflowsMiddleware var responseContent = JsonSerializer.Serialize(new { errorMessage = "The call is ambiguous and matches multiple workflows.", - workflows = pendingWorkflows + workflows = workflowMatches }); await httpContext.Response.WriteAsync(responseContent, cancellationToken); return true; } - private async Task HandleWorkflowFaultAsync(HttpContext httpContext, CollectedWorkflow pendingWorkflow, CancellationToken cancellationToken) + private async Task HandleWorkflowFaultAsync(HttpContext httpContext, WorkflowMatch workflowMatch, CancellationToken cancellationToken) { - var instanceFilter = new WorkflowInstanceFilter { Id = pendingWorkflow.WorkflowInstanceId }; + var instanceFilter = new WorkflowInstanceFilter { Id = workflowMatch.WorkflowInstanceId }; var workflowInstance = await _workflowInstanceStore.FindAsync(instanceFilter, cancellationToken); if (workflowInstance is not null @@ -219,38 +217,36 @@ public class WorkflowsMiddleware private async Task AuthorizeAsync( HttpContext httpContext, - CollectedWorkflow pendingWorkflow, + WorkflowMatch pendingWorkflowMatch, HttpEndpointBookmarkPayload bookmarkPayload, - CancellationToken cancellationToken = default) + CancellationToken cancellationToken) { - var hash = _hasher.Hash(_activityTypeName, bookmarkPayload); - var payload = default(HttpEndpointBookmarkPayload); - - if (pendingWorkflow is CollectedStartableWorkflow) - { - var triggerFilter = new TriggerFilter() { Hash = hash }; - var triggers = (await _triggerStore.FindManyAsync(triggerFilter, cancellationToken)) - .Select(x => _serializer.Deserialize(x.Data!)).ToList(); - payload = triggers.Single(); - } - else - { - var bookmarkFilter = new BookmarkFilter() { Hash = hash }; - var bookmarks = (await _bookmarkStore.FindManyAsync(bookmarkFilter, cancellationToken)) - .Select(x => _serializer.Deserialize(x.Data!)).ToList(); - payload = bookmarks.Single(); - } - + var payload = await GetBookmarkPayloadAsync(pendingWorkflowMatch, bookmarkPayload, cancellationToken); + if (!(payload.Authorize ?? false)) return false; - var authorized = await _httpEndpointAuthorizationHandler.AuthorizeAsync(new AuthorizeHttpEndpointContext(httpContext, pendingWorkflow.WorkflowInstanceId, payload.Policy)); + var authorized = await _httpEndpointAuthorizationHandler.AuthorizeAsync(new AuthorizeHttpEndpointContext(httpContext, pendingWorkflowMatch.WorkflowInstanceId, payload.Policy)); - if (!authorized) - { + if (!authorized) httpContext.Response.StatusCode = (int)HttpStatusCode.Unauthorized; - } return !authorized; } + + private async Task GetBookmarkPayloadAsync(WorkflowMatch workflowMatch,HttpEndpointBookmarkPayload bookmarkPayload, CancellationToken cancellationToken) + { + var hash = _hasher.Hash(_activityTypeName, bookmarkPayload); + + if (workflowMatch is StartableWorkflowMatch) + { + var triggerFilter = new TriggerFilter { Hash = hash }; + var trigger = (await _triggerStore.FindManyAsync(triggerFilter, cancellationToken)).First(); + return _serializer.Deserialize(trigger.Data!); + } + + var bookmarkFilter = new BookmarkFilter { Hash = hash }; + var bookmark = (await _bookmarkStore.FindManyAsync(bookmarkFilter, cancellationToken)).First(); + return _serializer.Deserialize(bookmark.Data!); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Http/Services/RouteMatcher.cs b/src/modules/Elsa.Http/Services/RouteMatcher.cs index 2a029bf1c..eb24ce2a4 100644 --- a/src/modules/Elsa.Http/Services/RouteMatcher.cs +++ b/src/modules/Elsa.Http/Services/RouteMatcher.cs @@ -6,8 +6,10 @@ using Microsoft.AspNetCore.Routing.Template; namespace Elsa.Http.Services; +/// public class RouteMatcher : IRouteMatcher { + /// public RouteValueDictionary? Match(string routeTemplate, string requestPath) { var template = TemplateParser.Parse(routeTemplate); diff --git a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs index 9adafd3da..d86908219 100644 --- a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs @@ -165,9 +165,9 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime } /// - public async Task ExecutePendingWorkflowAsync(CollectedWorkflow collectedWorkflow, IDictionary? input = default, CancellationToken cancellationToken = default) + public async Task ExecuteWorkflowAsync(WorkflowMatch match, IDictionary? input = default, CancellationToken cancellationToken = default) { - if (collectedWorkflow is CollectedStartableWorkflow collectedStartableWorkflow) + if (match is StartableWorkflowMatch collectedStartableWorkflow) { var startOptions = new StartWorkflowRuntimeOptions(collectedStartableWorkflow.CorrelationId, input, VersionOptions.Published, collectedStartableWorkflow.ActivityId, collectedStartableWorkflow.WorkflowInstanceId); @@ -176,10 +176,10 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime } else { - var collectedResumableWorkflow = (collectedWorkflow as CollectedResumableWorkflow)!; + var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!; var runtimeOptions = new ResumeWorkflowRuntimeOptions(collectedResumableWorkflow.CorrelationId, Input: input); var resumeResult = await ResumeWorkflowAsync( - collectedWorkflow.WorkflowInstanceId, + match.WorkflowInstanceId, runtimeOptions with { BookmarkId = collectedResumableWorkflow.BookmarkId }, cancellationToken); @@ -188,10 +188,10 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime } /// - public async Task> FindWorkflowsAsync(WorkflowsQuery workflowsQuery, CancellationToken cancellationToken = default) + public async Task> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default) { - var startableWorkflows = await CollectStartableWorkflowsAsync(workflowsQuery, cancellationToken); - var resumableWorkflows = await CollectResumableWorkflowsAsync(workflowsQuery, cancellationToken); + var startableWorkflows = await CollectStartableWorkflowsAsync(filter, cancellationToken); + var resumableWorkflows = await CollectResumableWorkflowsAsync(filter, cancellationToken); var results = startableWorkflows.Concat(resumableWorkflows).ToList(); return results; } @@ -312,45 +312,45 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime x.AutoBurn, x.CallbackMethodName.NullIfEmpty())); - private async Task> CollectStartableWorkflowsAsync(WorkflowsQuery workflowsQuery, CancellationToken cancellationToken) + private async Task> CollectStartableWorkflowsAsync(WorkflowsFilter workflowsFilter, CancellationToken cancellationToken) { - var hash = _hasher.Hash(workflowsQuery.ActivityTypeName, workflowsQuery.BookmarkPayload); + var hash = _hasher.Hash(workflowsFilter.ActivityTypeName, workflowsFilter.BookmarkPayload); var filter = new TriggerFilter { Hash = hash }; var triggers = await _triggerStore.FindManyAsync(filter, cancellationToken); - var results = new List(); + var results = new List(); foreach (var trigger in triggers) { var definitionId = trigger.WorkflowDefinitionId; - var startOptions = new StartWorkflowRuntimeOptions(workflowsQuery.Options.CorrelationId, workflowsQuery.Options.Input, VersionOptions.Published, trigger.ActivityId); + var startOptions = new StartWorkflowRuntimeOptions(workflowsFilter.Options.CorrelationId, workflowsFilter.Options.Input, VersionOptions.Published, trigger.ActivityId); var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken); - var workflowInstance = await _workflowInstanceFactory.CreateAsync(definitionId, workflowsQuery.Options.CorrelationId, cancellationToken); + var workflowInstance = await _workflowInstanceFactory.CreateAsync(definitionId, workflowsFilter.Options.CorrelationId, cancellationToken); if (canStartResult.CanStart) { - results.Add(new CollectedStartableWorkflow(workflowInstance.Id, workflowInstance, workflowsQuery.Options.CorrelationId, trigger.ActivityId, definitionId)); + results.Add(new StartableWorkflowMatch(workflowInstance.Id, workflowInstance, workflowsFilter.Options.CorrelationId, trigger.ActivityId, definitionId)); } } return results; } - private async Task> CollectResumableWorkflowsAsync(WorkflowsQuery workflowsQuery, CancellationToken cancellationToken) + private async Task> CollectResumableWorkflowsAsync(WorkflowsFilter workflowsFilter, CancellationToken cancellationToken) { - var hash = _hasher.Hash(workflowsQuery.ActivityTypeName, workflowsQuery.BookmarkPayload); + var hash = _hasher.Hash(workflowsFilter.ActivityTypeName, workflowsFilter.BookmarkPayload); var client = _cluster.GetNamedBookmarkGrain(hash); var request = new ResolveBookmarksRequest { - ActivityTypeName = workflowsQuery.ActivityTypeName, - CorrelationId = workflowsQuery.Options.CorrelationId.EmptyIfNull() + ActivityTypeName = workflowsFilter.ActivityTypeName, + CorrelationId = workflowsFilter.Options.CorrelationId.EmptyIfNull() }; var bookmarksResponse = await client.Resolve(request, cancellationToken); var bookmarks = bookmarksResponse!.Bookmarks; - var collectedWorkflows = bookmarks.Select(b => new CollectedResumableWorkflow(b.WorkflowInstanceId, default, workflowsQuery.Options.CorrelationId, b.BookmarkId)).ToList(); + var collectedWorkflows = bookmarks.Select(b => new ResumableWorkflowMatch(b.WorkflowInstanceId, default, workflowsFilter.Options.CorrelationId, b.BookmarkId)).ToList(); return collectedWorkflows; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs index 2e49a9c1c..ea36003c7 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs @@ -55,19 +55,19 @@ public interface IWorkflowRuntime /// /// Executes a pending workflow. /// - /// + /// A workflow match to execute. /// /// /// - Task ExecutePendingWorkflowAsync(CollectedWorkflow collectedWorkflow, IDictionary? input = default, CancellationToken cancellationToken = default); + Task ExecuteWorkflowAsync(WorkflowMatch match, IDictionary? input = default, CancellationToken cancellationToken = default); /// /// Finds all the workflows that can be started or resumed based on a query model. /// - /// + /// /// /// - Task> FindWorkflowsAsync(WorkflowsQuery workflowsQuery, CancellationToken cancellationToken = default); + Task> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default); /// /// Exports the of the specified workflow instance. @@ -98,12 +98,12 @@ public record TriggerWorkflowsRuntimeOptions(string? CorrelationId = default, ID public record TriggerWorkflowsResult(ICollection TriggeredWorkflows); public record WorkflowExecutionResult(string InstanceId, ICollection Bookmarks, string? ActivityId = null); public record UpdateBookmarksContext(string InstanceId, Diff Diff, string? CorrelationId); -public record WorkflowsQuery(string ActivityTypeName, object BookmarkPayload, TriggerWorkflowsRuntimeOptions Options); -public record CollectedWorkflow(string WorkflowInstanceId, WorkflowInstance? WorkflowInstance, string? CorrelationId); -public record CollectedStartableWorkflow(string WorkflowInstanceId, WorkflowInstance? WorkflowInstance, string? CorrelationId, string? ActivityId, string? DefinitionId) - : CollectedWorkflow(WorkflowInstanceId, WorkflowInstance, CorrelationId); -public record CollectedResumableWorkflow(string WorkflowInstanceId, WorkflowInstance? WorkflowInstance, string? CorrelationId, string? BookmarkId) - : CollectedWorkflow(WorkflowInstanceId, WorkflowInstance, CorrelationId); +public record WorkflowsFilter(string ActivityTypeName, object BookmarkPayload, TriggerWorkflowsRuntimeOptions Options); +public record WorkflowMatch(string WorkflowInstanceId, WorkflowInstance? WorkflowInstance, string? CorrelationId); +public record StartableWorkflowMatch(string WorkflowInstanceId, WorkflowInstance? WorkflowInstance, string? CorrelationId, string? ActivityId, string? DefinitionId) + : WorkflowMatch(WorkflowInstanceId, WorkflowInstance, CorrelationId); +public record ResumableWorkflowMatch(string WorkflowInstanceId, WorkflowInstance? WorkflowInstance, string? CorrelationId, string? BookmarkId) + : WorkflowMatch(WorkflowInstanceId, WorkflowInstance, CorrelationId); /// /// Contains arguments to use for counting the number of workflow instances. diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs index 49abc6583..4bcab8a56 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs @@ -160,9 +160,9 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime return new TriggerWorkflowsResult(results); } - public async Task ExecutePendingWorkflowAsync(CollectedWorkflow collectedWorkflow, IDictionary? input = default, CancellationToken cancellationToken = default) + public async Task ExecuteWorkflowAsync(WorkflowMatch match, IDictionary? input = default, CancellationToken cancellationToken = default) { - if (collectedWorkflow is CollectedStartableWorkflow collectedStartableWorkflow) + if (match is StartableWorkflowMatch collectedStartableWorkflow) { var startOptions = new StartWorkflowRuntimeOptions(collectedStartableWorkflow.CorrelationId, input, VersionOptions.Published, collectedStartableWorkflow.ActivityId, collectedStartableWorkflow.WorkflowInstanceId); @@ -171,10 +171,10 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime } else { - var collectedResumableWorkflow = (collectedWorkflow as CollectedResumableWorkflow)!; + var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!; var runtimeOptions = new ResumeWorkflowRuntimeOptions(collectedResumableWorkflow.CorrelationId, Input: input); var resumeResult = await ResumeWorkflowAsync( - collectedWorkflow.WorkflowInstanceId, + match.WorkflowInstanceId, runtimeOptions with { BookmarkId = collectedResumableWorkflow.BookmarkId }, cancellationToken); @@ -182,10 +182,10 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime } } - public async Task> FindWorkflowsAsync(WorkflowsQuery workflowsQuery, CancellationToken cancellationToken = default) + public async Task> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default) { - var startableWorkflows = await CollectStartableWorkflowsAsync(workflowsQuery, cancellationToken); - var resumableWorkflows = await CollectResumableWorkflowsAsync(workflowsQuery, cancellationToken); + var startableWorkflows = await CollectStartableWorkflowsAsync(filter, cancellationToken); + var resumableWorkflows = await CollectResumableWorkflowsAsync(filter, cancellationToken); var results = startableWorkflows.Concat(resumableWorkflows).ToList(); return results; } @@ -258,12 +258,12 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime } } - private async Task> CollectStartableWorkflowsAsync( - WorkflowsQuery workflowsQuery, + private async Task> CollectStartableWorkflowsAsync( + WorkflowsFilter workflowsFilter, CancellationToken cancellationToken = default) { - var results = new List(); - var hash = _hasher.Hash(workflowsQuery.ActivityTypeName, workflowsQuery.BookmarkPayload); + var results = new List(); + var hash = _hasher.Hash(workflowsFilter.ActivityTypeName, workflowsFilter.BookmarkPayload); // Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances from being created. var sharedResource = $"{nameof(DefaultWorkflowRuntime)}__StartTriggeredWorkflows__{hash}"; @@ -275,14 +275,14 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime foreach (var trigger in triggers) { var definitionId = trigger.WorkflowDefinitionId; - var startOptions = new StartWorkflowRuntimeOptions(workflowsQuery.Options.CorrelationId, workflowsQuery.Options.Input, VersionOptions.Published, trigger.ActivityId); + var startOptions = new StartWorkflowRuntimeOptions(workflowsFilter.Options.CorrelationId, workflowsFilter.Options.Input, VersionOptions.Published, trigger.ActivityId); var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken); - var workflowInstance = await _workflowInstanceFactory.CreateAsync(definitionId, workflowsQuery.Options.CorrelationId, cancellationToken); + var workflowInstance = await _workflowInstanceFactory.CreateAsync(definitionId, workflowsFilter.Options.CorrelationId, cancellationToken); if (canStartResult.CanStart) { - results.Add(new CollectedStartableWorkflow(workflowInstance.Id, workflowInstance, workflowsQuery.Options.CorrelationId, trigger.ActivityId, definitionId)); + results.Add(new StartableWorkflowMatch(workflowInstance.Id, workflowInstance, workflowsFilter.Options.CorrelationId, trigger.ActivityId, definitionId)); } } } @@ -290,14 +290,14 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime return results; } - private async Task> CollectResumableWorkflowsAsync(WorkflowsQuery workflowsQuery, CancellationToken cancellationToken = default) + private async Task> CollectResumableWorkflowsAsync(WorkflowsFilter workflowsFilter, CancellationToken cancellationToken = default) { - var hash = _hasher.Hash(workflowsQuery.ActivityTypeName, workflowsQuery.BookmarkPayload); - var correlationId = workflowsQuery.Options.CorrelationId; + var hash = _hasher.Hash(workflowsFilter.ActivityTypeName, workflowsFilter.BookmarkPayload); + var correlationId = workflowsFilter.Options.CorrelationId; var filter = new BookmarkFilter { Hash = hash, CorrelationId = correlationId }; var bookmarks = await _bookmarkStore.FindManyAsync(filter, cancellationToken); - var collectedWorkflows = bookmarks.Select(b => new CollectedResumableWorkflow(b.WorkflowInstanceId, default, correlationId, b.BookmarkId)).ToList(); + var collectedWorkflows = bookmarks.Select(b => new ResumableWorkflowMatch(b.WorkflowInstanceId, default, correlationId, b.BookmarkId)).ToList(); return collectedWorkflows; } } \ No newline at end of file