Fix Race Condition in WorkflowDefinitionActivity (#5315)

* Refactor hash method in Hasher.cs

The Hash method in the Hasher.cs file has been refactored to use more efficient code. Instead of creating a SHA256 instance, it directly hashes the data using SHA256.HashData. As a result, the private Hash method and its usage, which only has relevance to the previous approach, have been removed.

* Replace GetHashCode with SHA256 for script caching

Two unused namespaces, Esprima and Esprima.Ast, have been removed, and two new ones, System.Security.Cryptography and System.Text, have been added. These changes in Elsa.JavaScript's JintJavaScriptEvaluator file were made to replace the basic GetHashCode method, which was used earlier to generate a cache key for JavaScript scripts, with SHA256 hash function to ensure uniqueness and avoid possible hash collisions. This greatly increases the reliability of the caching mechanism.

* Set HTTP client timeout in WorkflowServer

The client timeout in WorkflowServer has been set to 1 minute for `ResrService` and `client` methods. This was done to manage long running requests and prevent timeouts.

* Refactor code and add validation in ScheduleActivityAsync methods

An unnecessary line of code in the ScheduleActivityAsync method under the ActivityExecutionContext file has been removed to simplify the function. Meanwhile, validation has been added to ensure that specified activities are part of the workflow. This improves code clarity and prevents potential errors caused by incorrect activity scheduling.

* Refactor ActivityVisitor to use ActivityVisitorContext

The ActivityVisitor class in Elsa.Workflows.Core was refactored to use an ActivityVisitorContext class. The introduction of this context class replaced the multiple hashsets that were used within function signatures, consolidating them into a single object and reducing function complexity. This change improves readability and code organization.

* Refactor WorkflowDefinitionActivity to use WorkflowGraph

Adjusted the WorkflowDefinitionActivity class in the Elsa.Workflows.Management module. The changes include migration from using Workflow to WorkflowGraph objects and condensing code blocks for clearer readability. Additionally, a new 'IsInitialized' field is added to eliminate potential race conditions during the graph construction process.

* Add activity validation to WorkflowExecutionContext

A check has been implemented in the WorkflowExecutionContextExtensions to validate that a specified activity is part of the workflow. This prevents incorrect activity references when scheduling tasks in the workflow.

* Refactor to use WorkflowGraph instead of Workflow

The codebase has been refactored to utilize the WorkflowGraph instead of the Workflow while running and manipulating workflows. Additional changes include restructuring WorkflowRunner and WorkflowDefinitionService classes, introducing WorkflowGraphBuilder usage, and mapping updates in WorkflowDefinitionMapper. The WorkflowHost, WorkflowDispatcher, and WorkflowExecutionContext have also been updated accordingly.

* Add WorkflowGraph and related services

This commit introduces the IWorkflowGraphBuilder interface, the WorkflowGraph model, and an implementation of the interface in the WorkflowGraphBuilder class. The purpose of these additions is to establish the building and structure of a workflow graph. The WorkflowGraph model also includes activity node handling and hashing capabilities.

* Add tests to ensure exception when scheduling an activity not part of workflow

Several unit tests were written to ensure that the right behavior is exhibited when scheduling an activity that is not part of the workflow. An exception is expected to be thrown in this case. Additionally, new service definitions, workflow definitions and workflow queries were added for more comprehensive testing.

* Update WorkflowGraph class and add comment descriptions

This commit updates the WorkflowGraph class by extending its descriptions and implementing an explicit mention to the Workflow reference. Additionally, more attribute descriptions have been added to increase code readability and comprehension.

* Removed unnecessary brackets

* Remove unused services from workflow management

This commit deletes "ExpressionDescriptorRegistryPopulator.cs" and "ScopedWorkflowDefinitionLookup.cs" files from Elsa.Workflows.Management.Services. These files, containing obsolete services, are no longer used in the workflow management process. The services' registration has been removed as well from "WorkflowManagementFeature.cs".

* Remove TestWorkflowDefinitionService.cs from component tests

A file, TestWorkflowDefinitionService.cs, was removed under Elsa.Workflows.ComponentTests. This is part of the improvement process where inefficient or unnecessary test files are cleaned up.

* Add comments to WorkflowDefinitionActivityTests

This commit adds explanatory comments to the WorkflowDefinitionActivityTests file. The comments provide information about the purpose of these tests and a reference to a related issue in the project's Github repository.

* Updated duplicate and missing package references

---------

Co-authored-by: Raymond den Haan <raymond.den.haan@nexxbiz.io>
This commit is contained in:
Sipke Schoorstra 2024-05-01 14:53:37 +02:00 committed by GitHub
parent e966f62f29
commit b6689c7e2b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
51 changed files with 799 additions and 352 deletions

View file

@ -91,6 +91,7 @@
<PackageVersion Include="System.Linq.Async" Version="6.0.1" />
<PackageVersion Include="System.Linq.Dynamic.Core" Version="1.3.13" />
<PackageVersion Include="Testcontainers" Version="3.8.0" />
<PackageVersion Include="Testcontainers.RabbitMq" Version="3.8.0" />
<PackageVersion Include="Testcontainers.Redis" Version="3.8.0" />
<PackageVersion Include="Testcontainers.PostgreSql" Version="3.8.0" />
<PackageVersion Include="ThrottleDebounce" Version="2.0.0" />

View file

@ -1,6 +1,7 @@
using System.Text.Encodings.Web;
using Elsa.Alterations.Extensions;
using Elsa.Alterations.MassTransit.Extensions;
using Elsa.Caching.Options;
using Elsa.Common.DistributedLocks.Noop;
using Elsa.Dapper.Extensions;
using Elsa.Dapper.Services;
@ -358,6 +359,8 @@ services
ConfigureForTest?.Invoke(elsa);
});
services.Configure<CachingOptions>(options => options.CacheDuration = TimeSpan.FromDays(1));
services.AddHealthChecks();
services.AddControllers();
services.AddCors(cors => cors.AddDefaultPolicy(policy => policy.AllowAnyHeader().AllowAnyMethod().AllowAnyOrigin().WithExposedHeaders("*")));

View file

@ -30,7 +30,7 @@ public class MigrateHandler : AlterationHandlerBase<Migrate>
}
var targetWorkflow = await workflowDefinitionService.MaterializeWorkflowAsync(targetWorkflowDefinition, cancellationToken);
await context.WorkflowExecutionContext.SetWorkflowAsync(targetWorkflow);
await context.WorkflowExecutionContext.SetWorkflowGraphAsync(targetWorkflow);
context.Succeed();
}

View file

@ -61,29 +61,24 @@ public class DefaultAlterationRunner : IAlterationRunner
{
var log = new AlterationLog(_systemClock);
var result = new RunAlterationsResult(workflowInstanceId, log);
// Load workflow instance.
var workflowState = await _workflowRuntime.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken);
// If the workflow instance is not found, log an error and continue.
if (workflowState == null)
{
log.Add($"Workflow instance with ID '{workflowInstanceId}' not found.", LogLevel.Error);
return result;
}
// Load workflow definition.
var workflow = await _workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionVersionId, cancellationToken);
var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionVersionId, cancellationToken);
// If the workflow definition is not found, log an error and continue.
if (workflow == null)
if (workflowGraph == null)
{
log.Add($"Workflow definition with ID '{workflowState.DefinitionVersionId}' not found.", LogLevel.Error);
return result;
}
// Create workflow execution context.
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflow, workflowState, cancellationTokens: cancellationToken);
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflowGraph, workflowState, cancellationTokens: cancellationToken);
workflowExecutionContext.TransientProperties.Add(RunAlterationsMiddleware.AlterationsPropertyKey, alterations);
workflowExecutionContext.TransientProperties.Add(RunAlterationsMiddleware.AlterationsLogPropertyKey, log);

View file

@ -87,8 +87,8 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions<HttpActivity
var trigger = triggers?.FirstOrDefault();
if (trigger != null)
{
var workflow = lookupResult.Workflow!;
await StartWorkflowAsync(httpContext, trigger, workflow, input);
var workflowGraph = lookupResult.WorkflowGraph!;
await StartWorkflowAsync(httpContext, trigger, workflowGraph, input);
return;
}
}
@ -123,20 +123,20 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions<HttpActivity
await next(httpContext);
}
private async Task<Workflow?> FindWorkflowAsync(IServiceProvider serviceProvider, StoredTrigger trigger, CancellationToken cancellationToken)
private async Task<WorkflowGraph?> FindWorkflowGraphAsync(IServiceProvider serviceProvider, StoredTrigger trigger, CancellationToken cancellationToken)
{
var workflowDefinitionService = serviceProvider.GetRequiredService<IWorkflowDefinitionService>();
var workflowDefinitionId = trigger.WorkflowDefinitionVersionId;
return await workflowDefinitionService.FindWorkflowAsync(workflowDefinitionId, cancellationToken);
return await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, cancellationToken);
}
private async Task StartWorkflowAsync(HttpContext httpContext, StoredTrigger trigger, Workflow workflow, IDictionary<string, object> input)
private async Task StartWorkflowAsync(HttpContext httpContext, StoredTrigger trigger, WorkflowGraph workflowGraph, IDictionary<string, object> input)
{
var serviceProvider = httpContext.RequestServices;
var cancellationToken = httpContext.RequestAborted;
var bookmarkPayload = trigger.GetPayload<HttpEndpointBookmarkPayload>();
var workflowHostFactory = serviceProvider.GetRequiredService<IWorkflowHostFactory>();
var workflowHost = await workflowHostFactory.CreateAsync(workflow, cancellationToken);
var workflowHost = await workflowHostFactory.CreateAsync(workflowGraph, cancellationToken);
if (await AuthorizeAsync(serviceProvider, httpContext, workflowHost.Workflow, bookmarkPayload, cancellationToken))
return;
@ -175,7 +175,7 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions<HttpActivity
}
var workflowDefinitionService = serviceProvider.GetRequiredService<IWorkflowDefinitionService>();
var workflow = await workflowDefinitionService.FindWorkflowAsync(workflowInstance.DefinitionVersionId, cancellationToken);
var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(workflowInstance.DefinitionVersionId, cancellationToken);
if (workflow == null)
{

View file

@ -1,6 +1,9 @@
using Elsa.Workflows.Activities;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Http.Models;
public record HttpWorkflowLookupResult(Workflow? Workflow, ICollection<StoredTrigger> Triggers);
/// <summary>
/// Represents the result of a workflow lookup.
/// </summary>
public record HttpWorkflowLookupResult(WorkflowGraph? WorkflowGraph, ICollection<StoredTrigger> Triggers);

View file

@ -29,8 +29,8 @@ public class CachingHttpWorkflowLookupService(
if (result == null)
return null;
var workflow = result.Workflow!;
var changeTokenKey = cacheManager.GetWorkflowChangeTokenKey(workflow.Identity.DefinitionId);
var workflowGraph = result.WorkflowGraph!;
var changeTokenKey = cacheManager.GetWorkflowChangeTokenKey(workflowGraph.Workflow.Identity.DefinitionId);
var changeToken = cache.GetToken(changeTokenKey);
entry.AddExpirationToken(changeToken);

View file

@ -80,7 +80,7 @@ public class HttpBookmarkProcessor : IHttpBookmarkProcessor
continue;
}
var workflow = await _workflowDefinitionService.FindWorkflowAsync(
var workflow = await _workflowDefinitionService.FindWorkflowGraphAsync(
workflowState.DefinitionId,
VersionOptions.SpecificVersion(workflowState.DefinitionVersion),
systemCancellationToken);

View file

@ -1,7 +1,7 @@
using Elsa.Http.Contracts;
using Elsa.Http.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Filters;
@ -25,12 +25,12 @@ public class HttpWorkflowLookupService(ITriggerStore triggerStore, IWorkflowDefi
if (trigger == null)
return default;
var workflow = await FindWorkflowAsync(trigger, cancellationToken);
var workflowGraph = await FindWorkflowGraphAsync(trigger, cancellationToken);
if (workflow == null)
if (workflowGraph == null)
return default;
return new(workflow, triggers);
return new(workflowGraph, triggers);
}
private async Task<IEnumerable<StoredTrigger>> FindTriggersAsync(string bookmarkHash, CancellationToken cancellationToken)
@ -42,9 +42,9 @@ public class HttpWorkflowLookupService(ITriggerStore triggerStore, IWorkflowDefi
return await triggerStore.FindManyAsync(triggerFilter, cancellationToken);
}
private async Task<Workflow?> FindWorkflowAsync(StoredTrigger trigger, CancellationToken cancellationToken)
private async Task<WorkflowGraph?> FindWorkflowGraphAsync(StoredTrigger trigger, CancellationToken cancellationToken)
{
var workflowDefinitionVersionId = trigger.WorkflowDefinitionVersionId;
return await workflowDefinitionService.FindWorkflowAsync(workflowDefinitionVersionId, cancellationToken);
return await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionVersionId, cancellationToken);
}
}

