Add support for attribute-based workflow selection & execution (#478)

This commit is contained in:
Sipke Schoorstra 2020-11-28 05:11:14 +01:00 committed by GitHub
parent 2ff8d006f3
commit 5bab43b8d2
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
15 changed files with 233 additions and 5 deletions

View file

@ -160,6 +160,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.Timers", "src\
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.WhileLoopWorker", "src\samples\worker\Elsa.Samples.WhileLoopWorker\Elsa.Samples.WhileLoopWorker.csproj", "{EA3832C4-8079-4E84-AF8C-12744E4EFD44}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.CustomAttributesChildWorker", "src\samples\worker\Elsa.Samples.CustomAttributesChildWorker\Elsa.Samples.CustomAttributesChildWorker.csproj", "{AEA00027-F280-43A4-AC60-E2AF3014CF34}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -387,6 +389,10 @@ Global
{EA3832C4-8079-4E84-AF8C-12744E4EFD44}.Debug|Any CPU.Build.0 = Debug|Any CPU
{EA3832C4-8079-4E84-AF8C-12744E4EFD44}.Release|Any CPU.ActiveCfg = Release|Any CPU
{EA3832C4-8079-4E84-AF8C-12744E4EFD44}.Release|Any CPU.Build.0 = Release|Any CPU
{AEA00027-F280-43A4-AC60-E2AF3014CF34}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{AEA00027-F280-43A4-AC60-E2AF3014CF34}.Debug|Any CPU.Build.0 = Debug|Any CPU
{AEA00027-F280-43A4-AC60-E2AF3014CF34}.Release|Any CPU.ActiveCfg = Release|Any CPU
{AEA00027-F280-43A4-AC60-E2AF3014CF34}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -463,6 +469,7 @@ Global
{A4DF29BD-95CC-4A86-96EC-FB6692270E68} = {E42743A0-FBDD-4150-9D53-6000496D9B87}
{FD2D1BD0-1229-4DCF-BE70-6BFD396DD489} = {E42743A0-FBDD-4150-9D53-6000496D9B87}
{EA3832C4-8079-4E84-AF8C-12744E4EFD44} = {E42743A0-FBDD-4150-9D53-6000496D9B87}
{AEA00027-F280-43A4-AC60-E2AF3014CF34} = {E42743A0-FBDD-4150-9D53-6000496D9B87}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158}

View file

@ -1,8 +1,9 @@
using System.Threading;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Extensions;
using Elsa.Models;
using Elsa.Services;
using Elsa.Services.Models;
@ -26,11 +27,12 @@ namespace Elsa.Activities.Workflows
_workflowRegistry = workflowRegistry;
}
[ActivityProperty] public string WorkflowDefinitionId { get; set; } = default!;
[ActivityProperty] public string? WorkflowDefinitionId { get; set; } = default!;
[ActivityProperty] public string? TenantId { get; set; } = default!;
[ActivityProperty] public object? Input { get; set; }
[ActivityProperty] public string? CorrelationId { get; set; }
[ActivityProperty] public string? ContextId { get; set; }
[ActivityProperty] public Variables? CustomAttributes { get; set; } = default!;
[ActivityProperty] public RunWorkflowMode Mode { get; set; }
public string ChildWorkflowInstanceId
@ -41,8 +43,8 @@ namespace Elsa.Activities.Workflows
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
var workflowBlueprint = (await _workflowRegistry.GetWorkflowAsync(WorkflowDefinitionId, TenantId, VersionOptions.Published, cancellationToken))!;
var workflowInstance = await _workflowScheduler.RunWorkflowAsync(workflowBlueprint, TenantId, Input, CorrelationId, ContextId, cancellationToken);
var workflowBlueprint = await FindWorkflowBlueprintAsync(cancellationToken);
var workflowInstance = await _workflowScheduler.RunWorkflowAsync(workflowBlueprint!, TenantId, Input, CorrelationId, ContextId, cancellationToken);
ChildWorkflowInstanceId = workflowInstance.WorkflowInstanceId;
return Mode switch
@ -59,6 +61,25 @@ namespace Elsa.Activities.Workflows
return Done(input);
}
private async Task<IWorkflowBlueprint?> FindWorkflowBlueprintAsync(CancellationToken cancellationToken)
{
var query = (IEnumerable<IWorkflowBlueprint>)(await _workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken));
query = query.Where(x => x.WithVersion(VersionOptions.Published));
if (WorkflowDefinitionId != null)
query = query.Where(x => x.Id == WorkflowDefinitionId);
if (TenantId != null)
query = query.Where(x => x.TenantId == TenantId);
if (CustomAttributes != null)
foreach (var customAttribute in CustomAttributes.Data)
query = query.Where(x => Equals(x.CustomAttributes.Get(customAttribute.Key), customAttribute.Value));
return query.FirstOrDefault();
}
public enum RunWorkflowMode
{
/// <summary>

View file

@ -1,5 +1,6 @@
using System;
using Elsa.Builders;
using Elsa.Models;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
@ -37,5 +38,14 @@ namespace Elsa.Activities.Workflows
public static IActivityBuilder RunWorkflow(this IBuilder builder, string workflowDefinitionId, string tenantId, RunWorkflow.RunWorkflowMode mode, string correlationId) =>
builder.RunWorkflow(activity => activity.WithWorkflow(workflowDefinitionId).WithMode(mode).WithCorrelationId(correlationId).WithTenantId(tenantId));
public static IActivityBuilder RunWorkflow(this IBuilder builder, string workflowDefinitionId, Variables customAttributes, RunWorkflow.RunWorkflowMode mode, string correlationId) =>
builder.RunWorkflow(activity => activity.WithWorkflow(workflowDefinitionId).WithMode(mode).WithCorrelationId(correlationId).WithCustomAttributes(customAttributes));
public static IActivityBuilder RunWorkflow(this IBuilder builder, Variables customAttributes, RunWorkflow.RunWorkflowMode mode, string correlationId) =>
builder.RunWorkflow(activity => activity.WithMode(mode).WithCorrelationId(correlationId).WithCustomAttributes(customAttributes));
public static IActivityBuilder RunWorkflow(this IBuilder builder, Variables customAttributes, RunWorkflow.RunWorkflowMode mode) =>
builder.RunWorkflow(activity => activity.WithMode(mode).WithCustomAttributes(customAttributes));
}
}

