diff --git a/doc/changelogs/3.6.0.md b/doc/changelogs/3.6.0.md index 434719dba..a1f3575bd 100644 --- a/doc/changelogs/3.6.0.md +++ b/doc/changelogs/3.6.0.md @@ -47,6 +47,10 @@ Compare: [`3.5.3...3.6.0`](https://github.com/elsa-workflows/elsa-core/compare/3 - **`WorkflowDefinitionsReloaded` notification**: `DefaultRegistriesPopulator` now dispatches a `WorkflowDefinitionsReloaded` notification after repopulating the workflow definition store, enabling subscriber nodes to synchronize their registries. ([f5505d66c9](https://github.com/elsa-workflows/elsa-core/commit/f5505d66c9)) ([#7293](https://github.com/elsa-workflows/elsa-core/pull/7293)) +### OpenTelemetry workflow instrumentation + +- **Workflow and activity telemetry**: Elsa now emits `Elsa.Workflows` `ActivitySource` spans for workflow execution cycles and activity execution, plus optional `Elsa.Workflows` meter instruments for workflow started/completed/faulted counts and activity duration. `SendHttpRequest` and `FlowSendHttpRequest` use the configured `HttpClient`, so standard .NET HTTP client instrumentation can create child HTTP spans and propagate W3C trace context when tracing is enabled. ([#7489](https://github.com/elsa-workflows/elsa-core/issues/7489)) + ### Consumers API & recursive export - **`GET /workflow-definitions/{definitionId}/consumers`**: New endpoint returns all recursive consuming workflow definitions for a given definition. Built on a new `IWorkflowReferenceGraphBuilder` / `WorkflowReferenceGraph` model that replaces the previous ad-hoc recursive query approach. ([0a2e99ad42](https://github.com/elsa-workflows/elsa-core/commit/0a2e99ad42)) ([#7309](https://github.com/elsa-workflows/elsa-core/pull/7309)) diff --git a/doc/wiki/opentelemetry-workflows.md b/doc/wiki/opentelemetry-workflows.md new file mode 100644 index 000000000..f41449514 --- /dev/null +++ b/doc/wiki/opentelemetry-workflows.md @@ -0,0 +1,32 @@ +# OpenTelemetry Workflow Instrumentation + +Elsa emits first-party OpenTelemetry-compatible workflow and activity instrumentation through `System.Diagnostics`. + +To collect workflow traces, configure OpenTelemetry to listen to the `Elsa.Workflows` activity source. To collect workflow metrics, configure it to listen to the `Elsa.Workflows` meter. + +```csharp +services.AddOpenTelemetry() + .WithTracing(builder => builder.AddSource("Elsa.Workflows")) + .WithMetrics(builder => builder.AddMeter("Elsa.Workflows")); +``` + +If you previously enabled workflow tracing through the `Elsa.OpenTelemetry` extension package, avoid enabling both the extension tracing middleware and the first-party workflow spans for the same host unless duplicate workflow and activity spans are acceptable. Both integrations publish to the `Elsa.Workflows` activity source so existing collectors can keep the same source configuration. + +## Traces + +Elsa creates spans around workflow execution cycles and activity execution. The spans include workflow and activity identifiers, definition metadata, status, tenant ID when available, and fault status. Workflow input, activity input, output payloads, headers, and variable values are not added as span attributes. + +Faulted workflow and activity spans use `ActivityStatusCode.Error` and record the exception type as `exception.type` when an exception is available. Exception messages and stack traces are not added to spans or exception events by Elsa workflow instrumentation. + +Outbound `SendHttpRequest` and `FlowSendHttpRequest` calls inject the current W3C trace context headers when an active workflow/activity span exists, so downstream services can continue the same trace without Elsa-specific middleware. Enable standard .NET HTTP client instrumentation in your OpenTelemetry setup when you want outbound HTTP spans for these calls. + +## Metrics + +The `Elsa.Workflows` meter emits: + +- `elsa.workflow.started` +- `elsa.workflow.completed` +- `elsa.workflow.faulted` +- `elsa.activity.duration` in seconds + +Metric tags use the same low-cardinality workflow and activity metadata as the spans where practical. The started counter omits execution status tags because it records the pre-execution boundary; completed and faulted counters include terminal workflow status tags. diff --git a/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs b/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs index bb8930ab5..49c3310bb 100644 --- a/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs +++ b/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs @@ -266,6 +266,8 @@ public abstract class SendHttpRequestBase(string? source = null, int? line = nul foreach (var header in headers) request.Headers.Add(header.Key, header.Value.AsEnumerable()); + InjectTraceContext(request); + var contentType = ContentType.GetOrDefault(context); var content = Content.GetOrDefault(context); @@ -279,6 +281,23 @@ public abstract class SendHttpRequestBase(string? source = null, int? line = nul return request; } + private static void InjectTraceContext(HttpRequestMessage request) + { + var activity = System.Diagnostics.Activity.Current; + + if (activity == null) + return; + + System.Diagnostics.DistributedContextPropagator.Current.Inject(activity, request, static (carrier, key, value) => + { + if (carrier is not HttpRequestMessage requestMessage) + return; + + if (!requestMessage.Headers.Contains(key)) + requestMessage.Headers.TryAddWithoutValidation(key, value); + }); + } + private IHttpContentFactory SelectContentWriter(string? contentType, IEnumerable factories) { if (string.IsNullOrWhiteSpace(contentType)) @@ -329,4 +348,4 @@ public abstract class SendHttpRequestBase(string? source = null, int? line = nul _ => false // Other errors are not transient }; } -} \ No newline at end of file +} diff --git a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs index b062b6645..19373143f 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs @@ -301,6 +301,11 @@ public partial class WorkflowExecutionContext : IExecutionContext /// The date and time the workflow execution context has finished. public DateTimeOffset? FinishedAt { get; set; } + /// + /// Gets the exception that occurred during workflow execution, if any. + /// + internal Exception? Exception { get; set; } + /// Gets the clock used to determine the current time. public ISystemClock SystemClock { get; } @@ -713,4 +718,3 @@ public partial class WorkflowExecutionContext : IExecutionContext return _commitStateHandler.CommitAsync(this, CancellationToken); } } - diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Workflows/ExceptionHandlingMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Workflows/ExceptionHandlingMiddleware.cs index 3d6070c36..30997b1b1 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Workflows/ExceptionHandlingMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Workflows/ExceptionHandlingMiddleware.cs @@ -43,6 +43,11 @@ public class ExceptionHandlingMiddleware : IWorkflowExecutionMiddleware { await _next(context); } + catch (OperationCanceledException) + { + context.Cancel(); + throw; + } catch (Exception e) { _logger.LogWarning(e, "An exception was caught from a downstream middleware component"); @@ -50,9 +55,10 @@ public class ExceptionHandlingMiddleware : IWorkflowExecutionMiddleware var now = _systemClock.UtcNow; var activity = context.Workflow; var incident = new ActivityIncident(activity.Id, activity.NodeId ,activity.Type, e.Message, exceptionState, now); + context.Exception ??= e; context.Incidents.Add(incident); context.TransitionTo(WorkflowSubStatus.Faulted); context.AddExecutionLogEntry("Faulted", e.Message, exceptionState); } } -} \ No newline at end of file +} diff --git a/src/modules/Elsa.Workflows.Core/Services/ActivityInvoker.cs b/src/modules/Elsa.Workflows.Core/Services/ActivityInvoker.cs index 4560f6974..7cfa34976 100644 --- a/src/modules/Elsa.Workflows.Core/Services/ActivityInvoker.cs +++ b/src/modules/Elsa.Workflows.Core/Services/ActivityInvoker.cs @@ -1,4 +1,5 @@ using Elsa.Workflows.Options; +using Elsa.Workflows.Telemetry; using Microsoft.Extensions.Logging; namespace Elsa.Workflows; @@ -41,10 +42,25 @@ public class ActivityInvoker( /// public async Task InvokeAsync(ActivityExecutionContext activityExecutionContext) { - var loggerState = loggerStateGenerator.GenerateLoggerState(activityExecutionContext); - using var loggingScope = logger.BeginScope(loggerState); + var telemetryScope = WorkflowInstrumentation.StartActivity(activityExecutionContext); + Exception? exception = null; - // Execute the activity execution pipeline. - await pipeline.ExecuteAsync(activityExecutionContext); + try + { + var loggerState = loggerStateGenerator.GenerateLoggerState(activityExecutionContext); + using var loggingScope = logger.BeginScope(loggerState); + + // Execute the activity execution pipeline. + await pipeline.ExecuteAsync(activityExecutionContext); + } + catch (Exception e) + { + exception = e; + throw; + } + finally + { + WorkflowInstrumentation.StopActivity(telemetryScope, activityExecutionContext, exception); + } } -} \ No newline at end of file +} diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs index ec3ea136a..cbaa6ea80 100644 --- a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs @@ -7,6 +7,7 @@ using Elsa.Workflows.Models; using Elsa.Workflows.Notifications; using Elsa.Workflows.Options; using Elsa.Workflows.State; +using Elsa.Workflows.Telemetry; using Microsoft.Extensions.Logging; namespace Elsa.Workflows; @@ -213,13 +214,33 @@ public class WorkflowRunner( await notificationSender.SendAsync(new WorkflowExecuting(workflow, workflowExecutionContext), cancellationToken); // If the status is Pending, it means the workflow is started for the first time. - if (workflowExecutionContext.SubStatus == WorkflowSubStatus.Pending) + var isStarting = workflowExecutionContext.SubStatus == WorkflowSubStatus.Pending; + if (isStarting) { workflowExecutionContext.TransitionTo(WorkflowSubStatus.Executing); await notificationSender.SendAsync(new WorkflowStarted(workflow, workflowExecutionContext), cancellationToken); } - await pipeline.ExecuteAsync(workflowExecutionContext); + var telemetryScope = WorkflowInstrumentation.StartWorkflow(workflowExecutionContext, isStarting); + Exception? workflowExecutionException = null; + + try + { + await pipeline.ExecuteAsync(workflowExecutionContext); + } + catch (Exception e) + { + if (e is not OperationCanceledException) + workflowExecutionContext.Exception ??= e; + + workflowExecutionException = e; + throw; + } + finally + { + WorkflowInstrumentation.StopWorkflow(telemetryScope, workflowExecutionContext, workflowExecutionException); + } + var workflowState = workflowStateExtractor.Extract(workflowExecutionContext); if (workflowState.Status == WorkflowStatus.Finished) @@ -235,4 +256,4 @@ public class WorkflowRunner( await commitStateHandler.CommitAsync(workflowExecutionContext, workflowState, cancellationToken); return new(workflowExecutionContext, workflowState, workflowExecutionContext.Workflow, result, journal); } -} \ No newline at end of file +} diff --git a/src/modules/Elsa.Workflows.Core/Telemetry/WorkflowInstrumentation.cs b/src/modules/Elsa.Workflows.Core/Telemetry/WorkflowInstrumentation.cs new file mode 100644 index 000000000..b1e084eec --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Telemetry/WorkflowInstrumentation.cs @@ -0,0 +1,282 @@ +using System.Diagnostics; +using System.Diagnostics.Metrics; + +using DiagnosticsActivity = System.Diagnostics.Activity; +using DiagnosticsActivityKind = System.Diagnostics.ActivityKind; +using WorkflowActivityStatus = Elsa.Workflows.ActivityStatus; +using WorkflowInstanceStatus = Elsa.Workflows.WorkflowStatus; +using WorkflowInstanceSubStatus = Elsa.Workflows.WorkflowSubStatus; + +namespace Elsa.Workflows.Telemetry; + +/// +/// Provides tracing and metrics instrumentation for workflow and activity execution. +/// +public static class WorkflowInstrumentation +{ + public const string ActivitySourceName = "Elsa.Workflows"; + public const string MeterName = "Elsa.Workflows"; + + public const string WorkflowSystem = "workflow.system"; + public const string WorkflowOperationName = "workflow.operation.name"; + public const string WorkflowName = "workflow.name"; + public const string WorkflowInstanceId = "workflow.instance.id"; + public const string WorkflowDefinitionId = "workflow.definition.id"; + public const string WorkflowDefinitionVersion = "workflow.definition.version"; + public const string WorkflowDefinitionVersionId = "workflow.definition.version.id"; + public const string WorkflowStatus = "workflow.status"; + public const string WorkflowSubStatus = "workflow.substatus"; + public const string WorkflowFaulted = "workflow.faulted"; + public const string WorkflowParentInstanceId = "workflow.parent.instance.id"; + public const string WorkflowCorrelationId = "workflow.correlation.id"; + + public const string ActivityOperationName = "workflow.activity.operation.name"; + public const string ActivityId = "workflow.activity.id"; + public const string ActivityName = "workflow.activity.name"; + public const string ActivityType = "workflow.activity.type"; + public const string ActivityVersion = "workflow.activity.version"; + public const string ActivityExecutionId = "workflow.activity.execution.id"; + public const string ActivityStatus = "workflow.activity.status"; + public const string ActivityOutcome = "workflow.activity.outcome"; + public const string ActivityParentExecutionId = "workflow.activity.parent.execution.id"; + public const string ActivityScheduledByExecutionId = "workflow.activity.scheduled.by.execution.id"; + public const string ActivityFaulted = "workflow.activity.faulted"; + + public const string TenantId = "elsa.tenant.id"; + public const string ExceptionType = "exception.type"; + + [Obsolete("Use ExceptionType. Workflow spans follow OpenTelemetry exception semantic conventions.")] + public const string ErrorType = ExceptionType; + + private const string SystemName = "elsa"; + private const string WorkflowExecuteOperation = "workflow.execute"; + private const string ActivityExecuteOperation = "activity.execute"; + + private static readonly string? Version = typeof(WorkflowInstrumentation).Assembly.GetName().Version?.ToString(); + private static readonly ActivitySource Source = new(ActivitySourceName, Version); + private static readonly Meter Meter = new(MeterName, Version); + private static readonly Counter WorkflowStartedCounter = Meter.CreateCounter("elsa.workflow.started", description: "Number of workflow instances started."); + private static readonly Counter WorkflowCompletedCounter = Meter.CreateCounter("elsa.workflow.completed", description: "Number of workflow instances completed."); + private static readonly Counter WorkflowFaultedCounter = Meter.CreateCounter("elsa.workflow.faulted", description: "Number of workflow execution cycles that faulted."); + private static readonly Histogram ActivityDuration = Meter.CreateHistogram("elsa.activity.duration", "s", "Duration of activity execution."); + + internal static WorkflowInstrumentationScope StartWorkflow(WorkflowExecutionContext context, bool? isStarting = null) + { + var shouldRecordStarted = isStarting ?? context.SubStatus == WorkflowInstanceSubStatus.Pending; + var activity = Source.StartActivity(WorkflowExecuteOperation, DiagnosticsActivityKind.Internal); + + if (activity != null) + { + SetWorkflowTags(activity, context); + activity.SetTag(WorkflowOperationName, WorkflowExecuteOperation); + } + + if (shouldRecordStarted && WorkflowStartedCounter.Enabled) + WorkflowStartedCounter.Add(1, CreateWorkflowTags(context, false)); + + return new WorkflowInstrumentationScope(activity); + } + + internal static ActivityInstrumentationScope StartActivity(ActivityExecutionContext context) + { + var activity = Source.StartActivity(ActivityExecuteOperation, DiagnosticsActivityKind.Internal); + + if (activity != null) + { + SetWorkflowTags(activity, context.WorkflowExecutionContext); + SetActivityTags(activity, context); + activity.SetTag(ActivityOperationName, ActivityExecuteOperation); + } + + return new ActivityInstrumentationScope(activity, Stopwatch.GetTimestamp()); + } + + internal static void StopWorkflow(WorkflowInstrumentationScope scope, WorkflowExecutionContext context, Exception? exception) + { + var activity = scope.Activity; + var workflowException = exception ?? (context.SubStatus == WorkflowInstanceSubStatus.Faulted ? context.Exception : null); + var cancelled = exception is OperationCanceledException || (exception == null && context.SubStatus == WorkflowInstanceSubStatus.Cancelled); + var faulted = !cancelled && (workflowException != null || context.SubStatus == WorkflowInstanceSubStatus.Faulted); + var workflowSubStatus = cancelled ? WorkflowInstanceSubStatus.Cancelled : faulted ? WorkflowInstanceSubStatus.Faulted : (WorkflowInstanceSubStatus?)null; + + if (activity != null) + { + SetWorkflowTags(activity, context, workflowSubStatus); + activity.SetTag(WorkflowFaulted, faulted); + SetError(activity, workflowException, faulted); + activity.Dispose(); + } + + if (faulted && WorkflowFaultedCounter.Enabled) + WorkflowFaultedCounter.Add(1, CreateWorkflowTags(context, workflowSubStatusOverride: workflowSubStatus)); + else if (context.SubStatus == WorkflowInstanceSubStatus.Finished && WorkflowCompletedCounter.Enabled) + WorkflowCompletedCounter.Add(1, CreateWorkflowTags(context)); + } + + internal static void StopActivity(ActivityInstrumentationScope scope, ActivityExecutionContext context, Exception? exception) + { + var activity = scope.Activity; + var cancelled = exception is OperationCanceledException || (exception == null && context.Status == WorkflowActivityStatus.Canceled); + var faulted = !cancelled && (exception != null || context.Status == WorkflowActivityStatus.Faulted); + var activityStatus = cancelled ? WorkflowActivityStatus.Canceled : context.Status; + var duration = Stopwatch.GetElapsedTime(scope.StartTimestamp).TotalSeconds; + + if (ActivityDuration.Enabled) + ActivityDuration.Record(duration, CreateActivityTags(context, activityStatus, faulted)); + + if (activity != null) + { + SetActivityTags(activity, context, activityStatus); + activity.SetTag(ActivityFaulted, faulted); + if (faulted) + activity.SetTag(ActivityStatus, WorkflowActivityStatus.Faulted.ToString()); + SetActivityOutcome(activity, context); + SetError(activity, exception ?? context.Exception, faulted); + activity.Dispose(); + } + } + + private static void SetWorkflowTags(DiagnosticsActivity activity, WorkflowExecutionContext context, WorkflowInstanceSubStatus? workflowSubStatusOverride = null) + { + var workflow = context.Workflow; + var identity = workflow.Identity; + var workflowSubStatus = workflowSubStatusOverride ?? context.SubStatus; + var workflowStatus = GetWorkflowStatus(context, workflowSubStatusOverride); + + activity.SetTag(WorkflowSystem, SystemName); + activity.SetTag(WorkflowInstanceId, context.Id); + activity.SetTag(WorkflowDefinitionId, identity.DefinitionId); + activity.SetTag(WorkflowDefinitionVersion, identity.Version); + activity.SetTag(WorkflowDefinitionVersionId, identity.Id); + activity.SetTag(WorkflowStatus, workflowStatus.ToString()); + activity.SetTag(WorkflowSubStatus, workflowSubStatus.ToString()); + AddIfNotNull(activity, WorkflowName, workflow.WorkflowMetadata.Name ?? workflow.Name); + AddIfNotNull(activity, WorkflowParentInstanceId, context.ParentWorkflowInstanceId); + AddIfNotNull(activity, WorkflowCorrelationId, context.CorrelationId); + AddIfNotNull(activity, TenantId, identity.TenantId); + } + + private static void SetActivityTags(DiagnosticsActivity activity, ActivityExecutionContext context, WorkflowActivityStatus? activityStatusOverride = null) + { + var currentActivity = context.Activity; + var activityStatus = activityStatusOverride ?? context.Status; + + activity.SetTag(ActivityId, currentActivity.Id); + activity.SetTag(ActivityType, currentActivity.Type); + activity.SetTag(ActivityVersion, currentActivity.Version); + activity.SetTag(ActivityExecutionId, context.Id); + activity.SetTag(ActivityStatus, activityStatus.ToString()); + AddIfNotNull(activity, ActivityName, currentActivity.Name ?? context.ActivityDescriptor.DisplayName ?? context.ActivityDescriptor.Name); + AddIfNotNull(activity, ActivityParentExecutionId, context.ParentActivityExecutionContext?.Id); + AddIfNotNull(activity, ActivityScheduledByExecutionId, context.SchedulingActivityExecutionId); + } + + private static void SetActivityOutcome(DiagnosticsActivity activity, ActivityExecutionContext context) + { + if (!context.JournalData.TryGetValue("Outcomes", out var outcomes)) + return; + + var outcome = outcomes switch + { + IEnumerable names => string.Join(",", names), + _ => outcomes?.ToString() + }; + + AddIfNotNull(activity, ActivityOutcome, outcome); + } + + private static void SetError(DiagnosticsActivity activity, Exception? exception, bool faulted) + { + if (!faulted) + { + activity.SetStatus(ActivityStatusCode.Ok); + return; + } + + activity.SetStatus(ActivityStatusCode.Error, "Faulted"); + + if (exception != null) + { + var exceptionType = exception.GetType().FullName; + activity.SetTag(ExceptionType, exceptionType); + activity.AddEvent(new("exception", tags: new ActivityTagsCollection + { + { ExceptionType, exceptionType } + })); + } + } + + private static TagList CreateWorkflowTags(WorkflowExecutionContext context, bool includeExecutionStatus = true, WorkflowInstanceSubStatus? workflowSubStatusOverride = null) + { + var workflow = context.Workflow; + var identity = workflow.Identity; + var workflowSubStatus = workflowSubStatusOverride ?? context.SubStatus; + var workflowStatus = GetWorkflowStatus(context, workflowSubStatusOverride); + var tags = new TagList + { + { WorkflowSystem, SystemName }, + { WorkflowDefinitionId, identity.DefinitionId }, + { WorkflowDefinitionVersion, identity.Version } + }; + + if (includeExecutionStatus) + { + tags.Add(WorkflowStatus, workflowStatus.ToString()); + tags.Add(WorkflowSubStatus, workflowSubStatus.ToString()); + } + + AddIfNotNull(ref tags, WorkflowName, workflow.WorkflowMetadata.Name ?? workflow.Name); + AddIfNotNull(ref tags, TenantId, identity.TenantId); + return tags; + } + + private static TagList CreateActivityTags(ActivityExecutionContext context, WorkflowActivityStatus activityStatus, bool faulted) + { + var currentActivity = context.Activity; + var tags = new TagList + { + { WorkflowSystem, SystemName }, + { WorkflowDefinitionId, context.WorkflowExecutionContext.Workflow.Identity.DefinitionId }, + { ActivityType, currentActivity.Type }, + { ActivityVersion, currentActivity.Version }, + { ActivityStatus, faulted ? WorkflowActivityStatus.Faulted.ToString() : activityStatus.ToString() }, + { ActivityFaulted, faulted } + }; + + AddIfNotNull(ref tags, TenantId, context.WorkflowExecutionContext.Workflow.Identity.TenantId); + return tags; + } + + private static WorkflowInstanceStatus GetWorkflowStatus(WorkflowExecutionContext context, WorkflowInstanceSubStatus? workflowSubStatusOverride) + { + if (workflowSubStatusOverride == null) + return context.Status; + + return workflowSubStatusOverride.Value switch + { + WorkflowInstanceSubStatus.Cancelled or WorkflowInstanceSubStatus.Faulted or WorkflowInstanceSubStatus.Finished => WorkflowInstanceStatus.Finished, + _ => WorkflowInstanceStatus.Running + }; + } + + private static void AddIfNotNull(ref TagList tags, string key, object? value) + { + if (ShouldAddTag(value)) + tags.Add(key, value!); + } + + private static void AddIfNotNull(DiagnosticsActivity activity, string key, object? value) + { + if (ShouldAddTag(value)) + activity.SetTag(key, value); + } + + private static bool ShouldAddTag(object? value) + { + return value is not null && (value is not string stringValue || !string.IsNullOrWhiteSpace(stringValue)); + } +} + +internal readonly record struct WorkflowInstrumentationScope(DiagnosticsActivity? Activity); + +internal readonly record struct ActivityInstrumentationScope(DiagnosticsActivity? Activity, long StartTimestamp); diff --git a/test/unit/Elsa.Activities.UnitTests/Http/FlowSendHttpRequestTests.cs b/test/unit/Elsa.Activities.UnitTests/Http/FlowSendHttpRequestTests.cs index 5dfa4a3f4..e3306345f 100644 --- a/test/unit/Elsa.Activities.UnitTests/Http/FlowSendHttpRequestTests.cs +++ b/test/unit/Elsa.Activities.UnitTests/Http/FlowSendHttpRequestTests.cs @@ -57,6 +57,16 @@ public class FlowSendHttpRequestTests Assert.Equal(authorizationHeader, requestCapture.CapturedRequest.Headers.Authorization.ToString()); } + [Fact] + public async Task Should_Propagate_Current_Trace_Context() + { + await SendHttpRequestTestHelpers.AssertPropagatesCurrentTraceContextAsync(async (url, responseHandler) => + { + var flowSendHttpRequest = CreateFlowSendHttpRequest(url); + await ExecuteActivityAsync(flowSendHttpRequest, responseHandler); + }); + } + [Theory] [MemberData(nameof(StatusCodeOutcomeTestCases))] public async Task Should_Return_Outcome_Based_On_Status_Code( diff --git a/test/unit/Elsa.Activities.UnitTests/Http/Helpers/SendHttpRequestTestHelpers.cs b/test/unit/Elsa.Activities.UnitTests/Http/Helpers/SendHttpRequestTestHelpers.cs index f654201ac..a8453039d 100644 --- a/test/unit/Elsa.Activities.UnitTests/Http/Helpers/SendHttpRequestTestHelpers.cs +++ b/test/unit/Elsa.Activities.UnitTests/Http/Helpers/SendHttpRequestTestHelpers.cs @@ -1,4 +1,6 @@ +using System.Diagnostics; using System.Net; +using Xunit; namespace Elsa.Activities.UnitTests.Http.Helpers; @@ -7,6 +9,8 @@ namespace Elsa.Activities.UnitTests.Http.Helpers; /// public static class SendHttpRequestTestHelpers { + private const string TestActivitySourceName = "Elsa.Tests"; + /// /// Creates a response handler that returns a specific HTTP status code and optional content. /// @@ -33,6 +37,35 @@ public static class SendHttpRequestTestHelpers return (_, _) => throw ((TException)Activator.CreateInstance(typeof(TException), message)!); } + public static async Task AssertPropagatesCurrentTraceContextAsync( + Func>, Task> executeActivityAsync) + { + using var listener = new ActivityListener + { + ShouldListenTo = source => source.Name == TestActivitySourceName, + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded, + SampleUsingParentId = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded + }; + ActivitySource.AddActivityListener(listener); + + using var source = new ActivitySource(TestActivitySourceName); + using var parentActivity = source.StartActivity("parent"); + Assert.NotNull(parentActivity); + + var requestCapture = new RequestCapture(); + var responseHandler = CreateResponseHandler(HttpStatusCode.OK, "{}", requestCapture); + + await executeActivityAsync(new Uri("https://api.example.com/traced"), responseHandler); + + Assert.NotNull(requestCapture.CapturedRequest); + Assert.True(requestCapture.CapturedRequest.Headers.TryGetValues("traceparent", out var traceParents)); + var traceParent = Assert.Single(traceParents); + var traceParentParts = traceParent.Split('-'); + Assert.Equal(4, traceParentParts.Length); + Assert.Equal(2, traceParentParts[0].Length); + Assert.Equal(parentActivity.TraceId.ToString(), traceParentParts[1]); + } + /// /// Captures HTTP request details during test execution. /// diff --git a/test/unit/Elsa.Activities.UnitTests/Http/SendHttpRequestTests.cs b/test/unit/Elsa.Activities.UnitTests/Http/SendHttpRequestTests.cs index 461b29c83..61940ddc8 100644 --- a/test/unit/Elsa.Activities.UnitTests/Http/SendHttpRequestTests.cs +++ b/test/unit/Elsa.Activities.UnitTests/Http/SendHttpRequestTests.cs @@ -58,6 +58,16 @@ public class SendHttpRequestTests Assert.Equal(authorizationHeader, requestCapture.CapturedRequest.Headers.Authorization.ToString()); } + [Fact] + public async Task Should_Propagate_Current_Trace_Context() + { + await SendHttpRequestTestHelpers.AssertPropagatesCurrentTraceContextAsync(async (url, responseHandler) => + { + var sendHttpRequest = CreateSendHttpRequest(url); + await ExecuteActivityAsync(sendHttpRequest, responseHandler); + }); + } + [Theory] [InlineData(new[]{200, 404}, new[]{"mockActivity200", "mockActivity404"}, "mockUnmatchedActivity", HttpStatusCode.NotFound, "mockActivity404")] [InlineData(new[]{200, 404}, new[]{"mockActivity200", "mockActivity404"}, "mockUnmatchedActivity", HttpStatusCode.InternalServerError, "mockUnmatchedActivity")] @@ -166,7 +176,7 @@ public class SendHttpRequestTests expectedCategory: "HTTP", expectedDisplayName: "HTTP Request", expectedDescription: "Send an HTTP request.", - expectedKind: ActivityKind.Task + expectedKind: Elsa.Workflows.ActivityKind.Task ); } @@ -181,6 +191,7 @@ public class SendHttpRequestTests { if (requestCapture != null) requestCapture.CapturedRequest = request; + return Task.FromResult(ActivityTestFixtureHttpExtensions.CreateHttpResponse(statusCode, content, additionalHeaders)); }; } @@ -285,4 +296,4 @@ public class SendHttpRequestTests { return (_, _) => throw ((TException)Activator.CreateInstance(typeof(TException), message)!); } -} \ No newline at end of file +} diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/Telemetry/WorkflowInstrumentationTests.cs b/test/unit/Elsa.Workflows.Core.UnitTests/Telemetry/WorkflowInstrumentationTests.cs new file mode 100644 index 000000000..840e170d0 --- /dev/null +++ b/test/unit/Elsa.Workflows.Core.UnitTests/Telemetry/WorkflowInstrumentationTests.cs @@ -0,0 +1,691 @@ +using System.Collections.Concurrent; +using System.Diagnostics; +using System.Diagnostics.Metrics; +using System.Linq; +using Elsa.Common; +using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Testing.Shared; +using Elsa.Workflows; +using Elsa.Workflows.CommitStates; +using Elsa.Workflows.Middleware.Workflows; +using Elsa.Workflows.Models; +using Elsa.Workflows.Notifications; +using Elsa.Workflows.Options; +using Elsa.Workflows.Pipelines.ActivityExecution; +using Elsa.Workflows.Pipelines.WorkflowExecution; +using Elsa.Workflows.Services; +using Elsa.Workflows.State; +using Elsa.Workflows.Telemetry; +using Microsoft.Extensions.Logging.Abstractions; +using NSubstitute; + +namespace Elsa.Workflows.Core.UnitTests.Telemetry; + +using DiagnosticsActivity = System.Diagnostics.Activity; + +[Collection(nameof(WorkflowInstrumentationTestCollection))] +public class WorkflowInstrumentationTests +{ + [Fact] + public async Task ActivityInvoker_Should_Emit_Activity_Span_And_Duration_Metric() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activity = new TestActivity { Name = "Test activity" }; + var context = await new ActivityTestFixture(activity).BuildAsync(); + var invoker = new ActivityInvoker(new CompletingActivityExecutionPipeline(), new ActivityLoggerStateGenerator(), NullLogger.Instance); + + await invoker.InvokeAsync(context); + + var span = GetStoppedActivity(activityCapture, "activity.execute", WorkflowInstrumentation.ActivityExecutionId, context.Id); + Assert.Equal("activity.execute", span.OperationName); + Assert.Equal(activity.Type, GetTag(span.TagObjects, WorkflowInstrumentation.ActivityType)); + Assert.Equal(context.WorkflowExecutionContext.Workflow.Identity.Id, GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowDefinitionVersionId)); + Assert.Equal(ActivityStatus.Completed.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.ActivityStatus)); + var activityDuration = GetActivityDuration(meterCapture, context); + Assert.False(activityDuration.Tags.ContainsKey(WorkflowInstrumentation.WorkflowDefinitionVersionId)); + Assert.False(activityDuration.Tags.ContainsKey(WorkflowInstrumentation.ActivityName)); + Assert.Equal(false, activityDuration.Tags[WorkflowInstrumentation.ActivityFaulted]); + } + + [Fact] + public async Task ActivityInvoker_Should_Not_Record_Faulted_Metric_When_Pipeline_Cancels() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var context = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var invoker = new ActivityInvoker(new CancellingActivityExecutionPipeline(), new ActivityLoggerStateGenerator(), NullLogger.Instance); + + await Assert.ThrowsAsync(() => invoker.InvokeAsync(context)); + + var span = GetStoppedActivity(activityCapture, "activity.execute", WorkflowInstrumentation.ActivityExecutionId, context.Id); + Assert.Equal(ActivityStatusCode.Ok, span.Status); + var activityDuration = GetActivityDuration(meterCapture, context); + Assert.Equal(false, activityDuration.Tags[WorkflowInstrumentation.ActivityFaulted]); + Assert.Equal(ActivityStatus.Canceled.ToString(), activityDuration.Tags[WorkflowInstrumentation.ActivityStatus]); + } + + [Fact] + public async Task ActivityInvoker_Should_Record_Canceled_Status_When_Pipeline_Cancels_Before_Status_Transition() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var context = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var invoker = new ActivityInvoker(new NonMutatingCancellingActivityExecutionPipeline(), new ActivityLoggerStateGenerator(), NullLogger.Instance); + + await Assert.ThrowsAsync(() => invoker.InvokeAsync(context)); + + var span = GetStoppedActivity(activityCapture, "activity.execute", WorkflowInstrumentation.ActivityExecutionId, context.Id); + Assert.Equal(ActivityStatusCode.Ok, span.Status); + Assert.Equal(ActivityStatus.Canceled.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.ActivityStatus)); + Assert.Equal(false, GetTag(span.TagObjects, WorkflowInstrumentation.ActivityFaulted)); + var activityDuration = GetActivityDuration(meterCapture, context); + Assert.Equal(ActivityStatus.Canceled.ToString(), activityDuration.Tags[WorkflowInstrumentation.ActivityStatus]); + Assert.Equal(false, activityDuration.Tags[WorkflowInstrumentation.ActivityFaulted]); + } + + [Fact] + public async Task WorkflowRunner_Should_Record_Completed_Metric_When_Pipeline_Finishes() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + context.Workflow.WorkflowMetadata = new("Test workflow"); + context.ParentWorkflowInstanceId = "parent-instance-id"; + context.CorrelationId = ""; + context.Workflow.Identity = context.Workflow.Identity with { TenantId = " " }; + var runner = CreateWorkflowRunner(context, new CompletingWorkflowExecutionPipeline()); + + await runner.RunAsync(context); + + var span = GetStoppedActivity(activityCapture, "workflow.execute", WorkflowInstrumentation.WorkflowInstanceId, context.Id); + Assert.Equal(ActivityStatusCode.Ok, span.Status); + Assert.Equal(context.Workflow.Identity.Id, GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowDefinitionVersionId)); + var workflowSubStatus = GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowSubStatus); + Assert.Equal(WorkflowSubStatus.Finished.ToString(), workflowSubStatus); + Assert.Equal("parent-instance-id", GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowParentInstanceId)); + Assert.False(HasTag(span.TagObjects, WorkflowInstrumentation.WorkflowCorrelationId)); + Assert.False(HasTag(span.TagObjects, WorkflowInstrumentation.TenantId)); + Assert.False(HasTag(span.TagObjects, "workflow.parent_instance.id")); + var started = GetWorkflowMeasurement(meterCapture, "elsa.workflow.started", context); + var completed = GetWorkflowMeasurement(meterCapture, "elsa.workflow.completed", context); + Assert.False(started.Tags.ContainsKey(WorkflowInstrumentation.WorkflowDefinitionVersionId)); + Assert.False(started.Tags.ContainsKey(WorkflowInstrumentation.TenantId)); + Assert.False(completed.Tags.ContainsKey(WorkflowInstrumentation.WorkflowDefinitionVersionId)); + Assert.False(started.Tags.ContainsKey(WorkflowInstrumentation.WorkflowSubStatus)); + Assert.Equal(WorkflowSubStatus.Finished.ToString(), completed.Tags[WorkflowInstrumentation.WorkflowSubStatus]); + } + + [Fact] + public async Task WorkflowRunner_Should_Record_Faulted_Span_And_Metric_When_Pipeline_Throws() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + context.Workflow.WorkflowMetadata = new("Test workflow"); + var runner = CreateWorkflowRunner(context, new ThrowingWorkflowExecutionPipeline()); + + await Assert.ThrowsAsync(() => runner.RunAsync(context)); + + var span = GetStoppedActivity(activityCapture, "workflow.execute", WorkflowInstrumentation.WorkflowInstanceId, context.Id); + Assert.Equal("workflow.execute", span.OperationName); + Assert.Equal(ActivityStatusCode.Error, span.Status); + Assert.Equal(WorkflowSubStatus.Executing, context.SubStatus); + Assert.Equal(context.Id, GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowInstanceId)); + Assert.Equal(WorkflowStatus.Finished.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowStatus)); + Assert.Equal(WorkflowSubStatus.Faulted.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowSubStatus)); + Assert.Equal(typeof(InvalidOperationException).FullName, GetTag(span.TagObjects, WorkflowInstrumentation.ExceptionType)); + Assert.False(HasTag(span.TagObjects, "exception.message")); + Assert.False(HasTag(span.TagObjects, "exception.stacktrace")); + var exceptionEvent = Assert.Single(span.Events, x => x.Name == "exception"); + Assert.Equal(typeof(InvalidOperationException).FullName, GetTag(exceptionEvent.Tags, WorkflowInstrumentation.ExceptionType)); + Assert.False(HasTag(exceptionEvent.Tags, "exception.message")); + Assert.False(HasTag(exceptionEvent.Tags, "exception.stacktrace")); + var started = GetWorkflowMeasurement(meterCapture, "elsa.workflow.started", context); + var faulted = GetWorkflowMeasurement(meterCapture, "elsa.workflow.faulted", context); + Assert.False(started.Tags.ContainsKey(WorkflowInstrumentation.WorkflowDefinitionVersionId)); + Assert.False(faulted.Tags.ContainsKey(WorkflowInstrumentation.WorkflowDefinitionVersionId)); + Assert.Equal("Test workflow", started.Tags[WorkflowInstrumentation.WorkflowName]); + Assert.Equal("Test workflow", faulted.Tags[WorkflowInstrumentation.WorkflowName]); + Assert.Equal(WorkflowStatus.Finished.ToString(), faulted.Tags[WorkflowInstrumentation.WorkflowStatus]); + Assert.Equal(WorkflowSubStatus.Faulted.ToString(), faulted.Tags[WorkflowInstrumentation.WorkflowSubStatus]); + } + + [Fact] + public async Task WorkflowRunner_Should_Not_Record_Faulted_Span_Or_Metric_When_Pipeline_Cancels() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + var runner = CreateWorkflowRunner(context, new CancellingWorkflowExecutionPipeline()); + + await Assert.ThrowsAsync(() => runner.RunAsync(context)); + + var span = GetStoppedActivity(activityCapture, "workflow.execute", WorkflowInstrumentation.WorkflowInstanceId, context.Id); + Assert.Equal(ActivityStatusCode.Ok, span.Status); + Assert.Equal(WorkflowSubStatus.Cancelled.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowSubStatus)); + Assert.Equal(false, GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowFaulted)); + Assert.DoesNotContain(meterCapture.LongMeasurements, x => IsWorkflowMeasurement(x, "elsa.workflow.faulted", context)); + } + + [Fact] + public async Task WorkflowRunner_Should_Record_Cancelled_SubStatus_When_Pipeline_Cancels_Before_Status_Transition() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + var runner = CreateWorkflowRunner(context, new NonMutatingCancellingWorkflowExecutionPipeline()); + + await Assert.ThrowsAsync(() => runner.RunAsync(context)); + + var span = GetStoppedActivity(activityCapture, "workflow.execute", WorkflowInstrumentation.WorkflowInstanceId, context.Id); + Assert.Equal(ActivityStatusCode.Ok, span.Status); + Assert.Equal(WorkflowStatus.Finished.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowStatus)); + Assert.Equal(WorkflowSubStatus.Cancelled.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowSubStatus)); + Assert.Equal(false, GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowFaulted)); + Assert.DoesNotContain(meterCapture.LongMeasurements, x => IsWorkflowMeasurement(x, "elsa.workflow.faulted", context)); + } + + [Fact] + public async Task WorkflowRunner_Should_Record_Cancel_When_ExceptionHandlingMiddleware_Catches_Cancellation() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + var runner = CreateWorkflowRunner(context, new ExceptionHandlingCancellingWorkflowExecutionPipeline()); + + await Assert.ThrowsAsync(() => runner.RunAsync(context)); + + var span = GetStoppedActivity(activityCapture, "workflow.execute", WorkflowInstrumentation.WorkflowInstanceId, context.Id); + Assert.Equal(ActivityStatusCode.Ok, span.Status); + Assert.Equal(WorkflowSubStatus.Cancelled.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowSubStatus)); + Assert.Equal(false, GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowFaulted)); + Assert.DoesNotContain(meterCapture.LongMeasurements, x => IsWorkflowMeasurement(x, "elsa.workflow.faulted", context)); + } + + [Fact] + public async Task WorkflowRunner_Should_Not_Start_When_WorkflowExecuting_Handler_Changes_SubStatus() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + var notificationSender = Substitute.For(); + notificationSender + .SendAsync(Arg.Any(), Arg.Any()) + .Returns(callInfo => + { + if (callInfo.Arg() is WorkflowExecuting) + context.Cancel(); + + return Task.CompletedTask; + }); + var runner = CreateWorkflowRunner(context, new NoopWorkflowExecutionPipeline(), notificationSender); + + await runner.RunAsync(context); + + Assert.Equal(WorkflowSubStatus.Cancelled, context.SubStatus); + Assert.DoesNotContain(meterCapture.LongMeasurements, x => IsWorkflowMeasurement(x, "elsa.workflow.started", context)); + await notificationSender + .DidNotReceive() + .SendAsync(Arg.Is(notification => notification is WorkflowStarted), Arg.Any()); + } + + [Fact] + public async Task ActivityInvoker_Should_Not_Replace_Pipeline_Exception_When_Outcome_Is_Null() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var context = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var invoker = new ActivityInvoker(new ThrowingActivityExecutionPipeline(), new ActivityLoggerStateGenerator(), NullLogger.Instance); + + var exception = await Assert.ThrowsAsync(() => invoker.InvokeAsync(context)); + + Assert.Equal("Pipeline failed", exception.Message); + var span = GetStoppedActivity(activityCapture, "activity.execute", WorkflowInstrumentation.ActivityExecutionId, context.Id); + Assert.Equal(ActivityStatus.Faulted.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.ActivityStatus)); + var activityDuration = GetActivityDuration(meterCapture, context); + Assert.Equal(true, activityDuration.Tags[WorkflowInstrumentation.ActivityFaulted]); + Assert.Equal(ActivityStatus.Faulted.ToString(), activityDuration.Tags[WorkflowInstrumentation.ActivityStatus]); + } + + [Fact] + public async Task ActivityInvoker_Should_Record_Fault_When_Cancelled_Context_Throws_NonCancellation_Exception() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var context = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var invoker = new ActivityInvoker(new CancelledThenThrowingActivityExecutionPipeline(), new ActivityLoggerStateGenerator(), NullLogger.Instance); + + await Assert.ThrowsAsync(() => invoker.InvokeAsync(context)); + + var span = GetStoppedActivity(activityCapture, "activity.execute", WorkflowInstrumentation.ActivityExecutionId, context.Id); + Assert.Equal(ActivityStatusCode.Error, span.Status); + Assert.Equal(ActivityStatus.Faulted.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.ActivityStatus)); + Assert.Equal(true, GetTag(span.TagObjects, WorkflowInstrumentation.ActivityFaulted)); + var activityDuration = GetActivityDuration(meterCapture, context); + Assert.Equal(ActivityStatus.Faulted.ToString(), activityDuration.Tags[WorkflowInstrumentation.ActivityStatus]); + Assert.Equal(true, activityDuration.Tags[WorkflowInstrumentation.ActivityFaulted]); + } + + [Fact] + public async Task ActivityInvoker_Should_Not_Emit_ExceptionType_When_Faulted_Without_Exception() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var context = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var invoker = new ActivityInvoker(new FaultingActivityExecutionPipeline(), new ActivityLoggerStateGenerator(), NullLogger.Instance); + + await invoker.InvokeAsync(context); + + var span = GetStoppedActivity(activityCapture, "activity.execute", WorkflowInstrumentation.ActivityExecutionId, context.Id); + Assert.Equal(ActivityStatusCode.Error, span.Status); + Assert.Equal(ActivityStatus.Faulted.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.ActivityStatus)); + Assert.False(HasTag(span.TagObjects, WorkflowInstrumentation.ExceptionType)); + } + + [Fact] + public async Task WorkflowRunner_Should_Use_Context_Exception_When_Faulted_Without_Thrown_Exception() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + var exception = new InvalidOperationException("Faulted by middleware"); + var runner = CreateWorkflowRunner(context, new FaultingWorkflowExecutionPipeline(exception)); + + await runner.RunAsync(context); + + var span = GetStoppedActivity(activityCapture, "workflow.execute", WorkflowInstrumentation.WorkflowInstanceId, context.Id); + Assert.Equal(ActivityStatusCode.Error, span.Status); + Assert.Equal(typeof(InvalidOperationException).FullName, GetTag(span.TagObjects, WorkflowInstrumentation.ExceptionType)); + _ = GetWorkflowMeasurement(meterCapture, "elsa.workflow.faulted", context); + } + + [Fact] + public async Task WorkflowRunner_Should_Record_Fault_When_Cancelled_Context_Throws_NonCancellation_Exception() + { + using var activityCapture = new ActivityCapture(); + using var meterCapture = new MeterCapture(); + var activityExecutionContext = await new ActivityTestFixture(new TestActivity()).BuildAsync(); + var context = activityExecutionContext.WorkflowExecutionContext; + var runner = CreateWorkflowRunner(context, new CancelledThenThrowingWorkflowExecutionPipeline()); + + await Assert.ThrowsAsync(() => runner.RunAsync(context)); + + var span = GetStoppedActivity(activityCapture, "workflow.execute", WorkflowInstrumentation.WorkflowInstanceId, context.Id); + Assert.Equal(ActivityStatusCode.Error, span.Status); + Assert.Equal(WorkflowStatus.Finished.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowStatus)); + Assert.Equal(WorkflowSubStatus.Faulted.ToString(), GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowSubStatus)); + Assert.Equal(true, GetTag(span.TagObjects, WorkflowInstrumentation.WorkflowFaulted)); + var faulted = GetWorkflowMeasurement(meterCapture, "elsa.workflow.faulted", context); + Assert.Equal(WorkflowSubStatus.Faulted.ToString(), faulted.Tags[WorkflowInstrumentation.WorkflowSubStatus]); + } + + private static WorkflowRunner CreateWorkflowRunner(WorkflowExecutionContext context, IWorkflowExecutionPipeline pipeline, INotificationSender? notificationSender = null) + { + if (notificationSender == null) + { + notificationSender = Substitute.For(); + notificationSender + .SendAsync(Arg.Any(), Arg.Any()) + .Returns(Task.CompletedTask); + } + + var workflowStateExtractor = Substitute.For(); + workflowStateExtractor.Extract(context).Returns(_ => new WorkflowState + { + Id = context.Id, + DefinitionId = context.Workflow.Identity.DefinitionId, + DefinitionVersionId = context.Workflow.Identity.Id, + DefinitionVersion = context.Workflow.Identity.Version, + Status = context.Status, + SubStatus = context.SubStatus + }); + + return new( + context.ServiceProvider, + pipeline, + workflowStateExtractor, + Substitute.For(), + Substitute.For(), + Substitute.For(), + notificationSender, + new WorkflowLoggerStateGenerator(), + Substitute.For(), + NullLogger.Instance); + } + + private static DiagnosticsActivity GetStoppedActivity(ActivityCapture capture, string operationName, string tagKey, object? tagValue) + { + return Assert.Single(capture.StoppedActivities, activity => + activity.OperationName == operationName && + activity.TagObjects.Any(tag => tag.Key == tagKey && Equals(tag.Value, tagValue))); + } + + private static CapturedMeasurement GetActivityDuration(MeterCapture capture, ActivityExecutionContext context) + { + return Assert.Single(capture.DoubleMeasurements, measurement => + measurement.InstrumentName == "elsa.activity.duration" && + measurement.Value >= 0 && + HasTag(measurement.Tags, WorkflowInstrumentation.WorkflowDefinitionId, context.WorkflowExecutionContext.Workflow.Identity.DefinitionId) && + HasTag(measurement.Tags, WorkflowInstrumentation.ActivityType, context.Activity.Type)); + } + + private static CapturedMeasurement GetWorkflowMeasurement(MeterCapture capture, string instrumentName, WorkflowExecutionContext context) + { + return Assert.Single(capture.LongMeasurements, measurement => IsWorkflowMeasurement(measurement, instrumentName, context)); + } + + private static bool IsWorkflowMeasurement(CapturedMeasurement measurement, string instrumentName, WorkflowExecutionContext context) + { + return measurement.InstrumentName == instrumentName && + measurement.Value == 1 && + HasTag(measurement.Tags, WorkflowInstrumentation.WorkflowDefinitionId, context.Workflow.Identity.DefinitionId); + } + + private static bool HasTag(IReadOnlyDictionary tags, string key, object? value) + { + return tags.TryGetValue(key, out var tagValue) && Equals(tagValue, value); + } + + private static object? GetTag(IEnumerable> tags, string key) + { + object? value = null; + var found = false; + + foreach (var tag in tags.Where(tag => tag.Key == key)) + { + value = tag.Value; + found = true; + } + + Assert.True(found, $"Expected tag '{key}' to be present."); + return value; + } + + private static bool HasTag(IEnumerable> tags, string key) + { + return tags.Any(tag => tag.Key == key); + } + + private sealed class CompletingActivityExecutionPipeline : IActivityExecutionPipeline + { + public ActivityMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public ActivityMiddlewareDelegate Setup(Action setup) => Pipeline; + + public async Task ExecuteAsync(ActivityExecutionContext context) + { + context.TransitionTo(ActivityStatus.Running); + await context.CompleteActivityAsync(); + } + } + + private sealed class ThrowingActivityExecutionPipeline : IActivityExecutionPipeline + { + public ActivityMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public ActivityMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(ActivityExecutionContext context) + { + context.JournalData["Outcomes"] = null!; + throw new InvalidOperationException("Pipeline failed"); + } + } + + private sealed class CancellingActivityExecutionPipeline : IActivityExecutionPipeline + { + public ActivityMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public ActivityMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(ActivityExecutionContext context) + { + context.TransitionTo(ActivityStatus.Canceled); + throw new OperationCanceledException(); + } + } + + private sealed class NonMutatingCancellingActivityExecutionPipeline : IActivityExecutionPipeline + { + public ActivityMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public ActivityMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(ActivityExecutionContext context) => throw new OperationCanceledException(); + } + + private sealed class CancelledThenThrowingActivityExecutionPipeline : IActivityExecutionPipeline + { + public ActivityMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public ActivityMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(ActivityExecutionContext context) + { + context.TransitionTo(ActivityStatus.Canceled); + throw new InvalidOperationException("Pipeline failed after cancellation."); + } + } + + private sealed class CompletingWorkflowExecutionPipeline : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public async Task ExecuteAsync(WorkflowExecutionContext context) + { + context.ScheduleWorkflow(); + + var middleware = new DefaultActivitySchedulerMiddleware( + _ => ValueTask.CompletedTask, + new CompletingActivityInvoker(), + Substitute.For(), + Microsoft.Extensions.Options.Options.Create(new CommitStateOptions())); + + await middleware.InvokeAsync(context); + } + } + + private sealed class CompletingActivityInvoker : IActivityInvoker + { + public async Task InvokeAsync(WorkflowExecutionContext workflowExecutionContext, IActivity activity, ActivityInvocationOptions? options = null) + { + var activityExecutionContext = options?.ExistingActivityExecutionContext ?? await workflowExecutionContext.CreateActivityExecutionContextAsync(activity, options); + + if (!workflowExecutionContext.ActivityExecutionContexts.Any(x => x.Id == activityExecutionContext.Id)) + workflowExecutionContext.AddActivityExecutionContext(activityExecutionContext); + + await InvokeAsync(activityExecutionContext); + return activityExecutionContext; + } + + public async Task InvokeAsync(ActivityExecutionContext activityExecutionContext) + { + activityExecutionContext.TransitionTo(ActivityStatus.Running); + await activityExecutionContext.CompleteActivityAsync(); + } + } + + private sealed class ThrowingWorkflowExecutionPipeline : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(WorkflowExecutionContext context) => throw new InvalidOperationException("Pipeline failed"); + } + + private sealed class CancellingWorkflowExecutionPipeline : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(WorkflowExecutionContext context) + { + context.Cancel(); + throw new OperationCanceledException(); + } + } + + private sealed class NonMutatingCancellingWorkflowExecutionPipeline : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(WorkflowExecutionContext context) => throw new OperationCanceledException(); + } + + private sealed class CancelledThenThrowingWorkflowExecutionPipeline : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(WorkflowExecutionContext context) + { + context.Cancel(); + throw new InvalidOperationException("Pipeline failed after cancellation."); + } + } + + private sealed class FaultingWorkflowExecutionPipeline(Exception exception) : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public async Task ExecuteAsync(WorkflowExecutionContext context) + { + var middleware = new ExceptionHandlingMiddleware( + _ => throw exception, + context.SystemClock, + NullLogger.Instance); + + await middleware.InvokeAsync(context); + } + } + + private sealed class ExceptionHandlingCancellingWorkflowExecutionPipeline : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public async Task ExecuteAsync(WorkflowExecutionContext context) + { + var middleware = new ExceptionHandlingMiddleware( + _ => throw new OperationCanceledException(), + context.SystemClock, + NullLogger.Instance); + + await middleware.InvokeAsync(context); + } + } + + private sealed class NoopWorkflowExecutionPipeline : IWorkflowExecutionPipeline + { + public Action ConfigurePipelineBuilder => _ => { }; + public WorkflowMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public WorkflowMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(WorkflowExecutionContext context) => Task.CompletedTask; + } + + private sealed class TestActivity : Elsa.Workflows.Activity + { + } + + private sealed class ActivityCapture : IDisposable + { + private readonly ActivityListener _listener; + + public ActivityCapture() + { + _listener = new() + { + ShouldListenTo = source => source.Name == WorkflowInstrumentation.ActivitySourceName, + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded, + ActivityStopped = activity => StoppedActivities.Enqueue(activity) + }; + + ActivitySource.AddActivityListener(_listener); + } + + public ConcurrentQueue StoppedActivities { get; } = new(); + + public void Dispose() => _listener.Dispose(); + } + + private sealed class MeterCapture : IDisposable + { + private readonly MeterListener _listener = new(); + + public MeterCapture() + { + _listener.InstrumentPublished = (instrument, listener) => + { + if (instrument.Meter.Name == WorkflowInstrumentation.MeterName) + listener.EnableMeasurementEvents(instrument); + }; + + _listener.SetMeasurementEventCallback((instrument, measurement, tags, _) => LongMeasurements.Enqueue(new(instrument.Name, measurement, CaptureTags(tags)))); + _listener.SetMeasurementEventCallback((instrument, measurement, tags, _) => DoubleMeasurements.Enqueue(new(instrument.Name, measurement, CaptureTags(tags)))); + _listener.Start(); + } + + public ConcurrentQueue> LongMeasurements { get; } = new(); + public ConcurrentQueue> DoubleMeasurements { get; } = new(); + + public void Dispose() => _listener.Dispose(); + + private static IReadOnlyDictionary CaptureTags(ReadOnlySpan> tags) + { + var result = new Dictionary(); + + foreach (var tag in tags) + result[tag.Key] = tag.Value; + + return result; + } + } + + private sealed class FaultingActivityExecutionPipeline : IActivityExecutionPipeline + { + public ActivityMiddlewareDelegate Pipeline => _ => ValueTask.CompletedTask; + + public ActivityMiddlewareDelegate Setup(Action setup) => Pipeline; + + public Task ExecuteAsync(ActivityExecutionContext context) + { + context.TransitionTo(ActivityStatus.Faulted); + return Task.CompletedTask; + } + } + + private readonly record struct CapturedMeasurement(string InstrumentName, T Value, IReadOnlyDictionary Tags); +} + +[CollectionDefinition(nameof(WorkflowInstrumentationTestCollection), DisableParallelization = true)] +public sealed class WorkflowInstrumentationTestCollection +{ +}