fix Signal activities, triggers and add sample

This commit is contained in:
Sipke Schoorstra 2020-11-03 16:58:00 +01:00
parent 3fc2e773b2
commit 3a9eaeb799
31 changed files with 263 additions and 144 deletions

View file

@ -116,6 +116,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.DeclarativeCom
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.GoBackConsole", "src\samples\Elsa.Samples.GoBackConsole\Elsa.Samples.GoBackConsole.csproj", "{F4454030-8295-4575-A54C-D816235173C5}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.SignalingConsole", "src\samples\Elsa.Samples.SignalingConsole\Elsa.Samples.SignalingConsole.csproj", "{3529F81C-9256-4C24-A1E6-ADC27093AB49}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -279,6 +281,10 @@ Global
{F4454030-8295-4575-A54C-D816235173C5}.Debug|Any CPU.Build.0 = Debug|Any CPU
{F4454030-8295-4575-A54C-D816235173C5}.Release|Any CPU.ActiveCfg = Release|Any CPU
{F4454030-8295-4575-A54C-D816235173C5}.Release|Any CPU.Build.0 = Release|Any CPU
{3529F81C-9256-4C24-A1E6-ADC27093AB49}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{3529F81C-9256-4C24-A1E6-ADC27093AB49}.Debug|Any CPU.Build.0 = Debug|Any CPU
{3529F81C-9256-4C24-A1E6-ADC27093AB49}.Release|Any CPU.ActiveCfg = Release|Any CPU
{3529F81C-9256-4C24-A1E6-ADC27093AB49}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -333,6 +339,7 @@ Global
{CC39E9F8-72D2-4A5B-846B-E2764DFE1C19} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
{E422689B-9F82-45B6-BE14-62B4EEF43B57} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
{F4454030-8295-4575-A54C-D816235173C5} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
{3529F81C-9256-4C24-A1E6-ADC27093AB49} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158}

View file

@ -23,6 +23,5 @@ namespace Elsa.Builders
IActivityBuilder WithId(string? id);
IActivityBuilder WithName(string? name);
Func<ActivityExecutionContext, CancellationToken, ValueTask<IActivity>> BuildActivityAsync();
//IWorkflowBlueprint Build();
}
}

View file

@ -27,7 +27,7 @@ namespace Elsa.Builders
where T : class, IActivity;
IActivityBuilder Add<T>(
Action<ISetupActivity<T>>? setup = default,
Action<ISetupActivity<T>>? setup,
Action<IActivityBuilder>? branch = default) where T : class, IActivity;
IActivityBuilder Add<T>(

View file

@ -1,3 +1,4 @@
using System;
using Elsa.Services.Models;
namespace Elsa.Builders
@ -8,6 +9,7 @@ namespace Elsa.Builders
IActivityBuilder Source { get; }
string? Outcome { get; }
IConnectionBuilder Then(string activityName);
IConnectionBuilder Then(IActivityBuilder targetActivity, Action<IActivityBuilder>? branch = default);
IWorkflowBlueprint Build();
}
}

View file

