diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs index ae041e601..7e56d5267 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs @@ -17,6 +17,7 @@ using Elsa.Workflows.Core.Notifications; using Elsa.Workflows.Core.Services; using Elsa.Workflows.Core.Signals; using JetBrains.Annotations; +using Microsoft.Extensions.Logging; // ReSharper disable once CheckNamespace namespace Elsa.Extensions; @@ -238,7 +239,7 @@ public static class ActivityExecutionContextExtensions context.ActivityState[inputDescriptor.Name] = value!; return value; } - + /// /// Returns the outcome name for the specified port property name. /// @@ -295,7 +296,7 @@ public static class ActivityExecutionContextExtensions /// public static IEnumerable GetActiveChildren(this ActivityExecutionContext context) => context.WorkflowExecutionContext.ActiveActivityExecutionContexts.Where(x => x.ParentActivityExecutionContext == context); - + /// /// Returns a flattened list of the current context's immediate children. /// @@ -308,11 +309,11 @@ public static class ActivityExecutionContextExtensions public static IEnumerable GetDescendants(this ActivityExecutionContext context) { var children = context.WorkflowExecutionContext.ActiveActivityExecutionContexts.Where(x => x.ParentActivityExecutionContext == context).ToList(); - + foreach (var child in children) { yield return child; - + foreach (var descendant in child.GetDescendants()) yield return descendant; } @@ -324,7 +325,8 @@ public static class ActivityExecutionContextExtensions public static async ValueTask SendSignalAsync(this ActivityExecutionContext context, object signal) { var receivingContexts = new[] { context }.Concat(context.GetAncestors()).ToList(); - var capturingContexts = receivingContexts.AsEnumerable().Reverse().ToList(); + var capturingContexts = receivingContexts.AsEnumerable().Reverse().ToList(); + var logger = context.GetRequiredService>(); // Let all ancestors capture the signal. foreach (var ancestorContext in capturingContexts) @@ -334,12 +336,16 @@ public static class ActivityExecutionContextExtensions if (ancestorContext.Activity is not ISignalHandler handler) continue; + logger.LogDebug("Capturing signal {SignalType} on activity {ActivityId} of type {ActivityType}", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type); await handler.CaptureSignalAsync(signal, signalContext); if (signalContext.StopPropagationRequested) + { + logger.LogDebug("Propagation of signal {SignalType} on activity {ActivityId} of type {ActivityType} was stopped", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type); return; + } } - + // Let all ancestors receive the signal. foreach (var ancestorContext in receivingContexts) { @@ -348,10 +354,14 @@ public static class ActivityExecutionContextExtensions if (ancestorContext.Activity is not ISignalHandler handler) continue; + logger.LogDebug("Receiving signal {SignalType} on activity {ActivityId} of type {ActivityType}", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type); await handler.ReceiveSignalAsync(signal, signalContext); if (signalContext.StopPropagationRequested) + { + logger.LogDebug("Propagation of signal {SignalType} on activity {ActivityId} of type {ActivityType} was stopped", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type); return; + } } } @@ -362,7 +372,7 @@ public static class ActivityExecutionContextExtensions { await context.TargetContext.CompleteActivityAsync(result); } - + /// /// Complete the current activity. This should only be called by activities that explicitly suppress automatic-completion. /// @@ -374,8 +384,8 @@ public static class ActivityExecutionContextExtensions // Update all child contexts. var childContexts = context.WorkflowExecutionContext.ActiveActivityExecutionContexts.Where(x => x.ParentActivityExecutionContext == context).ToList(); - - foreach (var childContext in childContexts) + + foreach (var childContext in childContexts) await childContext.CancelActivityAsync(); // Mark the activity as complete. @@ -424,7 +434,7 @@ public static class ActivityExecutionContextExtensions // Update the completed at timestamp. context.CompletedAt = context.WorkflowExecutionContext.SystemClock.UtcNow; } - + /// /// Complete the current activity. This should only be called by activities that explicitly suppress automatic-completion. /// @@ -452,7 +462,7 @@ public static class ActivityExecutionContextExtensions var serializedOutputValue = serializer.Serialize(outputValue); context.JournalData[outputName] = serializedOutputValue; } - + // Send a signal. await context.SendSignalAsync(new ScheduleActivityOutcomes(outcomes)); } @@ -461,7 +471,7 @@ public static class ActivityExecutionContextExtensions /// Complete the current activity with the specified outcome. /// public static ValueTask CompleteActivityWithOutcomesAsync(this ActivityCompletedContext context, params string[] outcomes) => context.CompleteActivityAsync(new Outcomes(outcomes)); - + /// /// Complete the current activity with the specified outcome. /// diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs index bcbc70bac..f63244dfc 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs @@ -11,6 +11,7 @@ using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; @@ -34,7 +35,7 @@ public class WorkflowDefinitionActivity : Composite, IInitializable /// The latest published version number set by the provider. This is used by tooling to let the user know that a newer version is available. /// public int LatestAvailablePublishedVersion { get; set; } - + /// /// The latest published version ID set by the provider. This is used by tooling to let the user know that a newer version is available. /// @@ -115,14 +116,18 @@ public class WorkflowDefinitionActivity : Composite, IInitializable private async ValueTask OnChildCompletedAsync(ActivityCompletedContext context) { var targetContext = context.TargetContext; - + // Do we have a "complete composite" signal that triggered the completion? var completeCompositeSignal = context.WorkflowExecutionContext.TransientProperties.TryGetValue(nameof(CompleteCompositeSignal), out var signal) ? (CompleteCompositeSignal)signal : default; // If we do, make sure to remove it from the transient properties. if (completeCompositeSignal != null) + { + var logger = context.GetRequiredService>(); + logger.LogDebug("Received a complete composite signal and removing it from the transient properties"); context.WorkflowExecutionContext.TransientProperties.Remove(nameof(CompleteCompositeSignal)); - + } + await targetContext.CompleteActivityAsync(completeCompositeSignal?.Value); } @@ -150,7 +155,7 @@ public class WorkflowDefinitionActivity : Composite, IInitializable // Find the latest version. workflowDefinition = await workflowDefinitionStore.FindAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Latest }, cancellationToken); } - + if (workflowDefinition == null) throw new Exception($"Could not find workflow definition with ID {WorkflowDefinitionId}."); }