Stash RunWorkflow activity

This commit is contained in:
Sipke Schoorstra 2020-11-01 11:31:53 +01:00
parent ce6509baff
commit 31ca0780e2
15 changed files with 327 additions and 16 deletions

View file

@ -110,6 +110,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.CorrelationHtt
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.ContextualWorkflowHttp", "src\samples\Elsa.Samples.ContextualWorkflowHttp\Elsa.Samples.ContextualWorkflowHttp.csproj", "{37753BF9-9D9B-4264-AB5B-D057E467D3C5}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.RunChildWorkflowWorker", "src\samples\Elsa.Samples.RunChildWorkflowWorker\Elsa.Samples.RunChildWorkflowWorker.csproj", "{1E5E1669-AEC9-4DA9-9A3D-3D405101A246}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -261,6 +263,10 @@ Global
{37753BF9-9D9B-4264-AB5B-D057E467D3C5}.Debug|Any CPU.Build.0 = Debug|Any CPU
{37753BF9-9D9B-4264-AB5B-D057E467D3C5}.Release|Any CPU.ActiveCfg = Release|Any CPU
{37753BF9-9D9B-4264-AB5B-D057E467D3C5}.Release|Any CPU.Build.0 = Release|Any CPU
{1E5E1669-AEC9-4DA9-9A3D-3D405101A246}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{1E5E1669-AEC9-4DA9-9A3D-3D405101A246}.Debug|Any CPU.Build.0 = Debug|Any CPU
{1E5E1669-AEC9-4DA9-9A3D-3D405101A246}.Release|Any CPU.ActiveCfg = Release|Any CPU
{1E5E1669-AEC9-4DA9-9A3D-3D405101A246}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -312,6 +318,7 @@ Global
{E508C1B6-916F-40F4-8731-58279B8F45B2} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
{565E2F7D-1F5F-45EF-B130-BC09CB3465B1} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
{37753BF9-9D9B-4264-AB5B-D057E467D3C5} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
{1E5E1669-AEC9-4DA9-9A3D-3D405101A246} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158}

View file

@ -1,6 +1,8 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Models;
using Elsa.Triggers;
namespace Elsa.Services
@ -11,17 +13,17 @@ namespace Elsa.Services
/// Schedules new and existing instances matching the specified trigger.
/// </summary>
Task TriggerWorkflowsAsync<TTrigger>(Func<TTrigger, bool> predicate, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default) where TTrigger : ITrigger;
/// <summary>
/// Schedules the specified workflow instance for execution.
/// </summary>
Task ScheduleWorkflowInstanceAsync(string instanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default);
/// <summary>
/// Creates a new workflow instance of the specified workflow definition and schedules it for execution.
/// </summary>
Task ScheduleWorkflowDefinitionAsync(string definitionId, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowInstance>> ScheduleWorkflowDefinitionAsync(string definitionId, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default);
/// <summary>
/// Schedules new workflows that start with the specified activity type or are blocked on the specified activity type.
/// </summary>

View file

@ -0,0 +1,70 @@
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
using Open.Linq.AsyncExtensions;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.Workflows
{
[ActivityDefinition(
Category = "Workflows",
Description = "Runs a child workflow."
)]
public class RunWorkflow : Activity
{
private readonly IWorkflowScheduler _workflowScheduler;
public RunWorkflow(IWorkflowScheduler workflowScheduler)
{
_workflowScheduler = workflowScheduler;
}
[ActivityProperty] public string WorkflowDefinitionId { get; set; } = default!;
[ActivityProperty] public object? Input { get; set; }
[ActivityProperty] public string? CorrelationId { get; set; }
[ActivityProperty] public string? ContextId { get; set; }
[ActivityProperty] public RunWorkflowMode Mode { get; set; }
public ICollection<string> WorkflowInstanceIds
{
get => GetState<ICollection<string>>();
set => SetState(value);
}
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
var workflowInstances = await _workflowScheduler.ScheduleWorkflowDefinitionAsync(WorkflowDefinitionId, Input, CorrelationId, ContextId, cancellationToken).ToList();
WorkflowInstanceIds = workflowInstances.Select(x => x.WorkflowInstanceId).ToList();
return Mode switch
{
RunWorkflowMode.FireAndForget => Done(),
RunWorkflowMode.Blocking => Suspend(),
_ => Suspend()
};
}
protected override ValueTask<IActivityExecutionResult> OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
return base.OnResumeAsync(context, cancellationToken);
}
public enum RunWorkflowMode
{
/// <summary>
/// Run the specified workflow and continue with the current one.
/// </summary>
FireAndForget,
/// <summary>
/// Run the specified workflow and continue once the child workflow finishes.
/// </summary>
Blocking
}
}
}