View file

@ -2,6 +2,7 @@ using System;
using System.Threading.Tasks;
using Elsa.Builders;
using Elsa.Extensions;
using Elsa.Models;
using Elsa.Services;
using Elsa.Services.Models;
@ -54,5 +55,11 @@ namespace Elsa.Activities.Workflows
public static ISetupActivity<RunWorkflow> WithTenantId(this ISetupActivity<RunWorkflow> activity, Func<ValueTask<string?>> value) => activity.Set(x => x.TenantId, value);
public static ISetupActivity<RunWorkflow> WithTenantId(this ISetupActivity<RunWorkflow> activity, Func<string?> value) => activity.Set(x => x.TenantId, value);
public static ISetupActivity<RunWorkflow> WithTenantId(this ISetupActivity<RunWorkflow> activity, string? value) => activity.Set(x => x.TenantId, value);
public static ISetupActivity<RunWorkflow> WithCustomAttributes(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, ValueTask<Variables?>> value) => activity.Set(x => x.CustomAttributes, value);
public static ISetupActivity<RunWorkflow> WithCustomAttributes(this ISetupActivity<RunWorkflow> activity, Func<ActivityExecutionContext, Variables?> value) => activity.Set(x => x.CustomAttributes, value);
public static ISetupActivity<RunWorkflow> WithCustomAttributes(this ISetupActivity<RunWorkflow> activity, Func<ValueTask<Variables?>> value) => activity.Set(x => x.CustomAttributes, value);
public static ISetupActivity<RunWorkflow> WithCustomAttributes(this ISetupActivity<RunWorkflow> activity, Func<Variables?> value) => activity.Set(x => x.CustomAttributes, value);
public static ISetupActivity<RunWorkflow> WithCustomAttributes(this ISetupActivity<RunWorkflow> activity, Variables? value) => activity.Set(x => x.CustomAttributes, value);
}
}

View file

@ -1,3 +1,4 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Runtime.CompilerServices;
@ -42,5 +43,11 @@ namespace Elsa.Services
.OrderByDescending(x => x.Version)
.FirstOrDefault();
}
public async Task<IEnumerable<IWorkflowBlueprint>> FindWorkflowsAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken) =>
await GetWorkflowsAsync(cancellationToken).Where(predicate).OrderByDescending(x => x.Version).ToListAsync(cancellationToken);
public async Task<IWorkflowBlueprint?> FindWorkflowAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken) =>
await GetWorkflowsAsync(cancellationToken).Where(predicate).OrderByDescending(x => x.Version).FirstOrDefaultAsync(cancellationToken);
}
}