View file

@ -3,8 +3,8 @@ using Elsa.Abstractions;
using Elsa.Common.Models;
using Elsa.JavaScript.TypeDefinitions.Contracts;
using Elsa.JavaScript.TypeDefinitions.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Models;
using JetBrains.Annotations;
namespace Elsa.JavaScript.Endpoints.TypeDefinitions;
@ -35,15 +35,16 @@ internal class Get : ElsaEndpoint<Request>
/// <inheritdoc />
public override async Task HandleAsync(Request request, CancellationToken cancellationToken)
{
var workflow = await GetWorkflowAsync(request.WorkflowDefinitionId, cancellationToken);
var workflowGraph = await GetWorkflowGraphAsync(request.WorkflowDefinitionId, cancellationToken);
if (workflow == null)
if (workflowGraph == null)
{
AddError($"Workflow definition {request.WorkflowDefinitionId} not found");
await SendErrorsAsync(cancellation: cancellationToken);
return;
}
var workflow = workflowGraph.Workflow;
var typeDefinitionContext = new TypeDefinitionContext(workflow, request.ActivityTypeName, request.PropertyName, cancellationToken);
var typeDefinitions = await _typeDefinitionService.GenerateTypeDefinitionsAsync(typeDefinitionContext);
var fileName = $"elsa.{request.WorkflowDefinitionId}.d.ts";
@ -52,9 +53,9 @@ internal class Get : ElsaEndpoint<Request>
await SendBytesAsync(data, fileName, "application/x-typescript", cancellation: cancellationToken);
}
private async Task<Workflow?> GetWorkflowAsync(string workflowDefinitionId, CancellationToken cancellationToken)
private async Task<WorkflowGraph?> GetWorkflowGraphAsync(string workflowDefinitionId, CancellationToken cancellationToken)
{
return await _workflowDefinitionService.FindWorkflowAsync(workflowDefinitionId, VersionOptions.Latest, cancellationToken);
return await _workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, VersionOptions.Latest, cancellationToken);
}
}

View file

