Make HTTP bookmark processing reusable
This commit is contained in:
parent
9c07d19903
commit
6dba654e9f
|
|
@ -27,8 +27,8 @@ public class HttpFeature : FeatureBase
|
|||
/// </summary>
|
||||
public Action<HttpActivityOptions>? ConfigureHttpOptions { get; set; }
|
||||
|
||||
public Func<IServiceProvider, IHttpEndpointAuthorizationHandler> HttpEndpointAuthorizationHandlerFactory { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance<AllowAnonymousHttpEndpointAuthorizationHandler>;
|
||||
public Func<IServiceProvider, IHttpEndpointWorkflowFaultHandler> HttpEndpointWorkflowFaultHandlerFactory { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance<DefaultHttpEndpointWorkflowFaultHandler>;
|
||||
public Func<IServiceProvider, IHttpEndpointAuthorizationHandler> HttpEndpointAuthorizationHandler { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance<AllowAnonymousHttpEndpointAuthorizationHandler>;
|
||||
public Func<IServiceProvider, IHttpEndpointWorkflowFaultHandler> HttpEndpointWorkflowFaultHandler { get; set; } = ActivatorUtilities.GetServiceOrCreateInstance<DefaultHttpEndpointWorkflowFaultHandler>;
|
||||
|
||||
/// <summary>
|
||||
/// A delegate to configure the <see cref="HttpClient"/> used when by the <see cref="SendHttpRequest"/> activity.
|
||||
|
|
@ -74,6 +74,7 @@ public class HttpFeature : FeatureBase
|
|||
.AddSingleton<IRouteMatcher, RouteMatcher>()
|
||||
.AddSingleton<IRouteTable, RouteTable>()
|
||||
.AddSingleton<IAbsoluteUrlProvider, DefaultAbsoluteUrlProvider>()
|
||||
.AddSingleton<IHttpBookmarkProcessor, HttpBookmarkProcessor>()
|
||||
.AddNotificationHandlersFrom<UpdateRouteTable>()
|
||||
.AddHttpContextAccessor()
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class HttpBookmarkProcessor : IHttpBookmarkProcessor
|
||||
{
|
||||
private readonly IWorkflowRuntime _workflowRuntime;
|
||||
private readonly IWorkflowDefinitionService _workflowDefinitionService;
|
||||
private readonly IWorkflowHostFactory _workflowHostFactory;
|
||||
private readonly IHttpContextAccessor _httpContextAccessor;
|
||||
|
||||
/// <summary>
|
||||
/// Constructor.
|
||||
/// </summary>
|
||||
public HttpBookmarkProcessor(
|
||||
IWorkflowRuntime workflowRuntime,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IWorkflowHostFactory workflowHostFactory,
|
||||
IHttpContextAccessor httpContextAccessor)
|
||||
{
|
||||
_workflowRuntime = workflowRuntime;
|
||||
_workflowDefinitionService = workflowDefinitionService;
|
||||
_workflowHostFactory = workflowHostFactory;
|
||||
_httpContextAccessor = httpContextAccessor;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task ProcessBookmarks(
|
||||
IEnumerable<WorkflowExecutionResult> executionResults,
|
||||
string? correlationId,
|
||||
IDictionary<string, object>? 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<HttpEndpoint>();
|
||||
var writeHttpResponseTypeName = ActivityTypeNameHelper.GenerateTypeName<WriteHttpResponse>();
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<HttpEndpoint>();
|
||||
|
||||
|
|
@ -33,12 +34,14 @@ public class WorkflowsMiddleware
|
|||
IWorkflowRuntime workflowRuntime,
|
||||
IWorkflowHostFactory workflowHostFactory,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IHttpBookmarkProcessor httpBookmarkProcessor,
|
||||
IOptions<HttpActivityOptions> 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<HttpEndpoint>();
|
||||
var writeHttpResponseTypeName = ActivityTypeNameHelper.GenerateTypeName<WriteHttpResponse>();
|
||||
|
||||
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)
|
||||
|
|
|
|||
18
src/modules/Elsa.Http/Services/IHttpBookmarkProcessor.cs
Normal file
18
src/modules/Elsa.Http/Services/IHttpBookmarkProcessor.cs
Normal file
|
|
@ -0,0 +1,18 @@
|
|||
using Elsa.Workflows.Runtime.Services;
|
||||
|
||||
namespace Elsa.Http.Services;
|
||||
|
||||
/// <summary>
|
||||
/// A helper service that can process <see cref="TriggerWorkflowsResult"/>s within the current HTTP context.
|
||||
/// </summary>
|
||||
public interface IHttpBookmarkProcessor
|
||||
{
|
||||
/// <summary>
|
||||
/// Processes the specified <see cref="executionResults"/> by resuming each HTTP bookmark while we are in an HTTP context.
|
||||
/// </summary>
|
||||
Task ProcessBookmarks(
|
||||
IEnumerable<WorkflowExecutionResult> executionResults,
|
||||
string? correlationId,
|
||||
IDictionary<string, object>? input,
|
||||
CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -66,7 +66,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
public async Task<WorkflowExecutionResult> 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);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -110,7 +110,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ICollection<ResumedWorkflow>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
public async Task<ICollection<WorkflowExecutionResult>> 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
|
|||
/// <inheritdoc />
|
||||
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var triggeredWorkflows = new List<TriggeredWorkflow>();
|
||||
var triggeredWorkflows = new List<WorkflowExecutionResult>();
|
||||
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<ICollection<ResumedWorkflow>> ResumeWorkflowsAsync(IEnumerable<StoredBookmark> bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default)
|
||||
private async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(IEnumerable<StoredBookmark> bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var resumedWorkflows = new List<ResumedWorkflow>();
|
||||
var resumedWorkflows = new List<WorkflowExecutionResult>();
|
||||
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -14,13 +14,10 @@
|
|||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\common\Elsa.Api.Common\Elsa.Api.Common.csproj" />
|
||||
<ProjectReference Include="..\Elsa.ActivityDefinitions\Elsa.ActivityDefinitions.csproj" />
|
||||
<ProjectReference Include="..\Elsa.Http\Elsa.Http.csproj" />
|
||||
<ProjectReference Include="..\Elsa.JavaScript\Elsa.JavaScript.csproj" />
|
||||
<ProjectReference Include="..\Elsa.Workflows.Management\Elsa.Workflows.Management.csproj" />
|
||||
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<Folder Include="Endpoints\Scripting" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
|
|||
|
|
@ -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<Request, Response>
|
|||
{
|
||||
private readonly IWorkflowDefinitionStore _store;
|
||||
private readonly IWorkflowRuntime _workflowRuntime;
|
||||
private readonly IHttpBookmarkProcessor _httpBookmarkProcessor;
|
||||
|
||||
/// <inheritdoc />
|
||||
public Execute(IWorkflowDefinitionStore store, IWorkflowRuntime workflowRuntime)
|
||||
public Execute(IWorkflowDefinitionStore store, IWorkflowRuntime workflowRuntime, IHttpBookmarkProcessor httpBookmarkProcessor)
|
||||
{
|
||||
_store = store;
|
||||
_workflowRuntime = workflowRuntime;
|
||||
_httpBookmarkProcessor = httpBookmarkProcessor;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -45,6 +48,9 @@ public class Execute : ElsaEndpoint<Request, Response>
|
|||
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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -38,7 +38,7 @@ internal class Import : ElsaEndpoint<WorkflowDefinitionRequest, WorkflowDefiniti
|
|||
public override void Configure()
|
||||
{
|
||||
Routes("workflow-definitions/import", "workflow-definitions/{definitionId}/import");
|
||||
Verbs(Http.POST, Http.PUT);
|
||||
Verbs(FastEndpoints.Http.POST, FastEndpoints.Http.PUT);
|
||||
ConfigurePermissions("write:workflow-definitions");
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -54,7 +54,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
public async Task<WorkflowExecutionResult> 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);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -104,7 +104,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ICollection<ResumedWorkflow>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
public async Task<ICollection<WorkflowExecutionResult>> 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<TriggeredWorkflow>();
|
||||
var triggeredWorkflows = new List<WorkflowExecutionResult>();
|
||||
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<ICollection<ResumedWorkflow>> ResumeWorkflowsAsync(IEnumerable<StoredBookmark> bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default)
|
||||
private async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(IEnumerable<StoredBookmark> bookmarks, ResumeWorkflowRuntimeOptions runtimeOptions, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var resumedWorkflows = new List<ResumedWorkflow>();
|
||||
var resumedWorkflows = new List<WorkflowExecutionResult>();
|
||||
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -21,7 +21,7 @@ public interface IWorkflowRuntime
|
|||
/// <param name="definitionId">The workflow definition ID to run.</param>
|
||||
/// <param name="options"></param>
|
||||
/// <param name="cancellationToken"></param>
|
||||
Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
|
||||
Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Resumes an existing workflow instance.
|
||||
|
|
@ -34,7 +34,7 @@ public interface IWorkflowRuntime
|
|||
/// <summary>
|
||||
/// Resumes all workflows that are bookmarked on the specified activity type.
|
||||
/// </summary>
|
||||
Task<ICollection<ResumedWorkflow>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
|
||||
Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// 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<string, object>? Input = default, VersionOptions VersionOptions = default, string? TriggerActivityId = default);
|
||||
public record ResumeWorkflowRuntimeOptions(string? CorrelationId = default, string? BookmarkId = default, string? ActivityId = default, IDictionary<string, object>? Input = default);
|
||||
public record CanStartWorkflowResult(string? InstanceId, bool CanStart);
|
||||
public record StartWorkflowResult(string InstanceId, ICollection<Bookmark> Bookmarks);
|
||||
public record ResumeWorkflowResult(ICollection<Bookmark> Bookmarks);
|
||||
public record TriggerWorkflowsRuntimeOptions(string? CorrelationId = default, IDictionary<string, object>? Input = default);
|
||||
public record TriggerWorkflowsResult(ICollection<TriggeredWorkflow> TriggeredWorkflows);
|
||||
public record ResumedWorkflow(string InstanceId, ICollection<Bookmark> Bookmarks);
|
||||
public record TriggeredWorkflow(string InstanceId, ICollection<Bookmark> Bookmarks);
|
||||
public record TriggerWorkflowsResult(ICollection<WorkflowExecutionResult> TriggeredWorkflows);
|
||||
public record WorkflowExecutionResult(string InstanceId, ICollection<Bookmark> Bookmarks);
|
||||
public record UpdateBookmarksContext(string InstanceId, Diff<Bookmark> Diff, string? CorrelationId);
|
||||
|
||||
/// <summary>
|
||||
|
|
|
|||
Loading…
Reference in a new issue