diff --git a/src/ConductorSharp.Engine/Builders/ForkJoinTaskBuilder.cs b/src/ConductorSharp.Engine/Builders/ForkJoinTaskBuilder.cs new file mode 100644 index 00000000..b179fce0 --- /dev/null +++ b/src/ConductorSharp.Engine/Builders/ForkJoinTaskBuilder.cs @@ -0,0 +1,112 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Linq.Expressions; +using ConductorSharp.Client.Generated; +using ConductorSharp.Engine.Interface; +using ConductorSharp.Engine.Model; +using ConductorSharp.Engine.Util; +using ConductorSharp.Engine.Util.Builders; + +namespace ConductorSharp.Engine.Builders +{ + public static class ForkJoinTaskExtensions + { + public static ITaskOptionsBuilder AddTask( + this ITaskSequenceBuilder builder, + Expression> reference, + Expression> input, + params Action>[] branches + ) + where TWorkflow : ITypedWorkflow + { + if (branches == null || branches.Length == 0) + throw new InvalidOperationException("FORK_JOIN task requires at least one branch."); + + var taskBuilder = new ForkJoinTaskBuilder( + reference.Body, + input.Body, + builder.BuildConfiguration, + builder.WorkflowBuildRegistry, + builder.ConfigurationProperties, + builder.BuildContext + ); + + foreach (var branch in branches) + { + taskBuilder.AddBranch(); + branch(taskBuilder); + } + + builder.AddTaskBuilderToSequence(taskBuilder); + return taskBuilder; + } + } + + public class ForkJoinTaskBuilder( + Expression taskExpression, + Expression inputExpression, + BuildConfiguration buildConfiguration, + WorkflowBuildItemRegistry workflowBuildItemRegistry, + IEnumerable configurationProperties, + BuildContext buildContext + ) : BaseTaskBuilder(taskExpression, inputExpression, buildConfiguration), ITaskSequenceBuilder + where TWorkflow : ITypedWorkflow + { + private readonly List> _branches = []; + + public BuildContext BuildContext { get; } = buildContext; + public BuildConfiguration BuildConfiguration { get; } = buildConfiguration; + public WorkflowBuildItemRegistry WorkflowBuildRegistry { get; } = workflowBuildItemRegistry; + public IEnumerable ConfigurationProperties { get; } = configurationProperties; + + public void AddBranch() => _branches.Add([]); + + public override WorkflowTask[] Build() + { + var forkTaskName = $"FORK_JOIN_{_taskRefferenceName}"; + var joinTaskName = $"JOIN_{_taskRefferenceName}"; + + var builtBranches = _branches + .Select(branch => (ICollection)branch.SelectMany(taskBuilder => taskBuilder.Build()).ToList()) + .ToList(); + + for (var i = 0; i < builtBranches.Count; i++) + { + if (builtBranches[i].Count == 0) + throw new InvalidOperationException($"FORK_JOIN branch {i} must contain at least one task."); + } + + var joinOn = builtBranches.Select(branch => branch.Last().TaskReferenceName).ToList(); + + return + [ + new() + { + Name = forkTaskName, + TaskReferenceName = forkTaskName, + WorkflowTaskType = WorkflowTaskType.FORK_JOIN, + Type = WorkflowTaskType.FORK_JOIN.ToString(), + InputParameters = _inputParameters.ToObject>(), + ForkTasks = builtBranches, + }, + new() + { + Name = joinTaskName, + TaskReferenceName = joinTaskName, + WorkflowTaskType = WorkflowTaskType.JOIN, + Type = WorkflowTaskType.JOIN.ToString(), + JoinOn = joinOn, + }, + ]; + } + + public void AddTaskBuilderToSequence(ITaskBuilder builder) + { + if (_branches.Count == 0) + throw new InvalidOperationException("Cannot add a task to a FORK_JOIN branch before the branch has been started."); + + _branches[^1].Add(builder); + } + } +} diff --git a/src/ConductorSharp.Engine/Model/ForkJoinTaskModel.cs b/src/ConductorSharp.Engine/Model/ForkJoinTaskModel.cs new file mode 100644 index 00000000..f962b692 --- /dev/null +++ b/src/ConductorSharp.Engine/Model/ForkJoinTaskModel.cs @@ -0,0 +1,8 @@ +using MediatR; + +namespace ConductorSharp.Engine.Model +{ + public class ForkJoinInput : IRequest { } + + public class ForkJoinTaskModel : TaskModel { } +} diff --git a/test/ConductorSharp.Engine.Tests/ConductorSharp.Engine.Tests.csproj b/test/ConductorSharp.Engine.Tests/ConductorSharp.Engine.Tests.csproj index a67cc2a7..59cd4f5a 100644 --- a/test/ConductorSharp.Engine.Tests/ConductorSharp.Engine.Tests.csproj +++ b/test/ConductorSharp.Engine.Tests/ConductorSharp.Engine.Tests.csproj @@ -43,6 +43,7 @@ + diff --git a/test/ConductorSharp.Engine.Tests/Integration/WorkflowBuilderTests.cs b/test/ConductorSharp.Engine.Tests/Integration/WorkflowBuilderTests.cs index 7c1dcc9d..3908e650 100644 --- a/test/ConductorSharp.Engine.Tests/Integration/WorkflowBuilderTests.cs +++ b/test/ConductorSharp.Engine.Tests/Integration/WorkflowBuilderTests.cs @@ -175,6 +175,27 @@ public void BuilderReturnsCorrectDefinitionDoWhileTask() Assert.Equal(expectedDefinition, definition); } + [Fact] + public void BuilderReturnsCorrectDefinitionForkJoinTask() + { + var definition = GetDefinitionFromWorkflow(); + var expectedDefinition = EmbeddedFileHelper.GetLinesFromEmbeddedFile("~/Samples/Workflows/ForkJoinTask.json"); + + Assert.Equal(expectedDefinition, definition); + } + + [Fact] + public void BuilderThrowsInvalidOperationExceptionForForkJoinTaskWithoutBranches() + { + Assert.Throws(GetDefinitionFromWorkflow); + } + + [Fact] + public void BuilderThrowsInvalidOperationExceptionForForkJoinTaskWithEmptyBranch() + { + Assert.Throws(GetDefinitionFromWorkflow); + } + [Fact] public void BuilderReturnsCorrectDefinitionSwitchTask() { diff --git a/test/ConductorSharp.Engine.Tests/Samples/Workflows/EmptyBranchForkJoinTask.cs b/test/ConductorSharp.Engine.Tests/Samples/Workflows/EmptyBranchForkJoinTask.cs new file mode 100644 index 00000000..b1687bee --- /dev/null +++ b/test/ConductorSharp.Engine.Tests/Samples/Workflows/EmptyBranchForkJoinTask.cs @@ -0,0 +1,37 @@ +using ConductorSharp.Client.Generated; + +namespace ConductorSharp.Engine.Tests.Samples.Workflows +{ + public sealed class EmptyBranchForkJoinTaskInput : WorkflowInput { } + + public sealed class EmptyBranchForkJoinTaskOutput : WorkflowOutput { } + + public sealed class EmptyBranchForkJoinTask : Workflow + { + public ForkJoinTaskModel ForkJoin { get; set; } + + public EmptyBranchForkJoinTask( + WorkflowDefinitionBuilder builder + ) + : base(builder) { } + + public override void BuildDefinition() + { + _builder.AddTask( + wf => wf.ForkJoin, + wf => new(), + branch => + branch.AddTasks( + new WorkflowTask + { + Name = "task_a1", + TaskReferenceName = "task_a1", + Type = WorkflowTaskType.SIMPLE.ToString(), + WorkflowTaskType = WorkflowTaskType.SIMPLE, + } + ), + branch => { } + ); + } + } +} diff --git a/test/ConductorSharp.Engine.Tests/Samples/Workflows/EmptyForkJoinTask.cs b/test/ConductorSharp.Engine.Tests/Samples/Workflows/EmptyForkJoinTask.cs new file mode 100644 index 00000000..3912dc5b --- /dev/null +++ b/test/ConductorSharp.Engine.Tests/Samples/Workflows/EmptyForkJoinTask.cs @@ -0,0 +1,19 @@ +namespace ConductorSharp.Engine.Tests.Samples.Workflows +{ + public sealed class EmptyForkJoinTaskInput : WorkflowInput { } + + public sealed class EmptyForkJoinTaskOutput : WorkflowOutput { } + + public sealed class EmptyForkJoinTask : Workflow + { + public ForkJoinTaskModel ForkJoin { get; set; } + + public EmptyForkJoinTask(WorkflowDefinitionBuilder builder) + : base(builder) { } + + public override void BuildDefinition() + { + _builder.AddTask(wf => wf.ForkJoin, wf => new()); + } + } +} diff --git a/test/ConductorSharp.Engine.Tests/Samples/Workflows/ForkJoinTask.cs b/test/ConductorSharp.Engine.Tests/Samples/Workflows/ForkJoinTask.cs new file mode 100644 index 00000000..d2116dd5 --- /dev/null +++ b/test/ConductorSharp.Engine.Tests/Samples/Workflows/ForkJoinTask.cs @@ -0,0 +1,38 @@ +using ConductorSharp.Client.Generated; + +namespace ConductorSharp.Engine.Tests.Samples.Workflows +{ + public sealed class ForkJoinTaskInput : WorkflowInput { } + + public sealed class ForkJoinTaskOutput : WorkflowOutput { } + + public sealed class ForkJoinTask : Workflow + { + public ForkJoinTaskModel ForkJoin { get; set; } + + public ForkJoinTask(WorkflowDefinitionBuilder builder) + : base(builder) { } + + public override void BuildDefinition() + { + _builder.AddTask( + wf => wf.ForkJoin, + wf => new(), + // Branch A has two tasks, to verify branch task ordering and that the JOIN's + // joinOn resolves to the branch's last task ("task_a2"), not its first ("task_a1"). + branch => branch.AddTasks(NewSimpleTask("task_a1"), NewSimpleTask("task_a2")), + // Branch B has a single task, to verify a simple one-task branch still works alongside a multi-task one. + branch => branch.AddTasks(NewSimpleTask("task_b1")) + ); + } + + private static WorkflowTask NewSimpleTask(string referenceName) => + new() + { + Name = referenceName, + TaskReferenceName = referenceName, + Type = WorkflowTaskType.SIMPLE.ToString(), + WorkflowTaskType = WorkflowTaskType.SIMPLE, + }; + } +} diff --git a/test/ConductorSharp.Engine.Tests/Samples/Workflows/ForkJoinTask.json b/test/ConductorSharp.Engine.Tests/Samples/Workflows/ForkJoinTask.json new file mode 100644 index 00000000..7cb81d1e --- /dev/null +++ b/test/ConductorSharp.Engine.Tests/Samples/Workflows/ForkJoinTask.json @@ -0,0 +1,51 @@ +{ + "name": "fork_join_task", + "version": 1, + "tasks": [ + { + "name": "FORK_JOIN_fork_join", + "taskReferenceName": "FORK_JOIN_fork_join", + "inputParameters": {}, + "type": "FORK_JOIN", + "forkTasks": [ + [ + { + "name": "task_a1", + "taskReferenceName": "task_a1", + "type": "SIMPLE", + "workflowTaskType": "SIMPLE" + }, + { + "name": "task_a2", + "taskReferenceName": "task_a2", + "type": "SIMPLE", + "workflowTaskType": "SIMPLE" + } + ], + [ + { + "name": "task_b1", + "taskReferenceName": "task_b1", + "type": "SIMPLE", + "workflowTaskType": "SIMPLE" + } + ] + ], + "workflowTaskType": "FORK_JOIN" + }, + { + "name": "JOIN_fork_join", + "taskReferenceName": "JOIN_fork_join", + "type": "JOIN", + "joinOn": [ + "task_a2", + "task_b1" + ], + "workflowTaskType": "JOIN" + } + ], + "inputParameters": [], + "outputParameters": {}, + "schemaVersion": 2, + "timeoutSeconds": 0 +} \ No newline at end of file