View file

@ -0,0 +1,16 @@
<Project Sdk="Microsoft.NET.Sdk.Worker">
<PropertyGroup>
<TargetFramework>net5.0</TargetFramework>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.Hosting" Version="5.0.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\..\activities\Elsa.Activities.Rebus\Elsa.Activities.Rebus.csproj" />
<ProjectReference Include="..\..\..\core\Elsa\Elsa.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,7 @@
namespace Elsa.Samples.CustomAttributesChildWorker.Messages
{
public class OrderReceived
{
public string CustomerId { get; set; }
}
}

View file

@ -0,0 +1,32 @@
using Elsa.Samples.CustomAttributesChildWorker.Messages;
using Elsa.Samples.CustomAttributesChildWorker.Workflows;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
namespace Elsa.Samples.CustomAttributesChildWorker
{
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()
.AddTimerActivities()
.AddConsoleActivities()
.AddRebusActivities<OrderReceived>()
.AddWorkflow<GenerateOrdersWorkflow>()
.AddWorkflow<OrderReceivedWorkflow>()
.AddWorkflow<Customer1Workflow>()
.AddWorkflow<Customer2Workflow>()
;
});
}
}

View file

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

View file

@ -0,0 +1,15 @@
using Elsa.Activities.Console;
using Elsa.Builders;
namespace Elsa.Samples.CustomAttributesChildWorker.Workflows
{
public class Customer1Workflow : IWorkflow
{
public void Build(IWorkflowBuilder workflow)
{
workflow
.WithCustomAttribute("Customer","Customer1")
.WriteLine("Specialized workflow for Customer 1");
}
}
}

View file

@ -0,0 +1,15 @@
using Elsa.Activities.Console;
using Elsa.Builders;
namespace Elsa.Samples.CustomAttributesChildWorker.Workflows
{
public class Customer2Workflow : IWorkflow
{
public void Build(IWorkflowBuilder workflow)
{
workflow
.WithCustomAttribute("Customer","Customer2")
.WriteLine("Specialized workflow for Customer 2");
}
}
}

View file

@ -0,0 +1,36 @@
using System;
using Elsa.Activities.Console;
using Elsa.Activities.ControlFlow;
using Elsa.Activities.Rebus;
using Elsa.Activities.Timers;
using Elsa.Builders;
using Elsa.Samples.CustomAttributesChildWorker.Messages;
using NodaTime;
namespace Elsa.Samples.CustomAttributesChildWorker.Workflows
{
/// <summary>
/// Generate a new order for a random customer every 5 seconds.
/// </summary>
public class GenerateOrdersWorkflow : IWorkflow
{
private readonly Random _random;
public GenerateOrdersWorkflow() => _random = new Random();
public void Build(IWorkflowBuilder workflow)
{
workflow
.TimerEvent(Duration.FromSeconds(5))
.SetVariable("CustomerId", SelectRandomCustomerId)
.WriteLine(context => $"Creating a new order for customer {context.GetVariable<string>("CustomerId")}.")
.Then<SendRebusMessage>(message => message.Set(x => x.Message, context => new OrderReceived { CustomerId = context.GetVariable<string>("CustomerId") }));
}
private string SelectRandomCustomerId()
{
var customers = new[] { "Customer1", "Customer2" };
var index = _random.Next(customers.Length);
return customers[index];
}
}
}

View file

@ -0,0 +1,26 @@
using Elsa.Activities.Console;
using Elsa.Activities.ControlFlow;
using Elsa.Activities.Rebus;
using Elsa.Activities.Workflows;
using Elsa.Builders;
using Elsa.Models;
using Elsa.Samples.CustomAttributesChildWorker.Messages;
namespace Elsa.Samples.CustomAttributesChildWorker.Workflows
{
/// <summary>
/// Listen for new OrderReceived messages and kick off customer-specific child workflows.
/// </summary>
public class OrderReceivedWorkflow : IWorkflow
{
public void Build(IWorkflowBuilder workflow)
{
workflow
.StartWith<RebusMessageReceived>(activity => activity.Set(x => x.MessageType, typeof(OrderReceived)))
.SetVariable(context => context.GetInput<OrderReceived>())
.WriteLine(context => $"Received a new order for {context.GetVariable<OrderReceived>().CustomerId}.")
.RunWorkflow(activity => activity.WithCustomAttributes(context => new Variables().Set("Customer", context.GetVariable<OrderReceived>().CustomerId)))
.WriteLine("Returned back from child 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"
}
}
}