View file

@ -0,0 +1,24 @@
using System;
using Elsa.Activities.ControlFlow;
using Elsa.Builders;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.Workflows
{
public static class RunWorkflowBuilderExtensions
{
public static IActivityBuilder RunWorkflow(this IBuilder builder, Action<ISetupActivity<RunWorkflow>>? setup = default) => builder.Then(setup);
public static IActivityBuilder RunWorkflow<T>(this IBuilder builder, RunWorkflow.RunWorkflowMode mode) => builder.RunWorkflow(activity => activity.WithWorkflow<T>().WithMode(mode));
public static IActivityBuilder RunWorkflow<T>(this IBuilder builder, RunWorkflow.RunWorkflowMode mode, object input) => builder.RunWorkflow(activity => activity.WithWorkflow<T>().WithMode(mode).WithInput(input));
public static IActivityBuilder RunWorkflow<T>(this IBuilder builder, RunWorkflow.RunWorkflowMode mode, object input, string correlationId) =>
builder.RunWorkflow(activity => activity.WithWorkflow<T>().WithMode(mode).WithInput(input).WithCorrelationId(correlationId));
public static IActivityBuilder RunWorkflow<T>(this IBuilder builder, RunWorkflow.RunWorkflowMode mode, object input, string correlationId, string contextId) =>
builder.RunWorkflow(activity => activity.WithWorkflow<T>().WithMode(mode).WithInput(input).WithCorrelationId(correlationId).WithContextId(contextId));
public static IActivityBuilder RunWorkflow<T>(this IBuilder builder, RunWorkflow.RunWorkflowMode mode, string correlationId) =>
builder.RunWorkflow(activity => activity.WithWorkflow<T>().WithMode(mode).WithCorrelationId(correlationId));
}
}

View file