@ -1,4 +1,6 @@
using System.Diagnostics.CodeAnalysis;
using System.Security.Cryptography;
using System.Text;
using System.Text.Encodings.Web;
using System.Text.Json;
using System.Text.Json.Serialization;
@ -12,8 +14,6 @@ using Elsa.JavaScript.Helpers;
using Elsa.JavaScript.Notifications;
using Elsa.JavaScript.Options;
using Elsa.Mediator.Contracts;
using Esprima;
using Esprima.Ast;
using Humanizer;
using Jint;
using Jint.Runtime.Interop;
@ -182,7 +182,7 @@ public class JintJavaScriptEvaluator : IJavaScriptEvaluator
private object? ExecuteExpressionAndGetResult(Engine engine, string expression)
{
var cacheKey = "jint:script:" + expression.GetHashCode(StringComparison.Ordinal);
var cacheKey = "jint:script:" + Hash(expression);
var parsedScript = _memoryCache.GetOrCreate(cacheKey, entry =>
{
@ -214,4 +214,11 @@ public class JintJavaScriptEvaluator : IJavaScriptEvaluator
return JsonSerializer.Serialize(value, options);
}
private string Hash(string input)
{
var bytes = Encoding.UTF8.GetBytes(input);
var hash = SHA256.HashData(bytes);
return Convert.ToBase64String(hash);
}
}

View file

@ -32,11 +32,12 @@ public class MassTransitWorkflowDispatcher(
/// <inheritdoc />
public async Task<DispatchWorkflowResponse> DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
var workflow = await workflowDefinitionService.FindWorkflowAsync(request.DefinitionId, request.VersionOptions, cancellationToken);
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(request.DefinitionId, request.VersionOptions, cancellationToken);
if (workflow == null)
if (workflowGraph == null)
throw new Exception($"Workflow definition with definition ID '{request.DefinitionId} and version {request.VersionOptions}' not found");
var workflow = workflowGraph.Workflow;
var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest
{
Workflow = workflow,
@ -106,9 +107,9 @@ public class MassTransitWorkflowDispatcher(
foreach (var trigger in triggers)
{
var workflow = await workflowDefinitionService.FindWorkflowAsync(trigger.WorkflowDefinitionVersionId, cancellationToken);
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(trigger.WorkflowDefinitionVersionId, cancellationToken);
if (workflow == null)
if (workflowGraph == null)
{
logger.LogWarning("Workflow definition with ID '{WorkflowDefinitionId}' not found", trigger.WorkflowDefinitionVersionId);
continue;
@ -116,7 +117,7 @@ public class MassTransitWorkflowDispatcher(
var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest
{
Workflow = workflow,
Workflow = workflowGraph.Workflow,
WorkflowInstanceId = request.WorkflowInstanceId,
Input = request.Input,
Properties = request.Properties,

View file

@ -82,12 +82,14 @@ internal class WorkflowInstance : WorkflowInstanceBase
using var scope = _scopeFactory.CreateScope();
var workflowDefinitionService = scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionService>();
// Load the workflow definition.
var workflow = await workflowDefinitionService.FindWorkflowAsync(_definitionId, VersionOptions.SpecificVersion(_version), cancellationToken);
// Get the workflow.
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(_definitionId, VersionOptions.SpecificVersion(_version), cancellationToken);
if (workflow == null)
if (workflowGraph == null)
throw new Exception("Workflow definition is no longer available");
var workflow = workflowGraph.Workflow;
// Create an initial workflow state.
if (_workflowState == null!)
{
@ -100,7 +102,7 @@ internal class WorkflowInstance : WorkflowInstanceBase
// Create a workflow host.
var workflowHostFactory = scope.ServiceProvider.GetRequiredService<IWorkflowHostFactory>();
_workflowHost = await workflowHostFactory.CreateAsync(workflow, _workflowState, cancellationToken);
_workflowHost = await workflowHostFactory.CreateAsync(workflowGraph, _workflowState, cancellationToken);
}
/// <inheritdoc />
@ -416,7 +418,7 @@ internal class WorkflowInstance : WorkflowInstanceBase
{
using var scope = _scopeFactory.CreateScope();
var workflowDefinitionService = scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionService>();
var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken);
var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken);
if (workflow == null)
throw new Exception("Specified workflow definition and version does not exist");
@ -438,7 +440,7 @@ internal class WorkflowInstance : WorkflowInstanceBase
var versionOptions = VersionOptions.SpecificVersion(workflowState.DefinitionVersion);
using var scope = _scopeFactory.CreateScope();
var workflowDefinitionService = scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionService>();
var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken);
var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken);
if (workflow == null)
throw new Exception("Specified workflow definition and version does not exist");

View file

@ -386,14 +386,15 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime
};
var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions);
var workflow = await _workflowDefinitionService.FindWorkflowAsync(trigger.WorkflowDefinitionVersionId, cancellationToken);
var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(trigger.WorkflowDefinitionVersionId, cancellationToken);
if (workflow == null)
if (workflowGraph == null)
{
_logger.LogWarning("Workflow version ID {DefinitionVersionId} not found", trigger.WorkflowDefinitionVersionId);
continue;
}
var workflow = workflowGraph.Workflow;
var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest
{
Workflow = workflow,

View file

@ -224,8 +224,7 @@ public partial class ActivityExecutionContext : IExecutionContext
/// <param name="options">The options used to schedule the activity.</param>
public async ValueTask ScheduleActivityAsync(IActivity? activity, ScheduleWorkOptions? options = default)
{
var activityNode = activity != null ? WorkflowExecutionContext.FindNodeByActivity(activity) : default;
await ScheduleActivityAsync(activityNode, this, options);
await ScheduleActivityAsync(activity, this, options);
}
/// <summary>
@ -236,7 +235,9 @@ public partial class ActivityExecutionContext : IExecutionContext
/// <param name="options">The options used to schedule the activity.</param>
public async ValueTask ScheduleActivityAsync(IActivity? activity, ActivityExecutionContext? owner, ScheduleWorkOptions? options = default)
{
var activityNode = activity != null ? WorkflowExecutionContext.FindNodeByActivity(activity) : default;
var activityNode = activity != null
? WorkflowExecutionContext.FindNodeByActivity(activity) ?? throw new InvalidOperationException("The specified activity is not part of the workflow.")
: null;
await ScheduleActivityAsync(activityNode, owner, options);
}
@ -250,6 +251,10 @@ public partial class ActivityExecutionContext : IExecutionContext
{
if (this.GetIsBackgroundExecution())
{
// Validate that the specified activity is part of the workflow.
if (activityNode != null && !WorkflowExecutionContext.NodeActivityLookup.ContainsKey(activityNode.Activity))
throw new InvalidOperationException("The specified activity is not part of the workflow.");
var scheduledActivity = new ScheduledActivity
{
ActivityNodeId = activityNode?.NodeId,

View file

@ -44,6 +44,7 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// </summary>
private WorkflowExecutionContext(
IServiceProvider serviceProvider,
WorkflowGraph workflowGraph,
string id,
string? correlationId,
string? parentWorkflowInstanceId,
@ -75,7 +76,7 @@ public partial class WorkflowExecutionContext : IExecutionContext
UpdatedAt = createdAt;
CancellationTokens = cancellationTokens;
Incidents = incidents.ToList();
WorkflowGraph = workflowGraph;
var appSource = CancellationTokenSource.CreateLinkedTokenSource(CancellationTokens.ApplicationCancellationToken);
_cancellationTokenSources.Add(appSource);
var sysSource = CancellationTokenSource.CreateLinkedTokenSource(CancellationTokens.SystemCancellationToken);
@ -89,22 +90,22 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// </summary>
public static async Task<WorkflowExecutionContext> CreateAsync(
IServiceProvider serviceProvider,
Workflow workflow,
WorkflowGraph workflowGraph,
string id,
string? correlationId,
string? parentWorkflowInstanceId = default,
IDictionary<string, object>? input = default,
IDictionary<string, object>? properties = default,
ExecuteActivityDelegate? executeDelegate = default,
string? triggerActivityId = default,
Action<WorkflowExecutionContext>? statusUpdatedCallback = default,
string? correlationId = null,
string? parentWorkflowInstanceId = null,
IDictionary<string, object>? input = null,
IDictionary<string, object>? properties = null,
ExecuteActivityDelegate? executeDelegate = null,
string? triggerActivityId = null,
Action<WorkflowExecutionContext>? statusUpdatedCallback = null,
CancellationTokens cancellationTokens = default)
{
var systemClock = serviceProvider.GetRequiredService<ISystemClock>();
return await CreateAsync(
serviceProvider,
workflow,
workflowGraph,
id,
new List<ActivityIncident>(),
systemClock.UtcNow,
@ -124,20 +125,20 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// </summary>
public static async Task<WorkflowExecutionContext> CreateAsync(
IServiceProvider serviceProvider,
Workflow workflow,
WorkflowGraph workflowGraph,
WorkflowState workflowState,
string? correlationId = default,
string? parentWorkflowInstanceId = default,
IDictionary<string, object>? input = default,
IDictionary<string, object>? properties = default,
ExecuteActivityDelegate? executeDelegate = default,
string? triggerActivityId = default,
Action<WorkflowExecutionContext>? statusUpdatedCallback = default,
string? correlationId = null,
string? parentWorkflowInstanceId = null,
IDictionary<string, object>? input = null,
IDictionary<string, object>? properties = null,
ExecuteActivityDelegate? executeDelegate = null,
string? triggerActivityId = null,
Action<WorkflowExecutionContext>? statusUpdatedCallback = null,
CancellationTokens cancellationTokens = default)
{
var workflowExecutionContext = await CreateAsync(
serviceProvider,
workflow,
workflowGraph,
workflowState.Id,
workflowState.Incidents,
workflowState.CreatedAt,
@ -161,22 +162,23 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// </summary>
public static async Task<WorkflowExecutionContext> CreateAsync(
IServiceProvider serviceProvider,
Workflow workflow,
WorkflowGraph workflowGraph,
string id,
IEnumerable<ActivityIncident> incidents,
DateTimeOffset createdAt,
string? correlationId = default,
string? parentWorkflowInstanceId = default,
IDictionary<string, object>? input = default,
IDictionary<string, object>? properties = default,
ExecuteActivityDelegate? executeDelegate = default,
string? triggerActivityId = default,
Action<WorkflowExecutionContext>? statusUpdatedCallback = default,
string? correlationId = null,
string? parentWorkflowInstanceId = null,
IDictionary<string, object>? input = null,
IDictionary<string, object>? properties = null,
ExecuteActivityDelegate? executeDelegate = null,
string? triggerActivityId = null,
Action<WorkflowExecutionContext>? statusUpdatedCallback = null,
CancellationTokens cancellationTokens = default)
{
// Setup a workflow execution context.
var workflowExecutionContext = new WorkflowExecutionContext(
serviceProvider,
workflowGraph,
id,
correlationId,
parentWorkflowInstanceId,
@ -189,47 +191,31 @@ public partial class WorkflowExecutionContext : IExecutionContext
statusUpdatedCallback,
cancellationTokens)
{
MemoryRegister = workflow.CreateRegister()
MemoryRegister = workflowGraph.Workflow.CreateRegister()
};
workflowExecutionContext.ExpressionExecutionContext = new ExpressionExecutionContext(serviceProvider, workflowExecutionContext.MemoryRegister, cancellationToken: cancellationTokens.ApplicationCancellationToken);
await workflowExecutionContext.SetWorkflowAsync(workflow);
await workflowExecutionContext.SetWorkflowGraphAsync(workflowGraph);
return workflowExecutionContext;
}
/// <summary>
/// Assigns the specified workflow to this workflow execution context.
/// </summary>
/// <param name="workflow">The workflow to assign.</param>
public async Task SetWorkflowAsync(Workflow workflow)
/// <param name="workflowGraph">The workflow graph to assign.</param>
public async Task SetWorkflowGraphAsync(WorkflowGraph workflowGraph)
{
var activityVisitor = GetRequiredService<IActivityVisitor>();
var root = workflow;
var graph = await activityVisitor.VisitAsync(root, CancellationTokens.ApplicationCancellationToken);
var nodes = graph.Flatten().ToList();
WorkflowGraph = workflowGraph;
var nodes = workflowGraph.Nodes;
// Register activity types.
var activityTypes = nodes.Select(x => x.Activity.GetType()).Distinct().ToList();
await ActivityRegistry.RegisterAsync(activityTypes, CancellationTokens.ApplicationCancellationToken);
var needsIdentityAssignment = nodes.Any(x => string.IsNullOrEmpty(x.Activity.Id));
if (needsIdentityAssignment)
{
var identityGraphService = GetRequiredService<IIdentityGraphService>();
identityGraphService.AssignIdentities(nodes);
}
Workflow = workflow;
Graph = graph;
Nodes = nodes;
NodeIdLookup = nodes.ToDictionary(x => x.NodeId);
NodeHashLookup = nodes.ToDictionary(x => Hash(x.NodeId));
NodeActivityLookup = nodes.ToDictionary(x => x.Activity);
foreach (var activityExecutionContext in ActivityExecutionContexts)
activityExecutionContext.Activity = NodeIdLookup[activityExecutionContext.Activity.NodeId].Activity;
// Update the activity execution contexts with the actual activity instances.
foreach (var activityExecutionContext in ActivityExecutionContexts)
activityExecutionContext.Activity = workflowGraph.NodeIdLookup[activityExecutionContext.Activity.NodeId].Activity;
}
/// <summary>
@ -242,15 +228,20 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// </summary>
public IActivityRegistry ActivityRegistry { get; }
/// <summary>
/// Gets the workflow graph.
/// </summary>
public WorkflowGraph WorkflowGraph { get; private set; }
/// <summary>
/// The <see cref="Workflow"/> associated with the execution context.
/// </summary>
public Workflow Workflow { get; private set; } = default!;
public Workflow Workflow => WorkflowGraph.Workflow;
/// <summary>
/// A graph of the workflow structure.
/// </summary>
public ActivityNode Graph { get; private set; } = default!;
public ActivityNode Graph => WorkflowGraph.Root;
/// <summary>
/// The current status of the workflow.
@ -279,7 +270,7 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// An application-specific identifier associated with the execution context.
/// </summary>
public string? CorrelationId { get; set; }
/// <summary>
/// The ID of the workflow instance that triggered this instance.
/// </summary>
@ -289,12 +280,12 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// The date and time the workflow execution context was created.
/// </summary>
public DateTimeOffset CreatedAt { get; set; }
/// <summary>
/// The date and time the workflow execution context was last updated.
/// </summary>
public DateTimeOffset UpdatedAt { get; set; }
/// <summary>
/// The date and time the workflow execution context has finished.
/// </summary>
@ -308,22 +299,22 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// <summary>
/// A flattened list of <see cref="ActivityNode"/>s from the <see cref="Graph"/>.
/// </summary>
public IReadOnlyCollection<ActivityNode> Nodes { get; private set; } = default!;
public IReadOnlyCollection<ActivityNode> Nodes => WorkflowGraph.Nodes.ToList();
/// <summary>
/// A map between activity IDs and <see cref="ActivityNode"/>s in the workflow graph.
/// </summary>
public IDictionary<string, ActivityNode> NodeIdLookup { get; private set; } = default!;
public IDictionary<string, ActivityNode> NodeIdLookup => WorkflowGraph.NodeIdLookup;
/// <summary>
/// A map between hashed activity node IDs and <see cref="ActivityNode"/>s in the workflow graph.
/// </summary>
public IDictionary<string, ActivityNode> NodeHashLookup { get; private set; } = default!;
public IDictionary<string, ActivityNode> NodeHashLookup => WorkflowGraph.NodeHashLookup;
/// <summary>
/// A map between <see cref="IActivity"/>s and <see cref="ActivityNode"/>s in the workflow graph.
/// </summary>
public IDictionary<IActivity, ActivityNode> NodeActivityLookup { get; private set; } = default!;
public IDictionary<IActivity, ActivityNode> NodeActivityLookup => WorkflowGraph.NodeActivityLookup;
/// <summary>
/// The <see cref="IActivityScheduler"/> for the execution context.
@ -559,15 +550,15 @@ public partial class WorkflowExecutionContext : IExecutionContext
SubStatus = subStatus;
UpdatedAt = SystemClock.UtcNow;
if (Status == WorkflowStatus.Finished)
FinishedAt = UpdatedAt;
//For now only trigger on Cancelled, since the other statuses are handling via the host/runner
if (SubStatus == WorkflowSubStatus.Cancelled
&& _statusUpdatedCallback is not null)
_statusUpdatedCallback(this);
if (Status == WorkflowStatus.Finished
|| SubStatus == WorkflowSubStatus.Suspended)
{
@ -594,7 +585,10 @@ public partial class WorkflowExecutionContext : IExecutionContext
var id = IdentityGenerator.GenerateId();
var activityExecutionContext = new ActivityExecutionContext(id, this, parentContext, expressionExecutionContext, activity, activityDescriptor, now, tag, SystemClock, CancellationTokens.ApplicationCancellationToken);
var variablesToDeclare = options?.Variables ?? Array.Empty<Variable>();
var variableContainer = new[] { activityExecutionContext.ActivityNode }.Concat(activityExecutionContext.ActivityNode.Ancestors()).FirstOrDefault(x => x.Activity is IVariableContainer)?.Activity as IVariableContainer;
var variableContainer = new[]
{
activityExecutionContext.ActivityNode
}.Concat(activityExecutionContext.ActivityNode.Ancestors()).FirstOrDefault(x => x.Activity is IVariableContainer)?.Activity as IVariableContainer;
expressionExecutionContext.TransientProperties[ExpressionExecutionContextExtensions.ActivityExecutionContextKey] = activityExecutionContext;
if (variableContainer != null)
@ -674,7 +668,4 @@ public partial class WorkflowExecutionContext : IExecutionContext
var currentMainStatus = GetMainStatus(SubStatus);
return currentMainStatus != WorkflowStatus.Finished;
}
private string Hash(string nodeId) => _hasher.Hash(nodeId);
}

View file

@ -0,0 +1,15 @@
using Elsa.Workflows.Activities;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Contracts;
/// <summary>
/// Builds a workflow graph from a workflow.
/// </summary>
public interface IWorkflowGraphBuilder
{
/// <summary>
/// Builds a workflow graph from a workflow.
/// </summary>
Task<WorkflowGraph> BuildAsync(Workflow workflow, CancellationToken cancellationToken = default);
}

View file

@ -17,5 +17,7 @@ public interface IWorkflowRunner
Task<TResult> RunAsync<T, TResult>(RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) where T : WorkflowBase<TResult>, new();
Task<RunWorkflowResult> RunAsync(Workflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default);
Task<RunWorkflowResult> RunAsync(Workflow workflow, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default);
Task<RunWorkflowResult> RunAsync(WorkflowGraph workflowGraph, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default);
Task<RunWorkflowResult> RunAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default);
Task<RunWorkflowResult> RunAsync(WorkflowExecutionContext workflowExecutionContext);
}

View file

@ -99,6 +99,10 @@ public static class WorkflowExecutionContextExtensions
ActivityExecutionContext owner,
ScheduleWorkOptions? options = default)
{
// Validate that the specified activity is part of the workflow.
if (!workflowExecutionContext.NodeActivityLookup.ContainsKey(activityNode.Activity))
throw new InvalidOperationException("The specified activity is not part of the workflow.");
var scheduler = workflowExecutionContext.Scheduler;
if (options?.PreventDuplicateScheduling == true)

View file

@ -127,6 +127,7 @@ public class WorkflowsFeature : FeatureBase
.AddScoped<IWorkflowRunner, WorkflowRunner>()
.AddScoped<IActivityVisitor, ActivityVisitor>()
.AddScoped<IIdentityGraphService, IdentityGraphService>()
.AddScoped<IWorkflowGraphBuilder, WorkflowGraphBuilder>()
.AddScoped<IWorkflowStateExtractor, WorkflowStateExtractor>()
.AddScoped<IActivitySchedulerFactory, ActivitySchedulerFactory>()
.AddSingleton<IHasher, Hasher>()

View file

@ -0,0 +1,62 @@
using System.Security.Cryptography;
using System.Text;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
namespace Elsa.Workflows.Models;
/// <summary>
/// Hold a reference to a <see cref="Workflow"/> and a collection of <see cref="ActivityNode"/> instances representing the workflow.
/// </summary>
public record WorkflowGraph
{
/// <summary>
/// Initializes a new instance of the <see cref="WorkflowGraph"/> class.
/// </summary>
public WorkflowGraph(Workflow workflow, ActivityNode root, IEnumerable<ActivityNode> nodes)
{
using var hashAlgorithm = SHA256.Create();
Workflow = workflow;
Root = root;
Nodes = nodes.ToList();
NodeIdLookup = Nodes.ToDictionary(x => x.NodeId);
NodeHashLookup = Nodes.ToDictionary(x => Hash(hashAlgorithm, x.NodeId));
NodeActivityLookup = Nodes.ToDictionary(x => x.Activity);
}
/// <summary>
/// Gets the workflow.
/// </summary>
public Workflow Workflow { get; }
/// <summary>
/// Gets the root node.
/// </summary>
public ActivityNode Root { get; }
/// <summary>
/// Gets a flat collection of all nodes in the workflow.
/// </summary>
public ICollection<ActivityNode> Nodes { get; }
/// <summary>
/// Gets a lookup of nodes by their activity.
/// </summary>
public IDictionary<IActivity, ActivityNode> NodeActivityLookup { get; }
/// <summary>
/// Gets a lookup of nodes by their hash.
/// </summary>
public IDictionary<string, ActivityNode> NodeHashLookup { get; }
/// <summary>
/// Gets a lookup of nodes by their ID.
/// </summary>
public IDictionary<string, ActivityNode> NodeIdLookup { get; }
private static string Hash(HashAlgorithm hashAlgorithm, string input)
{
var data = hashAlgorithm.ComputeHash(Encoding.UTF8.GetBytes(input));
return Convert.ToHexString(data);
}
}

View file

@ -10,7 +10,7 @@ public class ActivityVisitor : IActivityVisitor
private readonly IEnumerable<IActivityResolver> _portResolvers;
/// <summary>
/// Constructor.
/// Initializes a new instance of the <see cref="ActivityVisitor"/> class.
/// </summary>
public ActivityVisitor(IEnumerable<IActivityResolver> portResolvers, IServiceProvider serviceProvider)
{
@ -21,14 +21,26 @@ public class ActivityVisitor : IActivityVisitor
/// <inheritdoc />
public async Task<ActivityNode> VisitAsync(IActivity activity, CancellationToken cancellationToken = default)
{
var collectedActivities = new HashSet<IActivity>(new[] { activity });
var graph = new ActivityNode(activity);
var collectedNodes = new HashSet<ActivityNode>(new[] { graph });
await VisitRecursiveAsync((graph, activity), collectedActivities, collectedNodes, cancellationToken);
var collectedNodes = new HashSet<ActivityNode>(new[]
{
graph
});
var collectedActivities = new HashSet<IActivity>(new[]
{
activity
});
var visitorContext = new ActivityVisitorContext
{
CollectedActivities = collectedActivities,
CollectedNodes = collectedNodes
};
await VisitRecursiveAsync((graph, activity), visitorContext, cancellationToken);
return graph;
}
private async Task VisitRecursiveAsync((ActivityNode Node, IActivity Activity) pair, HashSet<IActivity> collectedActivities, HashSet<ActivityNode> collectedNodes, CancellationToken cancellationToken)
private async Task VisitRecursiveAsync((ActivityNode Node, IActivity Activity) pair, ActivityVisitorContext visitorContext, CancellationToken cancellationToken)
{
if (pair.Activity is IInitializable initializable)
{
@ -36,10 +48,10 @@ public class ActivityVisitor : IActivityVisitor
await initializable.InitializeAsync(context);
}
await VisitPortsRecursiveAsync(pair, collectedActivities, collectedNodes, cancellationToken);
await VisitPortsRecursiveAsync(pair, visitorContext, cancellationToken);
}
private async Task VisitPortsRecursiveAsync((ActivityNode Node, IActivity Activity) pair, HashSet<IActivity> collectedActivities, HashSet<ActivityNode> collectedNodes, CancellationToken cancellationToken)
private async Task VisitPortsRecursiveAsync((ActivityNode Node, IActivity Activity) pair, ActivityVisitorContext visitorContext, CancellationToken cancellationToken)
{
var resolver = _portResolvers.FirstOrDefault(x => x.GetSupportsActivity(pair.Activity));
@ -47,6 +59,8 @@ public class ActivityVisitor : IActivityVisitor
return;
var activities = await resolver.GetActivitiesAsync(pair.Activity, cancellationToken);
var collectedActivities = visitorContext.CollectedActivities;
var collectedNodes = visitorContext.CollectedNodes;
foreach (var activity in activities)
{
@ -65,7 +79,13 @@ public class ActivityVisitor : IActivityVisitor
childNode.Parents.Add(pair.Node);
pair.Node.Children.Add(childNode);
collectedActivities.Add(activity);
await VisitRecursiveAsync((childNode, activity), collectedActivities, collectedNodes, cancellationToken);
await VisitRecursiveAsync((childNode, activity), visitorContext, cancellationToken);
}
}
private class ActivityVisitorContext
{
public HashSet<IActivity> CollectedActivities { get; set; } = new();
public HashSet<ActivityNode> CollectedNodes { get; set; } = new();
}
}

View file

@ -33,8 +33,8 @@ public class Hasher : IHasher
/// <inheritdoc />
public string Hash(string value)
{
using var sha = SHA256.Create();
return Hash(sha, value);
var data = SHA256.HashData(Encoding.UTF8.GetBytes(value));
return Convert.ToHexString(data);
}
/// <inheritdoc />
@ -45,12 +45,6 @@ public class Hasher : IHasher
var input = string.Join("|", strings);
return Hash(input);
}
private static string Hash(HashAlgorithm hashAlgorithm, string input)
{
var data = hashAlgorithm.ComputeHash(Encoding.UTF8.GetBytes(input));
return Convert.ToHexString(data);
}
[RequiresUnreferencedCode("Calls System.Text.Json.JsonSerializer.Serialize(Object, Type, JsonSerializerOptions)")]
private string Serialize(object? payload)

View file

@ -0,0 +1,20 @@
using Elsa.Extensions;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Services;
/// <inheritdoc />
public class WorkflowGraphBuilder(IActivityVisitor activityVisitor, IIdentityGraphService identityGraphService, IServiceProvider serviceProvider) : IWorkflowGraphBuilder
{
/// <inheritdoc />
public async Task<WorkflowGraph> BuildAsync(Workflow workflow, CancellationToken cancellationToken = default)
{
var graph = await activityVisitor.VisitAsync(workflow, cancellationToken);
var nodes = graph.Flatten().ToList();
identityGraphService.AssignIdentities(nodes);
return new WorkflowGraph(workflow, graph, nodes);
}
}

View file

@ -10,45 +10,28 @@ using Elsa.Workflows.State;
namespace Elsa.Workflows.Services;
/// <inheritdoc />
public class WorkflowRunner : IWorkflowRunner
public class WorkflowRunner(
IServiceProvider serviceProvider,
IWorkflowExecutionPipeline pipeline,
IWorkflowStateExtractor workflowStateExtractor,
IWorkflowBuilderFactory workflowBuilderFactory,
IWorkflowGraphBuilder workflowGraphBuilder,
IIdentityGenerator identityGenerator,
INotificationSender notificationSender)
: IWorkflowRunner
{
private readonly IServiceProvider _serviceProvider;
private readonly IWorkflowExecutionPipeline _pipeline;
private readonly IWorkflowStateExtractor _workflowStateExtractor;
private readonly IWorkflowBuilderFactory _workflowBuilderFactory;
private readonly IIdentityGenerator _identityGenerator;
private readonly INotificationSender _notificationSender;
/// <summary>
/// Constructor.
/// </summary>
public WorkflowRunner(
IServiceProvider serviceProvider,
IWorkflowExecutionPipeline pipeline,
IWorkflowStateExtractor workflowStateExtractor,
IWorkflowBuilderFactory workflowBuilderFactory,
IIdentityGenerator identityGenerator,
INotificationSender notificationSender)
{
_serviceProvider = serviceProvider;
_pipeline = pipeline;
_workflowStateExtractor = workflowStateExtractor;
_workflowBuilderFactory = workflowBuilderFactory;
_identityGenerator = identityGenerator;
_notificationSender = notificationSender;
}
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(IActivity activity, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
var workflow = Workflow.FromActivity(activity);
return await RunAsync(workflow, options, cancellationToken);
var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow, cancellationToken);
return await RunAsync(workflowGraph, options, cancellationToken);
}
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(IWorkflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
var builder = _workflowBuilderFactory.CreateBuilder();
var builder = workflowBuilderFactory.CreateBuilder();
var workflowDefinition = await builder.BuildWorkflowAsync(workflow, cancellationToken);
return await RunAsync(workflowDefinition, options, cancellationToken);
}
@ -65,7 +48,7 @@ public class WorkflowRunner : IWorkflowRunner
RunWorkflowOptions? options = default,
CancellationToken cancellationToken = default) where T : IWorkflow, new()
{
var builder = _workflowBuilderFactory.CreateBuilder();
var builder = workflowBuilderFactory.CreateBuilder();
var workflowDefinition = await builder.BuildWorkflowAsync<T>(cancellationToken);
return await RunAsync(workflowDefinition, options, cancellationToken);
}
@ -75,17 +58,17 @@ public class WorkflowRunner : IWorkflowRunner
RunWorkflowOptions? options = default,
CancellationToken cancellationToken = default) where T : WorkflowBase<TResult>, new()
{
var builder = _workflowBuilderFactory.CreateBuilder();
var builder = workflowBuilderFactory.CreateBuilder();
var workflow = await builder.BuildWorkflowAsync<T>(cancellationToken);
var result = await RunAsync(workflow, options, cancellationToken);
return (TResult)result.Result!;
}
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(Workflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
public async Task<RunWorkflowResult> RunAsync(WorkflowGraph workflowGraph, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
// Set up a workflow execution context.
var instanceId = options?.WorkflowInstanceId ?? _identityGenerator.GenerateId();
var instanceId = options?.WorkflowInstanceId ?? identityGenerator.GenerateId();
var input = options?.Input;
var properties = options?.Properties;
var correlationId = options?.CorrelationId;
@ -93,8 +76,8 @@ public class WorkflowRunner : IWorkflowRunner
var parentWorkflowInstanceId = options?.ParentWorkflowInstanceId;
var statusUpdatedCallback = options?.StatusUpdatedCallback;
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(
_serviceProvider,
workflow,
serviceProvider,
workflowGraph,
instanceId,
correlationId,
parentWorkflowInstanceId,
@ -111,8 +94,22 @@ public class WorkflowRunner : IWorkflowRunner
return await RunAsync(workflowExecutionContext);
}
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(Workflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow, cancellationToken);
return await RunAsync(workflowGraph, options, cancellationToken);
}
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(Workflow workflow, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow, cancellationToken);
return await RunAsync(workflowGraph, workflowState, options, cancellationToken);
}
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
// Create workflow execution context.
var input = options?.Input;
@ -122,8 +119,8 @@ public class WorkflowRunner : IWorkflowRunner
var parentWorkflowInstanceId = options?.ParentWorkflowInstanceId;
var statusUpdatedCallback = options?.StatusUpdatedCallback;
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(
_serviceProvider,
workflow,
serviceProvider,
workflowGraph,
workflowState,
correlationId,
parentWorkflowInstanceId,
@ -187,25 +184,25 @@ public class WorkflowRunner : IWorkflowRunner
var applicationCancellationToken = workflowExecutionContext.CancellationTokens.ApplicationCancellationToken;
var systemCancellationToken = workflowExecutionContext.CancellationTokens.SystemCancellationToken;
await _notificationSender.SendAsync(new WorkflowExecuting(workflow, workflowExecutionContext), applicationCancellationToken);
await notificationSender.SendAsync(new WorkflowExecuting(workflow, workflowExecutionContext), applicationCancellationToken);
// If the status is Pending, it means the workflow is started for the first time.
if (workflowExecutionContext.SubStatus == WorkflowSubStatus.Pending)
{
workflowExecutionContext.TransitionTo(WorkflowSubStatus.Executing);
await _notificationSender.SendAsync(new WorkflowStarted(workflow, workflowExecutionContext), applicationCancellationToken);
await notificationSender.SendAsync(new WorkflowStarted(workflow, workflowExecutionContext), applicationCancellationToken);
}
await _pipeline.ExecuteAsync(workflowExecutionContext);
var workflowState = _workflowStateExtractor.Extract(workflowExecutionContext);
await pipeline.ExecuteAsync(workflowExecutionContext);
var workflowState = workflowStateExtractor.Extract(workflowExecutionContext);
if (workflowState.Status == WorkflowStatus.Finished)
{
await _notificationSender.SendAsync(new WorkflowFinished(workflow, workflowState, workflowExecutionContext), applicationCancellationToken);
await notificationSender.SendAsync(new WorkflowFinished(workflow, workflowState, workflowExecutionContext), applicationCancellationToken);
}
var result = workflow.ResultVariable?.Get(workflowExecutionContext.MemoryRegister);
await _notificationSender.SendAsync(new WorkflowExecuted(workflow, workflowState, workflowExecutionContext), systemCancellationToken);
await notificationSender.SendAsync(new WorkflowExecuted(workflow, workflowState, workflowExecutionContext), systemCancellationToken);
return new RunWorkflowResult(workflowState, workflowExecutionContext.Workflow, result);
}
}

View file

@ -20,6 +20,8 @@ namespace Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
[Browsable(false)]
public class WorkflowDefinitionActivity : Composite, IInitializable
{
private bool IsInitialized => Root.Id != null;
/// <summary>
/// The definition ID of the workflow to schedule for execution.
/// </summary>
@ -74,7 +76,7 @@ public class WorkflowDefinitionActivity : Composite, IInitializable
.Any(x => x.ActivityDescriptor.Outputs.Any(y => y.Name == outputDescriptor.Name)) == true
? activityExecutionContext.ParentActivityExecutionContext ?? activityExecutionContext
: activityExecutionContext;
parentActivityExecutionContext.Set(output, value, outputDescriptor.Name);
}
@ -134,24 +136,35 @@ public class WorkflowDefinitionActivity : Composite, IInitializable
}
}
private async Task<Workflow?> FindWorkflowAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
private async Task<WorkflowGraph?> FindWorkflowGraphAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
{
var workflowDefinitionService = serviceProvider.GetRequiredService<IWorkflowDefinitionService>();
var filter = new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId };
var filter = new WorkflowDefinitionFilter
{
DefinitionId = WorkflowDefinitionId
};
if (!string.IsNullOrWhiteSpace(WorkflowDefinitionVersionId))
filter.Id = WorkflowDefinitionVersionId;
else
filter.VersionOptions = VersionOptions.SpecificVersion(Version);
var workflow =
await workflowDefinitionService.FindWorkflowAsync(filter, cancellationToken)
?? (await workflowDefinitionService.FindWorkflowAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Published }, cancellationToken)
?? await workflowDefinitionService.FindWorkflowAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Latest }, cancellationToken));
var workflowGraph =
await workflowDefinitionService.FindWorkflowGraphAsync(filter, cancellationToken)
?? (await workflowDefinitionService.FindWorkflowGraphAsync(new WorkflowDefinitionFilter
{
DefinitionId = WorkflowDefinitionId,
VersionOptions = VersionOptions.Published
}, cancellationToken)
?? await workflowDefinitionService.FindWorkflowGraphAsync(new WorkflowDefinitionFilter
{
DefinitionId = WorkflowDefinitionId,
VersionOptions = VersionOptions.Latest
}, cancellationToken));
return workflow;
return workflowGraph;
}
private ActivityDescriptor FindActivityDescriptor(IServiceProvider serviceProvider)
{
var activityRegistry = serviceProvider.GetRequiredService<IActivityRegistry>();
@ -160,20 +173,26 @@ public class WorkflowDefinitionActivity : Composite, IInitializable
async ValueTask IInitializable.InitializeAsync(InitializationContext context)
{
// This is not just for efficiency, but also a necessity to avoid potential race conditions.
// Such conditions can occur when multiple threads are simultaneously creating consuming workflows,
// especially when cached workflows are being updated during the graph construction process.
if (IsInitialized)
return;
var serviceProvider = context.ServiceProvider;
var cancellationToken = context.CancellationToken;
var workflow = await FindWorkflowAsync(serviceProvider, cancellationToken);
var workflowGraph = await FindWorkflowGraphAsync(serviceProvider, cancellationToken);
if (workflow == null)
if (workflowGraph == null)
throw new Exception($"Could not find workflow definition with ID {WorkflowDefinitionId}.");
var activityDescriptor = FindActivityDescriptor(serviceProvider);
// Declare input and output variables.
DeclareInputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable));
DeclareOutputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable));
// Set the root activity.
Root = workflow;
Root = workflowGraph.Workflow;
}
}