@ -1,12 +1,14 @@
using System;
using System.Collections.Generic;
using Elsa.Models;
using Elsa.Services.Models;
namespace Elsa.Services
{
public interface IActivityActivator
{
IActivity ActivateActivity(string activityTypeName, Action<IActivity>? setup = default);
IActivity ActivateActivity(IActivityBlueprint activityBlueprint);
T ActivateActivity<T>(Action<T>? configure = default) where T : class, IActivity;
IActivity ActivateActivity(ActivityDefinition activityDefinition);
IEnumerable<Type> GetActivityTypes();

View file

@ -43,6 +43,13 @@ namespace Elsa.Services.Models
public T GetVariable<T>() => GetVariable<T>(typeof(T).Name);
public T GetService<T>() => WorkflowExecutionContext.ServiceProvider.GetService<T>();
public async ValueTask<IActivity> ActivateActivityAsync(CancellationToken cancellationToken = default)
{
var activity = ActivateActivity();
await SetActivityPropertiesAsync(activity, cancellationToken);
return activity;
}
public async ValueTask SetActivityPropertiesAsync(
IActivity activity,
CancellationToken cancellationToken = default) =>
@ -51,10 +58,10 @@ namespace Elsa.Services.Models
this,
cancellationToken);
public IActivity ActivateActivity(string activityType, Action<IActivity>? setupActivity = default)
public IActivity ActivateActivity()
{
var activityActivator = ServiceProvider.GetRequiredService<IActivityActivator>();
var activity = activityActivator.ActivateActivity(activityType, setupActivity);
var activity = activityActivator.ActivateActivity(ActivityBlueprint);
activity.Data = ActivityInstance.Data;
return activity;
}

View file

@ -1,11 +1,12 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Models;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Localization;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
namespace Elsa.Services.Models
{
@ -15,8 +16,8 @@ namespace Elsa.Services.Models
IServiceProvider serviceProvider,
IWorkflowBlueprint workflowBlueprint,
WorkflowInstance workflowInstance,
object? input,
object? workflowContext
object? input = default,
object? workflowContext = default
)
{
ServiceProvider = serviceProvider;
@ -114,5 +115,11 @@ namespace Elsa.Services.Models
public T GetOutputFrom<T>(string activityName) => (T)GetOutputFrom(activityName)!;
public void SetWorkflowContext(object? value) => WorkflowContext = value;
public T GetWorkflowContext<T>() => (T)WorkflowContext!;
public async ValueTask<IEnumerable<IActivity>> ActivateActivitiesAsync(CancellationToken cancellationToken = default)
{
var activityExecutionContexts = WorkflowBlueprint.Activities.Select(x => new ActivityExecutionContext(this, ServiceProvider, x));
return await Task.WhenAll(activityExecutionContexts.Select(async x => await x.ActivateActivityAsync(cancellationToken)));
}
}
}

View file