@ -0,0 +1,53 @@
using System;
using System.Threading.Tasks;
using Elsa.Activities.Workflows;
using Elsa.Builders;
using Elsa.Extensions;
using Elsa.Services;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
public static class RunWorkflowExtensions
{
public static ISetupActivity<RunWorkflow> WithWorkflow(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, ValueTask<string>> value) => activity.Set(x => x.WorkflowDefinitionId, value);
public static ISetupActivity<RunWorkflow> WithWorkflow(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, string> value) => activity.Set(x => x.WorkflowDefinitionId, value);
public static ISetupActivity<RunWorkflow> WithWorkflow(this ISetupActivity<RunWorkflow> activity, Func<ValueTask<string>> value) => activity.Set(x => x.WorkflowDefinitionId, value);
public static ISetupActivity<RunWorkflow> WithWorkflow(this ISetupActivity<RunWorkflow> activity, Func<string> value) => activity.Set(x => x.WorkflowDefinitionId, value);
public static ISetupActivity<RunWorkflow> WithWorkflow(this ISetupActivity<RunWorkflow> activity, string value) => activity.Set(x => x.WorkflowDefinitionId, value);
public static ISetupActivity<RunWorkflow> WithWorkflow<T>(this ISetupActivity<RunWorkflow> activity) =>
activity.WithWorkflow(
async context =>
{
var workflowRegistry = context.GetService<IWorkflowRegistry>();
var workflow = (await workflowRegistry.GetWorkflowAsync<T>())!;
return workflow.Id;
});
public static ISetupActivity<RunWorkflow> WithInput(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, ValueTask<object?>> value) => activity.Set(x => x.Input, value);
public static ISetupActivity<RunWorkflow> WithInput(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, object?> value) => activity.Set(x => x.Input, value);
public static ISetupActivity<RunWorkflow> WithInput(this ISetupActivity<RunWorkflow> activity, Func<ValueTask<object?>> value) => activity.Set(x => x.Input, value);
public static ISetupActivity<RunWorkflow> WithInput(this ISetupActivity<RunWorkflow> activity, Func<object?> value) => activity.Set(x => x.Input, value);
public static ISetupActivity<RunWorkflow> WithInput(this ISetupActivity<RunWorkflow> activity, object? value) => activity.Set(x => x.Input, value);
public static ISetupActivity<RunWorkflow> WithCorrelationId(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, ValueTask<string?>> value) => activity.Set(x => x.CorrelationId, value);
public static ISetupActivity<RunWorkflow> WithCorrelationId(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, string?> value) => activity.Set(x => x.CorrelationId, value);
public static ISetupActivity<RunWorkflow> WithCorrelationId(this ISetupActivity<RunWorkflow> activity, Func<ValueTask<string?>> value) => activity.Set(x => x.CorrelationId, value);
public static ISetupActivity<RunWorkflow> WithCorrelationId(this ISetupActivity<RunWorkflow> activity, Func<string?> value) => activity.Set(x => x.CorrelationId, value);
public static ISetupActivity<RunWorkflow> WithCorrelationId(this ISetupActivity<RunWorkflow> activity, string? value) => activity.Set(x => x.CorrelationId, value);
public static ISetupActivity<RunWorkflow> WithContextId(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, ValueTask<string?>> value) => activity.Set(x => x.ContextId, value);
public static ISetupActivity<RunWorkflow> WithContextId(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, string?> value) => activity.Set(x => x.ContextId, value);
public static ISetupActivity<RunWorkflow> WithContextId(this ISetupActivity<RunWorkflow> activity, Func<ValueTask<string?>> value) => activity.Set(x => x.ContextId, value);
public static ISetupActivity<RunWorkflow> WithContextId(this ISetupActivity<RunWorkflow> activity, Func<string?> value) => activity.Set(x => x.ContextId, value);
public static ISetupActivity<RunWorkflow> WithContextId(this ISetupActivity<RunWorkflow> activity, string? value) => activity.Set(x => x.ContextId, value);
public static ISetupActivity<RunWorkflow> WithMode(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, ValueTask<RunWorkflow.RunWorkflowMode>> value) => activity.Set(x => x.Mode, value);
public static ISetupActivity<RunWorkflow> WithMode(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, RunWorkflow.RunWorkflowMode> value) => activity.Set(x => x.Mode, value);
public static ISetupActivity<RunWorkflow> WithMode(this ISetupActivity<RunWorkflow> activity, Func<ValueTask<RunWorkflow.RunWorkflowMode>> value) => activity.Set(x => x.Mode, value);
public static ISetupActivity<RunWorkflow> WithMode(this ISetupActivity<RunWorkflow> activity, Func<RunWorkflow.RunWorkflowMode> value) => activity.Set(x => x.Mode, value);
public static ISetupActivity<RunWorkflow> WithMode(this ISetupActivity<RunWorkflow> activity, RunWorkflow.RunWorkflowMode value) => activity.Set(x => x.Mode, value);
}
}

View file