View file

@ -19,6 +19,6 @@ public class WorkflowDefinitionActivityResolver : IActivityResolver
var definitionActivity = (WorkflowDefinitionActivity)activity;
var root = definitionActivity.Root;
return new(new[] { root });
return new([root]);
}
}

View file

@ -2,6 +2,7 @@ using Elsa.Common.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Contracts;
@ -13,7 +14,7 @@ public interface IWorkflowDefinitionService
/// <summary>
/// Constructs an executable <see cref="Workflow"/> from the specified <see cref="WorkflowDefinition"/>.
/// </summary>
Task<Workflow> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default);
Task<WorkflowGraph> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default);
/// <summary>
/// Looks for a <see cref="WorkflowDefinition"/> by the specified definition ID and <see cref="VersionOptions"/>.
@ -33,15 +34,15 @@ public interface IWorkflowDefinitionService
/// <summary>
/// Looks for a <see cref="Workflow"/> by the specified definition ID and <see cref="VersionOptions"/>.
/// </summary>
Task<Workflow?> FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default);
Task<WorkflowGraph?> FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default);
/// <summary>
/// Looks for a <see cref="Workflow"/> by the specified version ID.
/// </summary>
Task<Workflow?> FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default);
Task<WorkflowGraph?> FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default);
/// <summary>
/// Looks for a <see cref="Workflow"/> by the specified <see cref="WorkflowDefinitionFilter"/>.
/// </summary>
Task<Workflow?> FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
Task<WorkflowGraph?> FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
}

