diff --git a/.gitignore b/.gitignore index d63bbd6bb..8a28551f4 100644 --- a/.gitignore +++ b/.gitignore @@ -78,3 +78,5 @@ unlist.sh /artifacts /docker/data/ + +docker/docker-compose-datadog.yml diff --git a/Directory.Packages.props b/Directory.Packages.props index 626f995bf..8ef10f029 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -1,206 +1,254 @@ - - true - true - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + true + true + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/Elsa.sln b/Elsa.sln index f78fa138b..f5901cf8a 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -87,6 +87,7 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "docker", "docker", "{986E54 docker\ElsaStudio.Dockerfile = docker\ElsaStudio.Dockerfile docker\otel-collector-config.yaml = docker\otel-collector-config.yaml docker\init-db-postgres.sh = docker\init-db-postgres.sh + docker\docker-compose-datadog+otel-collector.yml = docker\docker-compose-datadog+otel-collector.yml EndProjectSection EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Elasticsearch", "src\modules\Elsa.Elasticsearch\Elsa.Elasticsearch.csproj", "{3246883E-2FA7-4B4A-BDC5-99039A2869BC}" diff --git a/docker/ElsaServer-Datadog.Dockerfile b/docker/ElsaServer-Datadog.Dockerfile index a74069cf2..7d5022c63 100644 --- a/docker/ElsaServer-Datadog.Dockerfile +++ b/docker/ElsaServer-Datadog.Dockerfile @@ -10,7 +10,7 @@ COPY ./NuGet.Config ./ COPY *.props ./ # Restore packages. -RUN dotnet restore "./src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj" +RUN dotnet restore "./src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj" # Build and publish (UseAppHost=false creates platform independent binaries). WORKDIR /source/src/bundles/Elsa.Server.Web diff --git a/docker/docker-compose-datadog.yml b/docker/docker-compose-datadog+otel-collector.yml similarity index 99% rename from docker/docker-compose-datadog.yml rename to docker/docker-compose-datadog+otel-collector.yml index e90991c45..20a056580 100644 --- a/docker/docker-compose-datadog.yml +++ b/docker/docker-compose-datadog+otel-collector.yml @@ -91,7 +91,7 @@ services: datadog-agent: image: datadog/agent:latest environment: - DD_API_KEY: "" + DD_API_KEY: "" DD_SITE: "datadoghq.eu" DD_HOSTNAME: "otel-collector" DD_LOGS_ENABLED: "true" diff --git a/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj b/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj index e48fa21f3..4e4903a35 100644 --- a/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj +++ b/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj @@ -59,6 +59,7 @@ + @@ -68,7 +69,13 @@ + + + + + + diff --git a/src/apps/Elsa.Server.Web/Program.cs b/src/apps/Elsa.Server.Web/Program.cs index cfb2d5a26..f5ac0c5cf 100644 --- a/src/apps/Elsa.Server.Web/Program.cs +++ b/src/apps/Elsa.Server.Web/Program.cs @@ -70,6 +70,10 @@ using Medallion.Threading.Redis; using Microsoft.Data.Sqlite; using Microsoft.Extensions.Options; using OpenTelemetry; +using OpenTelemetry.Exporter; +using OpenTelemetry.Logs; +using OpenTelemetry.Metrics; +using OpenTelemetry.Resources; using OpenTelemetry.Trace; using Proto.Cluster.Kubernetes; using Proto.Persistence.Sqlite; @@ -98,7 +102,7 @@ const bool useTenantsFromConfiguration = true; const bool useSecrets = false; const bool disableVariableWrappers = false; const bool disableVariableCopying = false; -const bool useOtel = false; +const bool useManualOtelInstrumentation = true; var builder = WebApplication.CreateBuilder(args); var services = builder.Services; @@ -125,14 +129,39 @@ var sqlDatabaseProvider = Enum.Parse(configuration["Databas TypeAliasRegistry.RegisterAlias("OrderReceivedProducerFactory", typeof(GenericProducerFactory)); TypeAliasRegistry.RegisterAlias("OrderReceivedConsumerFactory", typeof(GenericConsumerFactory)); -if (useOtel) +if (useManualOtelInstrumentation) { -// Configure OpenTelemetry Tracing - using var tracerProvider = Sdk.CreateTracerProviderBuilder() - .AddSource("Elsa.Workflows") // Match your ActivitySource name here - .SetSampler(new AlwaysOnSampler()) // Always record traces for testing - .AddConsoleExporter() // Export spans to the console (optional) - .Build(); + services.AddOpenTelemetry() + .ConfigureResource(resource => resource.AddService("elsa-workflows", serviceVersion: "3.4.0").AddTelemetrySdk()) + .WithTracing(tracing => + { + tracing + .AddSource("*") + .SetSampler(new AlwaysOnSampler()) + .AddAspNetCoreInstrumentation() + .AddHttpClientInstrumentation() + .AddSqlClientInstrumentation() + //.AddConsoleExporter() + .AddOtlpExporter() + ; + }) + .WithMetrics(metrics => + { + metrics + .AddAspNetCoreInstrumentation() + .AddHttpClientInstrumentation() + //.AddConsoleExporter() + .AddOtlpExporter() + ; + }); + + // Enable OpenTelemetry Logging (optional) + builder.Logging.AddOpenTelemetry(options => + { + options.IncludeFormattedMessage = true; + options.IncludeScopes = true; + options.ParseStateValues = true; + }); } // Add Elsa services. @@ -223,7 +252,9 @@ services ef.UsePostgreSql(cockroachDbConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.Oracle) ef.UseOracle(oracleConnectionString, new() - { SchemaName = "ELSA"}); + { + SchemaName = "ELSA" + }); else ef.UseSqlite(sp => sp.GetSqliteConnectionString()); @@ -271,7 +302,9 @@ services ef.UsePostgreSql(cockroachDbConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.Oracle) ef.UseOracle(oracleConnectionString, new() - { SchemaName = "ELSA"}); + { + SchemaName = "ELSA" + }); else ef.UseSqlite(sp => sp.GetSqliteConnectionString()); @@ -324,7 +357,9 @@ services ef.UsePostgreSql(cockroachDbConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.Oracle) ef.UseOracle(oracleConnectionString, new() - { SchemaName = "ELSA"}); + { + SchemaName = "ELSA" + }); else ef.UseSqlite(sp => sp.GetSqliteConnectionString()); @@ -476,7 +511,9 @@ services ef.UsePostgreSql(cockroachDbConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.Oracle) ef.UseOracle(oracleConnectionString, new() - { SchemaName = "ELSA"}); + { + SchemaName = "ELSA" + }); else ef.UseSqlite(sp => sp.GetSqliteConnectionString()); @@ -579,7 +616,7 @@ services } }); } - + if (useKafka) { elsa.UseKafka(kafka => @@ -587,7 +624,7 @@ services kafka.ConfigureOptions(options => configuration.GetSection("Kafka").Bind(options)); }); } - + if (useSecrets) { elsa @@ -663,8 +700,11 @@ services if (sqlDatabaseProvider == SqlDatabaseProvider.PostgreSql) ef.UsePostgreSql(postgresConnectionString); if (sqlDatabaseProvider == SqlDatabaseProvider.Citus) ef.UsePostgreSql(citusConnectionString); if (sqlDatabaseProvider == SqlDatabaseProvider.YugabyteDb) ef.UsePostgreSql(yugabyteDbConnectionString); - if (sqlDatabaseProvider == SqlDatabaseProvider.Oracle) ef.UseOracle(oracleConnectionString, new() - { SchemaName = "ELSA"}); + if (sqlDatabaseProvider == SqlDatabaseProvider.Oracle) + ef.UseOracle(oracleConnectionString, new() + { + SchemaName = "ELSA" + }); #if !NET9_0 if (sqlDatabaseProvider == SqlDatabaseProvider.MySql) ef.UseMySql(mySqlConnectionString); @@ -684,7 +724,7 @@ services tenantHttpRouting.WithTenantHeader("X-Tenant-ID"); }); } - + elsa.UseWebhooks(webhooks => webhooks.ConfigureSinks += options => builder.Configuration.GetSection("Webhooks").Bind(options)); elsa.InstallDropIns(options => options.DropInRootDirectory = Path.Combine(Directory.GetCurrentDirectory(), "App_Data", "DropIns")); elsa.AddSwagger(); diff --git a/src/apps/Elsa.Server.Web/appsettings.json b/src/apps/Elsa.Server.Web/appsettings.json index 3fc86acdb..91735a81b 100644 --- a/src/apps/Elsa.Server.Web/appsettings.json +++ b/src/apps/Elsa.Server.Web/appsettings.json @@ -2,6 +2,7 @@ "Logging": { "LogLevel": { "Default": "Warning", + "OpenTelemetry": "Debug", "Microsoft.Hosting.Lifetime": "Information", "Elsa": "Information", "Elsa.Workflows.Runtime.Middleware.Workflows.WorkflowHeartbeatMiddleware": "Debug" diff --git a/src/modules/Elsa.OpenTelemetry/Abstractions/ErrorSpanHandlerBase.cs b/src/modules/Elsa.OpenTelemetry/Abstractions/ErrorSpanHandlerBase.cs index 420936119..06e77bd4f 100644 --- a/src/modules/Elsa.OpenTelemetry/Abstractions/ErrorSpanHandlerBase.cs +++ b/src/modules/Elsa.OpenTelemetry/Abstractions/ErrorSpanHandlerBase.cs @@ -5,5 +5,7 @@ namespace Elsa.OpenTelemetry.Abstractions; public abstract class ErrorSpanHandlerBase : IErrorSpanHandler { + public virtual float Order => 0; + public abstract bool CanHandle(ErrorSpanContext context); public abstract void Handle(ErrorSpanContext context); } \ No newline at end of file diff --git a/src/modules/Elsa.OpenTelemetry/Contracts/IErrorSpanHandler.cs b/src/modules/Elsa.OpenTelemetry/Contracts/IErrorSpanHandler.cs index fd61ae265..70b16f645 100644 --- a/src/modules/Elsa.OpenTelemetry/Contracts/IErrorSpanHandler.cs +++ b/src/modules/Elsa.OpenTelemetry/Contracts/IErrorSpanHandler.cs @@ -4,6 +4,8 @@ namespace Elsa.OpenTelemetry.Contracts; public interface IErrorSpanHandler { + float Order { get; } + bool CanHandle(ErrorSpanContext context); void Handle(ErrorSpanContext context); } diff --git a/src/modules/Elsa.OpenTelemetry/Handlers/DefaultErrorSpanHandler.cs b/src/modules/Elsa.OpenTelemetry/Handlers/DefaultErrorSpanHandler.cs index 04bebef62..aac9fbae9 100644 --- a/src/modules/Elsa.OpenTelemetry/Handlers/DefaultErrorSpanHandler.cs +++ b/src/modules/Elsa.OpenTelemetry/Handlers/DefaultErrorSpanHandler.cs @@ -5,20 +5,12 @@ namespace Elsa.OpenTelemetry.Handlers; public class DefaultErrorSpanHandler : ErrorSpanHandlerBase { + public override float Order => 100000; + + public override bool CanHandle(ErrorSpanContext context) => context.Exception != null; + public override void Handle(ErrorSpanContext context) { - var span = context.Span; - var exception = context.Exception; - var errorMessage = string.IsNullOrWhiteSpace(exception?.Message) ? "Unknown error" : exception.Message; - span.SetTag("error", true); - span.SetTag("error.message", errorMessage); - - if (exception != null) - { - span.SetTag("error.exceptionType", exception.GetType().FullName); - - if (!string.IsNullOrEmpty(exception.StackTrace)) - span.SetTag("error.stackTrace", exception.StackTrace); - } + context.Span.AddException(context.Exception!); } } \ No newline at end of file diff --git a/src/modules/Elsa.OpenTelemetry/Handlers/FaultExceptionErrorSpanHandler.cs b/src/modules/Elsa.OpenTelemetry/Handlers/FaultExceptionErrorSpanHandler.cs index 14aa2860a..9403e20e9 100644 --- a/src/modules/Elsa.OpenTelemetry/Handlers/FaultExceptionErrorSpanHandler.cs +++ b/src/modules/Elsa.OpenTelemetry/Handlers/FaultExceptionErrorSpanHandler.cs @@ -6,14 +6,18 @@ namespace Elsa.OpenTelemetry.Handlers; public class FaultExceptionErrorSpanHandler : ErrorSpanHandlerBase { + public override bool CanHandle(ErrorSpanContext context) => context.Exception is FaultException; + public override void Handle(ErrorSpanContext context) { - if(context.Exception is not FaultException faultException) - return; - + var faultException = (FaultException)context.Exception!; var span = context.Span; - span.SetTag("error.code", faultException.Code); - span.SetTag("error.category", faultException.Category); - span.SetTag("error.faultType", faultException.Type); + var tags = new Dictionary + { + ["exception.code"] = faultException.Code, + ["exception.category"] = faultException.Category, + ["exception.type"] = faultException.Type + }; + span.AddException(context.Exception, new(tags.ToArray())); } } \ No newline at end of file diff --git a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs index 5410d70bf..6115f7e2d 100644 --- a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs +++ b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingActivityExecutionMiddleware.cs @@ -1,5 +1,6 @@ using System.Diagnostics; using Elsa.Common; +using Elsa.Extensions; using Elsa.OpenTelemetry.Contracts; using Elsa.OpenTelemetry.Helpers; using Elsa.OpenTelemetry.Models; @@ -19,7 +20,7 @@ public class OpenTelemetryTracingActivityExecutionMiddleware(ActivityMiddlewareD public async ValueTask InvokeAsync(ActivityExecutionContext context) { var activity = context.Activity; - using var span = ElsaOpenTelemetry.ActivitySource.StartActivity("ActivityExecution", ActivityKind.Internal, Activity.Current?.Context ?? default); + using var span = ElsaOpenTelemetry.ActivitySource.StartActivity($"execute activity {activity.Type}", ActivityKind.Internal, Activity.Current?.Context ?? default); if (span == null) { @@ -27,46 +28,47 @@ public class OpenTelemetryTracingActivityExecutionMiddleware(ActivityMiddlewareD return; } - span.SetTag("activity.nodeId", activity.NodeId); + span.SetTag("operation.name", "elsa.activity.execution"); + span.SetTag("activity.id", activity.NodeId); + span.SetTag("activity.node.id", activity.NodeId); span.SetTag("activity.type", activity.Type); span.SetTag("activity.name", activity.Name); - span.SetTag("activityInstance.id", context.Id); - span.SetTag("activityExecution.startTimeUtc", span.StartTimeUtc); - span.SetTag("tenantId", context.WorkflowExecutionContext.Workflow.Identity.TenantId); - - span.AddEvent(new("Executing", tags: CreateStatusTags(context))); + 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())) + span.SetTag("span.type", "job"); + + span.AddEvent(new("executing")); + await next(context); if (context.Status == ActivityStatus.Faulted) { - span.AddEvent(new("Faulted", tags: CreateStatusTags(context))); + span.AddEvent(new("faulted")); span.SetStatus(ActivityStatusCode.Error); - span.SetTag("activityInstance.hasIncidents", true); - var errorSpanHandlers = context.GetServices(); var errorSpanHandlerContext = new ErrorSpanContext(span, context.Exception); + var errorSpanHandler = context.GetServices() + .OrderBy(x => x.Order) + .FirstOrDefault(x => x.CanHandle(errorSpanHandlerContext)); - foreach (var handler in errorSpanHandlers) - handler.Handle(errorSpanHandlerContext); + errorSpanHandler?.Handle(errorSpanHandlerContext); } - else + else if (context.Status == ActivityStatus.Canceled) { - span.AddEvent(new("Executed", tags: CreateStatusTags(context))); - span.SetStatus(ActivityStatusCode.Ok); + span.AddEvent(new("canceled")); } - - var now = systemClock.UtcNow; - span.SetTag("activityExecution.endTimeUtc", now); - span.SetTag("activityExecution.durationMs", (now - span.StartTimeUtc).TotalMilliseconds); - } - - private ActivityTagsCollection CreateStatusTags(ActivityExecutionContext context) - { - return new(new Dictionary + else if (context.Status == ActivityStatus.Completed) { - ["activityInstance.status"] = context.Status.ToString() - }); + span.AddEvent(new("completed")); + } + else if (context.Status == ActivityStatus.Pending) + { + span.AddEvent(new("pending")); + } } } diff --git a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs index c0e88ef64..837226655 100644 --- a/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs +++ b/src/modules/Elsa.OpenTelemetry/Middleware/OpenTelemetryTracingWorkflowExecutionMiddleware.cs @@ -1,10 +1,13 @@ using System.Diagnostics; +using System.Runtime.InteropServices.Marshalling; using System.Text.Json; using Elsa.Common; using Elsa.Expressions.Services; using Elsa.Extensions; using Elsa.OpenTelemetry.Helpers; using Elsa.Workflows; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Models; using Elsa.Workflows.Pipelines.WorkflowExecution; using Elsa.Workflows.Serialization.Converters; using JetBrains.Annotations; @@ -19,14 +22,22 @@ namespace Elsa.OpenTelemetry.Middleware; [UsedImplicitly] public class OpenTelemetryTracingWorkflowExecutionMiddleware(WorkflowMiddlewareDelegate next, ISystemClock systemClock) : WorkflowExecutionMiddleware(next) { - private readonly JsonSerializerOptions? _incidentSerializerOptions = new JsonSerializerOptions().WithConverters(new TypeJsonConverter(WellKnownTypeRegistry.CreateDefault())); - /// public override async ValueTask InvokeAsync(WorkflowExecutionContext context) { var workflowInstanceId = context.Id; var workflow = context.Workflow; - using var span = ElsaOpenTelemetry.ActivitySource.StartActivity("WorkflowExecution", ActivityKind.Internal, Activity.Current?.Context ?? default); + var startNewTrace = context.Properties.TryGetValue("StartNewTrace", out var startNewTraceValue) && (bool)startNewTraceValue; + var parentTraceContext = startNewTrace ? default : Activity.Current?.Context ?? default; + var linkedTraceContext = startNewTrace ? Activity.Current : null; + + if(startNewTrace) + { + Activity.Current?.Stop(); + Activity.Current = null; + } + + using var span = ElsaOpenTelemetry.ActivitySource.StartActivity($"execute workflow {workflow.WorkflowMetadata.Name}", ActivityKind.Server, parentTraceContext); if (span == null) // No listener is registered. { @@ -34,61 +45,77 @@ public class OpenTelemetryTracingWorkflowExecutionMiddleware(WorkflowMiddlewareD return; } - span.SetTag("workflowInstance.id", workflowInstanceId); - span.SetTag("workflowDefinition.definitionId", workflow.Identity.DefinitionId); - span.SetTag("workflowDefinition.version", workflow.Identity.Version); - span.SetTag("workflowDefinition.name", workflow.WorkflowMetadata.Name); - span.SetTag("workflowExecution.startTimeUtc", span.StartTimeUtc); - span.SetTag("tenantId", workflow.Identity.TenantId); - - if(context.TriggerActivityId != null) + if(startNewTrace) { - var activity = context.FindActivityById(context.TriggerActivityId) ?? throw new Exception($"Trigger activity with ID {context.TriggerActivityId} not found. This should not happen."); - span.SetTag("workflowExecution.trigger.activityId", activity.Id); - span.SetTag("workflowExecution.trigger.activityName", activity.Name); - span.SetTag("workflowExecution.trigger.activityType", activity.Type); + if (linkedTraceContext != null) + span.AddLink(new(linkedTraceContext.Context)); } - - span.AddEvent(new ActivityEvent("Executing", tags: CreateStatusTags(context))); + + span.SetTag("operation.name", "elsa.workflow.execution"); + span.SetTag("span.type", "workflow"); + span.SetTag("workflow.definition.id", workflow.Identity.DefinitionId); + span.SetTag("workflow.definition.version", workflow.Identity.Version); + span.SetTag("workflow.definition.name", workflow.WorkflowMetadata.Name); + span.SetTag("workflow.definition.tenant.id", workflow.Identity.TenantId); + span.SetTag("workflow.instance.id", workflowInstanceId); + + if (context.TriggerActivityId != null) + { + var activity = context.FindActivityById(context.TriggerActivityId) ?? throw new($"Trigger activity with ID {context.TriggerActivityId} not found. This should not happen."); + span.SetTag("workflow.trigger.activity.id", activity.Id); + span.SetTag("workflow.trigger.activity.name", activity.Name); + span.SetTag("workflow.trigger.activity.type", activity.Type); + span.SetTag("workflow.trigger.activity.version", activity.Version); + } + + span.AddEvent(new("executing")); await Next(context); if (context.SubStatus == WorkflowSubStatus.Faulted) { - span.AddEvent(new ActivityEvent("Faulted", tags: CreateStatusTags(context))); - span.SetStatus(ActivityStatusCode.Error); - span.SetTag("error", true); + span.AddEvent(new("faulted")); + span.SetStatus(ActivityStatusCode.Error, "The workflow entered the Faulted state. See incidents for details."); } - else + else if (context.SubStatus == WorkflowSubStatus.Finished) { - span.AddEvent(new ActivityEvent("Executed", tags: CreateStatusTags(context))); - span.SetStatus(ActivityStatusCode.Ok); + span.AddEvent(new("finished")); } - - if(context.Incidents.Any()) + else if (context.SubStatus == WorkflowSubStatus.Cancelled) { - span.SetStatus(ActivityStatusCode.Error); - span.SetTag("workflowInstance.hasIncidents", true); - span.SetTag("error", true); + span.AddEvent(new("canceled")); + } + else if (context.SubStatus == WorkflowSubStatus.Suspended) + { + span.AddEvent(new("suspended")); + } - if (context.Incidents.Count > 0) - span.SetTag("error.message", JsonSerializer.Serialize(context.Incidents, _incidentSerializerOptions)); - } - - if (!string.IsNullOrWhiteSpace(context.CorrelationId)) - span.SetTag("workflowInstance.correlationId", context.CorrelationId); - - var now = systemClock.UtcNow; - span.SetTag("workflowExecution.endTimeUtc", now); - span.SetTag("workflowExecution.durationMs", (now - span.StartTimeUtc).TotalMilliseconds); - } - - private ActivityTagsCollection CreateStatusTags(WorkflowExecutionContext context) - { - return new ActivityTagsCollection(new Dictionary + if (context.Incidents.Any()) { - ["workflowInstance.status"] = context.Status.ToString(), - ["workflowInstance.subStatus"] = context.SubStatus.ToString() + span.SetTag("workflow.incidents.count", context.Incidents.Count); + + foreach (var incident in context.Incidents) + span.AddEvent(new("incident", incident.Timestamp, CreateIncidentTags(incident))); + } + + if (!string.IsNullOrWhiteSpace(context.CorrelationId)) + span.SetTag("workflow.correlation_id", context.CorrelationId); + } + + private ActivityTagsCollection CreateIncidentTags(ActivityIncident incident) + { + var tags = new ActivityTagsCollection(new Dictionary + { + ["incident.message"] = incident.Message, }); + + if (incident.Exception != null) + { + tags["incident.exception.message"] = incident.Exception.Message; + tags["incident.exception.stackTrace"] = incident.Exception.StackTrace; + tags["incident.exception.type"] = incident.Exception.Type.FullName; + } + + return tags; } } diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs b/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs index 8074be822..2207e05bd 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs @@ -78,6 +78,12 @@ public class BulkDispatchWorkflows : Activity Description = "Wait for the dispatched workflows to complete before completing this activity.", DefaultValue = true)] public Input WaitForCompletion { get; set; } = new(true); + + /// + /// Indicates whether a new trace context should be started for the workflow execution. + /// + [Input(Description = "Start a new trace context when using Open Telemetry.", Category = "Open Telemetry")] + public Input StartNewTrace { get; set; } /// /// The channel to dispatch the workflow to. @@ -104,12 +110,13 @@ public class BulkDispatchWorkflows : Activity protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { var waitForCompletion = WaitForCompletion.GetOrDefault(context); + var startNewTrace = StartNewTrace.GetOrDefault(context); var items = await context.GetItemSource(Items).ToListAsync(context.CancellationToken); var count = items.Count; // Dispatch the child workflows. foreach (var item in items) - await DispatchChildWorkflowAsync(context, item, waitForCompletion); + await DispatchChildWorkflowAsync(context, item, waitForCompletion, startNewTrace); // Store the number of dispatched instances for tracking. context.SetProperty(DispatchedInstancesCountKey, count); @@ -139,7 +146,7 @@ public class BulkDispatchWorkflows : Activity } } - private async ValueTask DispatchChildWorkflowAsync(ActivityExecutionContext context, object item, bool waitForCompletion) + private async ValueTask DispatchChildWorkflowAsync(ActivityExecutionContext context, object item, bool waitForCompletion, bool startNewTrace) { var workflowDefinitionId = WorkflowDefinitionId.Get(context); var workflowDefinitionService = context.GetRequiredService(); @@ -157,8 +164,8 @@ public class BulkDispatchWorkflows : Activity ["ParentInstanceId"] = parentInstanceId }; - if (waitForCompletion) - properties["WaitForCompletion"] = true; + if (waitForCompletion) properties["WaitForCompletion"] = true; + if (startNewTrace) properties["StartNewTrace"] = true; var itemDictionary = new Dictionary { diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs b/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs index f8c49ad83..9619e1d73 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs @@ -55,6 +55,12 @@ public class DispatchWorkflow : Activity [Input(Description = "Wait for the child workflow to complete before completing this activity.")] public Input WaitForCompletion { get; set; } = null!; + /// + /// Indicates whether a new trace context should be started for the workflow execution. + /// + [Input(Description = "Start a new trace context when using Open Telemetry.", Category = "Open Telemetry")] + public Input StartNewTrace { get; set; } + /// /// The channel to dispatch the workflow to. /// @@ -103,16 +109,17 @@ public class DispatchWorkflow : Activity var input = Input.GetOrDefault(context) ?? new Dictionary(); var channelName = ChannelName.GetOrDefault(context); + var startNewTrace = StartNewTrace.GetOrDefault(context); var parentInstanceId = context.WorkflowExecutionContext.Id; var properties = new Dictionary { - ["ParentInstanceId"] = parentInstanceId + ["ParentInstanceId"] = parentInstanceId, }; // If we need to wait for the child workflow to complete, set the property. This will be used by the ResumeDispatchWorkflowActivity handler. - if (waitForCompletion) - properties["WaitForCompletion"] = true; - + if (waitForCompletion) properties["WaitForCompletion"] = true; + if (startNewTrace) properties["StartNewTrace"] = true; + input["ParentInstanceId"] = parentInstanceId; var correlationId = CorrelationId.GetOrDefault(context); diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs index e9cbefc06..09e889628 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs @@ -22,7 +22,7 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher } /// - public async Task DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) + public async Task DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default) { var command = new DispatchWorkflowDefinitionCommand(request.DefinitionVersionId) { @@ -38,7 +38,7 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher } /// - public async Task DispatchAsync(DispatchWorkflowInstanceRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) + public async Task DispatchAsync(DispatchWorkflowInstanceRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default) { var command = new DispatchWorkflowInstanceCommand(request.InstanceId){ BookmarkId = request.BookmarkId, @@ -52,7 +52,7 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher } /// - public async Task DispatchAsync(DispatchTriggerWorkflowsRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) + public async Task DispatchAsync(DispatchTriggerWorkflowsRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default) { var command = new DispatchTriggerWorkflowsCommand(request.ActivityTypeName, request.BookmarkPayload) { @@ -67,7 +67,7 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher } /// - public async Task DispatchAsync(DispatchResumeWorkflowsRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) + public async Task DispatchAsync(DispatchResumeWorkflowsRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default) { var command = new DispatchResumeWorkflowsCommand(request.ActivityTypeName, request.BookmarkPayload) {