Send HTTP error when workflow faults (#3937)
* Implement workflow fault handling in HTTP context * Handle faulted workflows across process boundariesl
This commit is contained in:
parent
b53ef86128
commit
53dcd11f17
|
|
@ -62,13 +62,13 @@ public static class ServiceProviderExtensions
|
|||
while (bookmarks.TryPop(out var bookmark))
|
||||
{
|
||||
var resumeOptions = new ResumeWorkflowRuntimeOptions(BookmarkId: bookmark.Id);
|
||||
var resumeResult = await workflowRuntime.ResumeWorkflowAsync(result.InstanceId, resumeOptions);
|
||||
var resumeResult = await workflowRuntime.ResumeWorkflowAsync(result.WorkflowInstanceId, resumeOptions);
|
||||
|
||||
foreach (var newBookmark in resumeResult.Bookmarks)
|
||||
bookmarks.Push(newBookmark);
|
||||
}
|
||||
|
||||
// Return the workflow state.
|
||||
return (await workflowRuntime.ExportWorkflowStateAsync(result.InstanceId))!;
|
||||
return (await workflowRuntime.ExportWorkflowStateAsync(result.WorkflowInstanceId))!;
|
||||
}
|
||||
}
|
||||
|
|
@ -1,3 +1,4 @@
|
|||
using Elsa.Workflows.Core.State;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
|
||||
namespace Elsa.Http.Contracts;
|
||||
|
|
@ -10,7 +11,7 @@ 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(
|
||||
Task<IEnumerable<WorkflowState>> ProcessBookmarks(
|
||||
IEnumerable<WorkflowExecutionResult> executionResults,
|
||||
string? correlationId,
|
||||
IDictionary<string, object>? input,
|
||||
|
|
|
|||
|
|
@ -26,22 +26,16 @@ public class DefaultHttpEndpointWorkflowFaultHandler : IHttpEndpointWorkflowFaul
|
|||
public virtual async ValueTask HandleAsync(HttpEndpointFaultedWorkflowContext context)
|
||||
{
|
||||
var httpContext = context.HttpContext;
|
||||
var workflowInstance = context.WorkflowInstance;
|
||||
var fault = workflowInstance.WorkflowState.Fault!;
|
||||
var workflowState = context.WorkflowState;
|
||||
var fault = workflowState.Fault!;
|
||||
|
||||
httpContext.Response.ContentType = MediaTypeNames.Application.Json;
|
||||
httpContext.Response.StatusCode = StatusCodes.Status500InternalServerError;
|
||||
|
||||
var faultedResponse = _apiSerializer.Serialize(new
|
||||
{
|
||||
errorMessage = $"Workflow faulted at {workflowInstance.FaultedAt!} with error: {fault.Message}",
|
||||
exception = fault?.Exception,
|
||||
workflow = new
|
||||
{
|
||||
name = workflowInstance.Name,
|
||||
version = workflowInstance.Version,
|
||||
instanceId = workflowInstance.Id
|
||||
}
|
||||
errorMessage = $"Workflow faulted with error: {fault.Message}",
|
||||
workflowState = workflowState
|
||||
});
|
||||
|
||||
await httpContext.Response.WriteAsync(faultedResponse, context.CancellationToken);
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ using System.Net;
|
|||
using System.Net.Mime;
|
||||
using System.Text.Json;
|
||||
using Elsa.Extensions;
|
||||
using Elsa.Workflows.Core.State;
|
||||
|
||||
namespace Elsa.Http.Middleware;
|
||||
|
||||
|
|
@ -125,7 +126,13 @@ public class WorkflowsMiddleware
|
|||
return;
|
||||
|
||||
// Process the trigger result by resuming each HTTP bookmark, if any.
|
||||
await _httpBookmarkProcessor.ProcessBookmarks(new List<WorkflowExecutionResult> { executionResult }, correlationId, input, cancellationToken);
|
||||
var affectedWorkflowStates = await _httpBookmarkProcessor.ProcessBookmarks(new List<WorkflowExecutionResult> { executionResult }, correlationId, input, cancellationToken);
|
||||
|
||||
// Check if there were any errors.
|
||||
var faultedWorkflowState = affectedWorkflowStates.FirstOrDefault(x => x.SubStatus == WorkflowSubStatus.Faulted);
|
||||
|
||||
if (faultedWorkflowState != null)
|
||||
await HandleWorkflowFaultAsync(httpContext, faultedWorkflowState, cancellationToken);
|
||||
}
|
||||
|
||||
private string GetMatchingRoute(string path)
|
||||
|
|
@ -202,18 +209,19 @@ public class WorkflowsMiddleware
|
|||
|
||||
private async Task<bool> HandleWorkflowFaultAsync(HttpContext httpContext, WorkflowExecutionResult workflowExecutionResult, CancellationToken cancellationToken)
|
||||
{
|
||||
var instanceFilter = new WorkflowInstanceFilter { Id = workflowExecutionResult.InstanceId };
|
||||
var workflowInstance = await _workflowInstanceStore.FindAsync(instanceFilter, cancellationToken);
|
||||
var subStatus = workflowExecutionResult.SubStatus;
|
||||
|
||||
if (workflowInstance is not null
|
||||
&& workflowInstance.SubStatus == WorkflowSubStatus.Faulted
|
||||
&& !httpContext.Response.HasStarted)
|
||||
{
|
||||
await _httpEndpointWorkflowFaultHandler.HandleAsync(new HttpEndpointFaultedWorkflowContext(httpContext, workflowInstance, null, cancellationToken));
|
||||
return true;
|
||||
}
|
||||
if (subStatus != WorkflowSubStatus.Faulted || httpContext.Response.HasStarted)
|
||||
return false;
|
||||
|
||||
return false;
|
||||
var workflowState = (await _workflowRuntime.ExportWorkflowStateAsync(workflowExecutionResult.WorkflowInstanceId, cancellationToken))!;
|
||||
return await HandleWorkflowFaultAsync(httpContext, workflowState, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<bool> HandleWorkflowFaultAsync(HttpContext httpContext, WorkflowState workflowState, CancellationToken cancellationToken)
|
||||
{
|
||||
await _httpEndpointWorkflowFaultHandler.HandleAsync(new HttpEndpointFaultedWorkflowContext(httpContext, workflowState, cancellationToken));
|
||||
return true;
|
||||
}
|
||||
|
||||
private async Task<bool> AuthorizeAsync(
|
||||
|
|
|
|||
|
|
@ -1,6 +1,12 @@
|
|||
using Elsa.Workflows.Management.Entities;
|
||||
using Elsa.Workflows.Core.State;
|
||||
using Microsoft.AspNetCore.Http;
|
||||
|
||||
namespace Elsa.Http.Models;
|
||||
|
||||
public record HttpEndpointFaultedWorkflowContext(HttpContext HttpContext, WorkflowInstance WorkflowInstance, Exception? Exception, CancellationToken CancellationToken);
|
||||
/// <summary>
|
||||
/// Provides context about the faulted workflow.
|
||||
/// </summary>
|
||||
/// <param name="HttpContext">The HTTP context.</param>
|
||||
/// <param name="WorkflowState">The faulted workflow state.</param>
|
||||
/// <param name="CancellationToken">The cancellation token.</param>
|
||||
public record HttpEndpointFaultedWorkflowContext(HttpContext HttpContext, WorkflowState WorkflowState, CancellationToken CancellationToken);
|
||||
|
|
@ -1,6 +1,7 @@
|
|||
using Elsa.Common.Models;
|
||||
using Elsa.Http.Contracts;
|
||||
using Elsa.Workflows.Core.Helpers;
|
||||
using Elsa.Workflows.Core.State;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using Microsoft.AspNetCore.Http;
|
||||
|
||||
|
|
@ -30,7 +31,7 @@ public class HttpBookmarkProcessor : IHttpBookmarkProcessor
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task ProcessBookmarks(
|
||||
public async Task<IEnumerable<WorkflowState>> ProcessBookmarks(
|
||||
IEnumerable<WorkflowExecutionResult> executionResults,
|
||||
string? correlationId,
|
||||
IDictionary<string, object>? input,
|
||||
|
|
@ -51,9 +52,10 @@ public class HttpBookmarkProcessor : IHttpBookmarkProcessor
|
|||
from executionResult in executionResults
|
||||
from bookmark in executionResult.Bookmarks
|
||||
where bookmark.Name == writeHttpResponseTypeName || bookmark.Name == httpEndpointTypeName
|
||||
select (executionResult.InstanceId, bookmark.Id);
|
||||
select (InstanceId: executionResult.WorkflowInstanceId, bookmark.Id);
|
||||
|
||||
var workflowExecutionResults = new Stack<(string InstanceId, string BookmarkId)>(query);
|
||||
var workflowStates = new List<WorkflowState>();
|
||||
|
||||
while (workflowExecutionResults.TryPop(out var result))
|
||||
{
|
||||
|
|
@ -86,9 +88,13 @@ public class HttpBookmarkProcessor : IHttpBookmarkProcessor
|
|||
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.
|
||||
workflowState = workflowHost.WorkflowState;
|
||||
await _workflowRuntime.ImportWorkflowStateAsync(workflowState, cancellationToken);
|
||||
workflowStates.Add(workflowState);
|
||||
}
|
||||
|
||||
return workflowStates;
|
||||
}
|
||||
}
|
||||
|
|
@ -3,5 +3,5 @@ namespace Elsa.ProtoActor.Extensions;
|
|||
internal static class ProtoStringExtensions
|
||||
{
|
||||
public static string EmptyIfNull(this string? value) => value ?? "";
|
||||
public static string? NullIfEmpty(this string value) => value == "" ? default : value;
|
||||
public static string? NullIfEmpty(this string? value) => value == "" ? default : value;
|
||||
}
|
||||
|
|
@ -4,6 +4,7 @@ using Elsa.Features.Attributes;
|
|||
using Elsa.Features.Services;
|
||||
using Elsa.ProtoActor.Grains;
|
||||
using Elsa.ProtoActor.HostedServices;
|
||||
using Elsa.ProtoActor.Mappers;
|
||||
using Elsa.ProtoActor.Protos;
|
||||
using Elsa.ProtoActor.Services;
|
||||
using Elsa.Workflows.Core.Features;
|
||||
|
|
@ -39,7 +40,7 @@ public class ProtoActorFeature : FeatureBase
|
|||
{
|
||||
// Configure default workflow execution pipeline suitable for Proto Actor.
|
||||
Module.UseWorkflows(workflows => workflows.WithProtoActorRuntimeWorkflowExecutionPipeline());
|
||||
|
||||
|
||||
// Configure runtime with ProtoActor workflow runtime.
|
||||
Module.Configure<WorkflowRuntimeFeature>().WorkflowRuntime = sp => ActivatorUtilities.CreateInstance<ProtoActorWorkflowRuntime>(sp);
|
||||
}
|
||||
|
|
@ -48,22 +49,22 @@ public class ProtoActorFeature : FeatureBase
|
|||
/// The name of the cluster to configure.
|
||||
/// </summary>
|
||||
public string ClusterName { get; set; } = "elsa-cluster";
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// A delegate that returns an instance of a concrete implementation of <see cref="IClusterProvider"/>.
|
||||
/// </summary>
|
||||
public Func<IServiceProvider, IClusterProvider> ClusterProvider { get; set; } = _ => new TestProvider(new TestProviderOptions(), new InMemAgent());
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// A delegate that configures an instance of <see cref="ActorSystemConfig"/>.
|
||||
/// </summary>
|
||||
public Action<IServiceProvider, ActorSystemConfig> ActorSystemConfig { get; set; } = SetupDefaultConfig;
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// A delegate that configures an instance of an <see cref="ActorSystem"/>.
|
||||
/// </summary>
|
||||
public Action<IServiceProvider, ActorSystem> ActorSystem { get; set; } = (_, _) => { };
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// A delegate that returns an instance of <see cref="GrpcNetRemoteConfig"/> to be used by the actor system.
|
||||
/// </summary>
|
||||
|
|
@ -84,19 +85,19 @@ public class ProtoActorFeature : FeatureBase
|
|||
public override void Apply()
|
||||
{
|
||||
var services = Services;
|
||||
|
||||
|
||||
// Register ActorSystem.
|
||||
services.AddSingleton(sp =>
|
||||
{
|
||||
var systemConfig = Proto.ActorSystemConfig
|
||||
.Setup()
|
||||
.WithMetrics();
|
||||
|
||||
|
||||
var clusterProvider = ClusterProvider(sp);
|
||||
var system = new ActorSystem(systemConfig).WithServiceProvider(sp);
|
||||
var workflowGrainProps = system.DI().PropsFor<WorkflowGrainActor>();
|
||||
var workflowRegistryGrainProps = system.DI().PropsFor<RunningWorkflowsGrainActor>();
|
||||
|
||||
|
||||
var clusterConfig = ClusterConfig
|
||||
.Setup(ClusterName, clusterProvider, new PartitionIdentityLookup())
|
||||
.WithHeartbeatExpiration(TimeSpan.FromDays(1))
|
||||
|
|
@ -114,17 +115,26 @@ public class ProtoActorFeature : FeatureBase
|
|||
system
|
||||
.WithRemote(remoteConfig)
|
||||
.WithCluster(clusterConfig);
|
||||
|
||||
ActorSystem(sp, system);
|
||||
|
||||
ActorSystem(sp, system);
|
||||
return system;
|
||||
});
|
||||
|
||||
// Logging.
|
||||
Log.SetLoggerFactory(LoggerFactory.Create(l => l.AddConsole().SetMinimumLevel(LogLevel.Warning)));
|
||||
|
||||
|
||||
// Persistence.
|
||||
services.AddSingleton(PersistenceProvider);
|
||||
|
||||
|
||||
// Mappers.
|
||||
services
|
||||
.AddSingleton<BookmarkMapper>()
|
||||
.AddSingleton<ExceptionMapper>()
|
||||
.AddSingleton<WorkflowExecutionResultMapper>()
|
||||
.AddSingleton<WorkflowFaultStateMapper>()
|
||||
.AddSingleton<WorkflowStatusMapper>()
|
||||
.AddSingleton<WorkflowSubStatusMapper>();
|
||||
|
||||
// Mediator handlers.
|
||||
services.AddHandlersFrom<ProtoActorFeature>();
|
||||
|
||||
|
|
@ -137,11 +147,11 @@ public class ProtoActorFeature : FeatureBase
|
|||
.AddTransient(sp => new RunningWorkflowsGrainActor((context, _) => ActivatorUtilities.CreateInstance<RunningWorkflowsGrain>(sp, context)))
|
||||
;
|
||||
}
|
||||
|
||||
|
||||
private static void SetupDefaultConfig(IServiceProvider serviceProvider, ActorSystemConfig config)
|
||||
{
|
||||
}
|
||||
|
||||
|
||||
private static GrpcNetRemoteConfig CreateDefaultRemoteConfig(IServiceProvider serviceProvider) =>
|
||||
GrpcNetRemoteConfig.BindToLocalhost()
|
||||
.WithProtoMessages(MessagesReflection.Descriptor)
|
||||
|
|
|
|||
|
|
@ -1,13 +1,16 @@
|
|||
using Elsa.Common.Models;
|
||||
using Elsa.ProtoActor.Extensions;
|
||||
using Elsa.ProtoActor.Mappers;
|
||||
using Elsa.ProtoActor.Protos;
|
||||
using Elsa.Workflows.Core.Contracts;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using Elsa.Workflows.Core.State;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using Proto;
|
||||
using Proto.Cluster;
|
||||
using Proto.Persistence;
|
||||
using Exception = System.Exception;
|
||||
using WorkflowStatus = Elsa.Workflows.Core.Models.WorkflowStatus;
|
||||
using ProtoBookmark = Elsa.ProtoActor.Protos.Bookmark;
|
||||
|
||||
namespace Elsa.ProtoActor.Grains;
|
||||
|
||||
|
|
@ -19,8 +22,11 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
private const int MaxSnapshotsToKeep = 5;
|
||||
private readonly IWorkflowDefinitionService _workflowDefinitionService;
|
||||
private readonly IWorkflowHostFactory _workflowHostFactory;
|
||||
private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer;
|
||||
private readonly IWorkflowStateSerializer _workflowStateSerializer;
|
||||
private readonly BookmarkMapper _bookmarkMapper;
|
||||
private readonly WorkflowStatusMapper _workflowStatusMapper;
|
||||
private readonly WorkflowSubStatusMapper _workflowSubStatusMapper;
|
||||
private readonly WorkflowFaultStateMapper _workflowFaultStateMapper;
|
||||
private readonly Persistence _persistence;
|
||||
|
||||
private string _definitionId = default!;
|
||||
|
|
@ -34,15 +40,21 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
public WorkflowGrain(
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IWorkflowHostFactory workflowHostFactory,
|
||||
IBookmarkPayloadSerializer bookmarkPayloadSerializer,
|
||||
IWorkflowStateSerializer workflowStateSerializer,
|
||||
IProvider provider,
|
||||
IContext context) : base(context)
|
||||
IContext context,
|
||||
BookmarkMapper bookmarkMapper,
|
||||
WorkflowStatusMapper workflowStatusMapper,
|
||||
WorkflowSubStatusMapper workflowSubStatusMapper,
|
||||
WorkflowFaultStateMapper workflowFaultStateMapper) : base(context)
|
||||
{
|
||||
_workflowDefinitionService = workflowDefinitionService;
|
||||
_workflowHostFactory = workflowHostFactory;
|
||||
_bookmarkPayloadSerializer = bookmarkPayloadSerializer;
|
||||
_workflowStateSerializer = workflowStateSerializer;
|
||||
_bookmarkMapper = bookmarkMapper;
|
||||
_workflowStatusMapper = workflowStatusMapper;
|
||||
_workflowSubStatusMapper = workflowSubStatusMapper;
|
||||
_workflowFaultStateMapper = workflowFaultStateMapper;
|
||||
_persistence = Persistence.WithSnapshotting(provider, Context.ClusterIdentity()!.Identity, ApplySnapshot);
|
||||
}
|
||||
|
||||
|
|
@ -72,7 +84,6 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
{
|
||||
DefinitionId = workflow.Identity.DefinitionId,
|
||||
DefinitionVersion = workflow.Identity.Version,
|
||||
//Bookmarks = _bookmarks
|
||||
};
|
||||
}
|
||||
|
||||
|
|
@ -90,15 +101,15 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
var versionOptions = VersionOptions.FromString(request.VersionOptions);
|
||||
var cancellationToken = Context.CancellationToken;
|
||||
var startWorkflowOptions = new StartWorkflowHostOptions(instanceId, correlationId, input, request.TriggerActivityId);
|
||||
|
||||
|
||||
_workflowHost = await CreateWorkflowHostAsync(definitionId, versionOptions, cancellationToken);
|
||||
_version = _workflowHost.Workflow.Identity.Version;
|
||||
_definitionId = definitionId;
|
||||
_instanceId = instanceId;
|
||||
_input = input;
|
||||
|
||||
|
||||
var canStart = await _workflowHost.CanStartWorkflowAsync(startWorkflowOptions, cancellationToken);
|
||||
|
||||
|
||||
return new CanStartWorkflowResponse
|
||||
{
|
||||
CanStart = canStart
|
||||
|
|
@ -106,7 +117,7 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override async Task<StartWorkflowResponse> Start(StartWorkflowRequest request)
|
||||
public override async Task<WorkflowExecutionResponse> Start(StartWorkflowRequest request)
|
||||
{
|
||||
var definitionId = request.DefinitionId;
|
||||
var instanceId = request.InstanceId;
|
||||
|
|
@ -134,15 +145,20 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
|
||||
await SaveSnapshotAsync();
|
||||
|
||||
return new StartWorkflowResponse
|
||||
return new WorkflowExecutionResponse
|
||||
{
|
||||
Result = result,
|
||||
Bookmarks = { Map(workflowState.Bookmarks) }
|
||||
Bookmarks = { _bookmarkMapper.Map(workflowState.Bookmarks).ToList() },
|
||||
Status = _workflowStatusMapper.Map(workflowState.Status),
|
||||
SubStatus = _workflowSubStatusMapper.Map(workflowState.SubStatus),
|
||||
Fault = workflowState.Fault != null ? _workflowFaultStateMapper.Map(workflowState.Fault) : default,
|
||||
TriggeredActivityId = string.Empty,
|
||||
WorkflowInstanceId = instanceId
|
||||
};
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override async Task<ResumeWorkflowResponse> Resume(ResumeWorkflowRequest request)
|
||||
public override async Task<WorkflowExecutionResponse> Resume(ResumeWorkflowRequest request)
|
||||
{
|
||||
_input = request.Input?.Deserialize();
|
||||
var correlationId = request.CorrelationId;
|
||||
|
|
@ -152,26 +168,26 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
var activityInstanceId = request.ActivityInstanceId.NullIfEmpty();
|
||||
var activityHash = request.ActivityHash.NullIfEmpty();
|
||||
var cancellationToken = Context.CancellationToken;
|
||||
|
||||
|
||||
var resumeWorkflowHostOptions = new ResumeWorkflowHostOptions(
|
||||
correlationId,
|
||||
bookmarkId,
|
||||
activityId,
|
||||
correlationId,
|
||||
bookmarkId,
|
||||
activityId,
|
||||
activityNodeId,
|
||||
activityInstanceId,
|
||||
activityHash,
|
||||
_input);
|
||||
|
||||
|
||||
var definitionId = _definitionId;
|
||||
var versionOptions = VersionOptions.SpecificVersion(_version);
|
||||
|
||||
|
||||
// Only need to reconstruct a workflow host if not already done so during CanStart.
|
||||
if (_workflowHost == null!)
|
||||
{
|
||||
_workflowHost = await CreateWorkflowHostAsync(definitionId, versionOptions, cancellationToken);
|
||||
_version = _workflowHost.Workflow.Identity.Version;
|
||||
}
|
||||
|
||||
|
||||
await _workflowHost.ResumeWorkflowAsync(resumeWorkflowHostOptions, cancellationToken);
|
||||
var finished = _workflowHost.WorkflowState.Status == WorkflowStatus.Finished;
|
||||
|
||||
|
|
@ -179,10 +195,15 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
|
||||
await SaveSnapshotAsync();
|
||||
|
||||
return new ResumeWorkflowResponse
|
||||
return new WorkflowExecutionResponse
|
||||
{
|
||||
Result = finished ? Protos.RunWorkflowResult.Finished : Protos.RunWorkflowResult.Suspended,
|
||||
Bookmarks = { Map(_workflowHost.WorkflowState.Bookmarks) }
|
||||
Result = finished ? RunWorkflowResult.Finished : RunWorkflowResult.Suspended,
|
||||
Bookmarks = { _bookmarkMapper.Map(_workflowHost.WorkflowState.Bookmarks).ToList() },
|
||||
Fault = _workflowState.Fault != null ? _workflowFaultStateMapper.Map(_workflowState.Fault) : default,
|
||||
TriggeredActivityId = string.Empty,
|
||||
WorkflowInstanceId = _workflowState.Id,
|
||||
Status = _workflowStatusMapper.Map(_workflowState.Status),
|
||||
SubStatus = _workflowSubStatusMapper.Map(_workflowState.SubStatus)
|
||||
};
|
||||
}
|
||||
|
||||
|
|
@ -218,6 +239,7 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
}
|
||||
|
||||
private void ApplySnapshot(Snapshot snapshot) => (_definitionId, _instanceId, _version, _workflowState, _input) = (WorkflowSnapshot)snapshot.State;
|
||||
|
||||
private async Task SaveSnapshotAsync()
|
||||
{
|
||||
if (_workflowState.Status == WorkflowStatus.Finished)
|
||||
|
|
@ -240,23 +262,4 @@ public class WorkflowGrain : WorkflowGrainBase
|
|||
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
|
||||
return await _workflowHostFactory.CreateAsync(workflow, cancellationToken);
|
||||
}
|
||||
|
||||
private IEnumerable<BookmarkDto> Map(IEnumerable<Bookmark> bookmarks)
|
||||
{
|
||||
return bookmarks.Select(x =>
|
||||
{
|
||||
var payloadJson = x.Payload != null ? _bookmarkPayloadSerializer.Serialize(x.Payload) : "";
|
||||
return new BookmarkDto
|
||||
{
|
||||
Id = x.Id,
|
||||
Name = x.Name,
|
||||
ActivityNodeId = x.ActivityNodeId,
|
||||
ActivityInstanceId = x.ActivityInstanceId,
|
||||
Hash = x.Hash,
|
||||
Data = payloadJson,
|
||||
AutoBurn = x.AutoBurn,
|
||||
CallbackMethodName = x.CallbackMethodName.EmptyIfNull()
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
52
src/modules/Elsa.ProtoActor/Mappers/BookmarkMapper.cs
Normal file
52
src/modules/Elsa.ProtoActor/Mappers/BookmarkMapper.cs
Normal file
|
|
@ -0,0 +1,52 @@
|
|||
using Elsa.ProtoActor.Extensions;
|
||||
using Elsa.Workflows.Core.Contracts;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using ProtoBookmark = Elsa.ProtoActor.Protos.Bookmark;
|
||||
|
||||
namespace Elsa.ProtoActor.Mappers;
|
||||
|
||||
/// <summary>
|
||||
/// Maps between <see cref="Bookmark"/> and <see cref="ProtoBookmark"/>.
|
||||
/// </summary>
|
||||
public class BookmarkMapper
|
||||
{
|
||||
private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of <see cref="BookmarkMapper"/>.
|
||||
/// </summary>
|
||||
public BookmarkMapper(IBookmarkPayloadSerializer bookmarkPayloadSerializer)
|
||||
{
|
||||
_bookmarkPayloadSerializer = bookmarkPayloadSerializer;
|
||||
}
|
||||
|
||||
public Bookmark Map(ProtoBookmark bookmark) =>
|
||||
new(bookmark.Id, bookmark.Name, bookmark.Hash, bookmark.Payload, bookmark.ActivityNodeId, bookmark.ActivityInstanceId, bookmark.AutoBurn, bookmark.CallbackMethodName);
|
||||
|
||||
|
||||
public IEnumerable<ProtoBookmark> Map(IEnumerable<Bookmark> source) =>
|
||||
source.Select(x =>
|
||||
new ProtoBookmark
|
||||
{
|
||||
Id = x.Id,
|
||||
Name = x.Name,
|
||||
Hash = x.Hash,
|
||||
Payload = x.Payload != null ? _bookmarkPayloadSerializer.Serialize(x.Payload) : string.Empty,
|
||||
ActivityNodeId = x.ActivityNodeId,
|
||||
ActivityInstanceId = x.ActivityInstanceId,
|
||||
AutoBurn = x.AutoBurn,
|
||||
CallbackMethodName = x.CallbackMethodName.NullIfEmpty()
|
||||
});
|
||||
|
||||
public IEnumerable<Bookmark> Map(IEnumerable<ProtoBookmark> source) =>
|
||||
source.Select(x =>
|
||||
new Bookmark(
|
||||
x.Id,
|
||||
x.Name,
|
||||
x.Hash,
|
||||
x.Payload.NullIfEmpty(),
|
||||
x.ActivityNodeId,
|
||||
x.ActivityInstanceId,
|
||||
x.AutoBurn,
|
||||
x.CallbackMethodName.NullIfEmpty()));
|
||||
}
|
||||
19
src/modules/Elsa.ProtoActor/Mappers/ExceptionMapper.cs
Normal file
19
src/modules/Elsa.ProtoActor/Mappers/ExceptionMapper.cs
Normal file
|
|
@ -0,0 +1,19 @@
|
|||
using Elsa.Workflows.Core.State;
|
||||
using ProtoException = Elsa.ProtoActor.Protos.ExceptionState;
|
||||
|
||||
namespace Elsa.ProtoActor.Mappers;
|
||||
|
||||
public class ExceptionMapper
|
||||
{
|
||||
public ExceptionState Map(ProtoException exception) =>
|
||||
new(Type.GetType(exception.Type)!, exception.Message, exception.StackTrace, exception.InnerException != null ? Map(exception.InnerException) : null);
|
||||
|
||||
public ProtoException Map(ExceptionState exception) =>
|
||||
new()
|
||||
{
|
||||
Type = exception.Type.AssemblyQualifiedName,
|
||||
Message = exception.Message,
|
||||
StackTrace = exception.StackTrace,
|
||||
InnerException = exception.InnerException != null ? Map(exception.InnerException) : null
|
||||
};
|
||||
}
|
||||
|
|
@ -0,0 +1,59 @@
|
|||
using Elsa.ProtoActor.Extensions;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using ProtoWorkflowExecutionResponse = Elsa.ProtoActor.Protos.WorkflowExecutionResponse;
|
||||
|
||||
namespace Elsa.ProtoActor.Mappers;
|
||||
|
||||
/// <summary>
|
||||
/// Maps between <see cref="WorkflowExecutionResult"/> and <see cref="ProtoWorkflowExecutionResponse"/>.
|
||||
/// </summary>
|
||||
public class WorkflowExecutionResultMapper
|
||||
{
|
||||
private readonly WorkflowStatusMapper _workflowStatusMapper;
|
||||
private readonly WorkflowSubStatusMapper _workflowSubStatusMapper;
|
||||
private readonly BookmarkMapper _bookmarkMapper;
|
||||
private readonly WorkflowFaultStateMapper _workflowFaultStateMapper;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="WorkflowExecutionResultMapper"/> class.
|
||||
/// </summary>
|
||||
public WorkflowExecutionResultMapper(
|
||||
WorkflowStatusMapper workflowStatusMapper,
|
||||
WorkflowSubStatusMapper workflowSubStatusMapper,
|
||||
BookmarkMapper bookmarkMapper,
|
||||
WorkflowFaultStateMapper workflowFaultStateMapper)
|
||||
{
|
||||
_workflowStatusMapper = workflowStatusMapper;
|
||||
_workflowSubStatusMapper = workflowSubStatusMapper;
|
||||
_bookmarkMapper = bookmarkMapper;
|
||||
_workflowFaultStateMapper = workflowFaultStateMapper;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Maps a <see cref="ProtoWorkflowExecutionResponse"/> to a <see cref="WorkflowExecutionResult"/>.
|
||||
/// </summary>
|
||||
/// <param name="source">The source.</param>
|
||||
/// <returns>The mapped <see cref="WorkflowExecutionResult"/>.</returns>
|
||||
public WorkflowExecutionResult Map(ProtoWorkflowExecutionResponse source) => new(
|
||||
source.WorkflowInstanceId,
|
||||
_workflowStatusMapper.Map(source.Status),
|
||||
_workflowSubStatusMapper.Map(source.SubStatus),
|
||||
_bookmarkMapper.Map(source.Bookmarks).ToList(),
|
||||
source.TriggeredActivityId.NullIfEmpty(),
|
||||
source.Fault != null ? _workflowFaultStateMapper.Map(source.Fault) : default);
|
||||
|
||||
/// <summary>
|
||||
/// Maps a <see cref="WorkflowExecutionResult"/> to a <see cref="ProtoWorkflowExecutionResponse"/>.
|
||||
/// </summary>
|
||||
/// <param name="source">The source.</param>
|
||||
/// <returns>The mapped <see cref="ProtoWorkflowExecutionResponse"/>.</returns>
|
||||
public ProtoWorkflowExecutionResponse Map(WorkflowExecutionResult source) => new()
|
||||
{
|
||||
WorkflowInstanceId = source.WorkflowInstanceId,
|
||||
Status = _workflowStatusMapper.Map(source.Status),
|
||||
SubStatus = _workflowSubStatusMapper.Map(source.SubStatus),
|
||||
Bookmarks = {_bookmarkMapper.Map(source.Bookmarks)},
|
||||
TriggeredActivityId = source.TriggeredActivityId,
|
||||
Fault = source.Fault != null ? _workflowFaultStateMapper.Map(source.Fault) : default
|
||||
};
|
||||
}
|
||||
|
|
@ -0,0 +1,41 @@
|
|||
using Elsa.Workflows.Core.State;
|
||||
using ProtoWorkflowFault = Elsa.ProtoActor.Protos.WorkflowFault;
|
||||
|
||||
namespace Elsa.ProtoActor.Mappers;
|
||||
|
||||
/// <summary>
|
||||
/// Maps between <see cref="WorkflowFaultState"/> and <see cref="ProtoWorkflowFault"/>.
|
||||
/// </summary>
|
||||
public class WorkflowFaultStateMapper
|
||||
{
|
||||
private readonly ExceptionMapper _exceptionMapper;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="WorkflowFaultStateMapper"/> class.
|
||||
/// </summary>
|
||||
public WorkflowFaultStateMapper(ExceptionMapper exceptionMapper)
|
||||
{
|
||||
_exceptionMapper = exceptionMapper;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Maps a <see cref="ProtoWorkflowFault"/> to a <see cref="WorkflowFaultState"/>.
|
||||
/// </summary>
|
||||
/// <param name="source">The source.</param>
|
||||
/// <returns>The mapped <see cref="WorkflowFaultState"/>.</returns>
|
||||
public WorkflowFaultState Map(ProtoWorkflowFault source) =>
|
||||
new(_exceptionMapper.Map(source.Exception), source.Message, source.FaultedActivityId);
|
||||
|
||||
/// <summary>
|
||||
/// Maps a <see cref="WorkflowFaultState"/> to a <see cref="ProtoWorkflowFault"/>.
|
||||
/// </summary>
|
||||
/// <param name="source">The source.</param>
|
||||
/// <returns>The mapped <see cref="ProtoWorkflowFault"/>.</returns>
|
||||
public ProtoWorkflowFault Map(WorkflowFaultState source) =>
|
||||
new()
|
||||
{
|
||||
Exception = source.Exception != null ? _exceptionMapper.Map(source.Exception) : default,
|
||||
Message = source.Message,
|
||||
FaultedActivityId = source.FaultedActivityId
|
||||
};
|
||||
}
|
||||
23
src/modules/Elsa.ProtoActor/Mappers/WorkflowStatusMapper.cs
Normal file
23
src/modules/Elsa.ProtoActor/Mappers/WorkflowStatusMapper.cs
Normal file
|
|
@ -0,0 +1,23 @@
|
|||
using Elsa.Workflows.Core.Models;
|
||||
using ProtoWorkflowStatus = Elsa.ProtoActor.Protos.WorkflowStatus;
|
||||
|
||||
namespace Elsa.ProtoActor.Mappers;
|
||||
|
||||
public class WorkflowStatusMapper
|
||||
{
|
||||
public WorkflowStatus Map(ProtoWorkflowStatus status) =>
|
||||
status switch
|
||||
{
|
||||
ProtoWorkflowStatus.Finished => WorkflowStatus.Finished,
|
||||
ProtoWorkflowStatus.Running => WorkflowStatus.Running,
|
||||
_ => throw new ArgumentOutOfRangeException(nameof(status), status, null)
|
||||
};
|
||||
|
||||
public ProtoWorkflowStatus Map(WorkflowStatus status) =>
|
||||
status switch
|
||||
{
|
||||
WorkflowStatus.Finished => ProtoWorkflowStatus.Finished,
|
||||
WorkflowStatus.Running => ProtoWorkflowStatus.Running,
|
||||
_ => throw new ArgumentOutOfRangeException(nameof(status), status, null)
|
||||
};
|
||||
}
|
||||
|
|
@ -0,0 +1,29 @@
|
|||
using Elsa.Workflows.Core.Models;
|
||||
using ProtoWorkflowSubStatus = Elsa.ProtoActor.Protos.WorkflowSubStatus;
|
||||
|
||||
namespace Elsa.ProtoActor.Mappers;
|
||||
|
||||
public class WorkflowSubStatusMapper
|
||||
{
|
||||
public WorkflowSubStatus Map(ProtoWorkflowSubStatus subStatus) =>
|
||||
subStatus switch
|
||||
{
|
||||
ProtoWorkflowSubStatus.Faulted => WorkflowSubStatus.Faulted,
|
||||
ProtoWorkflowSubStatus.Finished => WorkflowSubStatus.Finished,
|
||||
ProtoWorkflowSubStatus.Cancelled => WorkflowSubStatus.Cancelled,
|
||||
ProtoWorkflowSubStatus.Executing => WorkflowSubStatus.Executing,
|
||||
ProtoWorkflowSubStatus.Suspended => WorkflowSubStatus.Suspended,
|
||||
_ => throw new ArgumentOutOfRangeException(nameof(subStatus), subStatus, null)
|
||||
};
|
||||
|
||||
public ProtoWorkflowSubStatus Map(WorkflowSubStatus subStatus) =>
|
||||
subStatus switch
|
||||
{
|
||||
WorkflowSubStatus.Faulted => ProtoWorkflowSubStatus.Faulted,
|
||||
WorkflowSubStatus.Finished => ProtoWorkflowSubStatus.Finished,
|
||||
WorkflowSubStatus.Cancelled => ProtoWorkflowSubStatus.Cancelled,
|
||||
WorkflowSubStatus.Executing => ProtoWorkflowSubStatus.Executing,
|
||||
WorkflowSubStatus.Suspended => ProtoWorkflowSubStatus.Suspended,
|
||||
_ => throw new ArgumentOutOfRangeException(nameof(subStatus), subStatus, null)
|
||||
};
|
||||
}
|
||||
|
|
@ -7,8 +7,8 @@ import "Messages.proto";
|
|||
|
||||
service WorkflowGrain {
|
||||
rpc CanStart (StartWorkflowRequest) returns (CanStartWorkflowResponse);
|
||||
rpc Start (StartWorkflowRequest) returns (StartWorkflowResponse);
|
||||
rpc Resume (ResumeWorkflowRequest) returns (ResumeWorkflowResponse);
|
||||
rpc Start (StartWorkflowRequest) returns (WorkflowExecutionResponse);
|
||||
rpc Resume (ResumeWorkflowRequest) returns (WorkflowExecutionResponse);
|
||||
rpc ExportState(ExportWorkflowStateRequest) returns (ExportWorkflowStateResponse);
|
||||
rpc ImportState(ImportWorkflowStateRequest) returns (ImportWorkflowStateResponse);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,9 +27,40 @@ message StartWorkflowRequest {
|
|||
optional string TriggerActivityId = 6;
|
||||
}
|
||||
|
||||
message StartWorkflowResponse {
|
||||
RunWorkflowResult Result = 1;
|
||||
repeated BookmarkDto Bookmarks = 2;
|
||||
message WorkflowExecutionResponse {
|
||||
string WorkflowInstanceId = 1;
|
||||
RunWorkflowResult Result = 2;
|
||||
WorkflowStatus Status = 3;
|
||||
WorkflowSubStatus SubStatus = 4;
|
||||
repeated Bookmark Bookmarks = 5;
|
||||
optional WorkflowFault Fault = 6;
|
||||
optional string TriggeredActivityId = 7;
|
||||
}
|
||||
|
||||
message WorkflowFault {
|
||||
optional ExceptionState Exception = 1;
|
||||
string Message = 2;
|
||||
optional string FaultedActivityId = 3;
|
||||
}
|
||||
|
||||
message ExceptionState {
|
||||
string Type = 1;
|
||||
string Message = 2;
|
||||
optional string StackTrace = 3;
|
||||
optional ExceptionState InnerException = 4;
|
||||
}
|
||||
|
||||
enum WorkflowStatus {
|
||||
WorkflowStatusRunning = 0;
|
||||
WorkflowStatusFinished = 1;
|
||||
}
|
||||
|
||||
enum WorkflowSubStatus {
|
||||
WorkflowSubStatusExecuting = 0;
|
||||
WorkflowSubStatusSuspended = 1;
|
||||
WorkflowSubStatusFinished = 2;
|
||||
WorkflowSubStatusCancelled = 3;
|
||||
WorkflowSubStatusFaulted = 4;
|
||||
}
|
||||
|
||||
message ResumeWorkflowRequest {
|
||||
|
|
@ -43,11 +74,6 @@ message ResumeWorkflowRequest {
|
|||
optional Input input = 8;
|
||||
}
|
||||
|
||||
message ResumeWorkflowResponse {
|
||||
RunWorkflowResult Result = 1;
|
||||
repeated BookmarkDto Bookmarks = 2;
|
||||
}
|
||||
|
||||
message ExportWorkflowStateRequest {}
|
||||
|
||||
message ExportWorkflowStateResponse {
|
||||
|
|
@ -65,11 +91,11 @@ enum RunWorkflowResult {
|
|||
RunWorkflowResultSuspended = 1;
|
||||
}
|
||||
|
||||
message BookmarkDto {
|
||||
message Bookmark {
|
||||
string Id = 1;
|
||||
string Name = 2;
|
||||
string Hash = 3;
|
||||
optional string Data = 4;
|
||||
optional string Payload = 4;
|
||||
string ActivityNodeId = 5;
|
||||
string ActivityInstanceId = 6;
|
||||
optional bool AutoBurn = 7;
|
||||
|
|
|
|||
|
|
@ -1,13 +1,20 @@
|
|||
using Elsa.Common.Models;
|
||||
using Elsa.Extensions;
|
||||
using Elsa.ProtoActor.Extensions;
|
||||
using Elsa.ProtoActor.Mappers;
|
||||
using Elsa.ProtoActor.Protos;
|
||||
using Elsa.Workflows.Core.Contracts;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using Elsa.Workflows.Core.State;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using Elsa.Workflows.Runtime.Entities;
|
||||
using Proto.Cluster;
|
||||
using Bookmark = Elsa.Workflows.Core.Models.Bookmark;
|
||||
using ProtoWorkflowStatus = Elsa.ProtoActor.Protos.WorkflowStatus;
|
||||
using ProtoWorkflowSubStatus = Elsa.ProtoActor.Protos.WorkflowSubStatus;
|
||||
using ProtoWorkflowFault = Elsa.ProtoActor.Protos.WorkflowFault;
|
||||
using ProtoException = Elsa.ProtoActor.Protos.ExceptionState;
|
||||
using ProtoWorkflowExecutionResponse = Elsa.ProtoActor.Protos.WorkflowExecutionResponse;
|
||||
using ProtoBookmark = Elsa.ProtoActor.Protos.Bookmark;
|
||||
|
||||
namespace Elsa.ProtoActor.Services;
|
||||
|
||||
|
|
@ -24,6 +31,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
private readonly IBookmarkHasher _hasher;
|
||||
private readonly IWorkflowDefinitionService _workflowDefinitionService;
|
||||
private readonly IWorkflowInstanceFactory _workflowInstanceFactory;
|
||||
private readonly WorkflowExecutionResultMapper _workflowExecutionResultMapper;
|
||||
|
||||
/// <summary>
|
||||
/// Constructor.
|
||||
|
|
@ -36,7 +44,8 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
IIdentityGenerator identityGenerator,
|
||||
IBookmarkHasher hasher,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IWorkflowInstanceFactory workflowInstanceFactory)
|
||||
IWorkflowInstanceFactory workflowInstanceFactory,
|
||||
WorkflowExecutionResultMapper workflowExecutionResultMapper)
|
||||
{
|
||||
_cluster = cluster;
|
||||
_workflowStateSerializer = workflowStateSerializer;
|
||||
|
|
@ -46,6 +55,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
_hasher = hasher;
|
||||
_workflowDefinitionService = workflowDefinitionService;
|
||||
_workflowInstanceFactory = workflowInstanceFactory;
|
||||
_workflowExecutionResultMapper = workflowExecutionResultMapper;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -104,9 +114,8 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
|
||||
var client = _cluster.GetNamedWorkflowGrain(workflowInstanceId);
|
||||
var response = await client.Start(request, cancellationToken);
|
||||
var bookmarks = Map(response!.Bookmarks).ToList();
|
||||
|
||||
return new WorkflowExecutionResult(workflowInstanceId, bookmarks);
|
||||
return _workflowExecutionResultMapper.Map(response!);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -128,15 +137,14 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
continue;
|
||||
|
||||
var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
|
||||
|
||||
results.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
|
||||
results.Add(startResult);
|
||||
}
|
||||
|
||||
return results;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
public async Task<WorkflowExecutionResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var request = new ResumeWorkflowRequest
|
||||
{
|
||||
|
|
@ -149,9 +157,8 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
|
||||
var client = _cluster.GetNamedWorkflowGrain(workflowInstanceId);
|
||||
var response = await client.Resume(request, cancellationToken);
|
||||
var bookmarks = Map(response!.Bookmarks).ToList();
|
||||
|
||||
return new ResumeWorkflowResult(bookmarks);
|
||||
return _workflowExecutionResultMapper.Map(response!);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -164,7 +171,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
var bookmarks = await _bookmarkStore.FindManyAsync(filter, cancellationToken);
|
||||
return await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(correlationId, Input: options.Input), cancellationToken);
|
||||
}
|
||||
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
{
|
||||
|
|
@ -182,19 +189,16 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
{
|
||||
var startOptions = new StartWorkflowRuntimeOptions(collectedStartableWorkflow.CorrelationId, input, VersionOptions.Published,
|
||||
collectedStartableWorkflow.ActivityId, collectedStartableWorkflow.WorkflowInstanceId);
|
||||
var startResult = await StartWorkflowAsync(collectedStartableWorkflow.DefinitionId!, startOptions, cancellationToken);
|
||||
return new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks);
|
||||
return await StartWorkflowAsync(collectedStartableWorkflow.DefinitionId!, startOptions, cancellationToken);
|
||||
}
|
||||
|
||||
var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!;
|
||||
var runtimeOptions = new ResumeWorkflowRuntimeOptions(collectedResumableWorkflow.CorrelationId, Input: input);
|
||||
|
||||
var resumeResult = await ResumeWorkflowAsync(
|
||||
return await ResumeWorkflowAsync(
|
||||
match.WorkflowInstanceId,
|
||||
runtimeOptions with { BookmarkId = collectedResumableWorkflow.BookmarkId },
|
||||
cancellationToken);
|
||||
|
||||
return new WorkflowExecutionResult(collectedResumableWorkflow.WorkflowInstanceId, resumeResult.Bookmarks);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -275,7 +279,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
runtimeOptions with { BookmarkId = bookmark.BookmarkId },
|
||||
cancellationToken);
|
||||
|
||||
resumedWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks));
|
||||
resumedWorkflows.Add(resumeResult);
|
||||
}
|
||||
|
||||
return resumedWorkflows;
|
||||
|
|
@ -299,18 +303,6 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
|||
}
|
||||
}
|
||||
|
||||
private static IEnumerable<Bookmark> Map(IEnumerable<BookmarkDto> source) =>
|
||||
source.Select(x =>
|
||||
new Bookmark(
|
||||
x.Id,
|
||||
x.Name,
|
||||
x.Hash,
|
||||
x.Data.NullIfEmpty(),
|
||||
x.ActivityNodeId,
|
||||
x.ActivityInstanceId,
|
||||
x.AutoBurn,
|
||||
x.CallbackMethodName.NullIfEmpty()));
|
||||
|
||||
private async Task<IEnumerable<WorkflowMatch>> FindStartableWorkflowsAsync(WorkflowsFilter workflowsFilter, CancellationToken cancellationToken)
|
||||
{
|
||||
var hash = _hasher.Hash(workflowsFilter.ActivityTypeName, workflowsFilter.BookmarkPayload);
|
||||
|
|
|
|||
|
|
@ -1,9 +1,14 @@
|
|||
using System.Net.Mime;
|
||||
using Elsa.Abstractions;
|
||||
using Elsa.Common.Models;
|
||||
using Elsa.Http.Contracts;
|
||||
using Elsa.Workflows.Core.Contracts;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using Elsa.Workflows.Core.State;
|
||||
using Elsa.Workflows.Management.Contracts;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using JetBrains.Annotations;
|
||||
using Microsoft.AspNetCore.Http;
|
||||
|
||||
namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Execute;
|
||||
|
||||
|
|
@ -16,13 +21,15 @@ public class Execute : ElsaEndpoint<Request, Response>
|
|||
private readonly IWorkflowDefinitionStore _store;
|
||||
private readonly IWorkflowRuntime _workflowRuntime;
|
||||
private readonly IHttpBookmarkProcessor _httpBookmarkProcessor;
|
||||
private IApiSerializer _apiSerializer;
|
||||
|
||||
/// <inheritdoc />
|
||||
public Execute(IWorkflowDefinitionStore store, IWorkflowRuntime workflowRuntime, IHttpBookmarkProcessor httpBookmarkProcessor)
|
||||
public Execute(IWorkflowDefinitionStore store, IWorkflowRuntime workflowRuntime, IHttpBookmarkProcessor httpBookmarkProcessor, IApiSerializer apiSerializer)
|
||||
{
|
||||
_store = store;
|
||||
_workflowRuntime = workflowRuntime;
|
||||
_httpBookmarkProcessor = httpBookmarkProcessor;
|
||||
_apiSerializer = apiSerializer;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -49,10 +56,39 @@ public class Execute : ElsaEndpoint<Request, Response>
|
|||
var startWorkflowOptions = new StartWorkflowRuntimeOptions(correlationId, input, VersionOptions.Published);
|
||||
var result = await _workflowRuntime.StartWorkflowAsync(definitionId, startWorkflowOptions, cancellationToken);
|
||||
|
||||
// If a workflow fault occurred, respond appropriately with a 500 internal server error.
|
||||
if (result.SubStatus == WorkflowSubStatus.Faulted)
|
||||
{
|
||||
var workflowState = await _workflowRuntime.ExportWorkflowStateAsync(result.WorkflowInstanceId, cancellationToken);
|
||||
await HandleFaultAsync(workflowState!, cancellationToken);
|
||||
return;
|
||||
}
|
||||
|
||||
// Resume any HTTP bookmarks.
|
||||
await _httpBookmarkProcessor.ProcessBookmarks(new[] { result }, correlationId, default, cancellationToken);
|
||||
|
||||
var workflowState2 = await _workflowRuntime.ExportWorkflowStateAsync(result.WorkflowInstanceId, cancellationToken);
|
||||
|
||||
if (workflowState2!.SubStatus == WorkflowSubStatus.Faulted)
|
||||
{
|
||||
await HandleFaultAsync(workflowState2, cancellationToken);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!HttpContext.Response.HasStarted)
|
||||
await SendOkAsync(new Response(result.InstanceId), cancellationToken);
|
||||
await SendOkAsync(new Response(result.WorkflowInstanceId), cancellationToken);
|
||||
}
|
||||
|
||||
private async Task HandleFaultAsync(WorkflowState workflowState, CancellationToken cancellationToken)
|
||||
{
|
||||
var faultedResponse = _apiSerializer.Serialize(new
|
||||
{
|
||||
errorMessage = $"Workflow faulted with error: {workflowState.Fault!.Message}",
|
||||
workflowState = workflowState
|
||||
});
|
||||
|
||||
HttpContext.Response.ContentType = MediaTypeNames.Application.Json;
|
||||
HttpContext.Response.StatusCode = StatusCodes.Status500InternalServerError;
|
||||
await HttpContext.Response.WriteAsync(faultedResponse, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,4 +1,5 @@
|
|||
using Elsa.Common.Models;
|
||||
using Elsa.Workflows.Core.Activities;
|
||||
using Elsa.Workflows.Core.Helpers;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using Elsa.Workflows.Core.State;
|
||||
|
|
@ -46,7 +47,7 @@ public interface IWorkflowRuntime
|
|||
/// <param name="workflowInstanceId">The ID of the workflow instance to resume.</param>
|
||||
/// <param name="options"></param>
|
||||
/// <param name="cancellationToken"></param>
|
||||
Task<ResumeWorkflowResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
|
||||
Task<WorkflowExecutionResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Resumes all workflows that are bookmarked on the specified activity type.
|
||||
|
|
@ -117,13 +118,11 @@ public record ResumeWorkflowRuntimeOptions(
|
|||
|
||||
public record CanStartWorkflowResult(string? InstanceId, bool CanStart);
|
||||
|
||||
public record ResumeWorkflowResult(ICollection<Bookmark> Bookmarks);
|
||||
|
||||
public record TriggerWorkflowsRuntimeOptions(string? CorrelationId = default, string? WorkflowInstanceId = default, IDictionary<string, object>? Input = default);
|
||||
|
||||
public record TriggerWorkflowsResult(ICollection<WorkflowExecutionResult> TriggeredWorkflows);
|
||||
|
||||
public record WorkflowExecutionResult(string InstanceId, ICollection<Bookmark> Bookmarks, string? ActivityId = null);
|
||||
public record WorkflowExecutionResult(string WorkflowInstanceId, WorkflowStatus Status, WorkflowSubStatus SubStatus, ICollection<Bookmark> Bookmarks, string? TriggeredActivityId = null, WorkflowFaultState? Fault = default);
|
||||
|
||||
public record UpdateBookmarksContext(string InstanceId, Diff<Bookmark> Diff, string? CorrelationId);
|
||||
|
||||
|
|
|
|||
|
|
@ -75,7 +75,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
|
||||
if (workflowDefinition == null)
|
||||
return null;
|
||||
|
||||
|
||||
return await StartWorkflowAsync(workflowDefinition, options, cancellationToken);
|
||||
}
|
||||
|
||||
|
|
@ -107,7 +107,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
continue;
|
||||
|
||||
var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
|
||||
results.Add(startResult with { ActivityId = trigger.ActivityId });
|
||||
results.Add(startResult with { TriggeredActivityId = trigger.ActivityId });
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -115,7 +115,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
public async Task<WorkflowExecutionResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await using (await _distributedLockProvider.AcquireLockAsync(workflowInstanceId, TimeSpan.FromMinutes(1), cancellationToken))
|
||||
{
|
||||
|
|
@ -134,8 +134,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
|
||||
if (workflowDefinition == null)
|
||||
{
|
||||
_logger.LogInformation("The workflow definition {DefinitionId} version {Version} was not found", definitionId, version);
|
||||
return new ResumeWorkflowResult(Array.Empty<Bookmark>());
|
||||
throw new Exception($"The workflow definition {definitionId} version {version} was not found");
|
||||
}
|
||||
|
||||
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
|
||||
|
|
@ -156,7 +155,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
|
||||
await SaveWorkflowStateAsync(workflowState, cancellationToken);
|
||||
|
||||
return new ResumeWorkflowResult(workflowState.Bookmarks);
|
||||
return new WorkflowExecutionResult(workflowState.Id, workflowState.Status, workflowState.SubStatus, workflowState.Bookmarks);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -188,17 +187,16 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
var startOptions = new StartWorkflowRuntimeOptions(collectedStartableWorkflow.CorrelationId, input, VersionOptions.Published,
|
||||
collectedStartableWorkflow.ActivityId, collectedStartableWorkflow.WorkflowInstanceId);
|
||||
var startResult = await StartWorkflowAsync(collectedStartableWorkflow.DefinitionId!, startOptions, cancellationToken);
|
||||
return startResult with { ActivityId = collectedStartableWorkflow.ActivityId };
|
||||
return startResult with { TriggeredActivityId = collectedStartableWorkflow.ActivityId };
|
||||
}
|
||||
|
||||
var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!;
|
||||
var runtimeOptions = new ResumeWorkflowRuntimeOptions(collectedResumableWorkflow.CorrelationId, Input: input);
|
||||
var resumeResult = await ResumeWorkflowAsync(
|
||||
|
||||
return await ResumeWorkflowAsync(
|
||||
match.WorkflowInstanceId,
|
||||
runtimeOptions with { BookmarkId = collectedResumableWorkflow.BookmarkId },
|
||||
cancellationToken);
|
||||
|
||||
return new WorkflowExecutionResult(collectedResumableWorkflow.WorkflowInstanceId, resumeResult.Bookmarks);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -248,14 +246,20 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
|
||||
await SaveWorkflowStateAsync(workflowState, cancellationToken);
|
||||
|
||||
return new WorkflowExecutionResult(workflowState.Id, workflowState.Bookmarks);
|
||||
return new WorkflowExecutionResult(
|
||||
workflowState.Id,
|
||||
workflowState.Status,
|
||||
workflowState.SubStatus,
|
||||
workflowState.Bookmarks,
|
||||
default,
|
||||
workflowState.Fault);
|
||||
}
|
||||
|
||||
|
||||
private async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken)
|
||||
{
|
||||
return await _workflowDefinitionService.FindAsync(definitionId, versionOptions, cancellationToken);
|
||||
}
|
||||
|
||||
|
||||
private async Task<IWorkflowHost> CreateWorkflowHostAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken)
|
||||
{
|
||||
var versionOptions = options.VersionOptions;
|
||||
|
|
@ -263,10 +267,10 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
|
||||
if (workflowDefinition == null)
|
||||
throw new Exception("Specified workflow definition and version does not exist");
|
||||
|
||||
|
||||
return await CreateWorkflowHostAsync(workflowDefinition, cancellationToken);
|
||||
}
|
||||
|
||||
|
||||
private async Task<IWorkflowHost> CreateWorkflowHostAsync(WorkflowDefinition workflowDefinition, CancellationToken cancellationToken)
|
||||
{
|
||||
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
|
||||
|
|
@ -286,7 +290,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
|
|||
runtimeOptions with { BookmarkId = bookmark.BookmarkId },
|
||||
cancellationToken);
|
||||
|
||||
resumedWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks));
|
||||
resumedWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Status, resumeResult.SubStatus, resumeResult.Bookmarks));
|
||||
}
|
||||
|
||||
return resumedWorkflows;
|
||||
|
|
|
|||
Loading…
Reference in a new issue