@ -1,9 +1,8 @@
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.Signaling
namespace Elsa.Activities.Signaling.Models
{
public class TriggeredSignal
public class Signal
{
public TriggeredSignal(string signalName, object? input)
public Signal(string signalName, object? input = default)
{
SignalName = signalName;
Input = input;

View file

@ -1,4 +1,5 @@
using System;
using Elsa.Activities.Signaling.Models;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
@ -8,29 +9,31 @@ using Elsa.Services.Models;
namespace Elsa.Activities.Signaling
{
/// <summary>
/// Halts workflow execution until the specified signal is received.
/// Suspends workflow execution until the specified signal is received.
/// </summary>
[ActivityDefinition(
Category = "Workflows",
Description = "Halt workflow execution until the specified signal is received.",
Icon = "fas fa-traffic-light"
)]
public class Signaled : Activity
public class ReceiveSignal : Activity
{
[ActivityProperty(Hint = "An expression that evaluates to the name of the signal to wait for.")]
[ActivityProperty(Hint = "The name of the signal to wait for.")]
public string Signal { get; set; } = default!;
protected override bool OnCanExecute(ActivityExecutionContext context)
{
var signal = Signal;
var triggeredSignal = (TriggeredSignal)context.Input!;
return string.Equals(triggeredSignal.SignalName, signal, StringComparison.OrdinalIgnoreCase);
if (context.Input is Signal triggeredSignal)
return string.Equals(triggeredSignal.SignalName, Signal, StringComparison.OrdinalIgnoreCase);
return false;
}
protected override IActivityExecutionResult OnExecute() => Suspend();
protected override IActivityExecutionResult OnResume(ActivityExecutionContext context)
{
var triggeredSignal = (TriggeredSignal)context.Input!;
var triggeredSignal = context.GetInput<Signal>();
return Done(triggeredSignal.Input);
}
}

View file

@ -0,0 +1,16 @@
using System;
using Elsa.Activities.Signaling;
using Elsa.Builders;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
public static class ReceiveSignalExtensions
{
public static IActivityBuilder ReceiveSignal(this IBuilder builder, Action<ISetupActivity<ReceiveSignal>> setup) => builder.Then(setup);
public static IActivityBuilder ReceiveSignal(this IBuilder builder, Func<ActivityExecutionContext, string> signal) => builder.ReceiveSignal(activity => activity.Set(x => x.Signal, signal));
public static IActivityBuilder ReceiveSignal(this IBuilder builder, Func<string> signal) => builder.ReceiveSignal(activity => activity.Set(x => x.Signal, signal));
public static IActivityBuilder ReceiveSignal(this IBuilder builder, string signal) => builder.ReceiveSignal(activity => activity.Set(x => x.Signal, signal));
}
}

View file

@ -0,0 +1,23 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Triggers;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.Signaling
{
public class ReceiveSignalTrigger : Trigger
{
public string Signal { get; set; } = default!;
public string? CorrelationId { get; set; }
}
public class ReceiveSignalTriggerProvider : TriggerProvider<ReceiveSignalTrigger, ReceiveSignal>
{
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<ReceiveSignal> context, CancellationToken cancellationToken) =>
new ReceiveSignalTrigger
{
Signal = await context.Activity.GetPropertyValueAsync(x => x.Signal, cancellationToken),
CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId
};
}
}

View file

@ -0,0 +1,10 @@
using System.Threading;
using System.Threading.Tasks;
namespace Elsa.Activities.Signaling.Services
{
public interface ISignaler
{
Task SendSignal(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default);
}
}

View file

@ -0,0 +1,25 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.Signaling.Models;
using Elsa.Services;
namespace Elsa.Activities.Signaling.Services
{
public class Signaler : ISignaler
{
private readonly IWorkflowScheduler _workflowScheduler;
public Signaler(IWorkflowScheduler workflowScheduler)
{
_workflowScheduler = workflowScheduler;
}
public async Task SendSignal(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default) =>
await _workflowScheduler.TriggerWorkflowsAsync<ReceiveSignalTrigger>(
x => x.Signal == signal && (x.CorrelationId == null || x.CorrelationId == correlationId),
new Signal(signal, input),
correlationId,
cancellationToken: cancellationToken
);
}
}

View file

@ -1,16 +0,0 @@
using System;
using Elsa.Activities.Signaling;
using Elsa.Builders;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
public static class SignaledBuilderExtensions
{
public static IActivityBuilder Signaled(this IBuilder builder, Action<ISetupActivity<Signaled>> setup) => builder.Then(setup);
public static IActivityBuilder Signaled(this IBuilder builder, Func<ActivityExecutionContext, string> signal) => builder.Signaled(activity => activity.Set(x => x.Signal, signal));
public static IActivityBuilder Signaled(this IBuilder builder, Func<string> signal) => builder.Signaled(activity => activity.Set(x => x.Signal, signal));
public static IActivityBuilder Signaled(this IBuilder builder, string signal) => builder.Signaled(activity => activity.Set(x => x.Signal, signal));
}
}

View file

@ -1,21 +0,0 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Triggers;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.Signaling
{
public class SignaledTrigger : Trigger
{
public string Signal { get; set; } = default!;
}
public class SignaledTriggerProvider : TriggerProvider<SignaledTrigger, Signaled>
{
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<Signaled> context, CancellationToken cancellationToken) =>
new SignaledTrigger
{
Signal = await context.Activity.GetPropertyValueAsync(x => x.Signal, cancellationToken)
};
}
}

View file

@ -1,48 +0,0 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.Signaling
{
[ActivityDefinition(
Category = "Workflows",
Description = "Trigger all workflows that start with or are blocked on the specified activity type.",
Icon = "fas fa-sitemap"
)]
public class TriggerEvent : Activity
{
private readonly IWorkflowScheduler _workflowScheduler;
public TriggerEvent(IWorkflowScheduler workflowScheduler)
{
_workflowScheduler = workflowScheduler;
}
[ActivityProperty(Hint = "An expression that evaluates to the activity type to use when triggering workflows.")]
public string ActivityType { get; set; } = default!;
[ActivityProperty(
Hint = "An expression that evaluates to a dictionary to be provided as input when triggering workflows."
)]
public object? Input { get; set; }
[ActivityProperty(Hint = "An expression that evaluates to the correlation ID to use when triggering workflows.")]
public string? CorrelationId { get; set; }
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
await _workflowScheduler.TriggerWorkflowsAsync(
ActivityType,
Input,
CorrelationId,
cancellationToken: cancellationToken
);
return Done();
}
}
}

View file