View file

@ -99,7 +99,8 @@ public class WorkflowDefinitionMapper
/// <returns>The mapped <see cref="Workflow"/>.</returns>
public async Task<WorkflowDefinitionModel> MapAsync(WorkflowDefinition workflowDefinition, CancellationToken cancellationToken = default)
{
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var workflow = workflowGraph.Workflow;
var variables = _variableDefinitionMapper.Map(workflow.Variables).ToList();
return new(
@ -120,7 +121,7 @@ public class WorkflowDefinitionMapper
workflowDefinition.IsLatest,
workflowDefinition.IsPublished,
workflow.Options,
default,
null,
workflow.Root);
}
@ -151,7 +152,7 @@ public class WorkflowDefinitionMapper
workflow.Publication.IsLatest,
workflow.Publication.IsPublished,
workflow.Options,
default,
null,
workflow.Root);
}
}

View file

@ -1,8 +1,8 @@
using Elsa.Common.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Models;
using JetBrains.Annotations;
using Microsoft.Extensions.Caching.Memory;
@ -15,7 +15,7 @@ namespace Elsa.Workflows.Management.Services;
public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager) : IWorkflowDefinitionService
{
/// <inheritdoc />
public async Task<Workflow> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
return await decoratedService.MaterializeWorkflowAsync(definition, cancellationToken);
}
@ -49,30 +49,30 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
}
/// <inheritdoc />
public async Task<Workflow?> FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph?> FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionId, versionOptions);
return await GetFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken),
x => x.Identity.DefinitionId);
() => decoratedService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken),
x => x.Workflow.Identity.DefinitionId);
}
/// <inheritdoc />
public async Task<Workflow?> FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph?> FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionVersionId);
return await GetFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowAsync(definitionVersionId, cancellationToken),
x => x.Identity.DefinitionId);
() => decoratedService.FindWorkflowGraphAsync(definitionVersionId, cancellationToken),
x => x.Workflow.Identity.DefinitionId);
}
/// <inheritdoc />
public async Task<Workflow?> FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph?> FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowFilterCacheKey(filter);
return await GetFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowAsync(filter, cancellationToken),
x => x.Identity.DefinitionId);
() => decoratedService.FindWorkflowGraphAsync(filter, cancellationToken),
x => x.Workflow.Identity.DefinitionId);
}
private async Task<T?> GetFromCacheAsync<T>(string cacheKey, Func<Task<T?>> getObjectFunc, Func<T, string> getChangeTokenKeyFunc)

