Composite Activity fixes and improvements (#3512)

* Implement consistent input/output/variable memory block identifiers

* Telnyx fixes and improvements

* Fix bookmark burning

* Refactor AnswerCall

* Incremental work on source + line registration

* Handle missing bookmarks

* Fix bookmark persisted callback mechanism

* Implement distributed locking for default workflow runtime

* Fix recording bookmark

* Increase source file and line coverage
This commit is contained in:
Sipke Schoorstra 2022-12-07 20:47:36 +01:00 committed by GitHub
parent 4ece4bb51e
commit 026430fd2d
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
99 changed files with 1472 additions and 472 deletions

View file

@ -10,6 +10,7 @@ namespace Elsa.ActivityDefinitions.Activities;
/// </summary>
public class ActivityDefinitionActivity : ActivityBase
{
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
// Construct the root activity stored in the activity definitions.
@ -17,8 +18,8 @@ public class ActivityDefinitionActivity : ActivityBase
var root = await materializer.MaterializeAsync(this, context.CancellationToken);
// Schedule the activity for execution.
await context.ScheduleActivityAsync(root, onChildCompletedAsync);
await context.ScheduleActivityAsync(root, OnChildCompletedAsync);
}
private async ValueTask onChildCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
private async ValueTask OnChildCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
}

View file

@ -32,7 +32,7 @@ public static class ObjectConverter
if (value is DahomeyJsonNode { ValueKind: JsonValueKind.Object } dahomyJsonObject)
return ToObject(dahomyJsonObject, targetType, options);
if (value is JsonElement { ValueKind: JsonValueKind.Object } jsonObject)
if (value is JsonElement { ValueKind: JsonValueKind.Object or JsonValueKind.Array } jsonObject)
return jsonObject.Deserialize(targetType, options);
var underlyingTargetType = Nullable.GetUnderlyingType(targetType) ?? targetType;

View file

@ -27,15 +27,7 @@ public class JavaScriptExpressionSyntaxProvider : IExpressionSyntaxProvider
Syntax = SyntaxName,
Type = typeof(JavaScriptExpression),
CreateExpression = CreateJavaScriptExpression,
CreateBlockReference = context =>
{
var reference = new JavaScriptExpressionBlockReference(context.GetExpression<JavaScriptExpression>());
if (string.IsNullOrWhiteSpace(reference.Id))
reference.Id = GenerateId();
return reference;
},
CreateBlockReference = context => new JavaScriptExpressionBlockReference(context.GetExpression<JavaScriptExpression>()),
CreateSerializableObject = context => new
{
Type = SyntaxName,

View file

@ -23,15 +23,7 @@ public class LiquidExpressionSyntaxProvider : IExpressionSyntaxProvider
Syntax = SyntaxName,
Type = typeof(LiquidExpression),
CreateExpression = CreateLiquidExpression,
CreateBlockReference = context =>
{
var reference = new LiquidExpressionBlockReference(context.GetExpression<LiquidExpression>());
if (string.IsNullOrWhiteSpace(reference.Id))
reference.Id = GenerateId();
return reference;
},
CreateBlockReference = context => new LiquidExpressionBlockReference(context.GetExpression<LiquidExpression>()),
CreateSerializableObject = context => new
{
Type = SyntaxName,

View file

@ -176,6 +176,7 @@ public class WorkflowGrain : WorkflowGrainBase
ActivityInstanceId = x.ActivityInstanceId,
Hash = x.Hash,
Data = x.Data.EmptyIfNull(),
AutoBurn = x.AutoBurn,
CallbackMethodName = x.CallbackMethodName.EmptyIfNull()
});
}

View file

@ -244,5 +244,6 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
x.Data.NullIfEmpty(),
x.ActivityId,
x.ActivityInstanceId,
x.AutoBurn,
x.CallbackMethodName.NullIfEmpty()));
}

View file

