|
|
|
|
@ -1,4 +1,8 @@
|
|
|
|
|
using Elsa.Models;
|
|
|
|
|
using System.Collections.Generic;
|
|
|
|
|
using System.Linq;
|
|
|
|
|
using System.Threading;
|
|
|
|
|
using System.Threading.Tasks;
|
|
|
|
|
using Elsa.Models;
|
|
|
|
|
using Elsa.Persistence;
|
|
|
|
|
using Elsa.Persistence.Specifications;
|
|
|
|
|
using Elsa.Services;
|
|
|
|
|
@ -24,7 +28,7 @@ namespace Elsa.WorkflowTesting.Services
|
|
|
|
|
_workflowFactory = workflowFactory;
|
|
|
|
|
_workflowRunner = workflowRunner;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
public async Task<RunWorkflowResult?> FindAndRestartTestWorkflowAsync(
|
|
|
|
|
string workflowDefinitionId,
|
|
|
|
|
string activityId,
|
|
|
|
|
@ -35,30 +39,30 @@ namespace Elsa.WorkflowTesting.Services
|
|
|
|
|
CancellationToken cancellationToken = default)
|
|
|
|
|
{
|
|
|
|
|
var workflowBlueprint = await _workflowRegistry.GetAsync(workflowDefinitionId, tenantId, VersionOptions.SpecificVersion(version), cancellationToken);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if (workflowBlueprint == null)
|
|
|
|
|
return null;
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
var lastWorkflowInstance = await _workflowInstanceStore.FindAsync(new EntityIdSpecification<WorkflowInstance>(lastWorkflowInstanceId), cancellationToken);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if (lastWorkflowInstance == null)
|
|
|
|
|
return null;
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
var startActivity = workflowBlueprint.Activities.First(x => x.Id == activityId);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
var startableWorkflowDefinition = new StartableWorkflowDefinition(workflowBlueprint, startActivity.Id);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
var workflow = await InstantiateStartableWorkflow(startableWorkflowDefinition, cancellationToken);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
var previousActivityData = GetActivityDataFromLastWorkflowInstance(workflow.WorkflowInstance, lastWorkflowInstance, workflowBlueprint, activityId);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
MergeActivityDataIntoInstance(workflow.WorkflowInstance, previousActivityData);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
SetMetadata(workflow.WorkflowInstance, signalRConnectionId);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
//if previousActivityOutput has any items, then the first one is from activity closest to the starting one
|
|
|
|
|
var previousActivityOutput = previousActivityData.Count == 0 ? null : previousActivityData.First().Value?.GetItem("Output");
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
return await ExecuteStartableWorkflowAsync(workflow, new WorkflowInput(previousActivityOutput), cancellationToken);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@ -67,7 +71,7 @@ namespace Elsa.WorkflowTesting.Services
|
|
|
|
|
workflowInstance.SetMetadata("isTest", true);
|
|
|
|
|
workflowInstance.SetMetadata("signalRConnectionId", signalRConnectionId);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private void MergeActivityDataIntoInstance(WorkflowInstance workflowInstance, IDictionary<string, IDictionary<string, object?>> activityData)
|
|
|
|
|
{
|
|
|
|
|
foreach (var (key, value) in activityData)
|
|
|
|
|
@ -75,24 +79,24 @@ namespace Elsa.WorkflowTesting.Services
|
|
|
|
|
workflowInstance.ActivityData[key] = value;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private IDictionary<string, IDictionary<string, object?>> GetActivityDataFromLastWorkflowInstance(WorkflowInstance currentWorkflowInstance, WorkflowInstance lastWorkflowInstance, IWorkflowBlueprint workflowBlueprint, string startingActivityId)
|
|
|
|
|
{
|
|
|
|
|
IDictionary<string, IDictionary<string, object?>> CollectSourceActivityData(string targetActivityId, IDictionary<string, IDictionary<string, object?>> activityDataAccumulator)
|
|
|
|
|
{
|
|
|
|
|
var sourceActivityId = workflowBlueprint.Connections.FirstOrDefault(x => x.Target.Activity.Id == targetActivityId)?.Source.Activity.Id;
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if (sourceActivityId == null)
|
|
|
|
|
return activityDataAccumulator;
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
activityDataAccumulator.Add(sourceActivityId, lastWorkflowInstance.ActivityData.GetItem(sourceActivityId));
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
return CollectSourceActivityData(sourceActivityId, activityDataAccumulator);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
return CollectSourceActivityData(startingActivityId, new Dictionary<string, IDictionary<string, object?>>());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private async Task<StartableWorkflow> InstantiateStartableWorkflow(StartableWorkflowDefinition startableWorkflowDefinition, CancellationToken cancellationToken)
|
|
|
|
|
{
|
|
|
|
|
var workflowInstance = await _workflowFactory.InstantiateAsync(
|
|
|
|
|
@ -100,11 +104,11 @@ namespace Elsa.WorkflowTesting.Services
|
|
|
|
|
startableWorkflowDefinition.CorrelationId,
|
|
|
|
|
startableWorkflowDefinition.ContextId,
|
|
|
|
|
cancellationToken: cancellationToken);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
await _workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken);
|
|
|
|
|
return new StartableWorkflow(startableWorkflowDefinition.WorkflowBlueprint, workflowInstance, startableWorkflowDefinition.ActivityId);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private async Task<RunWorkflowResult> ExecuteStartableWorkflowAsync(StartableWorkflow startableWorkflow, WorkflowInput? input, CancellationToken cancellationToken = default) =>
|
|
|
|
|
await _workflowRunner.RunWorkflowAsync(startableWorkflow.WorkflowBlueprint, startableWorkflow.WorkflowInstance, startableWorkflow.ActivityId, input, cancellationToken);
|
|
|
|
|
}
|
|
|
|
|
|