Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
112 changes: 112 additions & 0 deletions src/ConductorSharp.Engine/Builders/ForkJoinTaskBuilder.cs
Original file line number Diff line number Diff line change
@@ -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<TWorkflow>(
this ITaskSequenceBuilder<TWorkflow> builder,
Expression<Func<TWorkflow, ForkJoinTaskModel>> reference,
Expression<Func<TWorkflow, ForkJoinInput>> input,
params Action<ITaskSequenceBuilder<TWorkflow>>[] 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<TWorkflow>(
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<TWorkflow>(
Expression taskExpression,
Expression inputExpression,
BuildConfiguration buildConfiguration,
WorkflowBuildItemRegistry workflowBuildItemRegistry,
IEnumerable<ConfigurationProperty> configurationProperties,
BuildContext buildContext
) : BaseTaskBuilder<ForkJoinInput, NoOutput>(taskExpression, inputExpression, buildConfiguration), ITaskSequenceBuilder<TWorkflow>
where TWorkflow : ITypedWorkflow
{
private readonly List<List<ITaskBuilder>> _branches = [];

public BuildContext BuildContext { get; } = buildContext;
public BuildConfiguration BuildConfiguration { get; } = buildConfiguration;
public WorkflowBuildItemRegistry WorkflowBuildRegistry { get; } = workflowBuildItemRegistry;
public IEnumerable<ConfigurationProperty> 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<WorkflowTask>)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<IDictionary<string, object>>(),
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);
}
}
}
8 changes: 8 additions & 0 deletions src/ConductorSharp.Engine/Model/ForkJoinTaskModel.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
using MediatR;

namespace ConductorSharp.Engine.Model
{
public class ForkJoinInput : IRequest<NoOutput> { }

public class ForkJoinTaskModel : TaskModel<ForkJoinInput, NoOutput> { }
}
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
<ItemGroup>
<EmbeddedResource Include="Samples\Workflows\DictionaryInitializationWorkflow.json" />
<EmbeddedResource Include="Samples\Workflows\DoWhileTask.json" />
<EmbeddedResource Include="Samples\Workflows\ForkJoinTask.json" />
<EmbeddedResource Include="Samples\Workflows\FormatterWorkflow.json" />
<EmbeddedResource Include="Samples\Workflows\HumanTaskWorkflow.json" />
<EmbeddedResource Include="Samples\Workflows\ListInitializationWorkflow.json" />
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,27 @@ public void BuilderReturnsCorrectDefinitionDoWhileTask()
Assert.Equal(expectedDefinition, definition);
}

[Fact]
public void BuilderReturnsCorrectDefinitionForkJoinTask()
{
var definition = GetDefinitionFromWorkflow<ForkJoinTask>();
var expectedDefinition = EmbeddedFileHelper.GetLinesFromEmbeddedFile("~/Samples/Workflows/ForkJoinTask.json");

Assert.Equal(expectedDefinition, definition);
}

[Fact]
public void BuilderThrowsInvalidOperationExceptionForForkJoinTaskWithoutBranches()
{
Assert.Throws<InvalidOperationException>(GetDefinitionFromWorkflow<EmptyForkJoinTask>);
}

[Fact]
public void BuilderThrowsInvalidOperationExceptionForForkJoinTaskWithEmptyBranch()
{
Assert.Throws<InvalidOperationException>(GetDefinitionFromWorkflow<EmptyBranchForkJoinTask>);
}

[Fact]
public void BuilderReturnsCorrectDefinitionSwitchTask()
{
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
using ConductorSharp.Client.Generated;

namespace ConductorSharp.Engine.Tests.Samples.Workflows
{
public sealed class EmptyBranchForkJoinTaskInput : WorkflowInput<EmptyBranchForkJoinTaskOutput> { }

public sealed class EmptyBranchForkJoinTaskOutput : WorkflowOutput { }

public sealed class EmptyBranchForkJoinTask : Workflow<EmptyBranchForkJoinTask, EmptyBranchForkJoinTaskInput, EmptyBranchForkJoinTaskOutput>
{
public ForkJoinTaskModel ForkJoin { get; set; }

public EmptyBranchForkJoinTask(
WorkflowDefinitionBuilder<EmptyBranchForkJoinTask, EmptyBranchForkJoinTaskInput, EmptyBranchForkJoinTaskOutput> 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 => { }
);
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
namespace ConductorSharp.Engine.Tests.Samples.Workflows
{
public sealed class EmptyForkJoinTaskInput : WorkflowInput<EmptyForkJoinTaskOutput> { }

public sealed class EmptyForkJoinTaskOutput : WorkflowOutput { }

public sealed class EmptyForkJoinTask : Workflow<EmptyForkJoinTask, EmptyForkJoinTaskInput, EmptyForkJoinTaskOutput>
{
public ForkJoinTaskModel ForkJoin { get; set; }

public EmptyForkJoinTask(WorkflowDefinitionBuilder<EmptyForkJoinTask, EmptyForkJoinTaskInput, EmptyForkJoinTaskOutput> builder)
: base(builder) { }

public override void BuildDefinition()
{
_builder.AddTask(wf => wf.ForkJoin, wf => new());
}
}
}
38 changes: 38 additions & 0 deletions test/ConductorSharp.Engine.Tests/Samples/Workflows/ForkJoinTask.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
using ConductorSharp.Client.Generated;

namespace ConductorSharp.Engine.Tests.Samples.Workflows
{
public sealed class ForkJoinTaskInput : WorkflowInput<ForkJoinTaskOutput> { }

public sealed class ForkJoinTaskOutput : WorkflowOutput { }

public sealed class ForkJoinTask : Workflow<ForkJoinTask, ForkJoinTaskInput, ForkJoinTaskOutput>
{
public ForkJoinTaskModel ForkJoin { get; set; }

public ForkJoinTask(WorkflowDefinitionBuilder<ForkJoinTask, ForkJoinTaskInput, ForkJoinTaskOutput> 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,
};
}
}
Original file line number Diff line number Diff line change
@@ -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
}