@ -122,7 +122,7 @@ namespace Microsoft.Extensions.DependencyInjection
.AddStartupTask<StartServiceBusTask>()
.AddConsumer<RunWorkflow, RunWorkflowConsumer>()
.AddMetadataHandlers()
.AddPrimitiveActivities();
.AddCoreActivities();
return configuration;
}
@ -131,7 +131,7 @@ namespace Microsoft.Extensions.DependencyInjection
services
.AddSingleton<IActivityPropertyOptionsProvider, SelectOptionsProvider>();
private static IServiceCollection AddPrimitiveActivities(this IServiceCollection services) =>
private static IServiceCollection AddCoreActivities(this IServiceCollection services) =>
services
.AddActivity<Inline>()
.AddActivity<Finish>()
@ -148,7 +148,8 @@ namespace Microsoft.Extensions.DependencyInjection
.AddActivity<Signaled>()
.AddTriggerProvider<SignaledTriggerProvider>()
.AddActivity<TriggerEvent>()
.AddActivity<TriggerSignal>();
.AddActivity<TriggerSignal>()
.AddActivity<Elsa.Activities.Workflows.RunWorkflow>();
private static ElsaOptions AddServiceBus(this ElsaOptions options)
{

View file

@ -12,7 +12,7 @@ namespace Elsa.Extensions
{
public static Task<IWorkflowBlueprint?> GetWorkflowAsync<T>(
this IWorkflowRegistry workflowRegistry,
CancellationToken cancellationToken) =>
CancellationToken cancellationToken = default) =>
workflowRegistry.GetWorkflowAsync(typeof(T).Name, VersionOptions.Latest, cancellationToken);
public static async Task<IEnumerable<(IWorkflowBlueprint Workflow, IActivityBlueprint Activity)>>
@ -25,11 +25,10 @@ namespace Elsa.Extensions
return results.Select(x => (x.Workflow, x.Activity));
}
public static async Task<IEnumerable<(IWorkflowBlueprint Workflow, IActivityBlueprint Activity)>>
GetWorkflowsByStartActivityAsync(
this IWorkflowRegistry workflowRegistry,
string activityType,
CancellationToken cancellationToken = default)
public static async Task<IEnumerable<(IWorkflowBlueprint Workflow, IActivityBlueprint Activity)>> GetWorkflowsByStartActivityAsync(
this IWorkflowRegistry workflowRegistry,
string activityType,
CancellationToken cancellationToken = default)
{
var workflows = await workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken);

View file