@ -1,5 +1,6 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.Signaling.Services;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
@ -13,16 +14,16 @@ namespace Elsa.Activities.Signaling
/// </summary>
[ActivityDefinition(
Category = "Workflows",
Description = "Trigger the specified signal.",
Description = "Sends the specified signal.",
Icon = "fas fa-broadcast-tower"
)]
public class TriggerSignal : Activity
public class SendSignal : Activity
{
private readonly IWorkflowScheduler _workflowScheduler;
private readonly ISignaler _signaler;
public TriggerSignal(IWorkflowScheduler workflowScheduler)
public SendSignal(ISignaler signaler)
{
_workflowScheduler = workflowScheduler;
_signaler = signaler;
}
[ActivityProperty(Hint = "An expression that evaluates to the name of the signal to trigger.")]
@ -34,19 +35,9 @@ namespace Elsa.Activities.Signaling
[ActivityProperty(Hint = "An expression that evaluates to an input value when triggering the signal.")]
public object? Input { get; set; }
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(
ActivityExecutionContext context,
CancellationToken cancellationToken)
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
var triggeredSignal = new TriggeredSignal(Signal, Input);
await _workflowScheduler.TriggerWorkflowsAsync(
nameof(Signaled),
triggeredSignal,
CorrelationId,
cancellationToken: cancellationToken
);
await _signaler.SendSignal(Signal, Input, CorrelationId, cancellationToken);
return Done();
}
}

View file

@ -75,11 +75,10 @@ namespace Elsa.Builders
public Func<ActivityExecutionContext, CancellationToken, ValueTask<IActivity>> BuildActivityAsync() =>
async (context, cancellationToken) =>
{
var activity = context.ActivateActivity(context.ActivityBlueprint.Type);
var activity = await context.ActivateActivityAsync(cancellationToken);
activity.Id = ActivityId;
activity.Name = Name;
activity.Description = Description;
await context.SetActivityPropertiesAsync(activity, cancellationToken);
return activity;
};
}

View file

@ -20,11 +20,20 @@ namespace Elsa.Builders
public IActivityBuilder Then<T>(
Action<ISetupActivity<T>>? setup = default,
Action<IActivityBuilder>? branch = default) where T : class, IActivity =>
Then(WorkflowBuilder.Add(setup), branch);
Action<IActivityBuilder>? branch = default) where T : class, IActivity
{
var activityBuilder = WorkflowBuilder.Add(setup);
Then(activityBuilder, branch);
return activityBuilder;
}
public IActivityBuilder Then<T>(Action<IActivityBuilder>? branch = default)
where T : class, IActivity => Then(WorkflowBuilder.Add<T>(branch));
where T : class, IActivity
{
var activityBuilder = WorkflowBuilder.Add<T>(branch);
Then(activityBuilder);
return activityBuilder;
}
public IConnectionBuilder Then(string activityName)
{
@ -34,11 +43,10 @@ namespace Elsa.Builders
Outcome);
}
private IActivityBuilder Then(IActivityBuilder activityBuilder, Action<IActivityBuilder>? branch = default)
public IConnectionBuilder Then(IActivityBuilder activityBuilder, Action<IActivityBuilder>? branch = default)
{
branch?.Invoke(activityBuilder);
WorkflowBuilder.Connect(Source, activityBuilder, Outcome);
return activityBuilder;
return WorkflowBuilder.Connect(Source, activityBuilder, Outcome);
}
public IWorkflowBlueprint Build() => ((IWorkflowBuilder)WorkflowBuilder).BuildBlueprint();

View file

@ -1,3 +1,4 @@
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Ccontrolflow_005Cfor/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Csignaling_005Creceivesignal/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Csignaling_005Csignaled/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>

View file

@ -4,6 +4,7 @@ using Elsa;
using Elsa.Activities.ControlFlow;
using Elsa.Activities.Primitives;
using Elsa.Activities.Signaling;
using Elsa.Activities.Signaling.Services;
using Elsa.Activities.Workflows;
using Elsa.Builders;
using Elsa.Consumers;
@ -149,10 +150,10 @@ namespace Microsoft.Extensions.DependencyInjection
.AddActivity<While>()
.AddActivity<Correlate>()
.AddActivity<SetVariable>()
.AddActivity<Signaled>()
.AddTriggerProvider<SignaledTriggerProvider>()
.AddActivity<TriggerEvent>()
.AddActivity<TriggerSignal>()
.AddActivity<ReceiveSignal>()
.AddTriggerProvider<ReceiveSignalTriggerProvider>()
.AddActivity<SendSignal>()
.AddScoped<ISignaler, Signaler>()
.AddActivity<Elsa.Activities.Workflows.RunWorkflow>()
.AddTriggerProvider<RunWorkflowTriggerProvider>();

