diff --git a/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs b/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs
index d364e3cbd..459132897 100644
--- a/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs
+++ b/src/modules/Elsa.Jobs.Activities/Handlers/JobExecutedHandler.cs
@@ -7,15 +7,22 @@ using Elsa.Workflows.Runtime.Services;
namespace Elsa.Jobs.Activities.Handlers;
+///
+/// A handler that resumes workflows waiting for a given job to complete.
+///
public class JobExecutedHandler : INotificationHandler
{
private readonly IWorkflowRuntime _workflowRuntime;
+ ///
+ /// Constructor.
+ ///
public JobExecutedHandler(IWorkflowRuntime workflowRuntime)
{
_workflowRuntime = workflowRuntime;
}
+ ///
public async Task HandleAsync(JobExecuted notification, CancellationToken cancellationToken)
{
var job = notification.Job;
@@ -33,7 +40,7 @@ public class JobExecutedHandler : INotificationHandler
var jobType = notification.Job.GetType();
var payload = new EnqueuedJobPayload(job.Id);
var jobTypeName = JobTypeNameHelper.GenerateTypeName(jobType);
- await _workflowRuntime.ResumeWorkflowsAsync(jobTypeName, payload, new ResumeWorkflowRuntimeOptions(), cancellationToken);
+ await _workflowRuntime.ResumeWorkflowsAsync(jobTypeName, payload, new TriggerWorkflowsRuntimeOptions(), cancellationToken);
}
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs
index 15cc0adcb..a9112ac33 100644
--- a/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs
+++ b/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs
@@ -52,7 +52,7 @@ public class DispatchWorkflowRequestConsumer :
public async Task Consume(ConsumeContext context)
{
var message = context.Message;
- var options = new ResumeWorkflowRuntimeOptions(CorrelationId: message.CorrelationId, Input: message.Input);
+ var options = new TriggerWorkflowsRuntimeOptions(CorrelationId: message.CorrelationId, Input: message.Input);
await _workflowRuntime.ResumeWorkflowsAsync(message.ActivityTypeName, message.BookmarkPayload, options, context.CancellationToken);
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs
index 98f4b583c..605e7391d 100644
--- a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs
+++ b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs
@@ -89,6 +89,31 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
return new WorkflowExecutionResult(workflowInstanceId, bookmarks);
}
+ ///
+ public async Task> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
+ {
+ var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
+ var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
+ var results = new List();
+
+ foreach (var trigger in triggers)
+ {
+ var definitionId = trigger.WorkflowDefinitionId;
+ var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
+ var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
+
+ // If we can't start the workflow, don't try it.
+ if(!canStartResult.CanStart)
+ continue;
+
+ var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
+
+ results.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
+ }
+
+ return results;
+ }
+
///
public async Task ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
{
@@ -109,7 +134,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
}
///
- public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
+ public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
var client = _cluster.GetNamedBookmarkGrain(hash);
@@ -122,58 +147,17 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
var bookmarksResponse = await client.Resolve(request, cancellationToken);
var bookmarks = bookmarksResponse!.Bookmarks;
- return await ResumeWorkflowsAsync(bookmarks, options, cancellationToken);
+ return await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken);
}
///
public async Task TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
- var triggeredWorkflows = new List();
- var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
+ var startedWorkflows = await StartWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
+ var resumedWorkflows = await ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
+ var results = startedWorkflows.Concat(resumedWorkflows).ToList();
- // Start new workflows.
- var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
-
- foreach (var trigger in triggers)
- {
- var definitionId = trigger.WorkflowDefinitionId;
- var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
- var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
-
- // If we can't start the workflow, don't try it.
- if(!canStartResult.CanStart)
- continue;
-
- var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
-
- triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
- }
-
- // Resume existing workflow instances.
- var client = _cluster.GetNamedBookmarkGrain(hash);
-
- var request = new ResolveBookmarksRequest
- {
- ActivityTypeName = activityTypeName,
- CorrelationId = options.CorrelationId.EmptyIfNull()
- };
-
- var bookmarksResponse = await client.Resolve(request, cancellationToken);
- var bookmarks = bookmarksResponse!.Bookmarks;
-
- foreach (var bookmark in bookmarks)
- {
- var workflowInstanceId = bookmark.WorkflowInstanceId;
-
- var resumeResult = await ResumeWorkflowAsync(
- workflowInstanceId,
- new ResumeWorkflowRuntimeOptions(options.CorrelationId, bookmark.BookmarkId, null, options.Input),
- cancellationToken);
-
- triggeredWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks));
- }
-
- return new TriggerWorkflowsResult(triggeredWorkflows);
+ return new TriggerWorkflowsResult(results);
}
///
diff --git a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs
index 72339a8d8..0034329cc 100644
--- a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs
+++ b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs
@@ -36,7 +36,7 @@ internal class TriggerWebhookDrivenActivities : INotificationHandler FindActivityDescriptors(string eventType) =>
diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs
index 56cab6a03..98949b315 100644
--- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs
+++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Activities/Flowchart.cs
@@ -27,17 +27,25 @@ public class Flowchart : Container
}
[Port] [Browsable(false)] public IActivity? Start { get; set; }
+ ///
+ /// A list of connections between activities.
+ ///
public ICollection Connections { get; set; } = new List();
+ ///
protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context)
{
- if (Start == null!)
+ var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId;
+ var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : default;
+ var startActivity = triggerActivity ?? Start;
+
+ if (startActivity == null!)
{
await context.CompleteActivityAsync();
return;
}
- await context.ScheduleActivityAsync(Start);
+ await context.ScheduleActivityAsync(startActivity);
}
private async ValueTask OnDescendantCompletedAsync(ActivityCompleted signal, SignalContext context)
diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs
index 400619c4c..c30490bca 100644
--- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs
+++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/LoggingMiddleware.cs
@@ -6,12 +6,18 @@ using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Core.Middleware.Activities;
+///
+/// An activity execution middleware component that logs information about the activity being executed.
+///
public class LoggingMiddleware : IActivityExecutionMiddleware
{
private readonly ActivityMiddlewareDelegate _next;
private readonly ILogger _logger;
private readonly Stopwatch _stopwatch;
+ ///
+ /// Constructor.
+ ///
public LoggingMiddleware(ActivityMiddlewareDelegate next, ILogger logger)
{
_next = next;
@@ -19,18 +25,25 @@ public class LoggingMiddleware : IActivityExecutionMiddleware
_stopwatch = new Stopwatch();
}
+ ///
public async ValueTask InvokeAsync(ActivityExecutionContext context)
{
var activity = context.Activity;
- _logger.LogDebug("Executing activity {ActivityType}", activity.Type);
+ _logger.LogInformation("Executing activity {ActivityId}", activity.Id);
_stopwatch.Restart();
await _next(context);
_stopwatch.Stop();
- _logger.LogDebug("Executed activity {ActivityType} in {Elapsed}", activity.Type, _stopwatch.Elapsed);
+ _logger.LogInformation("Executed activity {ActivityId} in {Elapsed}", activity.Id, _stopwatch.Elapsed);
}
}
+///
+/// Extends to install the component.
+///
public static class LoggingMiddlewareExtensions
{
+ ///
+ /// Installs the component.
+ ///
public static IActivityExecutionPipelineBuilder UseLogging(this IActivityExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware();
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs
index 3ea5fdeac..a255f2314 100644
--- a/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs
+++ b/src/modules/Elsa.Workflows.Core/Models/WorkflowExecutionContext.cs
@@ -266,7 +266,7 @@ public class WorkflowExecutionContext
/// Returns the with the specified ID from the workflow graph.
///
public IActivity FindActivityByNodeId(string nodeId) => FindNodeById(nodeId).Activity;
-
+
///
/// Returns a custom property with the specified key from the dictionary.
///
diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs
index 354e7af85..f5a5646fa 100644
--- a/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs
@@ -46,7 +46,7 @@ internal class DispatchWorkflowRequestHandler :
public async Task HandleAsync(DispatchResumeWorkflowsCommand command, CancellationToken cancellationToken)
{
- var options = new ResumeWorkflowRuntimeOptions(CorrelationId: command.CorrelationId, Input: command.Input);
+ var options = new TriggerWorkflowsRuntimeOptions(CorrelationId: command.CorrelationId, Input: command.Input);
await _workflowRuntime.ResumeWorkflowsAsync(command.ActivityTypeName, command.BookmarkPayload, options, cancellationToken);
return Unit.Instance;
diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs
index e8f96472f..6ca90d31b 100644
--- a/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeDispatchWorkflowActivityHandler.cs
@@ -4,6 +4,7 @@ using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Notifications;
using Elsa.Workflows.Runtime.Activities;
using Elsa.Workflows.Runtime.Bookmarks;
+using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Services;
using JetBrains.Annotations;
@@ -15,9 +16,9 @@ namespace Elsa.Workflows.Runtime.Handlers;
[PublicAPI]
internal class ResumeDispatchWorkflowActivityHandler : INotificationHandler
{
- private readonly IWorkflowRuntime _workflowRuntime;
+ private readonly IWorkflowDispatcher _workflowRuntime;
- public ResumeDispatchWorkflowActivityHandler(IWorkflowRuntime workflowRuntime)
+ public ResumeDispatchWorkflowActivityHandler(IWorkflowDispatcher workflowRuntime)
{
_workflowRuntime = workflowRuntime;
}
@@ -29,6 +30,7 @@ internal class ResumeDispatchWorkflowActivityHandler : INotificationHandler();
- await _workflowRuntime.TriggerWorkflowsAsync(activityTypeName, bookmark, new TriggerWorkflowsRuntimeOptions(), cancellationToken);
+ var request = new DispatchResumeWorkflowsRequest(activityTypeName, bookmark);
+ await _workflowRuntime.DispatchAsync(request, cancellationToken);
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs
index 477a078f8..6592d0416 100644
--- a/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs
@@ -68,6 +68,40 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
return new WorkflowExecutionResult(workflowState.Id, workflowState.Bookmarks);
}
+ ///
+ public async Task> StartWorkflowsAsync(
+ string activityTypeName,
+ object bookmarkPayload,
+ TriggerWorkflowsRuntimeOptions options,
+ CancellationToken cancellationToken = default)
+ {
+ var results = new List();
+ var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
+
+ // Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances from being created.
+ var sharedResource = $"{nameof(DefaultWorkflowRuntime)}__StartTriggeredWorkflows__{hash}";
+ await using (await _distributedLockProvider.AcquireLockAsync(sharedResource, TimeSpan.FromMinutes(10), cancellationToken))
+ {
+ var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
+
+ foreach (var trigger in triggers)
+ {
+ var definitionId = trigger.WorkflowDefinitionId;
+ var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
+ var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
+
+ // If we can't start the workflow, don't try it.
+ if (!canStartResult.CanStart)
+ continue;
+
+ var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
+ results.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
+ }
+ }
+
+ return results;
+ }
+
///
public async Task ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
{
@@ -104,52 +138,21 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
}
///
- public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
+ public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
var correlationId = options.CorrelationId;
- var bookmarks = correlationId == null ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken);
- return await ResumeWorkflowsAsync(bookmarks, options, cancellationToken);
+ var bookmarks = string.IsNullOrWhiteSpace(correlationId) ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken);
+ return await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken);
}
///
- public async Task TriggerWorkflowsAsync(
- string activityTypeName,
- object bookmarkPayload,
- TriggerWorkflowsRuntimeOptions options,
- CancellationToken cancellationToken = default)
+ public async Task TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
- var triggeredWorkflows = new List();
- var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
-
- // Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances being created.
- const string sharedResource = $"{nameof(DefaultWorkflowRuntime)}__StartTriggeredWorkflows";
- await using (await _distributedLockProvider.AcquireLockAsync(sharedResource, TimeSpan.FromMinutes(10), cancellationToken))
- {
- var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
-
- foreach (var trigger in triggers)
- {
- var definitionId = trigger.WorkflowDefinitionId;
- var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
- var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
-
- // If we can't start the workflow, don't try it.
- if (!canStartResult.CanStart)
- continue;
-
- var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
- triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
- }
- }
-
- // Resume bookmarks.
- var correlationId = options.CorrelationId;
- var bookmarks = (string.IsNullOrEmpty(correlationId) ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken)).ToList();
- var resumedWorkflows = await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken);
-
- triggeredWorkflows.AddRange(resumedWorkflows.Select(x => new WorkflowExecutionResult(x.InstanceId, x.Bookmarks)));
- return new TriggerWorkflowsResult(triggeredWorkflows);
+ var startedWorkflows = await StartWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
+ var resumedWorkflows = await ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
+ var results = startedWorkflows.Concat(resumedWorkflows).ToList();
+ return new TriggerWorkflowsResult(results);
}
///
diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs
index 712d4d2a9..d95454d4d 100644
--- a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs
@@ -22,6 +22,16 @@ public interface IWorkflowRuntime
///
///
Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
+
+ ///
+ /// Starts all workflows with triggers matching the specified activity type and bookmark payload.
+ ///
+ ///
+ Task> StartWorkflowsAsync(
+ string activityTypeName,
+ object bookmarkPayload,
+ TriggerWorkflowsRuntimeOptions options,
+ CancellationToken cancellationToken = default);
///
/// Resumes an existing workflow instance.
@@ -34,7 +44,7 @@ public interface IWorkflowRuntime
///
/// Resumes all workflows that are bookmarked on the specified activity type.
///
- Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
+ Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default);
///
/// Starts all workflows and resumes existing workflow instances based on the specified activity type and bookmark payload.