@ -50,7 +50,7 @@ namespace Elsa.Services
CancellationToken cancellationToken = default) =>
await _serviceBus.Publish(new RunWorkflow(instanceId, activityId, input));
public async Task ScheduleWorkflowDefinitionAsync(
public async Task<IEnumerable<WorkflowInstance>> ScheduleWorkflowDefinitionAsync(
string definitionId,
object? input = default,
string? correlationId = default,
@ -66,9 +66,15 @@ namespace Elsa.Services
throw new WorkflowException($"No workflow definition found by ID {definitionId}");
var startActivities = workflow.GetStartActivities();
var workflowInstances = new List<WorkflowInstance>();
foreach (var activity in startActivities)
await ScheduleWorkflowAsync(workflow, activity, input, correlationId, contextId, cancellationToken);
{
var workflowInstance = await ScheduleWorkflowAsync(workflow, activity, input, correlationId, contextId, cancellationToken);
workflowInstances.Add(workflowInstance);
}
return workflowInstances;
}
public async Task TriggerWorkflowsAsync(
@ -188,7 +194,7 @@ namespace Elsa.Services
cancellationToken);
}
private async Task ScheduleWorkflowAsync(
private async Task<WorkflowInstance> ScheduleWorkflowAsync(
IWorkflowBlueprint workflowBlueprint,
IActivityBlueprint activity,
object? input,
@ -199,6 +205,7 @@ namespace Elsa.Services
var workflowInstance = await _workflowFactory.InstantiateAsync(workflowBlueprint, correlationId, contextId, cancellationToken);
await _workflowInstanceManager.SaveAsync(workflowInstance, cancellationToken);
await ScheduleWorkflowInstanceAsync(workflowInstance.WorkflowInstanceId, activity.Id, input, cancellationToken);
return workflowInstance;
}
private async Task<IEnumerable<IWorkflowBlueprint>> FilterRunningSingletonsAsync(IEnumerable<IWorkflowBlueprint> workflows)

View file

@ -0,0 +1,14 @@
<Project Sdk="Microsoft.NET.Sdk.Worker">
<PropertyGroup>
<TargetFramework>net5.0</TargetFramework>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.Hosting" Version="5.0.0-rc.2.20475.5" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa\Elsa.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,32 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Samples.RunChildWorkflowWorker.Workflows;
using Elsa.Services;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
namespace Elsa.Samples.RunChildWorkflowWorker.HostedServices
{
/// <summary>
/// Runs the parent workflow.
/// </summary>
public class RunParentWorkflow : IHostedService
{
private readonly IServiceProvider _serviceProvider;
public RunParentWorkflow(IServiceProvider serviceProvider)
{
_serviceProvider = serviceProvider;
}
public async Task StartAsync(CancellationToken cancellationToken)
{
using var scope = _serviceProvider.CreateScope();
var workflowRunner = scope.ServiceProvider.GetRequiredService<IWorkflowRunner>();
await workflowRunner.RunWorkflowAsync<ParentWorkflow>(cancellationToken: cancellationToken);
}
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
}
}

View file

@ -0,0 +1,31 @@
using System.Data;
using Elsa.Samples.RunChildWorkflowWorker.HostedServices;
using Elsa.Samples.RunChildWorkflowWorker.Workflows;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NodaTime;
using YesSql.Provider.Sqlite;
namespace Elsa.Samples.RunChildWorkflowWorker
{
public class Program
{
public static void Main(string[] args)
{
CreateHostBuilder(args).Build().Run();
}
public static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureServices(
(hostContext, services) =>
{
services
.AddElsa(options => options.UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared", IsolationLevel.ReadUncommitted)))
.AddConsoleActivities()
.AddHostedService<RunParentWorkflow>()
.AddWorkflow<ParentWorkflow>()
.AddWorkflow<ChildWorkflow>();
});
}
}

View file

@ -0,0 +1,11 @@
{
"profiles": {
"Elsa.Samples.WhileLoopWorker": {
"commandName": "Project",
"dotnetRunMessages": "true",
"environmentVariables": {
"DOTNET_ENVIRONMENT": "Development"
}
}
}
}

View file

@ -0,0 +1,42 @@
using Elsa.Activities.Console;
using Elsa.Activities.ControlFlow;
using Elsa.Activities.Timers;
using Elsa.Activities.Workflows;
using Elsa.Builders;
using NodaTime;
namespace Elsa.Samples.RunChildWorkflowWorker.Workflows
{
/// <summary>
/// Delegate the hard work of counting numbers to a child workflow.
/// </summary>
public class ParentWorkflow : IWorkflow
{
private const int Count = 3;
public void Build(IWorkflowBuilder workflow)
{
workflow
.WriteLine("This is the parent workflow.")
.WriteLine("Let's kick off the child workflow.")
.RunWorkflow<ChildWorkflow>(RunWorkflow.RunWorkflowMode.Blocking, Count)
.WriteLine("Parent finished.");
}
}
public class ChildWorkflow : IWorkflowV
{
public void Build(IWorkflowBuilder workflow)
{
workflow
.SetVariable("Count", context => (int)context.Input!)
.WriteLine(context => $"Child workflow counting down from {context.GetVariable<int>("Count")} to 0")
.For(context => context.GetVariable<int>("Count"), _ => 0,
iterate =>
{
iterate.WriteLine(context => $"{context.Input}");
})
.WriteLine("Done. Back to you, parent workflow!");
}
}
}

View file

@ -0,0 +1,9 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
}
}

View file

@ -0,0 +1,9 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
}
}