View file

@ -1,29 +0,0 @@
// using Elsa.Expressions.Contracts;
//
// namespace Elsa.Workflows.Management.Services;
//
// /// <inheritdoc />
// public class ExpressionDescriptorRegistryPopulator : IExpressionDescriptorRegistryPopulator
// {
// private readonly IEnumerable<IExpressionDescriptorProvider> _providers;
// private readonly IExpressionDescriptorRegistry _registry;
//
// /// <summary>
// /// Initializes a new instance of the <see cref="ExpressionDescriptorRegistryPopulator"/> class.
// /// </summary>
// public ExpressionDescriptorRegistryPopulator(IEnumerable<IExpressionDescriptorProvider> providers, IExpressionDescriptorRegistry registry)
// {
// _providers = providers;
// _registry = registry;
// }
//
// /// <inheritdoc />
// public async ValueTask PopulateRegistryAsync(CancellationToken cancellationToken)
// {
// foreach (var provider in _providers)
// {
// var descriptors = await provider.GetDescriptorsAsync(cancellationToken);
// _registry.AddRange(descriptors);
// }
// }
// }

View file

@ -85,8 +85,8 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
/// <inheritdoc />
public async Task<PublishWorkflowDefinitionResult> PublishAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
var responses = await _requestSender.SendAsync(new ValidateWorkflowRequest(workflow), cancellationToken);
var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
var responses = await _requestSender.SendAsync(new ValidateWorkflowRequest(workflowGraph.Workflow), cancellationToken);
var validationErrors = responses.SelectMany(r => r.ValidationErrors).ToList();
if (validationErrors.Any())

View file

@ -1,75 +1,61 @@
using Elsa.Common.Models;
using Elsa.Extensions;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Services;
/// <inheritdoc />
public class WorkflowDefinitionService : IWorkflowDefinitionService
public class WorkflowDefinitionService(
IWorkflowDefinitionStore workflowDefinitionStore,
IWorkflowGraphBuilder workflowGraphBuilder,
Func<IEnumerable<IWorkflowMaterializer>> materializers)
: IWorkflowDefinitionService
{
private readonly IWorkflowDefinitionStore _workflowDefinitionStore;
private readonly IActivityVisitor _activityVisitor;
private readonly IIdentityGraphService _identityGraphService;
private readonly Func<IEnumerable<IWorkflowMaterializer>> _materializers;
/// <summary>
/// Constructor.
/// </summary>
public WorkflowDefinitionService(
IWorkflowDefinitionStore workflowDefinitionStore,
IActivityVisitor activityVisitor,
IIdentityGraphService identityGraphService,
Func<IEnumerable<IWorkflowMaterializer>> materializers)
{
_workflowDefinitionStore = workflowDefinitionStore;
_activityVisitor = activityVisitor;
_identityGraphService = identityGraphService;
_materializers = materializers;
}
/// <inheritdoc />
public async Task<Workflow> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
var materializers = _materializers();
var materializer = materializers.FirstOrDefault(x => x.Name == definition.MaterializerName);
var workflowMaterializers = materializers();
var materializer = workflowMaterializers.FirstOrDefault(x => x.Name == definition.MaterializerName);
if (materializer == null)
throw new Exception("Provider not found");
var workflow = await materializer.MaterializeAsync(definition, cancellationToken);
var graph = (await _activityVisitor.VisitAsync(workflow, cancellationToken)).Flatten().ToList();
_identityGraphService.AssignIdentities(graph);
return workflow;
return await workflowGraphBuilder.BuildAsync(workflow, cancellationToken);
}
/// <inheritdoc />
public async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var filter = new WorkflowDefinitionFilter { DefinitionId = definitionId, VersionOptions = versionOptions };
return await _workflowDefinitionStore.FindAsync(filter, cancellationToken);
var filter = new WorkflowDefinitionFilter
{
DefinitionId = definitionId,
VersionOptions = versionOptions
};
return await workflowDefinitionStore.FindAsync(filter, cancellationToken);
}
/// <inheritdoc />
public async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(string definitionVersionId, CancellationToken cancellationToken = default)
{
var filter = new WorkflowDefinitionFilter { Id = definitionVersionId };
return await _workflowDefinitionStore.FindAsync(filter, cancellationToken);
var filter = new WorkflowDefinitionFilter
{
Id = definitionVersionId
};
return await workflowDefinitionStore.FindAsync(filter, cancellationToken);
}
/// <inheritdoc />
public async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
return await _workflowDefinitionStore.FindAsync(filter, cancellationToken);
return await workflowDefinitionStore.FindAsync(filter, cancellationToken);
}
/// <inheritdoc />
public async Task<Workflow?> FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph?> FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken);
@ -80,7 +66,7 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService
}
/// <inheritdoc />
public async Task<Workflow?> FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph?> FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken);
@ -91,7 +77,7 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService
}
/// <inheritdoc />
public async Task<Workflow?> FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
public async Task<WorkflowGraph?> FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken);

View file

@ -8,7 +8,6 @@ using Elsa.Workflows.Management.Mappers;
using Elsa.Workflows.Management.Notifications;
using Elsa.Workflows.Management.Requests;
using Elsa.Workflows.State;
using Exception = System.Exception;
namespace Elsa.Workflows.Management.Services;

View file

@ -11,11 +11,6 @@ namespace Elsa.Workflows.Runtime.Contracts;
/// </summary>
public interface ITriggerIndexer
{
/// <summary>
/// Removes triggers for the specified workflow.
/// </summary>
Task<IndexedWorkflowTriggers> DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default);
/// <summary>
/// Removes triggers matching the specified filter.
/// </summary>
@ -30,14 +25,6 @@ public interface ITriggerIndexer
/// Indexes triggers of the specified workflow.
/// </summary>
Task<IndexedWorkflowTriggers> IndexTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default);
/// <summary>
/// Returns triggers for the specified workflow definition.
/// </summary>
/// <param name="definition">The workflow definition.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
/// <returns>A collection of triggers.</returns>
Task<IEnumerable<StoredTrigger>> GetTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default);
/// <summary>
/// Returns triggers for the specified workflow definition.

View file

