From 17bb1d08c948d991e2fe677e61cd9fb0b5cab2ba Mon Sep 17 00:00:00 2001 From: MariusVuscanNx <96233009+MariusVuscanNx@users.noreply.github.com> Date: Thu, 4 Sep 2025 09:43:37 +0300 Subject: [PATCH] Otel improvements (#6890) * Otel improvements * Update src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs Co-authored-by: Sipke Schoorstra * Update src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs Co-authored-by: Sipke Schoorstra * Fixed comments --------- Co-authored-by: Sipke Schoorstra --- ...metryTracingActivityExecutionMiddleware.cs | 46 +++++++++++++++++-- ...metryTracingWorkflowExecutionMiddleware.cs | 12 +++-- .../Services/ResilientActivityInvoker.cs | 32 +++++++------ .../ActivityExecutionContextExtensions.cs | 25 ++++++++++ 4 files changed, 95 insertions(+), 20 deletions(-) diff --git a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs index 6143cd669..83fbfe75b 100644 --- a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs +++ b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs @@ -1,4 +1,5 @@ using System.Diagnostics; +using System.Text; using Elsa.Common; using Elsa.Extensions; using Elsa.OpenTelemetry.Contracts; @@ -6,6 +7,7 @@ using Elsa.OpenTelemetry.Helpers; using Elsa.OpenTelemetry.Models; using Elsa.Workflows; using Elsa.Workflows.Pipelines.ActivityExecution; +using Humanizer; using JetBrains.Annotations; using Activity = System.Diagnostics.Activity; using ActivityKind = System.Diagnostics.ActivityKind; @@ -36,9 +38,9 @@ public class OpenTelemetryTracingActivityExecutionMiddleware(ActivityMiddlewareD span.SetTag("activity.version", activity.Version); span.SetTag("activity.instance.id", context.Id); span.SetTag("activity.tenant.id", context.WorkflowExecutionContext.Workflow.Identity.TenantId); - + var activityKind = context.ActivityDescriptor.Kind; - if (activityKind == Elsa.Workflows.ActivityKind.Job || (activityKind == Workflows.ActivityKind.Task && activity.GetRunAsynchronously())) + if (activityKind == Elsa.Workflows.ActivityKind.Job || (activityKind == Workflows.ActivityKind.Task && activity.GetRunAsynchronously())) span.SetTag("span.type", "job"); span.AddEvent(new("executing")); @@ -54,7 +56,7 @@ public class OpenTelemetryTracingActivityExecutionMiddleware(ActivityMiddlewareD var errorSpanHandler = context.GetServices() .OrderBy(x => x.Order) .FirstOrDefault(x => x.CanHandle(errorSpanHandlerContext)); - + errorSpanHandler?.Handle(errorSpanHandlerContext); } else if (context.Status == ActivityStatus.Canceled) @@ -72,6 +74,44 @@ public class OpenTelemetryTracingActivityExecutionMiddleware(ActivityMiddlewareD if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId)) span.SetTag("workflow.correlation_id", context.WorkflowExecutionContext.CorrelationId); + + span.SetTag("workflow.instance.id", context.WorkflowExecutionContext.Id); + span.SetTag("workflow.definition.id", context.WorkflowExecutionContext.Workflow.Identity.DefinitionId); + span.SetTag("workflow.definition.version", context.WorkflowExecutionContext.Workflow.Identity.Version); + + SetExtensions(context, span); + } + + private void SetExtensions(ActivityExecutionContext context, Activity span) + { + var extensions = context.GetExtensionsMetadata(); + if (extensions is not null) + { + foreach (var key in extensions.Keys) + { + span.SetTag($"activity.extensions.{ToOtelFormat(key)}", extensions[key]); + } + } + } + + public string ToOtelFormat(string input) + { + if (string.IsNullOrWhiteSpace(input)) + return string.Empty; + + var humanized = input.Humanize(); + + var sb = new StringBuilder(humanized.Length); + + foreach (var c in humanized) + { + if (char.IsWhiteSpace(c)) + sb.Append('_'); + else + sb.Append(char.ToLowerInvariant(c)); + } + + return sb.ToString(); } } diff --git a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs index 5d3996ea2..937255a7e 100644 --- a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs +++ b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs @@ -96,10 +96,16 @@ public class OpenTelemetryTracingWorkflowExecutionMiddleware(WorkflowMiddlewareD if (context.Incidents.Any()) { - span.SetTag("workflow.incidents.count", context.Incidents.Count); - + var incidentTagsList = new List(); foreach (var incident in context.Incidents) - span.AddEvent(new("incident", incident.Timestamp, CreateIncidentTags(incident))); + { + var incidentTags = CreateIncidentTags(incident); + incidentTagsList.Add(incidentTags); + span.AddEvent(new("incident", incident.Timestamp, incidentTags)); + } + + span.SetTag("workflow.incidents.items", incidentTagsList); + span.SetTag("workflow.incidents.count", context.Incidents.Count); } if (!string.IsNullOrWhiteSpace(context.CorrelationId)) diff --git a/src/modules/Elsa.Resilience.Core/Services/ResilientActivityInvoker.cs b/src/modules/Elsa.Resilience.Core/Services/ResilientActivityInvoker.cs index 8bf2c250f..bc86d879e 100644 --- a/src/modules/Elsa.Resilience.Core/Services/ResilientActivityInvoker.cs +++ b/src/modules/Elsa.Resilience.Core/Services/ResilientActivityInvoker.cs @@ -1,5 +1,6 @@ using System.Text.Json; using Elsa.Expressions.Helpers; +using Elsa.Extensions; using Elsa.Resilience.Entities; using Elsa.Resilience.Extensions; using Elsa.Resilience.Models; @@ -12,55 +13,56 @@ using Polly.Telemetry; namespace Elsa.Resilience; public class ResilientActivityInvoker( - IResilienceStrategyConfigEvaluator resilienceStrategyConfigEvaluator, - IRetryAttemptRecorder retryAttemptRecorder, - IIdentityGenerator identityGenerator, + IResilienceStrategyConfigEvaluator resilienceStrategyConfigEvaluator, + IRetryAttemptRecorder retryAttemptRecorder, + IIdentityGenerator identityGenerator, ResilienceStrategySerializer resilienceStrategySerializer) : IResilientActivityInvoker { private const string ResilienceStrategyIdPropKey = "resilienceStrategy"; + private const string RetryAttemptsCountKey = "RetryAttemptsCount"; public async Task InvokeAsync(IResilientActivity activity, ActivityExecutionContext context, Func> action, CancellationToken cancellationToken = default) { // Get the resilience strategy. var strategyConfig = GetStrategyConfig(activity); var resilienceStrategy = await resilienceStrategyConfigEvaluator.EvaluateAsync(strategyConfig, context.ExpressionExecutionContext, cancellationToken); - + // If no resilience strategy is configured, execute the action as-is. if (resilienceStrategy == null) return await action(); - + // Record the applied strategy as part of the activity execution context for diagnostics. var resilienceStrategyModel = JsonSerializer.SerializeToNode(resilienceStrategy, resilienceStrategySerializer.SerializerOptions)!; context.SetResilienceStrategy(resilienceStrategyModel); - + // Create a resilience pipeline builder. var builder = CreateResiliencePipelineBuilder(); var retries = new List(); context.TransientProperties[RetryAttempt.RetriesKey] = retries; - + // Create a resilience context. var resilienceContext = ResilienceContextPool.Shared.Get(cancellationToken); resilienceContext.Properties.Set(new(nameof(ActivityExecutionContext)), context); - + try { // Configure the resilience pipeline. await resilienceStrategy.ConfigurePipeline(builder, resilienceContext); var pipeline = builder.Build(); - + // Execute the action within the resilience pipeline. var result = await pipeline.ExecuteAsync(async _ => await action(), resilienceContext); - + // Record the retry attempts. await RecordRetryAttempts(activity, context, retries, cancellationToken); - + return result; } finally { ResilienceContextPool.Shared.Return(resilienceContext); } - + } private async Task RecordRetryAttempts(IResilientActivity activity, ActivityExecutionContext context, ICollection attempts, CancellationToken cancellationToken = default) @@ -70,9 +72,11 @@ public class ResilientActivityInvoker( var records = Map(context, activity, attempts); var recordContext = new RecordRetryAttemptsContext(context, records, cancellationToken); await retryAttemptRecorder.RecordAsync(recordContext); - + // Propagate a flag that retries have occurred. This information can then be used to show the retry attempts in the workflow designer. context.SetRetriesAttemptedFlag(); + + context.SetExtensionsMetadata(RetryAttemptsCountKey, attempts.Count); } } @@ -89,7 +93,7 @@ public class ResilientActivityInvoker( ? null : value.ConvertTo(); } - + private ICollection Map(ActivityExecutionContext activityExecutionContext, IResilientActivity resilientActivity, ICollection attempts) { return attempts.Select(x => Map(activityExecutionContext, resilientActivity, x)).ToList(); diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs index 0e7715a5f..ceb0f01e5 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs @@ -24,6 +24,8 @@ namespace Elsa.Extensions; [PublicAPI] public static partial class ActivityExecutionContextExtensions { + private const string ExtensionsMetadataKey = "Extensions"; + /// /// Attempts to get a value from the input provided via . If a value was found, an attempt is made to convert it into the specified type T. /// @@ -436,6 +438,29 @@ public static partial class ActivityExecutionContextExtensions } } + /// + /// Sets extension data in the metadata. Represents specific data that is exposed generically for an activity. + /// + public static void SetExtensionsMetadata(this ActivityExecutionContext context, string key, object? value) + { + var extensionsDictionary = context.GetExtensionsMetadata(); + + if(extensionsDictionary == null) extensionsDictionary = new(); + + extensionsDictionary[key] = value; + + context.Metadata[ExtensionsMetadataKey] = extensionsDictionary; + + } + + /// + /// Retrives the extensin data from the metdata. Represents specific data that is exposed generically for an activity. + /// + public static Dictionary? GetExtensionsMetadata(this ActivityExecutionContext context) + { + return context.Metadata[ExtensionsMetadataKey] as Dictionary; + } + internal static bool GetHasEvaluatedProperties(this ActivityExecutionContext context) => context.TransientProperties.TryGetValue("HasEvaluatedProperties", out var value) && value; internal static void SetHasEvaluatedProperties(this ActivityExecutionContext context) => context.TransientProperties["HasEvaluatedProperties"] = true; } \ No newline at end of file