View file

@ -2,6 +2,7 @@ using System;
using System.Collections.Generic;
using System.Linq;
using Elsa.Models;
using Elsa.Services.Models;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Services
@ -48,6 +49,18 @@ namespace Elsa.Services
return activity;
}
public IActivity ActivateActivity(IActivityBlueprint activityBlueprint)
{
return ActivateActivity(
activityBlueprint.Type,
activity =>
{
activity.Id = activityBlueprint.Id;
activity.Name = activityBlueprint.Name;
activity.PersistWorkflow = activityBlueprint.PersistWorkflow;
});
}
public T ActivateActivity<T>(Action<T>? setup = null) where T : class, IActivity
{
var activity = ActivatorUtilities.GetServiceOrCreateInstance<T>(_serviceProvider);

View file

@ -98,7 +98,7 @@ namespace Elsa.Services
private static async ValueTask<IActivity> CreateActivityAsync(ActivityDefinition activityDefinition, ActivityExecutionContext context, CancellationToken cancellationToken)
{
var activity = context.ActivateActivity(activityDefinition.Type);
var activity = context.ActivateActivity();
activity.Description = activityDefinition.Description;
activity.Id = activityDefinition.ActivityId;
activity.Name = activityDefinition.Name;

View file

@ -2,7 +2,6 @@
using Elsa.Activities.ControlFlow;
using Elsa.Builders;
using Elsa.Samples.GoBackConsole.Activities;
using Elsa.Services.Models;
namespace Elsa.Samples.GoBackConsole.Workflows
{

View file

@ -0,0 +1,16 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>netcoreapp3.1</TargetFramework>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.DependencyInjection" Version="3.1.9" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa\Elsa.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,51 @@
using System;
using System.Threading.Tasks;
using Elsa.Activities.Signaling.Services;
using Elsa.Services;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Samples.SignalingConsole
{
/// <summary>
/// Demonstrates a workflow with a While looping construct.
/// </summary>
static class Program
{
private static async Task Main()
{
// Create a service container with Elsa services.
var services = new ServiceCollection()
.AddElsa()
.AddConsoleActivities()
.AddWorkflow<TrafficLightWorkflow>()
.BuildServiceProvider();
// Run startup actions (not needed when registering Elsa with a Host).
var startupRunner = services.GetRequiredService<IStartupRunner>();
await startupRunner.StartupAsync();
// Get a workflow runner.
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
// Define a couple of cars so we can correlate workflows with them.
var cars = new[] { "Car 1", "Car 2" };
// Execute a workflow for each car.
foreach (var car in cars)
await workflowRunner.RunWorkflowAsync<TrafficLightWorkflow>(correlationId: car);
Console.WriteLine("Hit enter to signal green light for Car 2.");
Console.ReadLine();
// The workflows are now suspended at the red light.
// Trigger a green light signal for the first car.
var signaler = services.GetRequiredService<ISignaler>();
await signaler.SendSignal("Green", correlationId: "Car 2");
// Notice that only the workflow correlated to the second car executed.
// Keep the application alive for the workflow scheduler to have enough time to resume the workflow.
Console.ReadLine();
}
}
}

View file

@ -0,0 +1,20 @@
using Elsa.Activities.Console;
using Elsa.Activities.ControlFlow;
using Elsa.Builders;
using Elsa.Services.Models;
namespace Elsa.Samples.SignalingConsole
{
public class TrafficLightWorkflow : IWorkflow
{
public void Build(IWorkflowBuilder workflow)
{
workflow
.WriteLine(context => $"{GetCarName(context)} is approaching red traffic light...")
.ReceiveSignal("Green")
.WriteLine(context => $"Light turned green for {GetCarName(context)}. Hit that power pedal!");
}
private string GetCarName(ActivityExecutionContext context) => context.WorkflowExecutionContext.CorrelationId;
}
}

View file

@ -23,7 +23,7 @@ namespace Elsa.Core.IntegrationTests.Workflows
_items,
iterate => iterate
.Then<WriteLine>(activity => activity.Set(x => x.Text, context => $"{context.Input}")).WithId("WriteLine")
.Then<Signaled>() /* Block workflow.*/
.Then<ReceiveSignal>() /* Block workflow.*/
.WriteLine("Resumed"))
.WriteLine("One iterations executing, rest is blocked");
}

View file

@ -20,9 +20,9 @@ namespace Elsa.Core.IntegrationTests.Workflows
activity => activity.Set(x => x.Branches, new HashSet<string>(new[] { "Branch 1", "Branch 2", "Branch 3" })),
fork =>
{
fork.When("Branch 1").Signaled("Signal1").WriteLine("Branch 1 executed", "WriteLine1").Then("Join");
fork.When("Branch 2").Signaled("Signal2").WriteLine("Branch 2 executed", "WriteLine2").Then("Join");
fork.When("Branch 3").Signaled("Signal3").WriteLine("Branch 3 executed", "WriteLine3").Then("Join");
fork.When("Branch 1").ReceiveSignal("Signal1").WriteLine("Branch 1 executed", "WriteLine1").Then("Join");
fork.When("Branch 2").ReceiveSignal("Signal2").WriteLine("Branch 2 executed", "WriteLine2").Then("Join");
fork.When("Branch 3").ReceiveSignal("Signal3").WriteLine("Branch 3 executed", "WriteLine3").Then("Join");
})
.Add<Join>(join => join.Set(x => x.Mode, _joinMode)).WithName("Join")
.WriteLine("Finished", "Finished");