@ -81,7 +81,8 @@ message BookmarkDto {
optional string Data = 4;
string ActivityId = 5;
string ActivityInstanceId = 6;
optional string CallbackMethodName = 7;
optional bool AutoBurn = 7;
optional string CallbackMethodName = 8;
}
message Json {

View file

@ -1,4 +1,5 @@
using System.Text.Json.Serialization;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Common.Services;
using Elsa.Expressions.Models;
using Elsa.Scheduling.Models;
@ -8,42 +9,72 @@ using Elsa.Workflows.Core.Models;
namespace Elsa.Scheduling.Activities;
/// <summary>
/// Delay execution for the specified amount of time.
/// </summary>
[Activity( "Elsa", "Scheduling", "Delay execution for the specified amount of time.")]
public class Delay : Activity
{
/// <inheritdoc />
[JsonConstructor]
public Delay()
public Delay([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
public Delay(Func<ExpressionExecutionContext, TimeSpan> timeSpan, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking) : this(new Input<TimeSpan>(timeSpan))
/// <inheritdoc />
public Delay(
Func<ExpressionExecutionContext, TimeSpan> timeSpan,
DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking,
[CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input<TimeSpan>(timeSpan), blockingStrategy, source, line)
{
}
public Delay(Func<ExpressionExecutionContext, ValueTask<TimeSpan>> timeSpan, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking) : this(new Input<TimeSpan>(timeSpan))
/// <inheritdoc />
public Delay(
Func<ExpressionExecutionContext, ValueTask<TimeSpan>> timeSpan,
DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking,
[CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input<TimeSpan>(timeSpan), blockingStrategy, source, line)
{
}
public Delay(Input<TimeSpan> timeSpan, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking)
/// <inheritdoc />
public Delay(
Input<TimeSpan> timeSpan,
DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking,
[CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
TimeSpan = timeSpan;
Strategy = blockingStrategy;
}
public Delay(TimeSpan timeSpan, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking)
/// <inheritdoc />
public Delay(
TimeSpan timeSpan,
DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking,
[CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
TimeSpan = new Input<TimeSpan>(timeSpan);
Strategy = blockingStrategy;
}
public Delay(Variable<TimeSpan> timeSpan, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking)
/// <inheritdoc />
public Delay(
Variable<TimeSpan> timeSpan,
DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking,
[CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
TimeSpan = new Input<TimeSpan>(timeSpan);
Strategy = blockingStrategy;
}
/// <summary>
/// The amount of time to delay execution.
/// </summary>
[Input] public Input<TimeSpan> TimeSpan { get; set; } = default!;
/// <summary>
/// A value controlling whether the delay should happen in-process (synchronously or out of process (asynchronously).
/// </summary>
[Input] public DelayBlockingStrategy Strategy { get; set; } = DelayBlockingStrategy.NonBlocking;
/// <summary>
@ -97,11 +128,33 @@ public class Delay : Activity
await NonBlockingStrategy(timeSpan, context);
}
/// <summary>
/// Creates a new <see cref="Delay"/> from the specified number of milliseconds.
/// </summary>
public static Delay FromMilliseconds(double value, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking) => new(System.TimeSpan.FromMilliseconds(value), blockingStrategy);
/// <summary>
/// Creates a new <see cref="Delay"/> from the specified number of seconds.
/// </summary>
public static Delay FromSeconds(double value, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking) => new(System.TimeSpan.FromSeconds(value), blockingStrategy);
/// <summary>
/// Creates a new <see cref="Delay"/> from the specified number of minutes.
/// </summary>
public static Delay FromMinutes(double value, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking) => new(System.TimeSpan.FromMinutes(value), blockingStrategy);
/// <summary>
/// Creates a new <see cref="Delay"/> from the specified number of hours.
/// </summary>
public static Delay FromHours(double value, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking) => new(System.TimeSpan.FromHours(value), blockingStrategy);
/// <summary>
/// Creates a new <see cref="Delay"/> from the specified number of days.
/// </summary>
public static Delay FromDays(double value, DelayBlockingStrategy blockingStrategy = DelayBlockingStrategy.NonBlocking) => new(System.TimeSpan.FromDays(value), blockingStrategy);
}
/// <summary>
/// A bookmark payload for <see cref="Delay"/>.
/// </summary>
public record DelayPayload(DateTimeOffset ResumeAt);

View file

@ -1,4 +1,6 @@
using Elsa.Common.Services;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Common.Services;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Attributes;
@ -7,44 +9,69 @@ using Microsoft.Extensions.Logging;
namespace Elsa.Scheduling.Activities;
/// <summary>
/// Triggers the workflow at a specific future timestamp.
/// </summary>
[Activity("Elsa", "Scheduling", "Trigger execution at a specific time in the future.")]
public class StartAt : Trigger
{
public const string InputKey = "ExecuteAt";
public StartAt()
private const string InputKey = "ExecuteAt";
/// <inheritdoc />
[JsonConstructor]
public StartAt([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
public StartAt(Input<DateTimeOffset> dateTime) => DateTime = dateTime;
/// <inheritdoc />
public StartAt(Input<DateTimeOffset> dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => DateTime = dateTime;
public StartAt(Func<ExpressionExecutionContext, DateTimeOffset> dateTime) : this(new Input<DateTimeOffset>(dateTime))
{
}
public StartAt(Func<ExpressionExecutionContext, ValueTask<DateTimeOffset>> dateTime) : this(new Input<DateTimeOffset>(dateTime))
{
}
public StartAt(Func<ValueTask<DateTimeOffset>> dateTime) : this(new Input<DateTimeOffset>(dateTime))
{
}
public StartAt(Func<DateTimeOffset> dateTime) : this(new Input<DateTimeOffset>(dateTime))
/// <inheritdoc />
public StartAt(Func<ExpressionExecutionContext, DateTimeOffset> dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<DateTimeOffset>(dateTime), source, line)
{
}
public StartAt(DateTimeOffset dateTime) => DateTime = new Input<DateTimeOffset>(dateTime);
public StartAt(Variable<DateTimeOffset> dateTime) => DateTime = new Input<DateTimeOffset>(dateTime);
/// <inheritdoc />
public StartAt(
Func<ExpressionExecutionContext, ValueTask<DateTimeOffset>> dateTime,
[CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input<DateTimeOffset>(dateTime), source, line)
{
}
/// <inheritdoc />
public StartAt(Func<ValueTask<DateTimeOffset>> dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<DateTimeOffset>(dateTime), source, line)
{
}
/// <inheritdoc />
public StartAt(Func<DateTimeOffset> dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<DateTimeOffset>(dateTime), source, line)
{
}
/// <inheritdoc />
public StartAt(DateTimeOffset dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) =>
DateTime = new Input<DateTimeOffset>(dateTime);
/// <inheritdoc />
public StartAt(Variable<DateTimeOffset> dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) =>
DateTime = new Input<DateTimeOffset>(dateTime);
/// <summary>
/// The timestamp at which the workflow should be triggered.
/// </summary>
[Input] public Input<DateTimeOffset> DateTime { get; set; } = default!;
/// <inheritdoc />
protected override object GetTriggerPayload(TriggerIndexingContext context)
{
var executeAt = context.ExpressionExecutionContext.Get(DateTime);
return new StartAtPayload(executeAt);
}
/// <inheritdoc />
protected override void Execute(ActivityExecutionContext context)
{
// If external input was received, it means this activity got triggered and does not need to create a bookmark.
@ -69,7 +96,10 @@ public class StartAt : Trigger
context.CreateBookmark(payload);
}
/// <summary>
/// Creates a new <see cref="StartAt"/> activity set to trigger at the specified timestamp.
/// </summary>
public static StartAt From(DateTimeOffset value) => new(value);
}
public record StartAtPayload(DateTimeOffset ExecuteAt);
internal record StartAtPayload(DateTimeOffset ExecuteAt);

View file

@ -1,4 +1,5 @@
using System.Text.Json.Serialization;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Common.Services;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Attributes;
@ -12,22 +13,29 @@ namespace Elsa.Scheduling.Activities;
[Activity( "Elsa", "Scheduling", "Trigger workflow execution at a specific interval.")]
public class Timer : EventGenerator
{
/// <inheritdoc />
[JsonConstructor]
public Timer()
public Timer([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
public Timer(TimeSpan interval) : this(new Input<TimeSpan>(interval))
/// <inheritdoc />
public Timer(TimeSpan interval, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input<TimeSpan>(interval), source, line)
{
}
public Timer(Input<TimeSpan> interval)
/// <inheritdoc />
public Timer(Input<TimeSpan> interval, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Interval = interval;
}
/// <summary>
/// Th interval at which the timer should execute.
/// </summary>
[Input] public Input<TimeSpan> Interval { get; set; } = default!;
/// <inheritdoc />
protected override object GetTriggerPayload(TriggerIndexingContext context)
{
var interval = context.ExpressionExecutionContext.Get(Interval);
@ -36,8 +44,15 @@ public class Timer : EventGenerator
return new TimerPayload(executeAt, interval);
}
/// <summary>
/// Creates a new <see cref="Timer"/> activity set to trigger at the specified interval.
/// </summary>
public static Timer FromTimeSpan(TimeSpan value) => new(value);
/// <summary>
/// Creates a new <see cref="Timer"/> activity set to trigger at the specified interval in seconds.
/// </summary>
public static Timer FromSeconds(double value) => FromTimeSpan(TimeSpan.FromSeconds(value));
}
public record TimerPayload(DateTimeOffset StartAt, TimeSpan Interval);
internal record TimerPayload(DateTimeOffset StartAt, TimeSpan Interval);

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Attributes;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
@ -9,19 +10,64 @@ using Elsa.Workflows.Core.Activities.Flowchart.Attributes;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Runtime.Services;
using Refit;
namespace Elsa.Telnyx.Activities;
/// <inheritdoc />
[FlowNode("Connected", "Disconnected")]
public class FlowAnswerCall : AnswerCallBase
{
/// <inheritdoc />
public FlowAnswerCall([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected override async ValueTask HandleConnectedAsync(ActivityExecutionContext context) => await context.CompleteActivityAsync(new Outcomes("Connected"));
/// <inheritdoc />
protected override async ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => await context.CompleteActivityAsync(new Outcomes("Disconnected"));
}
/// <inheritdoc />
public class AnswerCall : AnswerCallBase
{
/// <inheritdoc />
public AnswerCall([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The activity to schedule when the call was successfully answered.
/// </summary>
[Port]
public IActivity? Connected { get; set; }
/// <summary>
/// The activity to schedule when the call was no longer active.
/// </summary>
[Port]
public IActivity? Disconnected { get; set; }
protected override async ValueTask HandleConnectedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Connected);
protected override async ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Disconnected);
}
/// <summary>
/// Answer an incoming call. You must issue this command before executing subsequent commands on an incoming call.
/// </summary>
[Activity(Constants.Namespace, "Answer an incoming call. You must issue this command before executing subsequent commands on an incoming call.", Kind = ActivityKind.Task)]
[FlowNode("Connected", "Disconnected")]
[WebhookDriven(WebhookEventTypes.CallAnswered)]
public class AnswerCall : ActivityBase<CallAnsweredPayload>, IBookmarksPersistedHandler
public abstract class AnswerCallBase : ActivityBase<CallAnsweredPayload>, IBookmarksPersistedHandler
{
/// <inheritdoc />
protected AnswerCallBase(string? source = default, int? line = default) : base(source, line)
{
}
/// <summary>
/// The call control ID to answer. Leave blank when the workflow is driven by an incoming call and you wish to pick up that one.
/// </summary>
@ -38,14 +84,16 @@ public class AnswerCall : ActivityBase<CallAnsweredPayload>, IBookmarksPersisted
/// <summary>
/// Invokes Telnyx to answer the call.
/// </summary>
/// <param name="context"></param>
public async ValueTask BookmarksPersistedAsync(ActivityExecutionContext context) => await InvokeTelnyxAsync(context);
protected abstract ValueTask HandleConnectedAsync(ActivityExecutionContext context);
protected abstract ValueTask HandleDisconnectedAsync(ActivityExecutionContext context);
private async ValueTask ResumeAsync(ActivityExecutionContext context)
{
var payload = context.GetInput<CallAnsweredPayload>();
context.Set(Result, payload);
await context.CompleteActivityAsync(new Outcomes("Connected"));
await HandleConnectedAsync(context);
}
/// <summary>
@ -56,15 +104,15 @@ public class AnswerCall : ActivityBase<CallAnsweredPayload>, IBookmarksPersisted
var callControlId = context.GetPrimaryCallControlId(CallControlId) ?? throw new Exception("CallControlId is required.");
var request = new AnswerCallRequest();
var telnyxClient = context.GetRequiredService<ITelnyxClient>();
try
{
await telnyxClient.Calls.AnswerCallAsync(callControlId, request, context.CancellationToken);
await telnyxClient.Calls.AnswerCallAsync(callControlId, request, context.CancellationToken);
}
catch (ApiException e)
{
if (!await e.CallIsNoLongerActiveAsync()) throw;
await context.CompleteActivityAsync(new Outcomes("Disconnected"));
await HandleDisconnectedAsync(context);
}
}
}

View file

@ -1,30 +1,55 @@
using Elsa.Telnyx.Client.Models;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
using Elsa.Telnyx.Extensions;
using Elsa.Telnyx.Payloads.Call;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Activities.Flowchart.Attributes;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Runtime.Services;
using Refit;
namespace Elsa.Telnyx.Activities;
/// <inheritdoc />
[FlowNode("Bridged", "Disconnected")]
public class FlowBridgeCalls : BridgeCallsBase
{
/// <inheritdoc />
public FlowBridgeCalls([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
protected override ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => context.CompleteActivityAsync("Disconnected");
protected override ValueTask HandleBridgedAsync(ActivityExecutionContext context) => context.CompleteActivityAsync("Bridged");
}
/// <inheritdoc />
public class BridgeCalls : BridgeCallsBase
{
[Port]public IActivity? Disconnected { get; set; }
[Port]public IActivity? Bridged { get; set; }
/// <inheritdoc />
public BridgeCalls([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The <see cref="IActivity"/> to execute when the source leg call is no longer active.
/// </summary>
[Port]public IActivity? Disconnected { get; set; }
/// <summary>
/// The <see cref="IActivity"/> to execute when the two calls are bridged.
/// </summary>
[Port]public IActivity? Bridged { get; set; }
/// <inheritdoc />
protected override async ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Disconnected, OnCompleted);
/// <inheritdoc />
protected override async ValueTask HandleBridgedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Bridged, OnCompleted);
}
@ -32,9 +57,13 @@ public class BridgeCalls : BridgeCallsBase
/// Bridge two calls.
/// </summary>
[Activity(Constants.Namespace, "Bridge two calls.", Kind = ActivityKind.Task)]
[FlowNode("Bridged", "Disconnected")]
public abstract class BridgeCallsBase : ActivityBase<BridgedCallsOutput>
public abstract class BridgeCallsBase : ActivityBase<BridgedCallsOutput>, IBookmarksPersistedHandler
{
/// <inheritdoc />
protected BridgeCallsBase(string? source = default, int? line = default) : base(source, line)
{
}
/// <summary>
/// The source call control ID of one of the call to bridge with. Leave empty to use the ambient inbound call control Id, if there is one.
/// </summary>
@ -48,7 +77,7 @@ public abstract class BridgeCallsBase : ActivityBase<BridgedCallsOutput>
public Input<string?>? CallControlIdB { get; set; }
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
public async ValueTask BookmarksPersistedAsync(ActivityExecutionContext context)
{
var callControlIdA = context.GetPrimaryCallControlId(CallControlIdA) ?? throw new Exception("CallControlA is required");
var callControlIdB = context.GetSecondaryCallControlId(CallControlIdB) ?? throw new Exception("CallControlB is required");
@ -58,26 +87,34 @@ public abstract class BridgeCallsBase : ActivityBase<BridgedCallsOutput>
try
{
await telnyxClient.Calls.BridgeCallsAsync(callControlIdA, request, context.CancellationToken);
context.CreateBookmark(ResumeAsync);
}
catch (ApiException e)
{
if (!await e.CallIsNoLongerActiveAsync()) throw;
await context.CompleteActivityAsync("Disconnected");
await HandleDisconnectedAsync(context);
}
}
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
var callControlIdA = context.GetPrimaryCallControlId(CallControlIdA) ?? throw new Exception("CallControlA is required");
var callControlIdB = context.GetSecondaryCallControlId(CallControlIdB) ?? throw new Exception("CallControlB is required");
var bookmarkA = new CallBridgedBookmarkPayload(callControlIdA);
var bookmarkB = new CallBridgedBookmarkPayload(callControlIdB);
context.CreateBookmarks(new[]{ bookmarkA, bookmarkB }, ResumeAsync);
}
protected abstract ValueTask HandleDisconnectedAsync(ActivityExecutionContext context);
protected abstract ValueTask HandleBridgedAsync(ActivityExecutionContext context);
protected async ValueTask OnCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
private async ValueTask ResumeAsync(ActivityExecutionContext context)
{
var payload = context.GetInput<CallBridgedPayload>()!;
var callControlIdA = context.GetPrimaryCallControlId(CallControlIdA);
var callControlIdB = context.GetPrimaryCallControlId(CallControlIdA);
var callControlIdB = context.GetSecondaryCallControlId(CallControlIdB);
if (payload.CallControlId == callControlIdA) context.SetProperty("CallBridgedPayloadA", payload);
if (payload.CallControlId == callControlIdB) context.SetProperty("CallBridgedPayloadB", payload);
@ -88,11 +125,8 @@ public abstract class BridgeCallsBase : ActivityBase<BridgedCallsOutput>
if (callBridgedPayloadA != null && callBridgedPayloadB != null)
{
context.Set(Result, new BridgedCallsOutput(callBridgedPayloadA, callBridgedPayloadB));
await context.CompleteActivityAsync(new Outcomes("Bridged"));
return;
await HandleBridgedAsync(context);
}
context.CreateBookmark();
}
}

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Attributes;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
@ -23,6 +24,11 @@ namespace Elsa.Telnyx.Activities;
[WebhookDriven(WebhookEventTypes.CallAnswered, WebhookEventTypes.CallHangup, WebhookEventTypes.CallMachineGreetingEnded, WebhookEventTypes.CallMachinePremiumGreetingEnded)]
public abstract class DialBase : ActivityBase
{
/// <inheritdoc />
protected DialBase(string? source = default, int? line = default) : base(source, line)
{
}
[Input(Description = "The DID or SIP URI to dial out and bridge to the given call.")]
public Input<string> To { get; set; } = default!;
@ -117,32 +123,57 @@ public abstract class DialBase : ActivityBase
}
}
/// <inheritdoc />
[FlowNode("Answered", "Hangup", "Voicemail")]
public class FlowDial : DialBase
{
private FlowDial([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected override async ValueTask OnHandleAnsweredAsync(ActivityExecutionContext context, CallAnsweredPayload payload) => await context.CompleteActivityWithOutcomesAsync("Answered");
/// <inheritdoc />
protected override async ValueTask OnHandleHangupAsync(ActivityExecutionContext context, CallHangupPayload payload) => await context.CompleteActivityWithOutcomesAsync("Hangup");
/// <inheritdoc />
protected override async ValueTask OnHandleMachineGreetingEndedAsync(ActivityExecutionContext context, CallMachineGreetingEndedBase payload) => await context.CompleteActivityWithOutcomesAsync("Voicemail");
}
/// <summary>
/// Dial a phone number or SIP URI.
/// </summary>
public class Dial : DialBase
{
/// <inheritdoc />
public Dial([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The <see cref="IActivity"/> to execute when the call was answered.
/// </summary>
[Port] public IActivity? Answered { get; set; }
/// <summary>
/// The <see cref="IActivity"/> to execute when there is no reply.
/// </summary>
[Port] public IActivity? Hangup { get; set; }
/// <summary>
/// The <see cref="IActivity"/> to execute when a robot answered the call.
/// </summary>
[Port] public IActivity? Voicemail { get; set; }
protected override async ValueTask OnHandleAnsweredAsync(ActivityExecutionContext context, CallAnsweredPayload payload)
{
/// <inheritdoc />
protected override async ValueTask OnHandleAnsweredAsync(ActivityExecutionContext context, CallAnsweredPayload payload) =>
await context.ScheduleActivityAsync(Answered);
}
protected override async ValueTask OnHandleHangupAsync(ActivityExecutionContext context, CallHangupPayload payload)
{
await context.ScheduleActivityAsync(Hangup);
}
/// <inheritdoc />
protected override async ValueTask OnHandleHangupAsync(ActivityExecutionContext context, CallHangupPayload payload) => await context.ScheduleActivityAsync(Hangup);
protected override async ValueTask OnHandleMachineGreetingEndedAsync(ActivityExecutionContext context, CallMachineGreetingEndedBase payload)
{
/// <inheritdoc />
protected override async ValueTask OnHandleMachineGreetingEndedAsync(ActivityExecutionContext context, CallMachineGreetingEndedBase payload) =>
await context.ScheduleActivityAsync(Voicemail);
}
}

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Attributes;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
@ -21,6 +22,11 @@ namespace Elsa.Telnyx.Activities;
[WebhookDriven(WebhookEventTypes.CallGatherEnded)]
public class GatherUsingAudio : ActivityBase<CallGatherEndedPayload>, IBookmarksPersistedHandler
{
/// <inheritdoc />
public GatherUsingAudio([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The call control ID of the call from which to gather input. Leave empty to use the ambient call control ID, if there is any.
/// </summary>

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Attributes;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
@ -22,6 +23,11 @@ namespace Elsa.Telnyx.Activities;
[WebhookDriven(WebhookEventTypes.CallGatherEnded)]
public class GatherUsingSpeak : ActivityBase<CallGatherEndedPayload>, IBookmarksPersistedHandler
{
/// <inheritdoc />
public GatherUsingSpeak([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The call control ID of the call from which to gather input. Leave empty to use the ambient call control ID, if there is any.
/// </summary>

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Client.Models;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
using Elsa.Telnyx.Extensions;
using Elsa.Workflows.Core;
@ -10,20 +11,45 @@ using Refit;
namespace Elsa.Telnyx.Activities;
/// <inheritdoc />
[FlowNode("Done", "Disconnected")]
public class FlowHangupCall : HangupCallBase
{
/// <inheritdoc />
public FlowHangupCall([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected override ValueTask HandleDoneAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Done");
/// <inheritdoc />
protected override ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Disconnected");
}
/// <inheritdoc />
public class HangupCall : HangupCallBase
{
[Port] public IActivity? Done { get; set; }
/// <inheritdoc />
public HangupCall([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The <see cref="IActivity"/> to execute when the call was no longer active.
/// </summary>
[Port] public IActivity? Disconnected { get; set; }
protected override async ValueTask HandleDoneAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Done, OnCompletedAsync);
/// <inheritdoc />
protected override async ValueTask HandleDoneAsync(ActivityExecutionContext context) => await context.CompleteActivityAsync(OnCompletedAsync);
/// <inheritdoc />
protected override async ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Disconnected, OnCompletedAsync);
/// <summary>
/// Executed when any child activity completed.
/// </summary>
private async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
}
/// <summary>
@ -32,6 +58,11 @@ public class HangupCall : HangupCallBase
[Activity(Constants.Namespace, "Hang up the call.", Kind = ActivityKind.Task)]
public abstract class HangupCallBase : ActivityBase
{
/// <inheritdoc />
protected HangupCallBase([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// Unique identifier and token for controlling the call.
/// </summary>
@ -57,7 +88,13 @@ public abstract class HangupCallBase : ActivityBase
}
}
/// <summary>
/// Executed when the call was hangup.
/// </summary>
protected abstract ValueTask HandleDoneAsync(ActivityExecutionContext context);
/// <summary>
/// Executed when the call was no longer active.
/// </summary>
protected abstract ValueTask HandleDisconnectedAsync(ActivityExecutionContext context);
protected async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
}

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Extensions;
@ -11,6 +12,9 @@ using Elsa.Workflows.Management.Models;
namespace Elsa.Telnyx.Activities;
/// <summary>
/// Triggered when an inbound phone call is received for any of the specified source or destination phone numbers.
/// </summary>
[Activity(
"Telnyx",
"Telnyx",
@ -18,13 +22,27 @@ namespace Elsa.Telnyx.Activities;
Kind = ActivityKind.Trigger)]
public class IncomingCall : Trigger<CallInitiatedPayload>
{
/// <inheritdoc />
public IncomingCall([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// A list of destination numbers to respond to.
/// </summary>
[Input(Description = "A list of destination numbers to respond to.", UIHint = InputUIHints.MultiText)]
public Input<ICollection<string>> To { get; set; } = default!;
/// <summary>
/// A list of source numbers to respond to.
/// </summary>
[Input(Description = "A list of source numbers to respond to.", UIHint = InputUIHints.MultiText)]
public Input<ICollection<string>> From { get; set; } = default!;
[Input(Description = "Match any inbound calls")]
/// <summary>
/// Match any inbound calls.
/// </summary>
[Input(Description = "Match any inbound calls.")]
public Input<bool> CatchAll { get; set; } = default!;
/// <inheritdoc />

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Client.Models;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Attributes;
@ -13,6 +14,11 @@ namespace Elsa.Telnyx.Activities;
[Activity(Constants.Namespace, "Returns information about the provided phone number.", Kind = ActivityKind.Task)]
public class LookupNumber : Activity<NumberLookupResponse>
{
/// <inheritdoc />
public LookupNumber([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The phone number to be looked up.
/// </summary>

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Attributes;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
@ -14,20 +15,47 @@ using Refit;
namespace Elsa.Telnyx.Activities;
/// <inheritdoc />
[FlowNode("Playback started", "Disconnected")]
public class FlowPlayAudio : PlayAudioBase
{
/// <inheritdoc />
public FlowPlayAudio([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected override ValueTask HandlePlaybackStartedAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Playback started");
/// <inheritdoc />
protected override ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Disconnected");
}
/// <inheritdoc />
public class PlayAudio : PlayAudioBase
{
/// <inheritdoc />
public PlayAudio([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The <see cref="IActivity"/> to execute when audio playback has started.
/// </summary>
[Port] public IActivity? PlaybackStarted { get; set; }
/// <summary>
/// The <see cref="IActivity"/> to execute when the call was no longer active.
/// </summary>
[Port] public IActivity? Disconnected { get; set; }
/// <inheritdoc />
protected override async ValueTask HandlePlaybackStartedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(PlaybackStarted, OnCompletedAsync);
/// <inheritdoc />
protected override async ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Disconnected, OnCompletedAsync);
private async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
}
/// <summary>
@ -38,6 +66,11 @@ public class PlayAudio : PlayAudioBase
[WebhookDriven(WebhookEventTypes.CallPlaybackStarted)]
public abstract class PlayAudioBase : ActivityBase, IBookmarksPersistedHandler
{
/// <inheritdoc />
protected PlayAudioBase([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// Unique identifier and token for controlling the call.
/// </summary>
@ -119,8 +152,16 @@ public abstract class PlayAudioBase : ActivityBase, IBookmarksPersistedHandler
/// <inheritdoc />
protected override void Execute(ActivityExecutionContext context) => context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallPlaybackStarted), ResumeAsync);
/// <summary>
/// Called when playback has started.
/// </summary>
protected abstract ValueTask HandlePlaybackStartedAsync(ActivityExecutionContext context);
/// <summary>
/// Called when the call was no longer active.
/// </summary>
protected abstract ValueTask HandleDisconnectedAsync(ActivityExecutionContext context);
protected async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
private async ValueTask ResumeAsync(ActivityExecutionContext context) => await HandlePlaybackStartedAsync(context);
}

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Attributes;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
@ -13,10 +14,18 @@ using Refit;
namespace Elsa.Telnyx.Activities;
/// <summary>
/// Convert text to speech and play it back on the call.
/// </summary>
[Activity(Constants.Namespace, "Convert text to speech and play it back on the call.", Kind = ActivityKind.Task)]
[WebhookDriven(WebhookEventTypes.CallSpeakEnded)]
public abstract class SpeakTextBase : ActivityBase
{
/// <inheritdoc />
protected SpeakTextBase([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// Unique identifier and token for controlling the call.
/// </summary>
@ -105,31 +114,56 @@ public abstract class SpeakTextBase : ActivityBase
}
}
/// <summary>
/// Called when the call was no longer active.
/// </summary>
protected abstract ValueTask HandleDisconnected(ActivityExecutionContext context);
/// <summary>
/// Called when speaking has finished.
/// </summary>
protected abstract ValueTask HandleFinishedSpeaking(ActivityExecutionContext context);
private async ValueTask ResumeAsync(ActivityExecutionContext context) => await HandleFinishedSpeaking(context);
}
/// <inheritdoc />
[FlowNode("Finished speaking", "Disconnected")]
public class FlowSpeakText : SpeakTextBase
{
/// <inheritdoc />
public FlowSpeakText([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected override async ValueTask HandleDisconnected(ActivityExecutionContext context) => await context.CompleteActivityWithOutcomesAsync("Disconnected");
/// <inheritdoc />
protected override async ValueTask HandleFinishedSpeaking(ActivityExecutionContext context) => await context.CompleteActivityWithOutcomesAsync("Finished speaking");
}
/// <inheritdoc />
public class SpeakText : SpeakTextBase
{
[Port]public IActivity? FinishedSpeaking { get; set; }
[Port]public IActivity? Disconnected { get; set; }
/// <inheritdoc />
public SpeakText([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
protected override async ValueTask HandleDisconnected(ActivityExecutionContext context)
{
await context.ScheduleActivityAsync(Disconnected);
}
/// <summary>
/// The <see cref="IActivity"/> to execute when speaking has finished.
/// </summary>
[Port]public IActivity? FinishedSpeaking { get; set; }
/// <summary>
/// The <see cref="IActivity"/> to execute when the call was no longer active.
/// </summary>
[Port]public IActivity? Disconnected { get; set; }
protected override async ValueTask HandleFinishedSpeaking(ActivityExecutionContext context)
{
await context.ScheduleActivityAsync(FinishedSpeaking);
}
/// <inheritdoc />
protected override async ValueTask HandleDisconnected(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Disconnected);
/// <inheritdoc />
protected override async ValueTask HandleFinishedSpeaking(ActivityExecutionContext context) => await context.ScheduleActivityAsync(FinishedSpeaking);
}

View file

@ -1,4 +1,7 @@
using Elsa.Telnyx.Client.Models;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
using Elsa.Telnyx.Extensions;
using Elsa.Telnyx.Payloads.Call;
@ -12,28 +15,61 @@ using Refit;
namespace Elsa.Telnyx.Activities;
/// <inheritdoc />
[FlowNode("Recording finished", "Disconnected")]
public class FlowStartRecording : StartRecordingBase
{
/// <inheritdoc />
public FlowStartRecording([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected override ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Disconnected");
/// <inheritdoc />
protected override ValueTask HandleCallRecordingSavedAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Recording finished");
}
/// <inheritdoc />
public class StartRecording : StartRecordingBase
{
/// <inheritdoc />
public StartRecording([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The <see cref="IActivity"/> to execute when recording has finished.
/// </summary>
[Port] public IActivity? RecordingFinished { get; set; }
/// <summary>
/// The <see cref="IActivity"/> to executed when the call was no longer active.
/// </summary>
[Port] public IActivity? Disconnected { get; set; }
/// <inheritdoc />
protected override async ValueTask HandleCallRecordingSavedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(RecordingFinished, OnCompletedAsync);
/// <inheritdoc />
protected override async ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Disconnected, OnCompletedAsync);
private async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
}
/// <summary>
/// Start recording the call.
/// </summary>
[Activity(Constants.Namespace, "Start recording the call.", Kind = ActivityKind.Task)]
[WebhookDriven(WebhookEventTypes.CallRecordingSaved)]
public abstract class StartRecordingBase : ActivityBase<CallRecordingSavedPayload>
{
/// <inheritdoc />
protected StartRecordingBase(string? source = default, int? line = default) : base(source, line)
{
}
/// <summary>
/// Unique identifier and token for controlling the call.
/// </summary>
@ -87,7 +123,8 @@ public abstract class StartRecordingBase : ActivityBase<CallRecordingSavedPayloa
try
{
await telnyxClient.Calls.StartRecordingAsync(callControlId, request, context.CancellationToken);
context.CreateBookmark(ResumeAsync);
context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallRecordingSaved), ResumeAsync);
}
catch (ApiException e)
{
@ -96,9 +133,16 @@ public abstract class StartRecordingBase : ActivityBase<CallRecordingSavedPayloa
}
}
/// <summary>
/// Called when the recording was saved.
/// </summary>
protected abstract ValueTask HandleCallRecordingSavedAsync(ActivityExecutionContext context);
/// <summary>
/// Called when the call was no longer active.
/// </summary>
protected abstract ValueTask HandleDisconnectedAsync(ActivityExecutionContext context);
protected async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
private async ValueTask ResumeAsync(ActivityExecutionContext context)
{

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Client.Models;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
using Elsa.Telnyx.Extensions;
using Elsa.Workflows.Core;
@ -10,29 +11,55 @@ using Refit;
namespace Elsa.Telnyx.Activities;
[FlowNode("Playback ended", "Disconnected")]
/// <inheritdoc />
[FlowNode("Done", "Disconnected")]
public class FlowStopAudioPlayback : StopAudioPlaybackBase
{
protected override ValueTask HandlePlaybackEndedAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Playback ended");
/// <inheritdoc />
public FlowStopAudioPlayback([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected override ValueTask HandleDoneAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Done");
/// <inheritdoc />
protected override ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => context.CompleteActivityWithOutcomesAsync("Disconnected");
}
/// <inheritdoc />
public class StopAudioPlayback : StopAudioPlaybackBase
{
[Port] public IActivity? PlaybackEnded { get; set; }
/// <inheritdoc />
public StopAudioPlayback([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The <see cref="IActivity"/> to execute when the call was no longer active.
/// </summary>
[Port] public IActivity? Disconnected { get; set; }
protected override async ValueTask HandlePlaybackEndedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(PlaybackEnded, OnCompletedAsync);
/// <inheritdoc />
protected override async ValueTask HandleDoneAsync(ActivityExecutionContext context) => await context.CompleteActivityAsync();
/// <inheritdoc />
protected override async ValueTask HandleDisconnectedAsync(ActivityExecutionContext context) => await context.ScheduleActivityAsync(Disconnected, OnCompletedAsync);
private async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
}
/// <summary>
/// Stop audio playback.
/// </summary>
[Activity(Constants.Namespace, Description = "Stop audio playback.", Kind = ActivityKind.Task)]
[FlowNode("Playback ended", "Disconnected")]
public abstract class StopAudioPlaybackBase : ActivityBase
{
/// <inheritdoc />
protected StopAudioPlaybackBase(string? source = default, int? line = default) : base(source, line)
{
}
/// <summary>
/// Unique identifier and token for controlling the call.
/// </summary>
@ -53,6 +80,7 @@ public abstract class StopAudioPlaybackBase : ActivityBase
)]
public Input<string?> Stop { get; set; } = new("all");
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
var request = new StopAudioPlaybackRequest(Stop.Get(context));
@ -62,18 +90,21 @@ public abstract class StopAudioPlaybackBase : ActivityBase
try
{
await telnyxClient.Calls.StopAudioPlaybackAsync(callControlId, request, context.CancellationToken);
context.CreateBookmark(ResumeAsync);
await HandleDoneAsync(context);
}
catch (ApiException e)
{
if (!await e.CallIsNoLongerActiveAsync()) throw;
await context.CompleteActivityWithOutcomesAsync("Disconnected");
await HandleDisconnectedAsync(context);
}
}
protected abstract ValueTask HandlePlaybackEndedAsync(ActivityExecutionContext context);
/// <summary>
/// Called when audio playback is stopping.
/// </summary>
protected abstract ValueTask HandleDoneAsync(ActivityExecutionContext context);
/// <summary>
/// Called when the call was no longer active.
/// </summary>
protected abstract ValueTask HandleDisconnectedAsync(ActivityExecutionContext context);
protected async ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext) => await context.CompleteActivityAsync();
private async ValueTask ResumeAsync(ActivityExecutionContext context) => await context.CompleteActivityWithOutcomesAsync("Playback ended");
}

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Client.Models;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
using Elsa.Telnyx.Extensions;
using Elsa.Workflows.Core;
@ -16,6 +17,11 @@ namespace Elsa.Telnyx.Activities;
[FlowNode("Recording stopped", "Disconnected")]
public class StopRecording : ActivityBase
{
/// <inheritdoc />
public StopRecording([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// Unique identifier and token for controlling the call.
/// </summary>
@ -44,6 +50,4 @@ public class StopRecording : ActivityBase
await context.CompleteActivityWithOutcomesAsync("Disconnected");
}
}
private static string? EmptyToNull(string? value) => value is "" ? null : value;
}

View file

@ -1,4 +1,5 @@
using Elsa.Telnyx.Client.Models;
using System.Runtime.CompilerServices;
using Elsa.Telnyx.Client.Models;
using Elsa.Telnyx.Client.Services;
using Elsa.Telnyx.Extensions;
using Elsa.Telnyx.Payloads.Call;
@ -18,6 +19,11 @@ namespace Elsa.Telnyx.Activities;
[FlowNode("Transferred", "Hangup", "Disconnected")]
public class TransferCall : Activity
{
/// <inheritdoc />
public TransferCall([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// Unique identifier and token for controlling the call.
/// </summary>

View file

@ -1,5 +1,8 @@
using System.ComponentModel;
using System.Reflection;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Telnyx.Attributes;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Helpers;
using Elsa.Telnyx.Models;
@ -17,22 +20,25 @@ namespace Elsa.Telnyx.Activities;
[Browsable(false)]
public class WebhookEvent : ActivityBase<Payload>
{
/// <inheritdoc />
[JsonConstructor]
public WebhookEvent()
public WebhookEvent([CallerFilePath]string? source = default, [CallerLineNumber]int? line = default) : base(source, line)
{
}
public WebhookEvent(string eventType, Variable<Payload> result)
/// <inheritdoc />
public WebhookEvent(string eventType, string activityTypeName, Variable<Payload> result, int version = 1, [CallerFilePath]string? source = default, [CallerLineNumber]int? line = default)
: base(activityTypeName, version, source, line)
{
EventType = new (eventType);
EventType = eventType;
Result = new(result);
}
/// <summary>
/// The Telnyx webhook event type to listen for.
/// </summary>
[Input(Description = "The Telnyx webhook event type to listen for")]
public Input<string> EventType { get; set; } = default!;
[Description("The Telnyx webhook event type to listen for")]
public string EventType { get; set; } = default!;
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
@ -41,9 +47,10 @@ public class WebhookEvent : ActivityBase<Payload>
await Resume(context);
else
{
var eventType = context.Get(EventType)!;
var eventType = EventType;
var payload = new WebhookEventBookmarkPayload(eventType);
context.CreateBookmark(payload, Resume);
context.CreateBookmark(new CreateBookmarkOptions(payload, Resume, Type));
}
}

View file

@ -0,0 +1,3 @@
namespace Elsa.Telnyx.Bookmarks;
public record CallBridgedBookmarkPayload(string CallControlId);

View file

@ -1,4 +1,5 @@
using System.Text.Json.Serialization;
using Dahomey.Json.Util;
namespace Elsa.Telnyx.Client.Models;
@ -10,7 +11,13 @@ public record DialResponse(
string CallSessionId,
bool IsAlive,
string RecordType
);
)
{
[JsonConstructor]
public DialResponse() : this(default!, default!, default!, default, default!)
{
}
}
public record NumberLookupResponse(
CallerName CallerName,
@ -21,7 +28,13 @@ public record NumberLookupResponse(
string PhoneNumber,
Portability Portability,
string RecordType
);
)
{
[JsonConstructor]
public NumberLookupResponse() : this(default!, default!, default!, default!, default!, default!, default!, default!)
{
}
}
public record Portability(
string Altspid,
@ -37,7 +50,13 @@ public record Portability(
string SpidCarrierName,
string SpidCarrierType,
string State
);
)
{
[JsonConstructor]
public Portability() : this(default!, default!, default!, default!, default!, default!, default!, default, default!, default!, default!, default!, default!)
{
}
}
public record Carrier(
string ErrorCode,
@ -45,7 +64,13 @@ public record Carrier(
int MobileNetworkCode,
string Name,
string Type
);
)
{
[JsonConstructor]
public Carrier() : this(default!, default!, default, default!, default!)
{
}
}
public record CallerName
{

View file

@ -0,0 +1,43 @@
using Elsa.Mediator.Services;
using Elsa.Telnyx.Activities;
using Elsa.Telnyx.Bookmarks;
using Elsa.Telnyx.Events;
using Elsa.Telnyx.Extensions;
using Elsa.Telnyx.Payloads.Abstract;
using Elsa.Telnyx.Payloads.Call;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Helpers;
using Elsa.Workflows.Runtime.Services;
using Microsoft.Extensions.Logging;
namespace Elsa.Telnyx.Handlers;
/// <summary>
/// Triggers all workflows starting with or blocked on a <see cref="IncomingCall"/> activity.
/// </summary>
internal class TriggerBridgeCallActivities : INotificationHandler<TelnyxWebhookReceived>
{
private readonly IWorkflowRuntime _workflowRuntime;
private readonly ILogger _logger;
public TriggerBridgeCallActivities(IWorkflowRuntime workflowRuntime, ILogger<TriggerBridgeCallActivities> logger)
{
_workflowRuntime = workflowRuntime;
_logger = logger;
}
public async Task HandleAsync(TelnyxWebhookReceived notification, CancellationToken cancellationToken)
{
var webhook = notification.Webhook;
var payload = webhook.Data.Payload;
if (payload is not CallBridgedPayload callBridgedPayload)
return;
var correlationId = ((Payload)webhook.Data.Payload).GetCorrelationId();;
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<BridgeCalls>();
var input = new Dictionary<string, object>().AddInput(callBridgedPayload);
var callBridgedBookmarkPayload = new CallBridgedBookmarkPayload(callBridgedPayload.CallControlId);
await _workflowRuntime.TriggerWorkflowsAsync(activityTypeName, callBridgedBookmarkPayload, new TriggerWorkflowsRuntimeOptions(correlationId, input), cancellationToken);
}
}

View file

@ -34,6 +34,9 @@ internal class TriggerIncomingCallActivities : INotificationHandler<TelnyxWebhoo
if (payload is not CallInitiatedPayload callInitiatedPayload)
return;
if (callInitiatedPayload.Direction != "incoming")
return;
var correlationId = ((Payload)webhook.Data.Payload).GetCorrelationId();;
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<IncomingCall>();
var input = new Dictionary<string, object>().AddInput(webhook);

View file

@ -7,6 +7,7 @@ using Elsa.Telnyx.Events;
using Elsa.Telnyx.Extensions;
using Elsa.Telnyx.Payloads.Abstract;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Helpers;
using Elsa.Workflows.Runtime.Services;
using Microsoft.Extensions.Logging;
@ -31,8 +32,7 @@ internal class TriggerWebhookActivities : INotificationHandler<TelnyxWebhookRece
var webhook = notification.Webhook;
var eventType = webhook.Data.EventType;
var payload = webhook.Data.Payload;
var ns = Constants.Namespace;
var activityType = $"{ns}.{payload.GetType().GetCustomAttribute<WebhookAttribute>()?.ActivityType}";
var activityType = payload.GetType().GetCustomAttribute<WebhookAttribute>()?.ActivityType;
if (activityType == null)
return;

View file

@ -2,10 +2,9 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallAnswered, ActivityTypeName, "Call Answered", "Triggered when an incoming call is answered.")]
[Webhook(WebhookEventTypes.CallAnswered, WebhookActivityTypeNames.CallAnswered, "Call Answered", "Triggered when an incoming call is answered.")]
public sealed record CallAnsweredPayload : CallPayload
{
public const string ActivityTypeName = "CallAnswered";
public string From { get; init; } = default!;
public string To { get; init; } = default!;
public string State { get; init; } = default!;

View file

@ -2,10 +2,9 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallBridged, ActivityTypeName, "Call Bridged", "Triggered when an a call is bridged.")]
[Webhook(WebhookEventTypes.CallBridged, WebhookActivityTypeNames.CallBridged, "Call Bridged", "Triggered when an a call is bridged.")]
public sealed record CallBridgedPayload : CallPayload
{
public const string ActivityTypeName = "CallBridged";
public string From { get; init; } = default!;
public string To { get; init; } = default!;
public string State { get; init; } = default!;

View file

@ -2,10 +2,9 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallDtmfReceived, ActivityTypeName, "Call DTMF Received", "Triggered when DTMF input is received.")]
[Webhook(WebhookEventTypes.CallDtmfReceived, WebhookActivityTypeNames.CallDtmfReceived, "Call DTMF Received", "Triggered when DTMF input is received.")]
public sealed record CallDtmfReceivedPayload : CallPayload
{
public const string ActivityTypeName = "CallDtmfReceived";
public string Digit { get; set; } = default!;
public string From { get; set; } = default!;
public string To { get; set; } = default!;

View file

@ -2,10 +2,9 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallGatherEnded, ActivityTypeName, "Call Gather Ended", "Triggered when an call gather has ended.")]
[Webhook(WebhookEventTypes.CallGatherEnded, WebhookActivityTypeNames.CallGatherEnded, "Call Gather Ended", "Triggered when an call gather has ended.")]
public sealed record CallGatherEndedPayload : CallPayload
{
public const string ActivityTypeName = "CallGatherEnded";
public string Digits { get; set; } = default!;
public string From { get; set; } = default!;
public string To { get; set; } = default!;

View file

@ -2,10 +2,9 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallHangup, ActivityTypeName, "Call Hangup", "Triggered when an incoming call was hangup.")]
[Webhook(WebhookEventTypes.CallHangup, WebhookActivityTypeNames.CallHangup, "Call Hangup", "Triggered when an incoming call was hangup.")]
public sealed record CallHangupPayload : CallPayload
{
public const string ActivityTypeName = "CallHangup";
public DateTimeOffset StartTime { get; init; }
public DateTimeOffset EndTime { get; init; }
public string SipHangupCause { get; init; } = default!;

View file

@ -2,10 +2,9 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallInitiated, ActivityTypeName, "Call Initiated", "Triggered when an incoming call is received.")]
[Webhook(WebhookEventTypes.CallInitiated, WebhookActivityTypeNames.CallInitiated, "Call Initiated", "Triggered when an incoming call is received.")]
public sealed record CallInitiatedPayload : CallPayload
{
public const string ActivityTypeName = "CallInitiated";
{
public string Direction { get; init; } = default!;
public string State { get; init; } = default!;
public string To { get; init; } = default!;

View file

@ -2,8 +2,5 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallMachineGreetingEnded, ActivityTypeName, "Call Machine Greeting Ended", "Triggered when a machine greeting has ended.")]
public sealed record CallMachineGreetingEnded : CallMachineGreetingEndedBase
{
public const string ActivityTypeName = "CallMachineGreetingEnded";
}
[Webhook(WebhookEventTypes.CallMachineGreetingEnded, WebhookActivityTypeNames.CallMachineGreetingEnded, "Call Machine Greeting Ended", "Triggered when a machine greeting has ended.")]
public sealed record CallMachineGreetingEnded : CallMachineGreetingEndedBase;

View file

@ -2,8 +2,12 @@ using Elsa.Telnyx.Attributes;
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallMachinePremiumDetectionEnded, ActivityTypeName, "Call Machine Premium Detection Ended", "Triggered when machine detection has ended.")]
[Webhook(
WebhookEventTypes.CallMachinePremiumDetectionEnded,
WebhookActivityTypeNames.CallMachinePremiumDetectionEnded,
"Call Machine Premium Detection Ended",
"Triggered when machine detection has ended."
)]
public sealed record CallMachinePremiumDetectionEnded : CallMachineDetectionEndedBase
{
public const string ActivityTypeName = nameof(CallMachinePremiumDetectionEnded);
}

View file

@ -2,8 +2,10 @@ using Elsa.Telnyx.Attributes;
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallMachinePremiumGreetingEnded, ActivityTypeName, "Call Machine Premium Greeting Ended", "Triggered when a machine greeting has ended.")]
public sealed record CallMachinePremiumGreetingEnded : CallMachineGreetingEndedBase
{
public const string ActivityTypeName = nameof(CallMachinePremiumGreetingEnded);
}
[Webhook(
WebhookEventTypes.CallMachinePremiumGreetingEnded,
WebhookActivityTypeNames.CallMachinePremiumGreetingEnded,
"Call Machine Premium Greeting Ended",
"Triggered when a machine greeting has ended."
)]
public sealed record CallMachinePremiumGreetingEnded : CallMachineGreetingEndedBase;

View file

@ -2,9 +2,8 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallPlaybackEnded, ActivityTypeName, "Call Playback Ended", "Triggered when an audio playback has ended.")]
[Webhook(WebhookEventTypes.CallPlaybackEnded, WebhookActivityTypeNames.CallPlaybackEnded, "Call Playback Ended", "Triggered when an audio playback has ended.")]
public sealed record CallPlaybackEndedPayload : CallPlayback
{
public const string ActivityTypeName = "CallPlaybackEnded";
public string Status { get; set; } = default!;
}

View file

@ -2,8 +2,5 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallPlaybackStarted, ActivityTypeName, "Call Playback Started", "Triggered when an audio playback has started.")]
public sealed record CallPlaybackStartedPayload : CallPlayback
{
public const string ActivityTypeName = "CallPlaybackStarted";
}
[Webhook(WebhookEventTypes.CallPlaybackStarted, WebhookActivityTypeNames.CallPlaybackStarted, "Call Playback Started", "Triggered when an audio playback has started.")]
public sealed record CallPlaybackStartedPayload : CallPlayback;

View file

@ -2,10 +2,9 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallRecordingSaved, ActivityTypeName, "Call Recording Saved", "Triggered when a recording has been saved.")]
[Webhook(WebhookEventTypes.CallRecordingSaved, WebhookActivityTypeNames.CallRecordingSaved, "Call Recording Saved", "Triggered when a recording has been saved.")]
public sealed record CallRecordingSavedPayload : CallPayload
{
public const string ActivityTypeName = "CallRecordingSaved";
public string Channels { get; set; } = default!;
public CallRecordingUrls PublicRecordingUrls { get; set; } = default!;
public CallRecordingUrls RecordingUrls { get; set; } = default!;

View file

@ -2,8 +2,5 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallSpeakEnded, ActivityTypeName, "Call Speak Ended", "Triggered when speaking has ended.")]
public sealed record CallSpeakEnded : CallPayload
{
public const string ActivityTypeName = "CallSpeakEnded";
}
[Webhook(WebhookEventTypes.CallSpeakEnded, WebhookActivityTypeNames.CallSpeakEnded, "Call Speak Ended", "Triggered when speaking has ended.")]
public sealed record CallSpeakEnded : CallPayload;

View file

@ -2,8 +2,5 @@
namespace Elsa.Telnyx.Payloads.Call;
[Webhook(WebhookEventTypes.CallSpeakStarted, ActivityTypeName, "Call Speak Started", "Triggered when speaking has started.")]
public sealed record CallSpeakStarted : CallPayload
{
public const string ActivityTypeName = "CallSpeakStarted";
}
[Webhook(WebhookEventTypes.CallSpeakStarted, WebhookActivityTypeNames.CallSpeakStarted, "Call Speak Started", "Triggered when speaking has started.")]
public sealed record CallSpeakStarted : CallPayload;

View file

@ -39,9 +39,7 @@ public class WebhookEventActivityProvider : IActivityProvider
private ActivityDescriptor CreateDescriptor(Type payloadType)
{
var webhookAttribute = payloadType.GetCustomAttribute<WebhookAttribute>() ?? throw new Exception($"No WebhookAttribute found on payload type {payloadType}");
var ns = Constants.Namespace;
var typeName = webhookAttribute.ActivityType;
var fullTypeName = $"{ns}.{typeName}";
var displayNameAttr = payloadType.GetCustomAttribute<DisplayNameAttribute>();
var displayName = displayNameAttr?.DisplayName ?? webhookAttribute.DisplayName;
var categoryAttr = payloadType.GetCustomAttribute<CategoryAttribute>();
@ -51,7 +49,7 @@ public class WebhookEventActivityProvider : IActivityProvider
return new()
{
TypeName = fullTypeName,
TypeName = typeName,
Version = 1,
DisplayName = displayName,
Description = description,
@ -62,8 +60,8 @@ public class WebhookEventActivityProvider : IActivityProvider
Constructor = context =>
{
var activity = _activityFactory.Create<WebhookEvent>(context);
activity.Type = fullTypeName;
activity.EventType = new Input<string>(webhookAttribute!.EventType);
activity.Type = typeName;
activity.EventType = webhookAttribute!.EventType;
return activity;
}

View file

@ -0,0 +1,20 @@
namespace Elsa.Telnyx;
public static class WebhookActivityTypeNames
{
public const string CallAnswered = $"{Constants.Namespace}.{nameof(CallAnswered)}";
public const string CallBridged = $"{Constants.Namespace}.{nameof(CallBridged)}";
public const string CallDtmfReceived = $"{Constants.Namespace}.{nameof(CallDtmfReceived)}";
public const string CallGatherEnded = $"{Constants.Namespace}.{nameof(CallGatherEnded)}";
public const string CallHangup = $"{Constants.Namespace}.{nameof(CallHangup)}";
public const string CallInitiated = $"{Constants.Namespace}.{nameof(CallInitiated)}";
public const string CallMachineGreetingEnded = $"{Constants.Namespace}.{nameof(CallMachineGreetingEnded)}";
public const string CallMachinePremiumGreetingEnded = $"{Constants.Namespace}.{nameof(CallMachinePremiumGreetingEnded)}";
public const string CallMachineDetectionEnded = $"{Constants.Namespace}.{nameof(CallMachineDetectionEnded)}";
public const string CallMachinePremiumDetectionEnded = $"{Constants.Namespace}.{nameof(CallMachinePremiumDetectionEnded)}";
public const string CallPlaybackStarted = $"{Constants.Namespace}.{nameof(CallPlaybackStarted)}";
public const string CallPlaybackEnded = $"{Constants.Namespace}.{nameof(CallPlaybackEnded)}";
public const string CallRecordingSaved = $"{Constants.Namespace}.{nameof(CallRecordingSaved)}";
public const string CallSpeakStarted = $"{Constants.Namespace}.{nameof(CallSpeakStarted)}";
public const string CallSpeakEnded = $"{Constants.Namespace}.{nameof(CallSpeakEnded)}";
}

View file

@ -1,12 +1,21 @@
using System.Runtime.CompilerServices;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Signals;
namespace Elsa.Workflows.Core.Activities;
[Activity("Elsa", "Control Flow", "Break out of a loop")]
/// <summary>
/// Break out of a loop.
/// </summary>
[Activity("Elsa", "Control Flow", "Break out of a loop.")]
public class Break : Activity
{
/// <inheritdoc />
public Break([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
await context.SendSignalAsync(new BreakSignal());

View file

@ -1,4 +1,5 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
@ -16,27 +17,33 @@ public class Complete : Activity
{
/// <inheritdoc />
[JsonConstructor]
public Complete()
{
}
/// <inheritdoc />
public Complete(params string[] outcomes) : this(new Input<ICollection<string>>(outcomes))
public Complete([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
public Complete(Func<ExpressionExecutionContext, ICollection<string>> outcomes) : this(new Input<ICollection<string>>(outcomes))
public Complete(IEnumerable<string> outcomes, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<ICollection<string>>(outcomes.ToList()), source, line)
{
}
/// <inheritdoc />
public Complete(Func<ExpressionExecutionContext, string> outcome) : this(context => new[] { outcome(context) })
public Complete(Func<ExpressionExecutionContext, ICollection<string>> outcomes, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<ICollection<string>>(outcomes), source, line)
{
}
/// <inheritdoc />
public Complete(Input<ICollection<string>> outcomes) => Outcomes = outcomes;
public Complete(Func<ExpressionExecutionContext, string> outcome, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(context => new[] { outcome(context) }, source, line)
{
}
/// <inheritdoc />
public Complete(Input<ICollection<string>> outcomes, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Outcomes = outcomes;
}
/// <summary>
/// The outcome or set of outcomes to complete this activity with.

View file

@ -1,4 +1,5 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
@ -15,7 +16,7 @@ namespace Elsa.Workflows.Core.Activities;
public abstract class Composite : ActivityBase
{
/// <inheritdoc />
protected Composite()
protected Composite(string? source = default, int? line = default) : base(source, line)
{
OnSignalReceived<CompleteCompositeSignal>(OnCompleteCompositeSignal);
}
@ -70,29 +71,73 @@ public abstract class Composite : ActivityBase
private async ValueTask OnCompleteCompositeSignal(CompleteCompositeSignal signal, SignalContext context)
{
var activityExecutionContext = context.ReceiverActivityExecutionContext;
// Remove the existing completed handler.
activityExecutionContext.WorkflowExecutionContext.PopCompletionCallback(activityExecutionContext, Root);
// Complete this activity.
await activityExecutionContext.CompleteActivityAsync(signal.Result);
// Complete the sender first so that it notifies its parents to complete.
await context.SenderActivityExecutionContext.CompleteActivityAsync();
// Then complete this activity.
await context.ReceiverActivityExecutionContext.CompleteActivityAsync(signal.Result);
context.StopPropagation();
}
protected static Inline Inline(Func<ActivityExecutionContext, ValueTask> activity) => new(activity);
protected static Inline Inline(Func<ValueTask> activity) => new(activity);
protected static Inline Inline(Action<ActivityExecutionContext> activity) => new(activity);
protected static Inline Inline(Action activity) => new(activity);
protected static Inline<TResult> Inline<TResult>(Func<ActivityExecutionContext, ValueTask<TResult>> activity, MemoryBlockReference? output = default) => new(activity, output);
protected static Inline<TResult> Inline<TResult>(Func<ValueTask<TResult>> activity, MemoryBlockReference? output = default) => new(activity, output);
protected static Inline<TResult> Inline<TResult>(Func<ActivityExecutionContext, TResult> activity, MemoryBlockReference? output = default) => new(activity, output);
protected static Inline<TResult> Inline<TResult>(Func<TResult> activity, MemoryBlockReference? output = default) => new(activity, output);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Func<ActivityExecutionContext, ValueTask> activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, source, line);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Func<ValueTask> activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, source, line);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Action<ActivityExecutionContext> activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, source, line);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Action activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, source, line);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> From<TResult>(Func<ActivityExecutionContext, ValueTask<TResult>> activity, MemoryBlockReference? output = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, output, source, line);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> From<TResult>(Func<ValueTask<TResult>> activity, MemoryBlockReference? output = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, output, source, line);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> From<TResult>(Func<ActivityExecutionContext, TResult> activity, MemoryBlockReference? output, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, output, source, line);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> Inline<TResult>(Func<TResult> activity, MemoryBlockReference? output = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(activity, output, source, line);
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, T value) => new(variable, value);
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, Func<ExpressionExecutionContext, T> value) => new(variable, value);
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, Func<T> value) => new(variable, value);
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, Variable<T> value) => new(variable, value);
/// <summary>
/// Creates a new <see cref="Activities.SetVariable"/> activity.
/// </summary>
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, T value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(variable, value, source, line);
/// <summary>
/// Creates a new <see cref="Activities.SetVariable"/> activity.
/// </summary>
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, Func<ExpressionExecutionContext, T> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(variable, value, source, line);
/// <summary>
/// Creates a new <see cref="Activities.SetVariable"/> activity.
/// </summary>
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, Func<T> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(variable, value, source, line);
/// <summary>
/// Creates a new <see cref="Activities.SetVariable"/> activity.
/// </summary>
protected static SetVariable<T> SetVariable<T>(Variable<T> variable, Variable<T> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(variable, value, source, line);
}
/// <summary>
@ -101,7 +146,7 @@ public abstract class Composite : ActivityBase
public abstract class Composite<T> : ActivityBase<T>
{
/// <inheritdoc />
protected Composite()
protected Composite([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
OnSignalReceived<CompleteCompositeSignal>(OnCompleteCompositeSignal);
}
@ -128,18 +173,18 @@ public abstract class Composite<T> : ActivityBase<T>
{
}
private async ValueTask OnRootCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext)
{
await OnCompletedAsync(context, childContext);
await context.CompleteActivityAsync();
}
/// <summary>
/// Override this method to handle the completion event for this composite activity.
/// </summary>
protected virtual ValueTask OnCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext)
{
OnCompleted(context, childContext);
return new();
}
/// <summary>
/// Override this method to handle the completion event for this composite activity.
/// </summary>
protected virtual void OnCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext)
{
}
@ -154,6 +199,12 @@ public abstract class Composite<T> : ActivityBase<T>
/// </summary>
protected async Task CompleteAsync(ActivityExecutionContext context, params string[] outcomes) => await CompleteAsync(context, new Outcomes(outcomes));
private async ValueTask OnRootCompletedAsync(ActivityExecutionContext context, ActivityExecutionContext childContext)
{
await OnCompletedAsync(context, childContext);
await context.CompleteActivityAsync();
}
private async ValueTask OnCompleteCompositeSignal(CompleteCompositeSignal signal, SignalContext context)
{
var activityExecutionContext = context.ReceiverActivityExecutionContext;
@ -166,12 +217,43 @@ public abstract class Composite<T> : ActivityBase<T>
context.StopPropagation();
}
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Func<ActivityExecutionContext, ValueTask> activity) => new(activity);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Func<ValueTask> activity) => new(activity);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Action<ActivityExecutionContext> activity) => new(activity);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline From(Action activity) => new(activity);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> From<TResult>(Func<ActivityExecutionContext, ValueTask<TResult>> activity, MemoryBlockReference? output = default) => new(activity, output);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> From<TResult>(Func<ValueTask<TResult>> activity, MemoryBlockReference? output = default) => new(activity, output);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> From<TResult>(Func<ActivityExecutionContext, TResult> activity, MemoryBlockReference? output = default) => new(activity, output);
/// <summary>
/// Creates a new <see cref="Activities.Inline"/> activity.
/// </summary>
protected static Inline<TResult> From<TResult>(Func<TResult> activity, MemoryBlockReference? output = default) => new(activity, output);
}

View file

@ -1,6 +1,5 @@
using System.Collections.ObjectModel;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Behaviors;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
@ -11,23 +10,22 @@ namespace Elsa.Workflows.Core.Activities;
/// </summary>
public abstract class Container : ActivityBase, IContainer
{
protected Container()
/// <inheritdoc />
protected Container(string? source = default, int? line = default) : base(source, line)
{
}
protected Container(params IActivity[] activities)
{
Activities = activities;
}
protected Container(ICollection<Variable> variables, params IActivity[] activities) : this(activities)
{
Variables = variables;
}
/// <summary>
/// The <see cref="IActivity"/>s to execute.
/// </summary>
[Port] public ICollection<IActivity> Activities { get; set; } = new HashSet<IActivity>();
/// <summary>
/// The variables available to this scope.
/// </summary>
public ICollection<Variable> Variables { get; set; } = new Collection<Variable>();
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
// Register variables.
@ -37,12 +35,18 @@ public abstract class Container : ActivityBase, IContainer
await ScheduleChildrenAsync(context);
}
/// <summary>
/// Schedule the <see cref="Activities"/> for execution.
/// </summary>
protected virtual ValueTask ScheduleChildrenAsync(ActivityExecutionContext context)
{
ScheduleChildren(context);
return ValueTask.CompletedTask;
}
/// <summary>
/// Schedule the <see cref="Activities"/> for execution.
/// </summary>
protected virtual void ScheduleChildren(ActivityExecutionContext context)
{
}

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
@ -13,36 +14,41 @@ public class Event : Trigger<object?>
{
/// <inheritdoc />
[JsonConstructor]
public Event()
public Event([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
public Event(string eventName) : this(new Literal<string>(eventName))
public Event(string eventName, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Literal<string>(eventName), source, line)
{
}
/// <inheritdoc />
public Event(Func<string> text) : this(new DelegateBlockReference<string>(text))
public Event(Func<string> text, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new DelegateBlockReference<string>(text), source, line)
{
}
/// <inheritdoc />
public Event(Func<ExpressionExecutionContext, string?> text) : this(new DelegateBlockReference<string?>(text))
public Event(Func<ExpressionExecutionContext, string?> text, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new DelegateBlockReference<string?>(text), source, line)
{
}
/// <inheritdoc />
public Event(Variable<string> variable) => EventName = new Input<string>(variable);
public Event(Variable<string> variable, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) =>
EventName = new Input<string>(variable);
/// <inheritdoc />
public Event(Literal<string> literal) => EventName = new Input<string>(literal);
public Event(Literal<string> literal, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) =>
EventName = new Input<string>(literal);
/// <inheritdoc />
public Event(DelegateBlockReference delegateBlockExpression) => EventName = new Input<string>(delegateBlockExpression);
public Event(DelegateBlockReference delegateBlockExpression, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) =>
EventName = new Input<string>(delegateBlockExpression);
/// <inheritdoc />
public Event(Input<string> eventName) => EventName = eventName;
public Event(Input<string> eventName, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => EventName = eventName;
/// <summary>
/// The name of the event to listen for.

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Activities.Flowchart.Attributes;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
@ -6,10 +7,18 @@ using Elsa.Workflows.Core.Models;
namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
/// <summary>
/// Performs a boolean condition and returns an outcome based on the the result.
/// </summary>
[FlowNode("True", "False")]
[Activity("Elsa", "Flow", "Evaluate a Boolean condition to determine which path to execute next.")]
public class FlowDecision : ActivityBase
{
/// <inheritdoc />
public FlowDecision([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The condition to evaluate.
/// </summary>

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Common.Extensions;
using Elsa.Workflows.Core.Activities.Flowchart.Contracts;
using Elsa.Workflows.Core.Activities.Flowchart.Extensions;
@ -7,10 +8,18 @@ using Elsa.Workflows.Core.Models;
namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
/// <summary>
/// Merge multiple branches into a single branch of execution.
/// </summary>
[Activity("Elsa", "Flow", "Merge multiple branches into a single branch of execution.")]
public class FlowJoin : ActivityBase, IJoinNode
{
[Input] public Input<Models.FlowJoinMode> Mode { get; set; } = new(Models.FlowJoinMode.WaitAll);
/// <inheritdoc />
public FlowJoin([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
[Input] public Input<FlowJoinMode> Mode { get; set; } = new(FlowJoinMode.WaitAll);
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
@ -10,6 +11,11 @@ namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
[Activity("Elsa", "Flow", "A simple container that executes the specified activity.")]
public class FlowNode : ActivityBase
{
/// <inheritdoc />
public FlowNode([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The activity to execute.
/// </summary>

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions;
using Elsa.Expressions.Models;
@ -9,12 +10,21 @@ using Elsa.Workflows.Core.Models;
namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
/// <summary>
/// Evaluates the specified case conditions and schedules the one that evaluates to <code>true</code>.
/// </summary>
[FlowNode("Default")]
[Activity("Elsa", "Flow", "Evaluate a set of case conditions and schedule the activity for a matching case.")]
public class FlowSwitch : ActivityBase
{
/// <inheritdoc />
public FlowSwitch([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
[Input(UIHint = "flow-switch-editor")] public ICollection<FlowSwitchCase> Cases { get; set; } = new List<FlowSwitchCase>();
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
var matchingCase = await FindMatchingCaseAsync(context.ExpressionExecutionContext);

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Workflows.Core.Activities.Flowchart.Contracts;
using Elsa.Workflows.Core.Activities.Flowchart.Extensions;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
@ -8,12 +9,16 @@ using Elsa.Workflows.Core.Signals;
namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
/// <summary>
/// A flowchart consists of a collection of activities and connections between them.
/// </summary>
[Activity("Elsa", "Flow", "A flowchart is a collection of activities and connections between them.")]
public class Flowchart : Container
{
internal const string ScopeProperty = "Scope";
public Flowchart()
/// <inheritdoc />
public Flowchart([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
OnSignalReceived<ActivityCompleted>(OnDescendantCompletedAsync);
}

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
@ -12,13 +13,15 @@ public class For : ActivityBase
{
private const string CurrentStepProperty = "CurrentStep";
/// <inheritdoc />
[JsonConstructor]
public For()
public For([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
Behaviors.Add<BreakBehavior>(this);
}
public For(int start, int end) : this()
/// <inheritdoc />
public For(int start, int end, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Start = new Input<int>(start);
End = new Input<int>(end);

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
@ -7,17 +8,22 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// Iterate over a set of values.
/// </summary>
[Activity("Elsa", "Control Flow", "Iterate over a set of values.")]
public class ForEach : ActivityBase
{
private const string CurrentIndexProperty = "CurrentIndex";
public ForEach()
/// <inheritdoc />
public ForEach([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
{
Behaviors.Add<BreakBehavior>(this);
}
public ForEach(ICollection<object> items) : this()
/// <inheritdoc />
public ForEach(ICollection<object> items, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Items = new Input<ICollection<object>>(items);
}

View file

@ -1,4 +1,6 @@
using System.Collections.Immutable;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Common.Extensions;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
using Elsa.Workflows.Core.Attributes;
@ -8,14 +10,23 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// Branch execution into multiple branches.
/// </summary>
[Activity("Elsa", "Control Flow", "Branch execution into multiple branches.")]
public class Fork : ActivityBase
{
/// <inheritdoc />
[JsonConstructor]
public Fork([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// Controls when this activity yields control back to its parent activity.
/// </summary>
[Input]
public ForkJoinMode JoinMode { get; set; } = ForkJoinMode.WaitAny;
public ForkJoinMode JoinMode { get; set; } = ForkJoinMode.WaitAll;
/// <summary>
/// The branches to schedule.
@ -24,7 +35,7 @@ public class Fork : ActivityBase
public ICollection<IActivity> Branches { get; set; } = new List<IActivity>();
/// <inheritdoc />
protected override void Execute(ActivityExecutionContext context) => context.ScheduleActivities(Branches.Reverse(), CompleteChildAsync);
protected override ValueTask ExecuteAsync(ActivityExecutionContext context) => context.ScheduleActivities(Branches.Reverse(), CompleteChildAsync);
private async ValueTask CompleteChildAsync(ActivityExecutionContext context, ActivityExecutionContext childContext)
{

View file

@ -1,7 +1,17 @@
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// Controls when a <see cref="Fork"/> completes.
/// </summary>
public enum ForkJoinMode
{
/// <summary>
/// The <see cref="Fork"/> completes after all inbound activities have completed.
/// </summary>
WaitAll,
/// <summary>
/// The <see cref="Fork"/> completes as soon as any of its inbound activity completes.
/// </summary>
WaitAny
}

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
@ -11,18 +12,20 @@ public class If : ActivityBase<bool>
{
/// <inheritdoc />
[JsonConstructor]
public If()
public If([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
public If(Input<bool> condition) => Condition = condition;
public If(Input<bool> condition, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) =>
Condition = condition;
/// <inheritdoc />
public If(Func<ExpressionExecutionContext, bool> condition) => Condition = new Input<bool>(condition);
public If(Func<ExpressionExecutionContext, bool> condition, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) =>
Condition = new Input<bool>(condition);
/// <inheritdoc />
public If(Func<bool> condition) => Condition = new Input<bool>(condition);
public If(Func<bool> condition, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => Condition = new Input<bool>(condition);
/// <summary>
/// The condition to evaluate.
@ -47,7 +50,7 @@ public class If : ActivityBase<bool>
{
var result = context.Get(Condition);
var nextNode = result ? Then : Else;
context.Set(Result, result);
await context.ScheduleActivityAsync(nextNode, OnChildCompleted);
}

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
@ -12,25 +13,28 @@ public class Inline : Activity
{
private readonly Func<ActivityExecutionContext, ValueTask> _activity;
public Inline(Func<ActivityExecutionContext, ValueTask> activity) => _activity = activity;
public Inline(Func<ActivityExecutionContext, ValueTask> activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
_activity = activity;
}
public Inline(Func<ValueTask> activity) : this(_ => activity())
public Inline(Func<ValueTask> activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(_ => activity(), source, line)
{
}
public Inline(Action<ActivityExecutionContext> activity) : this(c =>
public Inline(Action<ActivityExecutionContext> activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(c =>
{
activity(c);
return new ValueTask();
})
}, source, line)
{
}
public Inline(Action activity) : this(c =>
public Inline(Action activity, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(c =>
{
activity();
return new ValueTask();
})
}, source, line)
{
}
@ -54,28 +58,32 @@ public class Inline<T> : Activity<T>
{
private readonly Func<ActivityExecutionContext, ValueTask<T>> _activity;
public Inline(Func<ActivityExecutionContext, ValueTask<T>> activity, MemoryBlockReference? output = default) : base(output)
public Inline(Func<ActivityExecutionContext, ValueTask<T>> activity, MemoryBlockReference? output = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: base(output, source, line)
{
_activity = activity;
}
public Inline(Func<ValueTask<T>> activity, MemoryBlockReference? output = default) : this(_ => activity(), output)
public Inline(Func<ValueTask<T>> activity, MemoryBlockReference? output = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(_ => activity(), output, source, line)
{
}
public Inline(Func<ActivityExecutionContext, T> activity, MemoryBlockReference? output = default) : this(c =>
public Inline(Func<ActivityExecutionContext, T> activity, MemoryBlockReference? output = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(c =>
{
var result = activity(c);
return new ValueTask<T>(result);
}, output)
}, output, source, line)
{
}
public Inline(Func<T> activity, MemoryBlockReference? output = default) : this(c =>
public Inline(Func<T> activity, MemoryBlockReference? output = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(c =>
{
var result = activity();
return new ValueTask<T>(result);
}, output)
}, output, source, line)
{
}

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
@ -5,11 +6,29 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// Schedule an activity for each item in parallel.
/// </summary>
/// <typeparam name="T"></typeparam>
[Activity("Elsa", "Control Flow", "Schedule an activity for each item in parallel.")]
public class ParallelForEach<T> : Activity
{
private const string CollectedCountProperty = nameof(CollectedCountProperty);
/// <inheritdoc />
public ParallelForEach([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The items to iterate.
/// </summary>
[Input] public Input<ICollection<T>> Items { get; set; } = new(Array.Empty<T>());
/// <summary>
/// The <see cref="IActivity"/> to execute each iteration.
/// </summary>
[Port] public IActivity Body { get; set; } = default!;
/// <summary>

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Implementations;
@ -6,24 +7,31 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// Read a line of text from the console
/// </summary>
[Activity("Elsa", "Console", "Read a line of text from the console.")]
public class ReadLine : Activity<string>
{
public ReadLine()
/// <inheritdoc />
public ReadLine([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
public ReadLine(MemoryBlockReference output) : base(output)
{
}
public ReadLine(Output<string>? output) : base(output)
/// <inheritdoc />
public ReadLine(MemoryBlockReference output, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(output, source, line)
{
}
/// <inheritdoc />
public ReadLine(Output<string>? output, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(output, source, line)
{
}
/// <inheritdoc />
protected override void Execute(ActivityExecutionContext context)
{
var provider = context.GetService<IStandardInStreamProvider>() ?? new StandardInStreamProvider(System.Console.In);
var provider = context.GetService<IStandardInStreamProvider>() ?? new StandardInStreamProvider(Console.In);
var reader = provider.GetTextReader();
var text = reader.ReadLine()!;
context.Set(Result, text);

View file

@ -1,4 +1,5 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
@ -6,25 +7,22 @@ using Elsa.Workflows.Core.Signals;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// Execute a set of activities in sequence.
/// </summary>
[Category("Workflows")]
[Activity("Elsa", "Workflows", "Execute a set of activities in sequence.")]
public class Sequence : Container
{
private const string CurrentIndexProperty = "CurrentIndex";
public Sequence()
/// <inheritdoc />
public Sequence([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
OnSignalReceived<BreakSignal>(OnBreak);
}
public Sequence(params IActivity[] activities) : base(activities)
{
}
public Sequence(ICollection<Variable> variables, params IActivity[] activities) : base(variables, activities)
{
}
/// <inheritdoc />
protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context)
{
await HandleItemAsync(context);
@ -53,7 +51,7 @@ public class Sequence : Container
private void OnBreak(BreakSignal signal, SignalContext context)
{
// Clear any scheduled child completion callbacks, since we no longer want to schedule any sibling.
// Clear any scheduled child completion callbacks, since we no longer want to schedule any siblings.
context.ReceiverActivityExecutionContext.ClearCompletionCallbacks();
}
}

View file

@ -1,4 +1,5 @@
using System.Text.Json.Serialization;
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
@ -14,11 +15,11 @@ public class SetName : Activity
internal static readonly object WorkflowInstanceNameKey = new();
[JsonConstructor]
public SetName()
public SetName([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
public SetName(Input<string> value)
public SetName(Input<string> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Value = value;
}

View file

@ -1,43 +1,65 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// Assign a workflow variable a value.
/// </summary>
[Browsable(false)]
[Activity("Elsa", "Primitives", "Set a workflow variable to a given value.")]
[Activity("Elsa", "Primitives", "Assign a workflow variable a value.")]
public class SetVariable<T> : Activity
{
public SetVariable()
/// <inheritdoc />
public SetVariable([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
public SetVariable(Variable<T> variable, Input<T> value)
/// <inheritdoc />
public SetVariable(Variable<T> variable, Input<T> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Variable = variable;
Value = value;
}
public SetVariable(Variable<T> variable, Variable<T> value) : this(variable, new Input<T>(value))
{
}
public SetVariable(Variable<T> variable, Func<ExpressionExecutionContext, T> value) : this(variable, new Input<T>(value))
{
}
public SetVariable(Variable<T> variable, Func<T> value) : this(variable, new Input<T>(value))
{
}
public SetVariable(Variable<T> variable, T value) : this(variable, new Input<T>(value))
/// <inheritdoc />
public SetVariable(Variable<T> variable, Variable<T> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(variable, new Input<T>(value), source, line)
{
}
[Input] public Input<T> Value { get; set; } = new(new Literal<T>());
/// <inheritdoc />
public SetVariable(Variable<T> variable, Func<ExpressionExecutionContext, T> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(variable, new Input<T>(value), source, line)
{
}
/// <inheritdoc />
public SetVariable(Variable<T> variable, Func<T> value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(variable, new Input<T>(value), source, line)
{
}
/// <inheritdoc />
public SetVariable(Variable<T> variable, T value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(variable, new Input<T>(value), source, line)
{
}
/// <summary>
/// The variable to assign the value to.
/// </summary>
public Variable<T> Variable { get; set; } = default!;
/// <summary>
/// The value to assign.
/// </summary>
[Input] public Input<T> Value { get; set; } = new(new Literal<T>());
/// <inheritdoc />
protected override void Execute(ActivityExecutionContext context)
{
var value = context.Get(Value);
@ -45,10 +67,25 @@ public class SetVariable<T> : Activity
}
}
[Activity("Elsa", "Primitives", "Set a workflow variable to a given value.")]
/// <summary>
/// Assign a workflow variable a value.
/// </summary>
[Activity("Elsa", "Primitives", "Assign a workflow variable a value.")]
public class SetVariable : Activity
{
/// <inheritdoc />
public SetVariable([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The variable to assign the value to.
/// </summary>
[Input] public Variable Variable { get; set; } = default!;
/// <summary>
/// The value to assign.
/// </summary>
[Input] public Input<object?> Value { get; set; } = new(default(object));
protected override void Execute(ActivityExecutionContext context)

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions;
using Elsa.Expressions.Models;
@ -15,6 +16,11 @@ namespace Elsa.Workflows.Core.Activities;
[Activity("Elsa", "Control Flow", "Evaluate a set of case conditions and schedule the activity for a matching case.")]
public class Switch : ActivityBase
{
/// <inheritdoc />
public Switch([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The value to switch on.
/// </summary>
@ -29,6 +35,7 @@ public class Switch : ActivityBase
[Input(UIHint = "switch-editor")] public ICollection<SwitchCase> Cases { get; set; } = new List<SwitchCase>();
public IActivity? Default { get; set; }
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
context.Set(Output, Expression);

View file

@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
@ -16,46 +17,52 @@ public class While : Activity
};
[JsonConstructor]
public While(IActivity? body = default)
public While(IActivity? body = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
Body = body!;
Behaviors.Add<BreakBehavior>(this);
Behaviors.Remove<AutoCompleteBehavior>();
}
public While(Input<bool> condition, IActivity? body = default) : this(body)
public While(Input<bool> condition, IActivity? body = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(body, source, line)
{
Condition = condition;
}
public While(Func<ExpressionExecutionContext, ValueTask<bool>> condition, IActivity? body = default) : this(new Input<bool>(condition), body)
public While(Func<ExpressionExecutionContext, ValueTask<bool>> condition, IActivity? body = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<bool>(condition), body, source, line)
{
}
public While(Func<ExpressionExecutionContext, bool> condition, IActivity? body = default) : this(new Input<bool>(condition), body)
public While(Func<ExpressionExecutionContext, bool> condition, IActivity? body = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<bool>(condition), body, source, line)
{
}
public While(Func<ValueTask<bool>> condition, IActivity? body = default) : this(new Input<bool>(condition), body)
public While(Func<ValueTask<bool>> condition, IActivity? body = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<bool>(condition), body, source, line)
{
}
public While(Func<bool> condition, IActivity? body = default) : this(new Input<bool>(condition), body)
public While(Func<bool> condition, IActivity? body = default, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new Input<bool>(condition), body, source, line)
{
}
/// <summary>
/// The condition to evaluate.
/// </summary>
[Input(AutoEvaluate = false)] public Input<bool> Condition { get; set; } = new(false);
/// <summary>
/// The <see cref="IActivity"/> to execute on every iteration.
/// </summary>
[Port] public IActivity Body { get; set; }
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
await HandleIterationAsync(context);
}
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) => await HandleIterationAsync(context);
private async ValueTask OnBodyCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext)
{
await HandleIterationAsync(context);
}
private async ValueTask OnBodyCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext) => await HandleIterationAsync(context);
private async ValueTask HandleIterationAsync(ActivityExecutionContext context)
{

View file

@ -1,4 +1,5 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Implementations;
@ -7,33 +8,53 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
///
/// </summary>
[Activity("Elsa", "Console", "Write a line of text to the console.")]
public class WriteLine : Activity
{
public WriteLine()
{
}
public WriteLine(string text) : this(new Literal<string>(text))
{
}
public WriteLine(Func<string> text) : this(new DelegateBlockReference<string>(text))
{
}
public WriteLine(Func<ExpressionExecutionContext, string?> text) : this(new DelegateBlockReference<string?>(text))
/// <inheritdoc />
private WriteLine([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
public WriteLine(Variable<string> variable) => Text = new Input<string>(variable);
public WriteLine(Literal<string> literal) => Text = new Input<string>(literal);
public WriteLine(DelegateBlockReference delegateBlockExpression) => Text = new Input<string>(delegateBlockExpression);
public WriteLine(Input<string> text) => Text = text;
/// <inheritdoc />
public WriteLine(string text, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Literal<string>(text), source, line)
{
}
/// <inheritdoc />
public WriteLine(Func<string> text, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new DelegateBlockReference<string>(text), source, line)
{
}
/// <inheritdoc />
public WriteLine(Func<ExpressionExecutionContext, string?> text, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: this(new DelegateBlockReference<string?>(text), source, line)
{
}
/// <inheritdoc />
public WriteLine(Variable<string> variable, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => Text = new Input<string>(variable);
/// <inheritdoc />
public WriteLine(Literal<string> literal, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => Text = new Input<string>(literal);
/// <inheritdoc />
public WriteLine(DelegateBlockReference delegateBlockExpression, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => Text = new Input<string>(delegateBlockExpression);
/// <inheritdoc />
public WriteLine(Input<string> text, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => Text = text;
/// <summary>
/// The text to write.
/// </summary>
[Description("The text to write.")]
public Input<string> Text { get; set; } = default!;
/// <inheritdoc />
protected override void Execute(ActivityExecutionContext context)
{
var text = context.Get(Text);

View file

@ -44,6 +44,10 @@ public static class ActivityExecutionContextExtensions
var parentActivityInstanceId = context.ParentActivityExecutionContext?.Id;
var workflowExecutionContext = context.WorkflowExecutionContext;
var now = context.GetRequiredService<ISystemClock>().UtcNow;
if (source == null && activity.Source != null)
source = $"{Path.GetFileName(activity.Source)}:{activity.Line}";
var logEntry = new WorkflowExecutionLogEntry(activityInstanceId, parentActivityInstanceId, activity.Id, activity.Type, now, eventName, message, source, payload);
workflowExecutionContext.ExecutionLog.Add(logEntry);
return logEntry;

View file

@ -1,3 +1,4 @@
using Elsa.Expressions.Helpers;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Models;
@ -29,11 +30,11 @@ public static class ExpressionExecutionContextExtensions
public static ActivityExecutionContext GetActivityExecutionContext(this ExpressionExecutionContext context) => (ActivityExecutionContext)context.TransientProperties[ActivityExecutionContextKey];
public static IDictionary<string, object> GetInput(this ExpressionExecutionContext context) => (IDictionary<string, object>)context.TransientProperties[InputKey];
public static T? Get<T>(this ExpressionExecutionContext context, Input<T>? input) => input != null ? (T?)context.GetBlock(input.MemoryBlockReference).Value : default;
public static T? Get<T>(this ExpressionExecutionContext context, Output output) => (T?)context.GetBlock(output.MemoryBlockReference).Value;
public static T? Get<T>(this ExpressionExecutionContext context, Input<T>? input) => input != null ? context.GetBlock(input.MemoryBlockReference).Value.ConvertTo<T>() : default;
public static T? Get<T>(this ExpressionExecutionContext context, Output output) => context.GetBlock(output.MemoryBlockReference).Value.ConvertTo<T>();
public static object? Get(this ExpressionExecutionContext context, Output output) => context.GetBlock(output.MemoryBlockReference).Value;
public static T? GetVariable<T>(this ExpressionExecutionContext context, string name) => (T?)context.GetVariable(name);
public static T? GetVariable<T>(this ExpressionExecutionContext context) => (T?)context.GetVariable(typeof(T).Name);
public static T? GetVariable<T>(this ExpressionExecutionContext context) => context.GetVariable(typeof(T).Name).ConvertTo<T>();
public static object? GetVariable(this ExpressionExecutionContext context, string name) => new Variable(name).Get(context);
public static Variable SetVariable<T>(this ExpressionExecutionContext context, T? value) => context.SetVariable(typeof(T).Name, value);
public static Variable SetVariable<T>(this ExpressionExecutionContext context, string name, T? value) => context.SetVariable(name, (object?)value);
@ -63,7 +64,7 @@ public static class ExpressionExecutionContextExtensions
foreach (var l in currentRegister.Blocks)
{
if (!memoryBlocks.ContainsKey(l.Key))
memoryBlocks.Add(l.Key, l.Value!.Value);
memoryBlocks.Add(l.Key, l.Value!.Value!);
}
currentRegister = currentRegister.Parent;

View file

@ -45,7 +45,11 @@ public static class WorkflowExecutionContextExtensions
public static void ScheduleBookmark(this WorkflowExecutionContext workflowExecutionContext, Bookmark bookmark)
{
// Construct bookmark.
var bookmarkedActivityContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.Id == bookmark.ActivityInstanceId);
var bookmarkedActivityContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Id == bookmark.ActivityInstanceId);
if(bookmarkedActivityContext == null)
return;
var bookmarkedActivity = bookmarkedActivityContext.Activity;
// Schedule the activity to resume.
@ -56,8 +60,8 @@ public static class WorkflowExecutionContextExtensions
// If no resumption point was specified, use "Complete" to prevent the regular "ExecuteAsync" method to be invoked and instead complete the activity.
workflowExecutionContext.ExecuteDelegate = bookmark.CallbackMethodName != null ? bookmarkedActivity.GetResumeActivityDelegate(bookmark.CallbackMethodName) : WorkflowExecutionContext.Complete;
// Remove the bookmark.
workflowExecutionContext.Bookmarks.Remove(bookmark);
// Store the bookmark to resume in the context.
workflowExecutionContext.ResumedBookmarkContext = new ResumedBookmarkContext(bookmark);
}
/// <summary>

View file

@ -17,5 +17,6 @@ public class BookmarkPayloadSerializer : IBookmarkPayloadSerializer
}
public T Deserialize<T>(string json) where T : notnull => JsonSerializer.Deserialize<T>(json, _settings)!;
public object Deserialize(string json, Type type) => JsonSerializer.Deserialize(json, type, _settings)!;
public string Serialize<T>(T payload) where T : notnull => JsonSerializer.Serialize(payload, payload.GetType(), _settings);
}

View file

@ -26,7 +26,8 @@ public class DefaultWorkflowExecutionContextFactory : IWorkflowExecutionContextF
_serviceProvider = serviceProvider;
}
public async Task<WorkflowExecutionContext> CreateAsync(Workflow workflow,
public async Task<WorkflowExecutionContext> CreateAsync(
Workflow workflow,
string instanceId,
WorkflowState? workflowState,
IDictionary<string, object>? input = default,
@ -55,9 +56,9 @@ public class DefaultWorkflowExecutionContextFactory : IWorkflowExecutionContextF
graph,
scheduler,
input,
executeActivityDelegate,
executeActivityDelegate,
triggerActivityId,
default,
default,
cancellationToken);
// Restore workflow execution context from state, if provided.

View file

@ -33,22 +33,17 @@ public class IdentityGraphService : IIdentityGraphService
}
}
private void AssignInputOutputs(IActivity activity)
public void AssignInputOutputs(IActivity activity)
{
var inputs = activity.GetInputs();
var assignedInputs = inputs.Where(x =>
{
var memoryBlockReference = x.MemoryBlockReference();
return memoryBlockReference.Id == null!;
}).ToList();
var seed = 0;
foreach (var input in assignedInputs)
foreach (var input in inputs)
{
var blockReference = input.MemoryBlockReference();
blockReference.Id = $"{activity.Id}:input-{++seed}";
if(string.IsNullOrEmpty(blockReference.Id))
blockReference.Id = $"{activity.Id}:input-{++seed}";
}
seed = 0;
@ -62,12 +57,14 @@ public class IdentityGraphService : IIdentityGraphService
foreach (var output in assignedOutputs)
{
var memoryReference = output.Value.MemoryBlockReference();
memoryReference.Id = $"{activity.Id}:output-{++seed}";
var blockReference = output.Value.MemoryBlockReference();
if(string.IsNullOrEmpty(blockReference.Id))
blockReference.Id = $"{activity.Id}:output-{++seed}";
}
}
private void AssignVariables(IActivity activity)
public void AssignVariables(IActivity activity)
{
var variables = activity.GetVariables();
var seed = 0;

View file

@ -84,7 +84,7 @@ public class WorkflowRunner : IWorkflowRunner
if (bookmarkId != null)
{
// Schedule the bookmark.
var bookmark = workflowExecutionContext.Bookmarks.FirstOrDefault(x => x.Id == bookmarkId);
var bookmark = workflowState.Bookmarks.FirstOrDefault(x => x.Id == bookmarkId);
if (bookmark != null)
workflowExecutionContext.ScheduleBookmark(bookmark);

View file

@ -105,6 +105,15 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
private void SerializeCompletionCallbacks(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
{
// Assert all referenced owner contexts exist.
foreach (var completionCallback in workflowExecutionContext.CompletionCallbacks)
{
var owmnerContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x == completionCallback.Owner);
if (owmnerContext == null)
throw new Exception("Lost an owner context");
}
var completionCallbacks = workflowExecutionContext.CompletionCallbacks.Select(x => new CompletionCallbackState(x.Owner.Id, x.Child.Id, x.CompletionCallback?.Method.Name));
state.CompletionCallbacks = completionCallbacks.ToList();
}
@ -120,8 +129,9 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
if (parentId != null)
{
var parentContext = activityExecutionContext.WorkflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Id == parentId);
Debug.Assert(parentContext != null);
if (parentContext == null)
throw new Exception("We lost a context");
}
var activityExecutionContextState = new ActivityExecutionContextState

View file

@ -36,6 +36,12 @@ public class DefaultActivityInvokerMiddleware : IActivityExecutionMiddleware
// Reset execute delegate.
workflowExecutionContext.ExecuteDelegate = null;
// If a bookmark was used to resume, burn it if not burnt by the activity.
var resumedBookmark = workflowExecutionContext.ResumedBookmarkContext?.Bookmark;
if (resumedBookmark is { AutoBurn: true })
workflowExecutionContext.Bookmarks.Remove(resumedBookmark);
// Update execution count.
context.IncrementExecutionCount();

View file

@ -3,58 +3,80 @@ using Elsa.Workflows.Core.Behaviors;
namespace Elsa.Workflows.Core.Models;
/// <summary>
/// Base class for custom activities with auto-complete behavior.
/// </summary>
public abstract class Activity : ActivityBase
{
protected Activity()
/// <inheritdoc />
protected Activity(string? source = default, int? line = default) : base(source, line)
{
Behaviors.Add<AutoCompleteBehavior>(this);
}
protected Activity(string activityType) : this()
/// <inheritdoc />
protected Activity(string activityType, int version = 1, string? source = default, int? line = default) : base(activityType, version, source, line)
{
Type = activityType;
}
}
/// <summary>
/// Base class for custom activities with auto-complete behavior that return a result.
/// </summary>
public abstract class ActivityWithResult : Activity
{
protected ActivityWithResult()
/// <inheritdoc />
protected ActivityWithResult(string? source = default, int? line = default) : base(source, line)
{
}
protected ActivityWithResult(string activityType) : base(activityType)
/// <inheritdoc />
protected ActivityWithResult(string activityType, int version = 1, string? source = default, int? line = default) : base(activityType, version, source, line)
{
}
protected ActivityWithResult(MemoryBlockReference? output)
/// <inheritdoc />
protected ActivityWithResult(MemoryBlockReference? output, string? source = default, int? line = default) : base(source, line)
{
if (output != null) Result = new Output(output);
}
protected ActivityWithResult(Output? output)
/// <inheritdoc />
protected ActivityWithResult(Output? output, string? source = default, int? line = default) : base(source, line)
{
Result = output;
}
/// <summary>
/// The result of the activity.
/// </summary>
public Output? Result { get; set; }
}
/// <summary>
/// Base class for custom activities that return a result.
/// </summary>
public abstract class ActivityBaseWithResult : ActivityBase
{
protected ActivityBaseWithResult()
/// <inheritdoc />
protected ActivityBaseWithResult(string? source = default, int? line = default) : base(source, line)
{
}
protected ActivityBaseWithResult(string activityType) : base(activityType)
/// <inheritdoc />
protected ActivityBaseWithResult(string activityType, int version = 1, string? source = default, int? line = default) : base(activityType, version, source, line)
{
}
protected ActivityBaseWithResult(MemoryBlockReference? output)
/// <inheritdoc />
protected ActivityBaseWithResult(MemoryBlockReference? output, string? source = default, int? line = default) : this(source, line)
{
if (output != null) Result = new Output(output);
}
protected ActivityBaseWithResult(Output? output)
/// <inheritdoc />
protected ActivityBaseWithResult(Output? output, string? source = default, int? line = default) : this(source, line)
{
Result = output;
}
@ -62,48 +84,68 @@ public abstract class ActivityBaseWithResult : ActivityBase
public Output? Result { get; set; }
}
/// <summary>
/// Base class for custom activities with auto-complete behavior that return a result.
/// </summary>
public abstract class Activity<T> : Activity
{
protected Activity()
/// <inheritdoc />
protected Activity(string? source = default, int? line = default) : base(source, line)
{
}
/// <inheritdoc />
protected Activity(string activityType, int version = 1, string? source = default, int? line = default) : base(activityType, version, source, line)
{
}
protected Activity(string activityType) : base(activityType)
{
}
protected Activity(MemoryBlockReference? output)
/// <inheritdoc />
protected Activity(MemoryBlockReference? output, string? source = default, int? line = default) : this(source, line)
{
if (output != null) Result = new Output<T>(output);
}
protected Activity(Output<T>? output)
/// <inheritdoc />
protected Activity(Output<T>? output, string? source = default, int? line = default) : this(source, line)
{
Result = output;
}
/// <summary>
/// The result of the activity.
/// </summary>
public Output<T>? Result { get; set; }
}
/// <summary>
/// Base class for custom activities that return a result.
/// </summary>
public abstract class ActivityBase<T> : ActivityBase
{
protected ActivityBase()
/// <inheritdoc />
protected ActivityBase(string? source = default, int? line = default) : base(source, line)
{
}
protected ActivityBase(string activityType) : base(activityType)
/// <inheritdoc />
protected ActivityBase(string activityType, int version = 1, string? source = default, int? line = default) : base(activityType, version, source, line)
{
}
protected ActivityBase(MemoryBlockReference? output)
/// <inheritdoc />
protected ActivityBase(MemoryBlockReference? output, string? source = default, int? line = default) : this(source, line)
{
if (output != null) Result = new Output<T>(output);
}
protected ActivityBase(Output<T>? output)
/// <inheritdoc />
protected ActivityBase(Output<T>? output, string? source = default, int? line = default) : this(source, line)
{
Result = output;
}
/// <summary>
/// The result of the activity.
/// </summary>
public Output<T>? Result { get; set; }
}

View file

@ -6,31 +6,81 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Models;
/// <summary>
/// Base class for custom activities.
/// </summary>
[DebuggerDisplay("{Type} - {Id}")]
public abstract class ActivityBase : IActivity, ISignalHandler
{
private readonly ICollection<SignalHandlerRegistration> _signalHandlers = new List<SignalHandlerRegistration>();
protected ActivityBase()
/// <summary>
/// Constructor.
/// </summary>
protected ActivityBase(string? source = default, int? line = default)
{
Source = source;
Line = line;
Type = ActivityTypeNameHelper.GenerateTypeName(GetType());
Version = 1;
Behaviors.Add<ExecutionLoggingBehavior>(this);
Behaviors.Add<ScheduledChildCallbackBehavior>(this);
}
protected ActivityBase(string activityType, int version = 1) : this()
/// <inheritdoc />
protected ActivityBase(string activityType, int version = 1, string? source = default, int? line = default) : this(source, line)
{
Type = activityType;
Version = version;
}
/// <summary>
/// The unique ID of this activity within the <see cref="Workflow"/>.
/// </summary>
public string Id { get; set; } = default!;
/// <summary>
/// The technical type name.
/// </summary>
public string Type { get; set; }
/// <summary>
/// The version number.
/// </summary>
public int Version { get; set; }
/// <summary>
/// A flag indicating whether this activity can be used for starting a workflow.
/// Usually used for triggers, but also used to disambiguate between two or more starting activities and no starting activity was specified.
/// </summary>
public bool CanStartWorkflow { get; set; }
/// <summary>
/// A flag indicating if this activity should execute synchronously or asynchronously.
/// By default, activities with an <see cref="ActivityKind"/> of <see cref="ActivityKind.Action"/>, <see cref="ActivityKind.Task"/> or <see cref="ActivityKind.Trigger"/>
/// will execute synchronously, while activities of the <see cref="ActivityKind.Job"/> kind will execute asynchronously.
/// </summary>
public bool RunAsynchronously { get; set; }
/// <summary>
/// A bag of properties that can be used by custom activities and other code such as middleware components to store additional values with the activity.
/// </summary>
public IDictionary<string, object> ApplicationProperties { get; set; } = new Dictionary<string, object>();
/// <summary>
/// Automatically set to the current source file name when instantiating this activity inside of a workflow class or composite activity class.
/// </summary>
public string? Source { get; set; }
/// <summary>
/// Automatically set to the current line of code when instantiating this activity inside of a workflow class or composite activity class.
/// </summary>
public int? Line { get; set; }
/// <summary>
/// Stores metadata such as x and y coordinates when created via the designer.
/// </summary>
public IDictionary<string, object> Metadata { get; set; } = new Dictionary<string, object>();
/// <summary>
@ -39,37 +89,51 @@ public abstract class ActivityBase : IActivity, ISignalHandler
[JsonIgnore]
public ICollection<IBehavior> Behaviors { get; } = new List<IBehavior>();
/// <summary>
/// Override this method to implement activity-specific logic.
/// </summary>
protected virtual ValueTask ExecuteAsync(ActivityExecutionContext context)
{
Execute(context);
return ValueTask.CompletedTask;
}
/// <summary>
/// Override this method to implement activity-specific logic.
/// </summary>
protected virtual void Execute(ActivityExecutionContext context)
{
}
/// <summary>
/// Override this method to handle any signals sent from downstream activities.
/// </summary>
protected virtual ValueTask OnSignalReceivedAsync(object signal, SignalContext context)
{
OnSignalReceived(signal, context);
return ValueTask.CompletedTask;
}
/// <summary>
/// Override this method to handle any signals sent from downstream activities.
/// </summary>
protected virtual void OnSignalReceived(object signal, SignalContext context)
{
}
protected virtual void Execute(ActivityExecutionContext context)
{
}
/// <summary>
/// Notify the system that this activity completed.
/// Register a signal handler delegate.
/// </summary>
protected async ValueTask CompleteAsync(ActivityExecutionContext context)
{
await context.CompleteActivityAsync();
}
protected void OnSignalReceived(Type signalType, Func<object, SignalContext, ValueTask> handler) => _signalHandlers.Add(new SignalHandlerRegistration(signalType, handler));
/// <summary>
/// Register a signal handler delegate.
/// </summary>
protected void OnSignalReceived<T>(Func<T, SignalContext, ValueTask> handler) => OnSignalReceived(typeof(T), (signal, context) => handler((T)signal, context));
/// <summary>
/// Register a signal handler delegate.
/// </summary>
protected void OnSignalReceived<T>(Action<T, SignalContext> handler)
{
OnSignalReceived<T>((signal, context) =>
@ -78,6 +142,14 @@ public abstract class ActivityBase : IActivity, ISignalHandler
return ValueTask.CompletedTask;
});
}
/// <summary>
/// Notify the workflow that this activity completed.
/// </summary>
protected async ValueTask CompleteAsync(ActivityExecutionContext context)
{
await context.CompleteActivityAsync();
}
async ValueTask IActivity.ExecuteAsync(ActivityExecutionContext context)
{

View file

@ -1,5 +1,4 @@
using System.Collections.ObjectModel;
using System.Runtime.CompilerServices;
using Elsa.Expressions.Helpers;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Services;
@ -83,6 +82,8 @@ public class ActivityExecutionContext
// ReSharper disable once CollectionNeverQueried.Global
public IDictionary<string, object?> JournalData { get; } = new Dictionary<string, object?>();
public ResumedBookmarkContext? ResumedBookmarkContext => WorkflowExecutionContext.ResumedBookmarkContext;
public async ValueTask ScheduleActivityAsync(IActivity? activity, ActivityCompletionCallback? completionCallback = default, IEnumerable<MemoryBlockReference>? references = default, object? tag = default)
{
await ScheduleActivityAsync(activity, this, completionCallback, references, tag);
@ -111,33 +112,39 @@ public class ActivityExecutionContext
public void CreateBookmarks(IEnumerable<object> payloads, ExecuteActivityDelegate? callback = default)
{
foreach (var payload in payloads)
CreateBookmark(payload, callback);
CreateBookmark(new CreateBookmarkOptions(payload, callback));
}
public void AddBookmarks(IEnumerable<Bookmark> bookmarks) => _bookmarks.AddRange(bookmarks);
public void AddBookmark(Bookmark bookmark) => _bookmarks.Add(bookmark);
public Bookmark CreateBookmark(ExecuteActivityDelegate callback) => CreateBookmark(default, callback);
public Bookmark CreateBookmark(ExecuteActivityDelegate callback) => CreateBookmark(new CreateBookmarkOptions(default, callback));
public Bookmark CreateBookmark(object payload, ExecuteActivityDelegate callback) => CreateBookmark(new CreateBookmarkOptions(payload, callback));
public Bookmark CreateBookmark(object payload) => CreateBookmark(new CreateBookmarkOptions(payload));
/// <summary>
/// Creates a bookmark so that this activity can be resumed at a later time.
/// Creating a bookmark will automatically suspend the workflow after all pending activities have executed.
/// </summary>
public Bookmark CreateBookmark(object? payload = default, ExecuteActivityDelegate? callback = default)
public Bookmark CreateBookmark(CreateBookmarkOptions? options = default)
{
var payload = options?.Payload;
var callback = options?.Callback;
var activityTypeName = options?.ActivityTypeName ?? Activity.Type;
var bookmarkHasher = GetRequiredService<IBookmarkHasher>();
var identityGenerator = GetRequiredService<IIdentityGenerator>();
var payloadSerializer = GetRequiredService<IBookmarkPayloadSerializer>();
var payloadJson = payload != null ? payloadSerializer.Serialize(payload) : default;
var hash = bookmarkHasher.Hash(Activity.Type, payloadJson);
var hash = bookmarkHasher.Hash(activityTypeName, payloadJson);
var bookmark = new Bookmark(
identityGenerator.GenerateId(),
Activity.Type,
activityTypeName,
hash,
payloadJson,
Activity.Id,
Id,
options?.AutoBurn ?? true,
callback?.Method.Name);
AddBookmark(bookmark);
@ -195,7 +202,7 @@ public class ActivityExecutionContext
public object? Get(MemoryBlockReference blockReference)
{
var location = GetBlock(blockReference) ?? throw new InvalidOperationException($"No location found with ID {blockReference.Id}. Did you forget to declare a variable with a container?");
var location = GetMemoryBlock(blockReference) ?? throw new InvalidOperationException($"No location found with ID {blockReference.Id}. Did you forget to declare a variable with a container?");
return location.Value;
}
@ -225,8 +232,8 @@ public class ActivityExecutionContext
internal void IncrementExecutionCount() => _executionCount++;
private MemoryBlock? GetBlock(MemoryBlockReference locationBlockReference) =>
ExpressionExecutionContext.Memory.TryGetBlock(locationBlockReference.Id, out var location)
? location
: ParentActivityExecutionContext?.GetBlock(locationBlockReference);
private MemoryBlock? GetMemoryBlock(MemoryBlockReference locationBlockReference) =>
ExpressionExecutionContext.Memory.TryGetBlock(locationBlockReference.Id, out var memoryBlock)
? memoryBlock
: ParentActivityExecutionContext?.GetMemoryBlock(locationBlockReference);
}

View file

@ -9,6 +9,7 @@ public record Bookmark(
string? Data,
string ActivityId,
string ActivityInstanceId,
bool AutoBurn = true,
string? CallbackMethodName = default
)
{

View file

@ -0,0 +1,5 @@
using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Models;
public record CreateBookmarkOptions(object? Payload = default, ExecuteActivityDelegate? Callback = default, string? ActivityTypeName = default, bool AutoBurn = true);

View file

@ -2,13 +2,18 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Models;
/// <summary>
/// Generates events on a workflow instance.
/// </summary>
public abstract class EventGenerator : Trigger, IEventGenerator
{
protected EventGenerator()
/// <inheritdoc />
protected EventGenerator(string? source = default, int? line = default) : base(source, line)
{
}
protected EventGenerator(string triggerType) : base(triggerType)
/// <inheritdoc />
protected EventGenerator(string triggerType, int version = 1, string? source = default, int? line = default) : base(triggerType, version, source, line)
{
}
}

View file

@ -0,0 +1,3 @@
namespace Elsa.Workflows.Core.Models;
public record ResumedBookmarkContext(Bookmark Bookmark);

View file

@ -2,13 +2,18 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Models;
/// <summary>
/// Represents an activity that acts as a workflow trigger.
/// </summary>
public abstract class Trigger : ActivityBase, ITrigger
{
protected Trigger()
/// <inheritdoc />
protected Trigger(string? source = default, int? line = default) : base(source, line)
{
}
protected Trigger(string activityType) : base(activityType)
/// <inheritdoc />
protected Trigger(string activityType, int version = 1, string? source = default, int? line = default) : base(activityType, version, source, line)
{
}
@ -36,11 +41,11 @@ public abstract class Trigger : ActivityBase, ITrigger
public abstract class Trigger<TResult> : ActivityBase<TResult>, ITrigger
{
protected Trigger()
protected Trigger(string? source = default, int? line = default) : base(source, line)
{
}
protected Trigger(string activityType) : base(activityType)
protected Trigger(string activityType, int version = 1, string? source = default, int? line = default) : base(activityType, version, source, line)
{
}

View file

@ -7,7 +7,6 @@ public class Variable : MemoryBlockReference
{
public Variable()
{
Id = Guid.NewGuid().ToString("N");
}
public Variable(string name)

View file

@ -1,5 +1,4 @@
using System.Collections.ObjectModel;
using System.Diagnostics;
using Elsa.Common.Extensions;
using Elsa.Expressions.Models;
using Elsa.Workflows.Core.Services;
@ -74,6 +73,7 @@ public class WorkflowExecutionContext
public IDictionary<object, object> TransientProperties { get; set; } = new Dictionary<object, object>();
public ExecuteActivityDelegate? ExecuteDelegate { get; set; }
public ResumedBookmarkContext? ResumedBookmarkContext { get; set; }
public string? TriggerActivityId { get; set; }
public CancellationToken CancellationToken { get; }
public ICollection<ActivityCompletionCallbackEntry> CompletionCallbacks => new ReadOnlyCollection<ActivityCompletionCallbackEntry>(_completionCallbackEntries);
@ -166,7 +166,14 @@ public class WorkflowExecutionContext
foreach (var childContext in childContexts) RemoveActivityExecutionContext(childContext);
// Remove the context.
_activityExecutionContexts.Remove(context);
// Remove all associated completion callbacks.
context.ClearCompletionCallbacks();
// Remove all associated bookmarks.
Bookmarks.RemoveWhere(x => x.ActivityInstanceId == context.Id);
}
public void AddActivityExecutionContext(ActivityExecutionContext context) => _activityExecutionContexts.Add(context);

View file

@ -25,17 +25,27 @@ public interface IActivity
/// <summary>
/// A value indicating whether this activity can start instances of the workflow it is a part of.
/// </summary>
public bool CanStartWorkflow { get; set; }
bool CanStartWorkflow { get; set; }
/// <summary>
/// A value indicating whether this activity can be executed asynchronously in the background.
/// </summary>
public bool RunAsynchronously { get; set; }
bool RunAsynchronously { get; set; }
/// <summary>
/// Can contain application-specific information about this activity.
/// </summary>
IDictionary<string, object> ApplicationProperties { get; set; }
/// <summary>
/// The source file where this activity was instantiated, if any.
/// </summary>
string? Source { get; set; }
/// <summary>
/// The source file line number where this activity was instantiated, if any.
/// </summary>
int? Line { get; set; }
/// <summary>
/// Invoked when the activity executes.

View file

@ -3,5 +3,6 @@ namespace Elsa.Workflows.Core.Services;
public interface IBookmarkPayloadSerializer
{
T Deserialize<T>(string json) where T : notnull;
object Deserialize(string json, Type type);
string Serialize<T>(T payload) where T : notnull;
}

View file

@ -7,4 +7,6 @@ public interface IIdentityGraphService
Task AssignIdentitiesAsync(Workflow workflow, CancellationToken cancellationToken = default);
Task AssignIdentitiesAsync(IActivity root, CancellationToken cancellationToken = default);
void AssignIdentities(ActivityNode root);
void AssignInputOutputs(IActivity activity);
void AssignVariables(IActivity activity);
}

View file

@ -53,15 +53,7 @@ public class DefaultExpressionSyntaxProvider : IExpressionSyntaxProvider
Syntax = syntax,
Type = typeof(TExpression),
CreateExpression = constructor,
CreateBlockReference = context =>
{
var reference = createBlockReference(context);
if (string.IsNullOrWhiteSpace(reference.Id))
reference.Id = context.MemoryReferenceId;
return reference;
},
CreateBlockReference = createBlockReference,
CreateSerializableObject = context => new
{
Type = syntax,

View file

@ -15,12 +15,14 @@ public class ActivityJsonConverter : JsonConverter<IActivity>
{
private readonly IActivityRegistry _activityRegistry;
private readonly IActivityFactory _activityFactory;
private readonly IIdentityGraphService _identityGraphService;
private readonly IServiceProvider _serviceProvider;
public ActivityJsonConverter(IActivityRegistry activityRegistry, IActivityFactory activityFactory, IServiceProvider serviceProvider)
public ActivityJsonConverter(IActivityRegistry activityRegistry, IActivityFactory activityFactory, IIdentityGraphService identityGraphService, IServiceProvider serviceProvider)
{
_activityRegistry = activityRegistry;
_activityFactory = activityFactory;
_identityGraphService = identityGraphService;
_serviceProvider = serviceProvider;
}
@ -56,6 +58,9 @@ public class ActivityJsonConverter : JsonConverter<IActivity>
var context = new ActivityConstructorContext(doc.RootElement, newOptions);
var activity = activityDescriptor.Constructor(context);
_identityGraphService.AssignInputOutputs(activity);
_identityGraphService.AssignVariables(activity);
return activity;
}

View file

@ -3,7 +3,6 @@ using System.Text.Json.Serialization;
using Elsa.Expressions.Models;
using Elsa.Expressions.Services;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Management.Serialization.Converters;
@ -13,17 +12,17 @@ namespace Elsa.Workflows.Management.Serialization.Converters;
public class InputJsonConverter<T> : JsonConverter<Input<T>>
{
private readonly IExpressionSyntaxRegistry _expressionSyntaxRegistry;
private readonly IIdentityGenerator _identityGenerator;
/// <inheritdoc />
public InputJsonConverter(IExpressionSyntaxRegistry expressionSyntaxRegistry, IIdentityGenerator identityGenerator)
public InputJsonConverter(IExpressionSyntaxRegistry expressionSyntaxRegistry)
{
_expressionSyntaxRegistry = expressionSyntaxRegistry;
_identityGenerator = identityGenerator;
}
/// <inheritdoc />
public override bool CanConvert(Type typeToConvert) => typeof(Input).IsAssignableFrom(typeToConvert);
/// <inheritdoc />
public override Input<T> Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
{
if (!JsonDocument.TryParseValue(ref reader, out var doc))
@ -32,12 +31,14 @@ public class InputJsonConverter<T> : JsonConverter<Input<T>>
if (doc.RootElement.ValueKind != JsonValueKind.Object)
return default!;
if (!doc.RootElement.TryGetProperty("typeName", out var inputTargetTypeElement))
if (!doc.RootElement.TryGetProperty("typeName", out _))
return default!;
var memoryReferenceId = doc.RootElement.TryGetProperty("memoryReference", out var memoryReferenceElement) && memoryReferenceElement.ValueKind == JsonValueKind.String
var memoryReferenceId = doc.RootElement.TryGetProperty("memoryReference", out var memoryReferenceElement) && memoryReferenceElement.ValueKind == JsonValueKind.Object
? memoryReferenceElement.GetProperty("id").GetString()!
: _identityGenerator.GenerateId();
: memoryReferenceElement.ValueKind == JsonValueKind.String
? memoryReferenceElement.GetString()!
: throw new Exception("No input ID specified");
var expressionElement = doc.RootElement.GetProperty("expression");
@ -52,11 +53,12 @@ public class InputJsonConverter<T> : JsonConverter<Input<T>>
var context = new ExpressionConstructorContext(expressionElement, options);
var expression = expressionSyntaxDescriptor.CreateExpression(context);
var locationReference = expressionSyntaxDescriptor.CreateBlockReference(new BlockReferenceConstructorContext(expression, memoryReferenceId));
var memoryBlockReference = expressionSyntaxDescriptor.CreateBlockReference(new BlockReferenceConstructorContext(expression, memoryReferenceId));
return (Input<T>)Activator.CreateInstance(typeof(Input<T>), expression, locationReference)!;
return (Input<T>)Activator.CreateInstance(typeof(Input<T>), expression, memoryBlockReference)!;
}
/// <inheritdoc />
public override void Write(Utf8JsonWriter writer, Input<T> value, JsonSerializerOptions options)
{
var expression = value.Expression;

View file

@ -18,6 +18,7 @@
</ItemGroup>
<ItemGroup>
<PackageReference Include="DistributedLock" Version="2.3.2" />
<PackageReference Include="Microsoft.Extensions.DependencyInjection" Version="6.0.0" />
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions" Version="6.0.0" />
<PackageReference Include="Open.Linq.AsyncExtensions" Version="1.2.0" />

View file

@ -13,6 +13,8 @@ using Elsa.Workflows.Runtime.Implementations;
using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Options;
using Elsa.Workflows.Runtime.Services;
using Medallion.Threading;
using Medallion.Threading.FileSystem;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.Runtime.Features;
@ -52,6 +54,8 @@ public class WorkflowRuntimeFeature : FeatureBase
public Func<IServiceProvider, ITriggerStore> WorkflowTriggerStore { get; set; } = sp => sp.GetRequiredService<MemoryTriggerStore>();
public Func<IServiceProvider, IWorkflowExecutionLogStore> WorkflowExecutionLogStore { get; set; } = sp => sp.GetRequiredService<MemoryWorkflowExecutionLogStore>();
public Func<IServiceProvider, IDistributedLockProvider> DistributedLockProvider { get; set; } = _ =>
new FileDistributedSynchronizationProvider(new DirectoryInfo( Path.Combine(Environment.CurrentDirectory, "App_Data/locks")));
public Func<IServiceProvider, IWorkflowStateExporter> WorkflowStateExporter { get; set; } =
sp => sp.GetRequiredService<NoopWorkflowStateExporter>();
@ -83,11 +87,14 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddSingleton(WorkflowTriggerStore)
.AddSingleton(WorkflowExecutionLogStore)
// Memory Stores
// Memory stores.
.AddMemoryStore<WorkflowState, MemoryWorkflowStateStore>()
.AddMemoryStore<StoredBookmark, MemoryBookmarkStore>()
.AddMemoryStore<StoredTrigger, MemoryTriggerStore>()
.AddMemoryStore<WorkflowExecutionLogRecord, MemoryWorkflowExecutionLogStore>()
// Distributed locking.
.AddSingleton(DistributedLockProvider)
// Workflow definition providers.
.AddWorkflowDefinitionProvider<ClrWorkflowDefinitionProvider>()

View file

@ -4,6 +4,7 @@ using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Core.State;
using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Services;
using Medallion.Threading;
namespace Elsa.Workflows.Runtime.Implementations;
@ -15,6 +16,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
private readonly ITriggerStore _triggerStore;
private readonly IBookmarkStore _bookmarkStore;
private readonly IBookmarkHasher _hasher;
private readonly IDistributedLockProvider _distributedLockProvider;
public DefaultWorkflowRuntime(
IWorkflowHostFactory workflowHostFactory,
@ -22,7 +24,8 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
IWorkflowStateStore workflowStateStore,
ITriggerStore triggerStore,
IBookmarkStore bookmarkStore,
IBookmarkHasher hasher)
IBookmarkHasher hasher,
IDistributedLockProvider distributedLockProvider)
{
_workflowHostFactory = workflowHostFactory;
_workflowDefinitionService = workflowDefinitionService;
@ -30,6 +33,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
_triggerStore = triggerStore;
_bookmarkStore = bookmarkStore;
_hasher = hasher;
_distributedLockProvider = distributedLockProvider;
}
public async Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
@ -56,32 +60,35 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
{
var workflowState = await _workflowStateStore.LoadAsync(workflowInstanceId, cancellationToken);
await using (await _distributedLockProvider.AcquireLockAsync(workflowInstanceId, TimeSpan.FromMinutes(1), cancellationToken))
{
var workflowState = await _workflowStateStore.LoadAsync(workflowInstanceId, cancellationToken);
if (workflowState == null)
throw new Exception($"Workflow instance {workflowInstanceId} not found");
if (workflowState == null)
throw new Exception($"Workflow instance {workflowInstanceId} not found");
var definitionId = workflowState.DefinitionId;
var version = workflowState.DefinitionVersion;
var definitionId = workflowState.DefinitionId;
var version = workflowState.DefinitionVersion;
var workflowDefinition = await _workflowDefinitionService.FindAsync(
definitionId,
VersionOptions.SpecificVersion(version),
cancellationToken);
var workflowDefinition = await _workflowDefinitionService.FindAsync(
definitionId,
VersionOptions.SpecificVersion(version),
cancellationToken);
if (workflowDefinition == null)
throw new Exception("Specified workflow definition and version does not exist");
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var workflowHost = await _workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken);
var resumeWorkflowOptions = new ResumeWorkflowHostOptions(options.CorrelationId, options.BookmarkId, options.ActivityId, options.Input);
if (workflowDefinition == null)
throw new Exception("Specified workflow definition and version does not exist");
await workflowHost.ResumeWorkflowAsync(resumeWorkflowOptions, cancellationToken);
workflowState = workflowHost.WorkflowState;
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var workflowHost = await _workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken);
var resumeWorkflowOptions = new ResumeWorkflowHostOptions(options.CorrelationId, options.BookmarkId, options.ActivityId, options.Input);
await SaveWorkflowStateAsync(workflowState, cancellationToken);
await workflowHost.ResumeWorkflowAsync(resumeWorkflowOptions, cancellationToken);
workflowState = workflowHost.WorkflowState;
return new ResumeWorkflowResult(workflowState.Bookmarks);
await SaveWorkflowStateAsync(workflowState, cancellationToken);
return new ResumeWorkflowResult(workflowState.Bookmarks);
}
}
public async Task<ICollection<ResumedWorkflow>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)

View file

@ -46,9 +46,9 @@ public class PersistBookmarkMiddleware : WorkflowExecutionMiddleware
await _eventPublisher.PublishAsync(new WorkflowBookmarksIndexed(new IndexedWorkflowBookmarks(context.Id, diff.Added, diff.Removed)), cancellationToken);
// Notify all interested activities that the bookmarks have been persisted.
var activityExecutionContexts = context.ActivityExecutionContexts.Where(x => x.Activity is IBookmarksPersistedHandler).ToList();
var activityExecutionContexts = context.ActivityExecutionContexts.Where(x => x.Activity is IBookmarksPersistedHandler && x.Bookmarks.Any()).ToList();
foreach (var activityExecutionContext in activityExecutionContexts)
foreach (var activityExecutionContext in activityExecutionContexts)
await ((IBookmarksPersistedHandler)activityExecutionContext.Activity).BookmarksPersistedAsync(activityExecutionContext);
}
}