@ -1,4 +1,5 @@
using Elsa.Workflows.Activities;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Parameters;
using Elsa.Workflows.Runtime.Results;
using Elsa.Workflows.State;
@ -11,9 +12,14 @@ namespace Elsa.Workflows.Runtime.Contracts;
public interface IWorkflowHost
{
/// <summary>
/// The workflow definition.
/// The workflow graph.
/// </summary>
Workflow Workflow { get; set; }
WorkflowGraph WorkflowGraph { get; }
/// <summary>
/// The workflow.
/// </summary>
Workflow Workflow { get; }
/// <summary>
/// The workflow state.

View file

@ -1,6 +1,5 @@
using Elsa.Common.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
using Elsa.Workflows.State;
namespace Elsa.Workflows.Runtime.Contracts;
@ -21,15 +20,15 @@ public interface IWorkflowHostFactory
/// <summary>
/// Creates a new <see cref="IWorkflowHost"/> object.
/// </summary>
/// <param name="workflow">The workflow.</param>
/// <param name="workflowGraph">The workflow.</param>
/// <param name="workflowState">The workflow state to initialize the workflow host with.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
Task<IWorkflowHost> CreateAsync(Workflow workflow, WorkflowState workflowState, CancellationToken cancellationToken = default);
Task<IWorkflowHost> CreateAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, CancellationToken cancellationToken = default);
/// <summary>
/// Creates a new <see cref="IWorkflowHost"/> object.
/// </summary>
/// <param name="workflow">The workflow.</param>
/// <param name="workflowGraph">The workflow.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
Task<IWorkflowHost> CreateAsync(Workflow workflow, CancellationToken cancellationToken = default);
Task<IWorkflowHost> CreateAsync(WorkflowGraph workflowGraph, CancellationToken cancellationToken = default);
}

View file

@ -58,12 +58,12 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker
if (workflowState == null)
throw new Exception("Workflow state not found");
var workflow = await _workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken);
var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken);
if (workflow == null)
if (workflowGraph == null)
throw new Exception("Workflow definition not found");
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflow, workflowState, cancellationTokens: cancellationToken);
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflowGraph, workflowState, cancellationTokens: cancellationToken);
var activityNodeId = scheduledBackgroundActivity.ActivityNodeId;
var activityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.NodeId == activityNodeId);

View file

@ -242,7 +242,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP
}
}
private async Task IndexTriggersAsync(MaterializedWorkflow workflow, CancellationToken cancellationToken) => await _triggerIndexer.IndexTriggersAsync(workflow.Workflow, cancellationToken);
private async Task IndexTriggersAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken) => await _triggerIndexer.IndexTriggersAsync(materializedWorkflow.Workflow, cancellationToken);
/// <summary>
/// Syncs the items in the primary list with existing items in the secondary list, even when the object instances are not the same (but their IDs are).

View file

@ -140,7 +140,7 @@ public class DefaultWorkflowRuntime(
var definitionId = workflowInstance.DefinitionId;
var version = workflowInstance.Version;
var workflow = await workflowDefinitionService.FindWorkflowAsync(
var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(
definitionId,
VersionOptions.SpecificVersion(version),
systemCancellationToken);
@ -267,7 +267,7 @@ public class DefaultWorkflowRuntime(
if (workflowState == null)
throw new Exception("Workflow state not found");
var workflow = await workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken);
var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken);
if (workflow == null)
throw new Exception("Workflow definition not found");
@ -385,7 +385,7 @@ public class DefaultWorkflowRuntime(
if (options?.IsExistingInstance == true)
{
var workflowState = await LoadWorkflowStateAsync(options.InstanceId!, cancellationToken);
var workflow = await workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionVersionId, cancellationToken) ?? throw new Exception("Specified workflow definition and version does not exist");
var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionVersionId, cancellationToken) ?? throw new Exception("Specified workflow definition and version does not exist");
return await workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken);
}
@ -465,14 +465,14 @@ public class DefaultWorkflowRuntime(
InstanceId = workflowsFilter.Options.WorkflowInstanceId
};
var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions);
var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, startOptions.VersionOptions, cancellationToken);
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, startOptions.VersionOptions, cancellationToken);
if (workflow == null)
if (workflowGraph == null)
throw new Exception($"The workflow definition {definitionId} was not found");
var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest
{
Workflow = workflow,
Workflow = workflowGraph.Workflow,
CorrelationId = workflowsFilter.Options.CorrelationId,
Properties = workflowsFilter.Options.Properties,
Input = workflowsFilter.Options.Input,

View file

@ -63,21 +63,20 @@ public class TriggerIndexer : ITriggerIndexer
foreach (string workflowDefinitionVersionId in workflowDefinitionVersionIds)
{
var workflowDefinition = await _workflowDefinitionService.FindWorkflowDefinitionAsync(workflowDefinitionVersionId, cancellationToken);
var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionVersionId, cancellationToken);
if (workflowDefinition == null)
if (workflowGraph == null)
continue;
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
await DeleteTriggersAsync(workflow, cancellationToken);
await DeleteTriggersAsync(workflowGraph.Workflow, cancellationToken);
}
}
/// <inheritdoc />
public async Task<IndexedWorkflowTriggers> IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
return await IndexTriggersAsync(workflow, cancellationToken);
var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
return await IndexTriggersAsync(workflowGraph.Workflow, cancellationToken);
}
/// <inheritdoc />
@ -104,21 +103,13 @@ public class TriggerIndexer : ITriggerIndexer
return indexedWorkflow;
}
/// <inheritdoc />
public async Task<IEnumerable<StoredTrigger>> GetTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
return await GetTriggersAsync(workflow, cancellationToken);
}
/// <inheritdoc />
public async Task<IEnumerable<StoredTrigger>> GetTriggersAsync(Workflow workflow, CancellationToken cancellationToken)
{
return await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken);
}
/// <inheritdoc />
public async Task<IndexedWorkflowTriggers> DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default)
private async Task DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default)
{
var emptyTriggerList = new List<StoredTrigger>(0);
var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList();
@ -126,7 +117,6 @@ public class TriggerIndexer : ITriggerIndexer
await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken);
var indexedWorkflow = new IndexedWorkflowTriggers(workflow, emptyTriggerList, currentTriggers, emptyTriggerList);
await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken);
return indexedWorkflow;
}
private async Task<IEnumerable<StoredTrigger>> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken)

View file

@ -27,12 +27,12 @@ public class WorkflowHost : IWorkflowHost
/// </summary>
public WorkflowHost(
IServiceScopeFactory serviceScopeFactory,
Workflow workflow,
WorkflowGraph workflowGraph,
WorkflowState workflowState,
IIdentityGenerator identityGenerator,
ILogger<WorkflowHost> logger)
{
Workflow = workflow;
WorkflowGraph = workflowGraph;
WorkflowState = workflowState;
_serviceScopeFactory = serviceScopeFactory;
_identityGenerator = identityGenerator;
@ -40,7 +40,10 @@ public class WorkflowHost : IWorkflowHost
}
/// <inheritdoc />
public Workflow Workflow { get; set; }
public WorkflowGraph WorkflowGraph { get; }
/// <inheritdoc />
public Workflow Workflow => WorkflowGraph.Workflow;
/// <inheritdoc />
public WorkflowState WorkflowState { get; set; }
@ -90,8 +93,8 @@ public class WorkflowHost : IWorkflowHost
using var scope = _serviceScopeFactory.CreateScope();
var workflowRunner = scope.ServiceProvider.GetRequiredService<IWorkflowRunner>();
var workflowResult = @params?.IsExistingInstance == true
? await workflowRunner.RunAsync(Workflow, WorkflowState, runOptions, cancellationToken)
: await workflowRunner.RunAsync(Workflow, runOptions, cancellationToken);
? await workflowRunner.RunAsync(WorkflowGraph, WorkflowState, runOptions, cancellationToken)
: await workflowRunner.RunAsync(WorkflowGraph, runOptions, cancellationToken);
WorkflowState = workflowResult.WorkflowState;
@ -133,7 +136,7 @@ public class WorkflowHost : IWorkflowHost
using var scope = _serviceScopeFactory.CreateScope();
var workflowRunner = scope.ServiceProvider.GetRequiredService<IWorkflowRunner>();
var workflowResult = await workflowRunner.RunAsync(Workflow, WorkflowState, runOptions, cancellationToken);
var workflowResult = await workflowRunner.RunAsync(WorkflowGraph, WorkflowState, runOptions, cancellationToken);
WorkflowState = workflowResult.WorkflowState;

View file

@ -1,8 +1,7 @@
using Elsa.Common.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.State;
using Microsoft.Extensions.DependencyInjection;
@ -31,7 +30,7 @@ public class WorkflowHostFactory : IWorkflowHostFactory
{
using var scope = _serviceScopeFactory.CreateScope();
var workflowDefinitionService = scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionService>();
var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken);
var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken);
if(workflow == null)
return default;
@ -40,22 +39,23 @@ public class WorkflowHostFactory : IWorkflowHostFactory
}
/// <inheritdoc />
public Task<IWorkflowHost> CreateAsync(Workflow workflow, WorkflowState workflowState, CancellationToken cancellationToken = default)
public Task<IWorkflowHost> CreateAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, CancellationToken cancellationToken = default)
{
var workflowHost = (IWorkflowHost)ActivatorUtilities.CreateInstance<WorkflowHost>(_serviceProvider, workflow, workflowState);
var workflowHost = (IWorkflowHost)ActivatorUtilities.CreateInstance<WorkflowHost>(_serviceProvider, workflowGraph, workflowState);
return Task.FromResult(workflowHost);
}
/// <inheritdoc />
public Task<IWorkflowHost> CreateAsync(Workflow workflow, CancellationToken cancellationToken = default)
public Task<IWorkflowHost> CreateAsync(WorkflowGraph workflowGraph, CancellationToken cancellationToken = default)
{
var workflow = workflowGraph.Workflow;
var workflowState = new WorkflowState
{
Id = _identityGenerator.GenerateId(),
DefinitionId = workflow.Identity.DefinitionId,
DefinitionId =workflow.Identity.DefinitionId,
DefinitionVersion = workflow.Identity.Version
};
return CreateAsync(workflow, workflowState, cancellationToken);
return CreateAsync(workflowGraph, workflowState, cancellationToken);
}
}

View file

@ -47,6 +47,15 @@
<None Update="Scenarios\WorkflowCompletion\hello-world.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\CachingAndWorkflowDefinitionActivity\Workflows\workflow-definition-child.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\CachingAndWorkflowDefinitionActivity\Workflows\workflow-definition-grand-child.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\CachingAndWorkflowDefinitionActivity\Workflows\workflow-definition-parent.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
</ItemGroup>
</Project>

View file

@ -22,6 +22,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
{
var client = CreateClient();
client.BaseAddress = new Uri(client.BaseAddress!, "/elsa/api");
client.Timeout = TimeSpan.FromMinutes(1);
return RestService.For<TClient>(client, CreateRefitSettings());
}
@ -29,6 +30,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
{
var client = CreateClient();
client.BaseAddress = new Uri(client.BaseAddress!, "/workflows/");
client.Timeout = TimeSpan.FromMinutes(1);
return client;
}