View file

@ -2,10 +2,10 @@ using System.Linq;
using System.Threading.Tasks;
using Elsa.Activities.ControlFlow;
using Elsa.Activities.Signaling;
using Elsa.Activities.Signaling.Models;
using Elsa.Models;
using Elsa.Services.Models;
using Elsa.Testing.Shared.Helpers;
using Open.Linq.AsyncExtensions;
using Xunit;
using Xunit.Abstractions;
@ -58,7 +58,7 @@ namespace Elsa.Core.IntegrationTests.Workflows
var workflow = new ForkJoinWorkflow(Join.JoinMode.WaitAny);
var workflowBlueprint = WorkflowBuilder.Build(workflow);
var workflowInstance = await WorkflowRunner.RunWorkflowAsync(workflowBlueprint);
bool GetActivityHasExecuted(string name) => (from entry in workflowInstance.ExecutionLog let activity = workflowBlueprint.Activities.First(x => x.Id == entry.ActivityId) where activity.Name == name select activity).Any();
bool GetIsFinished() => GetActivityHasExecuted("Finished");
@ -74,9 +74,14 @@ namespace Elsa.Core.IntegrationTests.Workflows
private async Task<WorkflowInstance> TriggerSignalAsync(IWorkflowBlueprint workflowBlueprint, WorkflowInstance workflowInstance, string signal)
{
var signaled = await WorkflowSelector.GetTriggersAsync<SignaledTrigger>(x => x.Signal == signal).First();
var triggeredSignal = new TriggeredSignal(signal, null);
return await WorkflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, signaled.Id, triggeredSignal);
var workflowExecutionContext = new WorkflowExecutionContext(ServiceProvider, workflowBlueprint, workflowInstance);
var activities = await workflowExecutionContext.ActivateActivitiesAsync();
var blockingActivityIds = workflowInstance.BlockingActivities.Where(x => x.ActivityType == nameof(ReceiveSignal)).Select(x => x.ActivityId);
var receiveSignalActivities = activities.Where(x => blockingActivityIds.Contains(x.Id)).Cast<ReceiveSignal>();
var receiveSignal = receiveSignalActivities.Single(x => x.Signal == signal);
var triggeredSignal = new Signal(signal);
return await WorkflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, receiveSignal.Id, triggeredSignal);
}
}
}

View file

@ -23,7 +23,7 @@ namespace Elsa.Core.IntegrationTests.Workflows
_items,
iterate => iterate
.Then<WriteLine>(activity => activity.Set(x => x.Text, context => $"{context.Input}")).WithId("WriteLine")
.Then<Signaled>() /* Block workflow.*/
.Then<ReceiveSignal>() /* Block workflow.*/
.WriteLine("Resumed"))
.WriteLine("All iterations executing in parallel");
}