diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/ActivityIncident.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/ActivityIncident.cs index 503b5f9ce..4c6029cce 100644 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/ActivityIncident.cs +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/ActivityIncident.cs @@ -40,6 +40,16 @@ public class ActivityIncident /// The Node ID of the activity that caused the incident. public string ActivityNodeId { get; init; } = default!; + /// + /// The ID of the individual activity execution that caused the incident. + /// + /// + /// identifies the static workflow node, which several executions can share when a node + /// is looped, retried or run concurrently. This tells those executions apart. Null for an incident raised outside an + /// activity execution, and for incidents persisted before the server recorded it. + /// + public string? ActivityInstanceId { get; init; } + /// The type of the activity that caused the incident. public string ActivityType { get; init; } = default!; diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs index b39f8aa6c..673f4246b 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs @@ -336,7 +336,7 @@ public static partial class ActivityExecutionContextExtensions var exceptionState = ExceptionState.FromException(e); var systemClock = activityExecutionContext.GetRequiredService(); var now = systemClock.UtcNow; - var incident = new ActivityIncident(activity.Id, activity.NodeId, activity.Type, e.Message, exceptionState, now); + var incident = new ActivityIncident(activity.Id, activity.NodeId, activity.Type, e.Message, exceptionState, now, activityExecutionContext.Id); activityExecutionContext.WorkflowExecutionContext.Incidents.Add(incident); activityExecutionContext.AggregateFaultCount++; @@ -347,13 +347,40 @@ public static partial class ActivityExecutionContextExtensions } /// - /// Recovers the current activity from a faulted state by resetting fault counts and transitioning the activity to the running status. + /// Recovers the current activity from a faulted state: undoes the fault counts, discards the incident and the recorded + /// exception, and transitions the activity back to the running status. /// /// - /// This method is asymmetric with Fault: it sets the current context's fault count to zero, which is idempotent, but - /// decrements the count on every ancestor, which is not. Calling it more than once per fault drives ancestor counts negative. - /// The status transition only applies to an activity that is still , so that a - /// handler that already terminalized the activity does not see its decision undone. + /// + /// This is the inverse of Fault, and is meant to leave no trace of a fault an enclosing container claimed. In + /// particular it removes the that Fault appended. A handled fault must not remain an + /// incident, because plenty of code treats a non-empty as "this workflow + /// failed" without looking further - the HTTP endpoint fault handler is one, and it would answer a caller with a fault + /// response for a workflow that caught its error and completed normally. The execution log still records the failure, so + /// nothing is hidden from anyone reading the journal. + /// + /// + /// The incident is matched on this execution's id rather than on the activity's node id, because a node inside + /// a loop, retried, or run concurrently has several executions that all raise incidents under the same node id, and + /// recovering one of them must not remove another's. Within a single execution the most recent is taken, which keeps + /// an activity that faults, recovers, and faults again correct: each recovery removes its own incident rather than + /// the whole execution's history. That last part relies on preserving + /// insertion order, which it does because it is list-backed; it is typed as a plain collection, so a future change to + /// an unordered one would silently pick an arbitrary incident of that execution rather than its newest. + /// + /// + /// Clearing the exception also clears it from the activity's execution record, which maps it from the live context. + /// That is intended and follows from the same reasoning as the incident: an execution whose failure a container + /// claimed is not a failed execution. The journal keeps the evidence either way, since the Faulted entry is + /// written by ExecutionLogMiddleware, which sits inside this middleware and logs on the way past. + /// + /// + /// It remains asymmetric with Fault in one respect: it sets the current context's fault count to zero, which is + /// idempotent, but decrements the count on every ancestor, which is not. Calling it more than once per fault drives + /// ancestor counts negative. The status transition only applies to an activity that is still + /// , so that a handler that already + /// terminalized the activity does not see its decision undone. + /// /// public void RecoverFromFault() { @@ -364,6 +391,14 @@ public static partial class ActivityExecutionContextExtensions foreach (var ancestor in ancestors) ancestor.AggregateFaultCount--; + var incidents = activityExecutionContext.WorkflowExecutionContext.Incidents; + var ownIncident = incidents.LastOrDefault(x => x.ActivityInstanceId == activityExecutionContext.Id); + + if (ownIncident != null) + incidents.Remove(ownIncident); + + activityExecutionContext.Exception = null; + if (activityExecutionContext.Status == ActivityStatus.Faulted) activityExecutionContext.TransitionTo(ActivityStatus.Running); } diff --git a/src/modules/Elsa.Workflows.Core/Models/ActivityIncident.cs b/src/modules/Elsa.Workflows.Core/Models/ActivityIncident.cs index c787f729a..a18dc6603 100644 --- a/src/modules/Elsa.Workflows.Core/Models/ActivityIncident.cs +++ b/src/modules/Elsa.Workflows.Core/Models/ActivityIncident.cs @@ -25,7 +25,8 @@ public class ActivityIncident /// The message of the incident. /// The exception that caused the incident. /// The timestamp of the incident. - public ActivityIncident(string activityId, string activityNodeId, string activityType, string message, ExceptionState? exception, DateTimeOffset timestamp) + /// The ID of the individual activity execution that caused the incident, if any. + public ActivityIncident(string activityId, string activityNodeId, string activityType, string message, ExceptionState? exception, DateTimeOffset timestamp, string? activityInstanceId = null) { ActivityId = activityId; ActivityNodeId = activityNodeId; @@ -33,6 +34,7 @@ public class ActivityIncident Message = message; Exception = exception; Timestamp = timestamp; + ActivityInstanceId = activityInstanceId; } /// The ID of the activity that caused the incident. @@ -40,6 +42,20 @@ public class ActivityIncident /// The Node ID of the activity that caused the incident. public string ActivityNodeId { get; init; } = default!; + /// + /// The ID of the individual activity execution that caused the incident. + /// + /// + /// identifies the static workflow node, and several executions can share one: a node + /// inside a loop, retried, or run concurrently raises an incident per execution, all carrying the same node id. This + /// tells them apart, so an incident can be tied back to the execution that raised it. + /// + /// Null for an incident raised outside an activity execution, such as one recorded against the workflow itself, and + /// for incidents persisted before this property existed. + /// + /// + public string? ActivityInstanceId { get; init; } + /// The type of the activity that caused the incident. public string ActivityType { get; init; } = default!; diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FaultSignals/FaultSignalTests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FaultSignals/FaultSignalTests.cs index 962e72a62..4eb648c54 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FaultSignals/FaultSignalTests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/FaultSignals/FaultSignalTests.cs @@ -33,8 +33,11 @@ public class FaultSignalTests(ITestOutputHelper testOutputHelper) Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); Assert.Equal(WorkflowSubStatus.Finished, result.WorkflowState.SubStatus); - // The incident stays on record even though the fault was handled: caught failures remain visible to operators. - Assert.Single(result.WorkflowState.Incidents); + // A fault a container claimed is not an incident. Plenty of code reads a non-empty Incidents as "this workflow + // failed" without looking further - the HTTP endpoint fault handler among them - so leaving one here would + // answer a caller with a fault response for a workflow that caught its error and finished normally. The + // execution log still records the failure for anyone reading the journal. + Assert.Empty(result.WorkflowState.Incidents); } [Theory(DisplayName = "A fault nobody handles is left to the incident strategy, exactly as before")] @@ -95,8 +98,9 @@ public class FaultSignalTests(ITestOutputHelper testOutputHelper) Assert.Equal(0, container.FaultsSeen); Assert.NotEqual(WorkflowSubStatus.Faulted, result.WorkflowState.SubStatus); - // The incident is still on record, and the fault bookkeeping was still recovered exactly once. - Assert.Single(result.WorkflowState.Incidents); + // The activity claimed the fault itself, so it is not an incident either, and the fault bookkeeping was still + // recovered exactly once. + Assert.Empty(result.WorkflowState.Incidents); var faultedContext = result.GetActivityContext(faultingActivity); Assert.NotNull(faultedContext); diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/Extensions/ActivityExecutionContextExtensions/FaultHandlingTests.cs b/test/unit/Elsa.Workflows.Core.UnitTests/Extensions/ActivityExecutionContextExtensions/FaultHandlingTests.cs index d8e15fe5f..8f66bf253 100644 --- a/test/unit/Elsa.Workflows.Core.UnitTests/Extensions/ActivityExecutionContextExtensions/FaultHandlingTests.cs +++ b/test/unit/Elsa.Workflows.Core.UnitTests/Extensions/ActivityExecutionContextExtensions/FaultHandlingTests.cs @@ -106,4 +106,130 @@ public class FaultHandlingTests Assert.Equal(ActivityStatus.Canceled, faultedContext.Status); Assert.Equal(new[] { 0, 0, 0 }, chain.Select(x => x.AggregateFaultCount)); } + + [Fact] + public async Task Fault_RecordsAnIncident() + { + // Arrange + var context = await CreateContextAsync(); + + // Act + context.Fault(new InvalidOperationException("Test error")); + + // Assert + var incident = Assert.Single(context.WorkflowExecutionContext.Incidents); + Assert.Equal(context.NodeId, incident.ActivityNodeId); + Assert.Equal("Test error", incident.Message); + } + + [Fact] + public async Task RecoverFromFault_RemovesTheIncidentAndTheException() + { + // A fault an enclosing container claimed is not an incident. Plenty of code reads + // WorkflowExecutionContext.Incidents as "this workflow failed" without looking further - the HTTP endpoint + // fault handler among them - and would otherwise answer a caller with a fault response for a workflow that + // caught its error and completed normally. + + // Arrange + var context = await CreateContextAsync(); + context.Fault(new InvalidOperationException("Test error")); + + // Act + context.RecoverFromFault(); + + // Assert + Assert.Empty(context.WorkflowExecutionContext.Incidents); + Assert.Null(context.Exception); + } + + [Fact] + public async Task RecoverFromFault_LeavesIncidentsFromAnotherExecutionOfTheSameNode() + { + // ActivityNodeId identifies the static workflow node, and one node can have several executions: inside a loop, + // retried, or run concurrently. Matching an incident on the node id alone would let one execution's recovery + // remove another execution's incident, so the match is on the execution id. + + // Arrange: two executions of one activity, hence one shared node id and two distinct execution ids. The ids are + // assigned here because this fixture substitutes the identity generator, which hands every context an empty id. + var first = await CreateContextAsync(); + var second = await first.WorkflowExecutionContext.CreateActivityExecutionContextAsync(first.Activity, new()); + first.Id = "execution-1"; + second.Id = "execution-2"; + + Assert.Equal(first.NodeId, second.NodeId); + + first.Fault(new InvalidOperationException("First execution")); + second.Fault(new InvalidOperationException("Second execution")); + + // Act: recover the earlier execution, whose incident is not the most recently appended. + first.RecoverFromFault(); + + // Assert: the other execution keeps its own. + var remaining = Assert.Single(first.WorkflowExecutionContext.Incidents); + Assert.Equal(second.Id, remaining.ActivityInstanceId); + Assert.Equal("Second execution", remaining.Message); + } + + [Fact] + public async Task Fault_StampsTheIncidentWithTheExecutionThatRaisedIt() + { + // Arrange + var context = await CreateContextAsync(); + + // Act + context.Fault(new InvalidOperationException("Test error")); + + // Assert + var incident = Assert.Single(context.WorkflowExecutionContext.Incidents); + Assert.Equal(context.Id, incident.ActivityInstanceId); + } + + [Fact] + public async Task RecoverFromFault_LeavesIncidentsBelongingToOtherActivities() + { + // Arrange + var chain = await CreateContextChainAsync(); + var other = chain[0]; + var faultedContext = chain[^1]; + other.Fault(new InvalidOperationException("Someone else's problem")); + faultedContext.Fault(new InvalidOperationException("Test error")); + + // Act + faultedContext.RecoverFromFault(); + + // Assert + var remaining = Assert.Single(faultedContext.WorkflowExecutionContext.Incidents); + Assert.Equal(other.NodeId, remaining.ActivityNodeId); + } + + [Fact] + public async Task RecoverFromFault_RemovesOneIncidentPerFault() + { + // An activity that faults, is recovered, and faults again keeps the incident that was never recovered. + // Recovery pairs with a single fault rather than wiping the activity's history wholesale. + + // Arrange + var context = await CreateContextAsync(); + context.Fault(new InvalidOperationException("First")); + context.RecoverFromFault(); + context.Fault(new InvalidOperationException("Second")); + + // Assert + var incident = Assert.Single(context.WorkflowExecutionContext.Incidents); + Assert.Equal("Second", incident.Message); + } + + [Fact] + public async Task RecoverFromFault_WithNoIncidentIsHarmless() + { + // Arrange + var context = await CreateContextAsync(); + + // Act + context.RecoverFromFault(); + + // Assert + Assert.Empty(context.WorkflowExecutionContext.Incidents); + Assert.Null(context.Exception); + } }