View file

@ -0,0 +1,53 @@
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Management.Filters;
using Microsoft.Extensions.DependencyInjection;
using Xunit.Abstractions;
namespace Elsa.Workflows.ComponentTests.Scenarios.CachingAndWorkflowDefinitionActivity;
// Tests the behavior of the WorkflowDefinitionActivity.
// See https://github.com/elsa-workflows/elsa-core/issues/5314
public class WorkflowDefinitionActivityTests : AppComponentTest
{
private readonly ITestOutputHelper _testOutputHelper;
private const string GrandChildDefinitionId = "29595e7b37a4836d";
private readonly IWorkflowDefinitionCacheManager _workflowDefinitionCacheManager;
private readonly IWorkflowInstanceStore _workflowInstanceStore;
private readonly HttpClient _httpWorkflowClient;
public WorkflowDefinitionActivityTests(App app, ITestOutputHelper testOutputHelper) : base(app)
{
_testOutputHelper = testOutputHelper;
_httpWorkflowClient = WorkflowServer.CreateHttpWorkflowClient();
_workflowDefinitionCacheManager = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionCacheManager>();
_workflowInstanceStore = Scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
}
[Fact]
public async Task SendHttpRequest_WhileEvictingCache_ShouldNotGenerateFaults()
{
var requestTasks = Enumerable.Range(0, 200).Select(SendRequestAsync).ToList();
await Task.WhenAll(requestTasks);
var filter = new WorkflowInstanceFilter
{
WorkflowSubStatus = WorkflowSubStatus.Faulted
};
var faultedWorkflows = (await _workflowInstanceStore.FindManyAsync(filter)).ToList();
var faultCount = faultedWorkflows.Count;
foreach (var faultedWorkflow in faultedWorkflows)
foreach (var incident in faultedWorkflow.WorkflowState.Incidents)
_testOutputHelper.WriteLine(incident.Message);
Assert.Equal(0, faultCount);
}
private async Task SendRequestAsync(int index = 0)
{
var requestTask = _httpWorkflowClient.PostAsync("parent", new StringContent("{}"));
var evictionTask = _workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(GrandChildDefinitionId);
await Task.WhenAll(requestTask, evictionTask);
}
}

View file

@ -0,0 +1,64 @@
{
"id": "26d186d68deb1250",
"definitionId": "a3390d1f4c2594a8",
"name": "Child",
"createdAt": "2024-04-30T08:22:44.543407+00:00",
"version": 1,
"toolVersion": "3.2.0.0",
"variables": [],
"inputs": [],
"outputs": [],
"outcomes": [],
"customProperties": {},
"isReadonly": false,
"isSystem": false,
"isLatest": true,
"isPublished": true,
"options": {
"usableAsActivity": true,
"autoUpdateConsumingWorkflows": true
},
"root": {
"type": "Elsa.Flowchart",
"version": 1,
"id": "e6c941b6f2292fcf",
"nodeId": "Workflow2:e6c941b6f2292fcf",
"metadata": {},
"customProperties": {
"source": "FlowchartJsonConverter.cs:45",
"notFoundConnections": [],
"canStartWorkflow": false,
"runAsynchronously": false
},
"activities": [
{
"workflowDefinitionId": "29595e7b37a4836d",
"workflowDefinitionVersionId": "f36dc264f14be3c6",
"latestAvailablePublishedVersion": 1,
"latestAvailablePublishedVersionId": "f36dc264f14be3c6",
"id": "377b23f4235ec26d",
"nodeId": "Workflow2:e6c941b6f2292fcf:377b23f4235ec26d",
"name": "GrandChild1",
"type": "GrandChild",
"version": 1,
"customProperties": {
"canStartWorkflow": false,
"runAsynchronously": false
},
"metadata": {
"designer": {
"position": {
"x": -193.95703125,
"y": 23
},
"size": {
"width": 115.1875,
"height": 50
}
}
}
}
],
"connections": []
}
}

View file

@ -0,0 +1,67 @@
{
"id": "f36dc264f14be3c6",
"definitionId": "29595e7b37a4836d",
"name": "Grand Child",
"createdAt": "2024-04-30T08:22:44.537237+00:00",
"version": 1,
"toolVersion": "3.2.0.0",
"variables": [],
"inputs": [],
"outputs": [],
"outcomes": [],
"customProperties": {},
"isReadonly": false,
"isSystem": false,
"isLatest": true,
"isPublished": true,
"options": {
"usableAsActivity": true,
"autoUpdateConsumingWorkflows": true
},
"root": {
"type": "Elsa.Flowchart",
"version": 1,
"id": "29044d8e2ba192c0",
"nodeId": "Workflow1:29044d8e2ba192c0",
"metadata": {},
"customProperties": {
"source": "FlowchartJsonConverter.cs:45",
"notFoundConnections": [],
"canStartWorkflow": false,
"runAsynchronously": false
},
"activities": [
{
"text": {
"typeName": "String",
"expression": {
"type": "Literal",
"value": "Grand Child"
}
},
"id": "fb00878500fde1d3",
"nodeId": "Workflow1:29044d8e2ba192c0:fb00878500fde1d3",
"name": "WriteLine1",
"type": "Elsa.WriteLine",
"version": 1,
"customProperties": {
"canStartWorkflow": false,
"runAsynchronously": false
},
"metadata": {
"designer": {
"position": {
"x": -234.796875,
"y": -191
},
"size": {
"width": 139.296875,
"height": 50
}
}
}
}
],
"connections": []
}
}

View file

@ -0,0 +1,142 @@
{
"id": "9178b94f961b1cf",
"definitionId": "189be5173f90b1f6",
"name": "Parent",
"createdAt": "2024-04-30T08:42:17.838484+00:00",
"version": 6,
"toolVersion": "3.2.0.0",
"variables": [],
"inputs": [],
"outputs": [],
"outcomes": [],
"customProperties": {
"Elsa:WorkflowContextProviderTypes": []
},
"isReadonly": false,
"isSystem": false,
"isLatest": true,
"isPublished": true,
"options": {
"autoUpdateConsumingWorkflows": false
},
"root": {
"type": "Elsa.Flowchart",
"version": 1,
"id": "96f60810318a6a89",
"nodeId": "Workflow3:96f60810318a6a89",
"metadata": {},
"customProperties": {
"source": "FlowchartJsonConverter.cs:45",
"notFoundConnections": [],
"canStartWorkflow": false,
"runAsynchronously": false
},
"activities": [
{
"path": {
"typeName": "String",
"expression": {
"type": "Literal",
"value": "parent"
}
},
"supportedMethods": {
"typeName": "String[]",
"expression": {
"type": "Object",
"value": "[\u0022POST\u0022]"
}
},
"authorize": {
"typeName": "Boolean",
"expression": {
"type": "Literal",
"value": false
}
},
"policy": {
"typeName": "String",
"expression": {
"type": "Literal"
}
},
"requestTimeout": null,
"requestSizeLimit": null,
"fileSizeLimit": null,
"allowedFileExtensions": null,
"blockedFileExtensions": null,
"allowedMimeTypes": null,
"exposeRequestTooLargeOutcome": false,
"exposeFileTooLargeOutcome": false,
"exposeInvalidFileExtensionOutcome": false,
"exposeInvalidFileMimeTypeOutcome": false,
"parsedContent": null,
"files": null,
"routeData": null,
"queryStringData": null,
"headers": null,
"result": null,
"id": "7e20498ffc94dbc6",
"nodeId": "Workflow3:96f60810318a6a89:7e20498ffc94dbc6",
"name": "HttpEndpoint1",
"type": "Elsa.HttpEndpoint",
"version": 1,
"customProperties": {
"canStartWorkflow": true,
"runAsynchronously": false
},
"metadata": {
"designer": {
"position": {
"x": -520,
"y": -120
},
"size": {
"width": 176.390625,
"height": 50
}
}
}
},
{
"workflowDefinitionId": "a3390d1f4c2594a8",
"workflowDefinitionVersionId": "26d186d68deb1250",
"latestAvailablePublishedVersion": 1,
"latestAvailablePublishedVersionId": "26d186d68deb1250",
"id": "ed6a03040f651906",
"nodeId": "Workflow3:96f60810318a6a89:ed6a03040f651906",
"name": "Child1",
"type": "Child",
"version": 1,
"customProperties": {
"canStartWorkflow": false,
"runAsynchronously": false
},
"metadata": {
"designer": {
"position": {
"x": -236.5810546875,
"y": -120
},
"size": {
"width": 67.7734375,
"height": 50
}
}
}
}
],
"connections": [
{
"source": {
"activity": "7e20498ffc94dbc6",
"port": "Done"
},
"target": {
"activity": "ed6a03040f651906",
"port": "In"
}
}
]
}
}

View file

@ -9,13 +9,6 @@
<IsTestProject>true</IsTestProject>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="coverlet.collector" Version="6.0.0"/>
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="17.8.0"/>
<PackageReference Include="xunit" Version="2.5.3"/>
<PackageReference Include="xunit.runner.visualstudio" Version="2.5.3"/>
</ItemGroup>
<ItemGroup>
<Using Include="Xunit"/>
</ItemGroup>

View file

@ -0,0 +1,30 @@
using Elsa.Testing.Shared;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
using Microsoft.Extensions.DependencyInjection;
using Xunit.Abstractions;
namespace Elsa.Workflows.Core.UnitTests;
public class ScheduleActivityExecutionContextTests(ITestOutputHelper testOutputHelper)
{
private readonly IServiceProvider _serviceProvider = new TestApplicationBuilder(testOutputHelper).Build();
[Fact(DisplayName = "Scheduling an activity that is not part of the workflow should throw an exception")]
public async Task ScheduleActivityAsync_WithActivityNotPartOfWorkflow_ShouldThrowException()
{
await _serviceProvider.PopulateRegistriesAsync();
var writeLineA = new WriteLine("Test");
var writeLineB = new WriteLine("Test");
var workflow = new Workflow
{
Root = writeLineA
};
var workflowGraphBuilder = _serviceProvider.GetRequiredService<IWorkflowGraphBuilder>();
var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow);
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflowGraph, "test");
var activityExecutionContext = workflowExecutionContext.CreateActivityExecutionContext(writeLineA);
await Assert.ThrowsAsync<InvalidOperationException>(() => activityExecutionContext.ScheduleActivityAsync(writeLineB).AsTask());
}
}