From 698d8ec9bf9157f6581ff5a17a7537577dd9faa9 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sun, 30 Aug 2026 11:37:10 -0400 Subject: [PATCH] fix(client): preserve interoperable trace context Signed-off-by: Yordis Prieto --- .config/mise/tasks/semconv/check | 22 ++ .config/mise/tasks/semconv/generate | 46 +++ .github/workflows/dotnet.yml | 17 + mise.toml | 2 + otel/semconv/registry-version | 1 + otel/semconv/registry/manifest.yaml | 6 + .../trogon/eventstore/client/spans.yaml | 60 ++++ .../csharp/telemetry-attributes.cs.j2 | 11 + .../csharp/trogon-telemetry-attributes.cs.j2 | 11 + .../templates/registry/csharp/weaver.yaml | 54 ++++ samples/diagnostics/Program.cs | 2 +- samples/diagnostics/diagnostics.csproj | 2 +- .../Diagnostics/ActivitySourceExtensions.cs | 73 +++-- .../ActivityTagsCollectionExtensions.cs | 52 ++-- .../Diagnostics/Core/ActivityExtensions.cs | 25 +- .../Core/ActivityStatusCodeHelper.cs | 23 -- .../Core/Telemetry/TelemetryTags.cs | 35 --- .../Core/Tracing/TracingConstants.cs | 10 - .../Core/Tracing/TracingMetadata.cs | 32 -- .../Diagnostics/EventMetadataExtensions.cs | 118 ++++--- .../Generated/TelemetryAttributes.g.cs | 22 ++ .../Generated/TrogonTelemetryAttributes.g.cs | 7 + .../Diagnostics/KurrentDBClientDiagnostics.cs | 4 +- .../Diagnostics/Telemetry/TelemetryTags.cs | 12 - .../Diagnostics/Tracing/TracingConstants.cs | 15 +- .../TracerProviderBuilderExtensions.cs | 4 +- ...entDBPersistentSubscriptionsClient.Read.cs | 5 +- .../Streams/KurrentDBClient.Append.cs | 10 +- .../Streams/KurrentDBClient.AppendRecords.cs | 6 +- .../Streams/KurrentDBClient.MultiAppend.cs | 12 +- .../Streams/KurrentDBClient.Subscriptions.cs | 5 +- .../Fixtures/DiagnosticsFixture.cs | 71 +++-- .../Fixtures/KurrentDBPermanentFixture.cs | 15 +- .../KurrentDB.Client.Tests.Common.csproj | 1 + .../Diagnostics/AppendRecordsTracingTests.cs | 23 +- .../Diagnostics/DiagnosticsCollection.cs | 6 + .../OpenTelemetryIntegrationTests.cs | 293 ++++++++++++++++++ ...ubscriptionsTracingInstrumentationTests.cs | 24 +- .../StreamsTracingInstrumentationTests.cs | 105 +++++-- .../KurrentDB.Client.Tests.csproj | 2 + 40 files changed, 913 insertions(+), 331 deletions(-) create mode 100755 .config/mise/tasks/semconv/check create mode 100755 .config/mise/tasks/semconv/generate create mode 100644 mise.toml create mode 100644 otel/semconv/registry-version create mode 100644 otel/semconv/registry/manifest.yaml create mode 100644 otel/semconv/registry/trogon/eventstore/client/spans.yaml create mode 100644 otel/semconv/templates/registry/csharp/telemetry-attributes.cs.j2 create mode 100644 otel/semconv/templates/registry/csharp/trogon-telemetry-attributes.cs.j2 create mode 100644 otel/semconv/templates/registry/csharp/weaver.yaml delete mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityStatusCodeHelper.cs delete mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Core/Telemetry/TelemetryTags.cs delete mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingConstants.cs delete mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingMetadata.cs create mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TelemetryAttributes.g.cs create mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TrogonTelemetryAttributes.g.cs delete mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Telemetry/TelemetryTags.cs create mode 100644 test/KurrentDB.Client.Tests/Diagnostics/DiagnosticsCollection.cs create mode 100644 test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs diff --git a/.config/mise/tasks/semconv/check b/.config/mise/tasks/semconv/check new file mode 100755 index 00000000..6956424c --- /dev/null +++ b/.config/mise/tasks/semconv/check @@ -0,0 +1,22 @@ +#!/bin/sh +#MISE description="Validate the client semantic convention registry and generated constants" + +set -eu + +root=$(CDPATH='' cd -- "$(dirname -- "$0")/../../../.." && pwd) +expected="$root/src/KurrentDB.Client/Core/Common/Diagnostics/Generated" +generated=$(mktemp -d) +trap 'rm -rf "$generated"' EXIT HUP INT TERM + +"$root/.config/mise/tasks/semconv/generate" "$generated" +diff -ru "$expected" "$generated" + +for generated_file in "$expected"/*.g.cs +do + last_byte=$(tail -c 1 "$generated_file" | od -An -t u1 | tr -d '[:space:]') + newline_count=$(tail -c 2 "$generated_file" | wc -l | tr -d '[:space:]') + if [ "$last_byte" != 10 ] || [ "$newline_count" != 1 ]; then + echo "$generated_file must end with exactly one newline." >&2 + exit 1 + fi +done diff --git a/.config/mise/tasks/semconv/generate b/.config/mise/tasks/semconv/generate new file mode 100755 index 00000000..9ec14464 --- /dev/null +++ b/.config/mise/tasks/semconv/generate @@ -0,0 +1,46 @@ +#!/bin/sh +#MISE description="Generate client semantic convention constants" + +set -eu + +root=$(CDPATH='' cd -- "$(dirname -- "$0")/../../../.." && pwd) +output=${1:-"$root/src/KurrentDB.Client/Core/Common/Diagnostics/Generated"} +registry_version=$(sed -n '1p' "$root/otel/semconv/registry-version") +registry="$root/otel/semconv/registry" +official_registry="https://github.com/open-telemetry/semantic-conventions@${registry_version}[model]" +staging=$(mktemp -d) +trap 'rm -rf "$staging"' EXIT HUP INT TERM + +grep -Fqx " registry_path: $official_registry" "$registry/manifest.yaml" + +weaver registry check \ + --future \ + --registry "$registry" + +weaver registry generate csharp "$staging" \ + --future \ + --registry "$registry" \ + --templates "$root/otel/semconv/templates" \ + -D custom_attributes=true \ + -D official_attributes=false + +weaver registry generate csharp "$staging" \ + --future \ + --registry "$official_registry" \ + --templates "$root/otel/semconv/templates" \ + -D custom_attributes=false \ + -D official_attributes=true + +set -- "$staging"/*.g.cs +[ -e "$1" ] +mkdir -p "$output" +for generated_file in "$output"/*.g.cs +do + [ -e "$generated_file" ] || break + rm "$generated_file" +done + +for generated_file +do + mv "$generated_file" "$output/" +done diff --git a/.github/workflows/dotnet.yml b/.github/workflows/dotnet.yml index 6fa92951..e407bcac 100644 --- a/.github/workflows/dotnet.yml +++ b/.github/workflows/dotnet.yml @@ -9,6 +9,23 @@ permissions: contents: read jobs: + semantic-conventions: + name: semantic-conventions + runs-on: ubuntu-latest + steps: + - name: Checkout + uses: actions/checkout@v5 + with: + persist-credentials: false + - name: Install semantic convention tooling + uses: jdx/mise-action@9e7f7633ff6f6d6048a9418a68d48f288f50eb14 # v4.2.3 + with: + version: 2026.8.2 + install_args: github:open-telemetry/weaver + cache: false + - name: Verify semantic conventions + run: mise run --skip-tools semconv:check + vulnerability-scan: name: scan-vulnerabilities runs-on: ubuntu-latest diff --git a/mise.toml b/mise.toml new file mode 100644 index 00000000..8d9f559d --- /dev/null +++ b/mise.toml @@ -0,0 +1,2 @@ +[tools] +"github:open-telemetry/weaver" = "0.24.2" diff --git a/otel/semconv/registry-version b/otel/semconv/registry-version new file mode 100644 index 00000000..518f2eaf --- /dev/null +++ b/otel/semconv/registry-version @@ -0,0 +1 @@ +v1.43.0 diff --git a/otel/semconv/registry/manifest.yaml b/otel/semconv/registry/manifest.yaml new file mode 100644 index 00000000..0e209965 --- /dev/null +++ b/otel/semconv/registry/manifest.yaml @@ -0,0 +1,6 @@ +name: trogon_eventstore_client +description: Semantic conventions for TrogonEventStore client telemetry. +schema_url: https://trogondb.com/schemas/client/0.1.0 +dependencies: + - schema_url: https://opentelemetry.io/schemas/1.43.0 + registry_path: https://github.com/open-telemetry/semantic-conventions@v1.43.0[model] diff --git a/otel/semconv/registry/trogon/eventstore/client/spans.yaml b/otel/semconv/registry/trogon/eventstore/client/spans.yaml new file mode 100644 index 00000000..09195778 --- /dev/null +++ b/otel/semconv/registry/trogon/eventstore/client/spans.yaml @@ -0,0 +1,60 @@ +groups: + - id: registry.trogon.eventstore.client.attributes + type: attribute_group + stability: development + brief: Attributes used by TrogonEventStore client spans. + attributes: + - id: trogon.eventstore.event.type + type: string + stability: development + brief: Event type processed by the client. + examples: [order-created] + - id: span.trogon.eventstore.client.database.operation + type: span + stability: development + span_kind: client + brief: Describes a database operation performed by the client. + attributes: + - ref: db.system.name + requirement_level: required + - ref: db.operation.name + requirement_level: required + - ref: db.collection.name + requirement_level: recommended + - ref: db.operation.batch.size + requirement_level: recommended + - ref: error.type + requirement_level: + conditionally_required: If the operation failed. + - ref: server.address + requirement_level: recommended + - ref: server.port + requirement_level: + conditionally_required: If the server port is available. + + - id: span.trogon.eventstore.client.process + type: span + stability: development + span_kind: consumer + brief: Describes processing an event delivered by a subscription. + attributes: + - ref: messaging.system + requirement_level: required + - ref: messaging.operation.name + requirement_level: required + - ref: messaging.operation.type + requirement_level: required + - ref: messaging.destination.name + requirement_level: required + - ref: messaging.message.id + requirement_level: recommended + - ref: messaging.consumer.group.name + requirement_level: + conditionally_required: If the event was delivered by a persistent subscription. + - ref: trogon.eventstore.event.type + requirement_level: recommended + - ref: server.address + requirement_level: recommended + - ref: server.port + requirement_level: + conditionally_required: If the server port is available. diff --git a/otel/semconv/templates/registry/csharp/telemetry-attributes.cs.j2 b/otel/semconv/templates/registry/csharp/telemetry-attributes.cs.j2 new file mode 100644 index 00000000..d9663edb --- /dev/null +++ b/otel/semconv/templates/registry/csharp/telemetry-attributes.cs.j2 @@ -0,0 +1,11 @@ +// + +namespace KurrentDB.Diagnostics.Telemetry; + +static class TelemetryAttributes { +{% for group in ctx %} +{% for attribute in group.attributes | sort(attribute="name") %} + public const string {{ attribute.name | pascal_case }} = "{{ attribute.name }}"; +{% endfor %} +{% endfor %} +}{{- "\n" -}} diff --git a/otel/semconv/templates/registry/csharp/trogon-telemetry-attributes.cs.j2 b/otel/semconv/templates/registry/csharp/trogon-telemetry-attributes.cs.j2 new file mode 100644 index 00000000..ad79aaf5 --- /dev/null +++ b/otel/semconv/templates/registry/csharp/trogon-telemetry-attributes.cs.j2 @@ -0,0 +1,11 @@ +// + +namespace KurrentDB.Diagnostics.Telemetry; + +static class TrogonTelemetryAttributes { +{% for group in ctx %} +{% for attribute in group.attributes | sort(attribute="name") %} + public const string {{ attribute.name | pascal_case | regex_replace("^TrogonEventstore", "") }} = "{{ attribute.name }}"; +{% endfor %} +{% endfor %} +}{{- "\n" -}} diff --git a/otel/semconv/templates/registry/csharp/weaver.yaml b/otel/semconv/templates/registry/csharp/weaver.yaml new file mode 100644 index 00000000..8b0b21b3 --- /dev/null +++ b/otel/semconv/templates/registry/csharp/weaver.yaml @@ -0,0 +1,54 @@ +whitespace_control: + trim_blocks: true + lstrip_blocks: true + +params: + custom_attributes: false + official_attributes: false + +templates: + - template: telemetry-attributes.cs.j2 + filter: > + if $official_attributes then + semconv_grouped_attributes + | map({ + root_namespace: .root_namespace, + attributes: [.attributes[] | select( + .name == "db.collection.name" or + .name == "db.operation.batch.size" or + .name == "db.operation.name" or + .name == "db.system.name" or + .name == "error.type" or + .name == "exception.message" or + .name == "exception.stacktrace" or + .name == "exception.type" or + .name == "messaging.consumer.group.name" or + .name == "messaging.destination.name" or + .name == "messaging.message.id" or + .name == "messaging.operation.name" or + .name == "messaging.operation.type" or + .name == "messaging.system" or + .name == "server.address" or + .name == "server.port" + )] + }) + | map(select(.attributes | length > 0)) + else + empty + end + application_mode: single + file_name: TelemetryAttributes.g.cs + - template: trogon-telemetry-attributes.cs.j2 + filter: > + if $custom_attributes then + semconv_grouped_attributes + | map({ + root_namespace: .root_namespace, + attributes: [.attributes[] | select(.name | startswith("trogon.eventstore."))] + }) + | map(select(.attributes | length > 0)) + else + empty + end + application_mode: single + file_name: TrogonTelemetryAttributes.g.cs diff --git a/samples/diagnostics/Program.cs b/samples/diagnostics/Program.cs index 357a0962..be69fd9b 100644 --- a/samples/diagnostics/Program.cs +++ b/samples/diagnostics/Program.cs @@ -11,7 +11,7 @@ /** # region import-required-packages // required -dotnet add package EventStore.Client.Extensions.OpenTelemetry +dotnet add package TrogonEventStore.Client // recommended dotnet add package OpenTelemetry.Exporter.OpenTelemetryProtocol diff --git a/samples/diagnostics/diagnostics.csproj b/samples/diagnostics/diagnostics.csproj index 6a740e16..86708e38 100644 --- a/samples/diagnostics/diagnostics.csproj +++ b/samples/diagnostics/diagnostics.csproj @@ -1,7 +1,7 @@  Exe - connecting_to_a_cluster + diagnostics diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs index d85d3261..aa792387 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs @@ -4,21 +4,38 @@ using System.Diagnostics; using KurrentDB.Diagnostics; using KurrentDB.Diagnostics.Telemetry; +using OpenTelemetry; using static KurrentDB.Diagnostics.Tracing.TracingConstants; namespace KurrentDB.Client.Diagnostics; static class ActivitySourceExtensions { - public static async ValueTask TraceClientOperation( + public static ValueTask TraceClientOperation( this ActivitySource source, Func> tracedOperation, string operationName, ActivityTagsCollection? tags = null + ) => source.TraceClientOperation(_ => tracedOperation(), operationName, tags); + + public static async ValueTask TraceClientOperation( + this ActivitySource source, + Func> tracedOperation, + string operationName, + ActivityTagsCollection? tags = null ) { - using var activity = StartActivity(source, operationName, ActivityKind.Client, tags, Activity.Current?.Context); + if (source.HasNoActiveListeners()) + return await tracedOperation(null).ConfigureAwait(false); + + (tags ??= new ActivityTagsCollection()) + .WithRequiredTag(TelemetryAttributes.DbSystemName, SystemName) + .WithRequiredTag(TelemetryAttributes.DbOperationName, operationName); + + var target = tags.FirstOrDefault(tag => tag.Key == TelemetryAttributes.DbCollectionName).Value as string; + var spanName = target is null ? operationName : $"{operationName} {target}"; + using var activity = StartActivity(source, spanName, ActivityKind.Client, tags, Activity.Current?.Context); try { - var res = await tracedOperation().ConfigureAwait(false); + var res = await tracedOperation(activity).ConfigureAwait(false); activity?.StatusOk(); return res; } catch (Exception ex) { @@ -29,34 +46,44 @@ public static async ValueTask TraceClientOperation( public static void TraceSubscriptionEvent( this ActivitySource source, - string? subscriptionId, + string? consumerGroupName, ResolvedEvent resolvedEvent, ChannelInfo channelInfo, - KurrentDBClientSettings settings, - UserCredentials? userCredentials + KurrentDBClientSettings settings ) { if (source.HasNoActiveListeners() || resolvedEvent.Event is null) return; - var parentContext = resolvedEvent.Event.Metadata.ExtractPropagationContext(); + var propagationContext = resolvedEvent.Event.Metadata.ExtractPropagationContext(); - if (parentContext == default(ActivityContext)) return; + if (propagationContext.ActivityContext == default) + return; + var destination = resolvedEvent.OriginalEvent.EventStreamId; var tags = new ActivityTagsCollection() - .WithRequiredTag(TelemetryTags.KurrentDB.Stream, resolvedEvent.OriginalEvent.EventStreamId) - .WithOptionalTag(TelemetryTags.KurrentDB.SubscriptionId, subscriptionId) - .WithRequiredTag(TelemetryTags.KurrentDB.EventId, resolvedEvent.OriginalEvent.EventId.ToString()) - .WithRequiredTag(TelemetryTags.KurrentDB.EventType, resolvedEvent.OriginalEvent.EventType) - // Ensure consistent server.address attribute when connecting to cluster via dns discovery + .WithRequiredTag(TelemetryAttributes.MessagingSystem, SystemName) + .WithRequiredTag(TelemetryAttributes.MessagingOperationName, Operations.Process) + .WithRequiredTag(TelemetryAttributes.MessagingOperationType, Operations.Process) + .WithRequiredTag(TelemetryAttributes.MessagingDestinationName, destination) + .WithOptionalTag(TelemetryAttributes.MessagingConsumerGroupName, consumerGroupName) + .WithRequiredTag(TelemetryAttributes.MessagingMessageId, resolvedEvent.OriginalEvent.EventId.ToString()) + .WithRequiredTag(TrogonTelemetryAttributes.EventType, resolvedEvent.OriginalEvent.EventType) .WithGrpcChannelServerTags(channelInfo) - .WithClientSettingsServerTags(settings) - .WithOptionalTag( - TelemetryTags.Database.User, - userCredentials?.Username ?? settings.DefaultCredentials?.Username - ); - - StartActivity(source, Operations.Subscribe, ActivityKind.Consumer, tags, parentContext) - ?.Dispose(); + .WithClientSettingsServerTags(settings); + + using var activity = StartActivity( + source, + $"{Operations.Process} {destination}", + ActivityKind.Consumer, + tags, + propagationContext.ActivityContext + ); + + if (activity is null) + return; + + foreach (var (name, value) in propagationContext.Baggage.GetBaggage()) + activity.AddBaggage(name, value); } static Activity? StartActivity( @@ -67,10 +94,6 @@ public static void TraceSubscriptionEvent( if (source.HasNoActiveListeners()) return null; - (tags ??= new ActivityTagsCollection()) - .WithRequiredTag(TelemetryTags.Database.System, KurrentDBClientDiagnostics.InstrumentationName) - .WithRequiredTag(TelemetryTags.Database.Operation, operationName); - return source .CreateActivity( operationName, diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/ActivityTagsCollectionExtensions.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/ActivityTagsCollectionExtensions.cs index d069f9d7..086f6745 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/ActivityTagsCollectionExtensions.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/ActivityTagsCollectionExtensions.cs @@ -6,30 +6,30 @@ namespace KurrentDB.Client.Diagnostics; static class ActivityTagsCollectionExtensions { - [MethodImpl(MethodImplOptions.AggressiveInlining)] - public static ActivityTagsCollection WithGrpcChannelServerTags(this ActivityTagsCollection tags, ChannelInfo? channelInfo) { - if (channelInfo is null) - return tags; - - var authorityParts = channelInfo.Channel.Target.Split(':'); - - tags = tags.WithRequiredTag(TelemetryTags.Server.Address, authorityParts[0]); - - if (authorityParts.Length > 1) - tags = tags.WithRequiredTag(TelemetryTags.Server.Port, int.Parse(authorityParts[1])); - - return tags; - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - public static ActivityTagsCollection WithClientSettingsServerTags(this ActivityTagsCollection source, KurrentDBClientSettings settings) { - if (settings.ConnectivitySettings.DnsGossipSeeds?.Length != 1) - return source; - - var gossipSeed = settings.ConnectivitySettings.DnsGossipSeeds[0]; - - return source - .WithRequiredTag(TelemetryTags.Server.Address, gossipSeed.Host) - .WithRequiredTag(TelemetryTags.Server.Port, gossipSeed.Port); - } + [MethodImpl(MethodImplOptions.AggressiveInlining)] + public static ActivityTagsCollection WithGrpcChannelServerTags(this ActivityTagsCollection tags, ChannelInfo? channelInfo) { + if (channelInfo is null) + return tags; + + var authorityParts = channelInfo.Channel.Target.Split(':'); + + tags = tags.WithRequiredTag(TelemetryAttributes.ServerAddress, authorityParts[0]); + + if (authorityParts.Length > 1) + tags = tags.WithRequiredTag(TelemetryAttributes.ServerPort, int.Parse(authorityParts[1])); + + return tags; + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + public static ActivityTagsCollection WithClientSettingsServerTags(this ActivityTagsCollection source, KurrentDBClientSettings settings) { + if (settings.ConnectivitySettings.DnsGossipSeeds?.Length != 1) + return source; + + var gossipSeed = settings.ConnectivitySettings.DnsGossipSeeds[0]; + + return source + .WithRequiredTag(TelemetryAttributes.ServerAddress, gossipSeed.Host) + .WithRequiredTag(TelemetryAttributes.ServerPort, gossipSeed.Port); + } } diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityExtensions.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityExtensions.cs index 81d23564..5789e8ba 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityExtensions.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityExtensions.cs @@ -3,15 +3,12 @@ using System.Diagnostics; using System.Runtime.CompilerServices; using KurrentDB.Diagnostics.Telemetry; -using KurrentDB.Diagnostics.Tracing; + +using static KurrentDB.Diagnostics.Tracing.TracingConstants; namespace KurrentDB.Diagnostics; static class ActivityExtensions { - [MethodImpl(MethodImplOptions.AggressiveInlining)] - public static TracingMetadata GetTracingMetadata(this Activity activity) => - new(activity.TraceId.ToString(), activity.SpanId.ToString()); - [MethodImpl(MethodImplOptions.AggressiveInlining)] public static Activity StatusOk(this Activity activity, string? description = null) => activity.SetActivityStatus(ActivityStatus.Ok(description)); @@ -22,30 +19,30 @@ public static Activity StatusError(this Activity activity, Exception exception) [MethodImpl(MethodImplOptions.AggressiveInlining)] static Activity RecordException(this Activity activity, Exception? exception) { - if (exception is null) return activity; + if (exception is null) + return activity; var ex = exception is AggregateException aex ? aex.Flatten() : exception; var tags = new ActivityTagsCollection { - { TelemetryTags.Exception.Type, ex.GetType().FullName }, - { TelemetryTags.Exception.Stacktrace, ex.ToInvariantString() } + { TelemetryAttributes.ExceptionType, ex.GetType().FullName }, + { TelemetryAttributes.ExceptionStacktrace, ex.ToInvariantString() } }; if (!string.IsNullOrWhiteSpace(exception.Message)) - tags.Add(TelemetryTags.Exception.Message, ex.Message); + tags.Add(TelemetryAttributes.ExceptionMessage, ex.Message); - activity.AddEvent(new ActivityEvent(TelemetryTags.Exception.EventName, default, tags)); + activity.AddEvent(new ActivityEvent(ExceptionEventName, default, tags)); return activity; } [MethodImpl(MethodImplOptions.AggressiveInlining)] static Activity SetActivityStatus(this Activity activity, ActivityStatus status) { - var statusCode = ActivityStatusCodeHelper.GetTagValueForStatusCode(status.StatusCode); - activity.SetStatus(status.StatusCode, status.Description); - activity.SetTag(TelemetryTags.Otel.StatusCode, statusCode); - activity.SetTag(TelemetryTags.Otel.StatusDescription, status.Description); + + if (status.Exception is { } exception) + activity.SetTag(TelemetryAttributes.ErrorType, exception.GetType().FullName); return activity.IsAllDataRequested ? activity.RecordException(status.Exception) : activity; } diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityStatusCodeHelper.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityStatusCodeHelper.cs deleted file mode 100644 index 657ce553..00000000 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/ActivityStatusCodeHelper.cs +++ /dev/null @@ -1,23 +0,0 @@ -// ReSharper disable CheckNamespace - -using System.Diagnostics; -using System.Runtime.CompilerServices; - -using static System.Diagnostics.ActivityStatusCode; - -namespace KurrentDB.Diagnostics; - -static class ActivityStatusCodeHelper { - public const string UnsetStatusCodeTagValue = "UNSET"; - public const string OkStatusCodeTagValue = "OK"; - public const string ErrorStatusCodeTagValue = "ERROR"; - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - public static string? GetTagValueForStatusCode(ActivityStatusCode statusCode) => - statusCode switch { - Unset => UnsetStatusCodeTagValue, - Error => ErrorStatusCodeTagValue, - Ok => OkStatusCodeTagValue, - _ => null - }; -} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Telemetry/TelemetryTags.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Telemetry/TelemetryTags.cs deleted file mode 100644 index e355c96f..00000000 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Telemetry/TelemetryTags.cs +++ /dev/null @@ -1,35 +0,0 @@ -// ReSharper disable CheckNamespace - -namespace KurrentDB.Diagnostics.Telemetry; - -// The attributes below match the specification of v1.24.0 of the Open Telemetry semantic conventions. -// Some attributes are ignored where not required or relevant. -// https://github.com/open-telemetry/semantic-conventions/blob/v1.24.0/docs/general/trace.md -// https://github.com/open-telemetry/semantic-conventions/blob/v1.24.0/docs/database/database-spans.md -// https://github.com/open-telemetry/semantic-conventions/blob/v1.24.0/docs/exceptions/exceptions-spans.md - -static partial class TelemetryTags { - public static class Database { - public const string User = "db.user"; - public const string System = "db.system"; - public const string Operation = "db.operation"; - } - - public static class Server { - public const string Address = "server.address"; - public const string Port = "server.port"; - public const string SocketAddress = "server.socket.address"; // replaces: "net.peer.ip" (AttributeNetPeerIp) - } - - public static class Exception { - public const string EventName = "exception"; - public const string Type = "exception.type"; - public const string Message = "exception.message"; - public const string Stacktrace = "exception.stacktrace"; - } - - public static class Otel { - public const string StatusCode = "otel.status_code"; - public const string StatusDescription = "otel.status_description"; - } -} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingConstants.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingConstants.cs deleted file mode 100644 index 355c6fca..00000000 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingConstants.cs +++ /dev/null @@ -1,10 +0,0 @@ -// ReSharper disable CheckNamespace - -namespace KurrentDB.Diagnostics.Tracing; - -static partial class TracingConstants { - public static class Metadata { - public const string TraceId = "$traceId"; - public const string SpanId = "$spanId"; - } -} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingMetadata.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingMetadata.cs deleted file mode 100644 index 87a1bbe7..00000000 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Core/Tracing/TracingMetadata.cs +++ /dev/null @@ -1,32 +0,0 @@ -// ReSharper disable CheckNamespace - -using System.Diagnostics; -using System.Text.Json.Serialization; - -namespace KurrentDB.Diagnostics.Tracing; - -readonly record struct TracingMetadata( - [property: JsonPropertyName(TracingConstants.Metadata.TraceId)] - string? TraceId, - [property: JsonPropertyName(TracingConstants.Metadata.SpanId)] - string? SpanId -) { - public static readonly TracingMetadata None = new(null, null); - - [JsonIgnore] public bool IsValid => TraceId != null && SpanId != null; - - public ActivityContext? ToActivityContext(bool isRemote = true) { - try { - return IsValid - ? new ActivityContext( - ActivityTraceId.CreateFromString(new ReadOnlySpan(TraceId!.ToCharArray())), - ActivitySpanId.CreateFromString(new ReadOnlySpan(SpanId!.ToCharArray())), - ActivityTraceFlags.Recorded, - isRemote: isRemote - ) - : default; - } catch (Exception) { - return default; - } - } -} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs index 06625bb7..b2391cb4 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs @@ -1,98 +1,126 @@ using System.Diagnostics; using System.Runtime.CompilerServices; using System.Text.Json; -using KurrentDB.Diagnostics; -using KurrentDB.Diagnostics.Tracing; +using OpenTelemetry; +using OpenTelemetry.Context.Propagation; namespace KurrentDB.Client.Diagnostics; static class EventMetadataExtensions { public static void InjectTracingContext(this Dictionary metadata, Activity? activity) { - if (!KurrentDBClientDiagnostics.ActivitySource.HasListeners() || activity is null) return; + if (activity is null) + return; + + Propagators.DefaultTextMapPropagator.Inject( + new PropagationContext(activity.Context, Baggage.Current), + metadata, + static (carrier, name, value) => SetPropagationField(carrier, name, value) + ); + } + + static void SetPropagationField(Dictionary metadata, string name, string value) { + while (true) { + string? existingName = null; + foreach (var key in metadata.Keys) { + if (!string.Equals(key, name, StringComparison.OrdinalIgnoreCase)) + continue; - if (!metadata.ContainsKey(TracingConstants.Metadata.TraceId)) - metadata[TracingConstants.Metadata.TraceId] = activity.TraceId.ToString(); + existingName = key; + break; + } - if (!metadata.ContainsKey(TracingConstants.Metadata.SpanId)) - metadata[TracingConstants.Metadata.SpanId] = activity.SpanId.ToString(); + if (existingName is null) + break; + + metadata.Remove(existingName); + } + + metadata[name] = value; } [MethodImpl(MethodImplOptions.AggressiveInlining)] public static ReadOnlySpan InjectTracingContext( this ReadOnlyMemory eventMetadata, Activity? activity - ) => - eventMetadata.InjectTracingMetadata(activity?.GetTracingMetadata() ?? TracingMetadata.None); + ) { + if (activity is null) + return eventMetadata.Span; - [MethodImpl(MethodImplOptions.AggressiveInlining)] - public static ActivityContext? ExtractPropagationContext(this ReadOnlyMemory eventMetadata) => - eventMetadata.ExtractTracingMetadata().ToActivityContext(isRemote: true); + var propagationMetadata = new Dictionary(StringComparer.OrdinalIgnoreCase); + Propagators.DefaultTextMapPropagator.Inject( + new PropagationContext(activity.Context, Baggage.Current), + propagationMetadata, + static (carrier, name, value) => carrier[name] = value + ); + + return eventMetadata.InjectPropagationMetadata(propagationMetadata); + } [MethodImpl(MethodImplOptions.AggressiveInlining)] - public static TracingMetadata ExtractTracingMetadata(this ReadOnlyMemory eventMetadata) { + public static PropagationContext ExtractPropagationContext(this ReadOnlyMemory eventMetadata) { if (eventMetadata.IsEmpty) - return TracingMetadata.None; + return default; var reader = new Utf8JsonReader(eventMetadata.Span); try { if (!JsonDocument.TryParseValue(ref reader, out var doc)) - return TracingMetadata.None; + return default; using (doc) { - if (!doc.RootElement.TryGetProperty(TracingConstants.Metadata.TraceId, out var traceId) - || !doc.RootElement.TryGetProperty(TracingConstants.Metadata.SpanId, out var spanId)) - return TracingMetadata.None; - - return new TracingMetadata(traceId.GetString(), spanId.GetString()); + if (doc.RootElement.ValueKind != JsonValueKind.Object) + return default; + + var propagationMetadata = new Dictionary(StringComparer.OrdinalIgnoreCase); + foreach (var property in doc.RootElement.EnumerateObject()) { + if (property.Value.ValueKind == JsonValueKind.String && property.Value.GetString() is { } value) + propagationMetadata[property.Name] = value; + } + + return Propagators.DefaultTextMapPropagator.Extract( + default, + propagationMetadata, + static (carrier, name) => + carrier.TryGetValue(name, out var value) ? [value] : [] + ); } } catch (Exception) { - return TracingMetadata.None; + return default; } } [MethodImpl(MethodImplOptions.AggressiveInlining)] - static ReadOnlySpan InjectTracingMetadata( - this ReadOnlyMemory eventMetadata, TracingMetadata tracingMetadata + static ReadOnlySpan InjectPropagationMetadata( + this ReadOnlyMemory eventMetadata, Dictionary propagationMetadata ) { - if (tracingMetadata == TracingMetadata.None || !tracingMetadata.IsValid) + if (propagationMetadata.Count == 0) return eventMetadata.Span; return eventMetadata.IsEmpty - ? JsonSerializer.SerializeToUtf8Bytes(tracingMetadata) - : TryInjectTracingMetadata(eventMetadata, tracingMetadata).ToArray(); + ? JsonSerializer.SerializeToUtf8Bytes(propagationMetadata) + : TryInjectPropagationMetadata(eventMetadata, propagationMetadata).ToArray(); } [MethodImpl(MethodImplOptions.AggressiveInlining)] - static ReadOnlyMemory TryInjectTracingMetadata( - this ReadOnlyMemory utf8Json, TracingMetadata tracingMetadata + static ReadOnlyMemory TryInjectPropagationMetadata( + this ReadOnlyMemory utf8Json, Dictionary propagationMetadata ) { try { - using var doc = JsonDocument.Parse(utf8Json); + using var doc = JsonDocument.Parse(utf8Json); using var stream = new MemoryStream(); using var writer = new Utf8JsonWriter(stream); if (doc.RootElement.ValueKind != JsonValueKind.Object) return utf8Json; - var hasTraceId = doc.RootElement.TryGetProperty(TracingConstants.Metadata.TraceId, out _); - var hasSpanId = doc.RootElement.TryGetProperty(TracingConstants.Metadata.SpanId, out _); - - if (hasTraceId && hasSpanId) - return utf8Json; - writer.WriteStartObject(); - foreach (var prop in doc.RootElement.EnumerateObject()) - prop.WriteTo(writer); - - if (!hasTraceId) { - writer.WritePropertyName(TracingConstants.Metadata.TraceId); - writer.WriteStringValue(tracingMetadata.TraceId); + foreach (var property in doc.RootElement.EnumerateObject()) { + if (!propagationMetadata.ContainsKey(property.Name)) + property.WriteTo(writer); } - if (!hasSpanId) { - writer.WritePropertyName(TracingConstants.Metadata.SpanId); - writer.WriteStringValue(tracingMetadata.SpanId); - } + + foreach (var (name, value) in propagationMetadata) + writer.WriteString(name, value); writer.WriteEndObject(); writer.Flush(); diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TelemetryAttributes.g.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TelemetryAttributes.g.cs new file mode 100644 index 00000000..a34ec240 --- /dev/null +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TelemetryAttributes.g.cs @@ -0,0 +1,22 @@ +// + +namespace KurrentDB.Diagnostics.Telemetry; + +static class TelemetryAttributes { + public const string DbCollectionName = "db.collection.name"; + public const string DbOperationBatchSize = "db.operation.batch.size"; + public const string DbOperationName = "db.operation.name"; + public const string DbSystemName = "db.system.name"; + public const string ErrorType = "error.type"; + public const string ExceptionMessage = "exception.message"; + public const string ExceptionStacktrace = "exception.stacktrace"; + public const string ExceptionType = "exception.type"; + public const string MessagingConsumerGroupName = "messaging.consumer.group.name"; + public const string MessagingDestinationName = "messaging.destination.name"; + public const string MessagingMessageId = "messaging.message.id"; + public const string MessagingOperationName = "messaging.operation.name"; + public const string MessagingOperationType = "messaging.operation.type"; + public const string MessagingSystem = "messaging.system"; + public const string ServerAddress = "server.address"; + public const string ServerPort = "server.port"; +} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TrogonTelemetryAttributes.g.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TrogonTelemetryAttributes.g.cs new file mode 100644 index 00000000..f785efd7 --- /dev/null +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/TrogonTelemetryAttributes.g.cs @@ -0,0 +1,7 @@ +// + +namespace KurrentDB.Diagnostics.Telemetry; + +static class TrogonTelemetryAttributes { + public const string EventType = "trogon.eventstore.event.type"; +} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/KurrentDBClientDiagnostics.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/KurrentDBClientDiagnostics.cs index be0115e0..dce45d72 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/KurrentDBClientDiagnostics.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/KurrentDBClientDiagnostics.cs @@ -3,6 +3,6 @@ namespace KurrentDB.Client.Diagnostics; public static class KurrentDBClientDiagnostics { - public const string InstrumentationName = "kurrentdb"; - public static readonly ActivitySource ActivitySource = new(InstrumentationName); + public const string InstrumentationName = "TrogonEventStore.Client"; + public static readonly ActivitySource ActivitySource = new(InstrumentationName); } diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Telemetry/TelemetryTags.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Telemetry/TelemetryTags.cs deleted file mode 100644 index 130c3fa8..00000000 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Telemetry/TelemetryTags.cs +++ /dev/null @@ -1,12 +0,0 @@ -// ReSharper disable CheckNamespace - -namespace KurrentDB.Diagnostics.Telemetry; - -static partial class TelemetryTags { - public static class KurrentDB { - public const string Stream = "db.kurrentdb.stream"; - public const string SubscriptionId = "db.kurrentdb.subscription.id"; - public const string EventId = "db.kurrentdb.event.id"; - public const string EventType = "db.kurrentdb.event.type"; - } -} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs index 5562748a..7dcd65fe 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs @@ -2,10 +2,13 @@ namespace KurrentDB.Diagnostics.Tracing; -static partial class TracingConstants { - public static class Operations { - public const string Append = "streams.append"; - public const string MultiAppend = "streams.multi-append"; - public const string Subscribe = "streams.subscribe"; - } +static class TracingConstants { + public const string SystemName = "trogoneventstore"; + public const string ExceptionEventName = "exception"; + + public static class Operations { + public const string Append = "append"; + public const string BatchAppend = "batch_append"; + public const string Process = "process"; + } } diff --git a/src/KurrentDB.Client/OpenTelemetry/TracerProviderBuilderExtensions.cs b/src/KurrentDB.Client/OpenTelemetry/TracerProviderBuilderExtensions.cs index 3844242c..85b04d56 100644 --- a/src/KurrentDB.Client/OpenTelemetry/TracerProviderBuilderExtensions.cs +++ b/src/KurrentDB.Client/OpenTelemetry/TracerProviderBuilderExtensions.cs @@ -5,12 +5,12 @@ namespace KurrentDB.Client.Extensions.OpenTelemetry; /// -/// Extension methods used to facilitate tracing instrumentation of the EventStore Client. +/// Extension methods used to facilitate tracing instrumentation of the TrogonEventStore client. /// [PublicAPI] public static class TracerProviderBuilderExtensions { /// - /// Adds the EventStore client ActivitySource name to the list of subscribed sources on the + /// Adds the TrogonEventStore client ActivitySource name to the list of subscribed sources on the /// /// being configured. /// The instance of to chain configuration. diff --git a/src/KurrentDB.Client/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs b/src/KurrentDB.Client/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs index 1a53cd65..2aa87ec3 100644 --- a/src/KurrentDB.Client/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs +++ b/src/KurrentDB.Client/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs @@ -271,11 +271,10 @@ async Task PumpMessages() { if (subscriptionMessage is PersistentSubscriptionMessage.Event evnt) KurrentDBClientDiagnostics.ActivitySource.TraceSubscriptionEvent( - SubscriptionId, + GroupName, evnt.ResolvedEvent, channelInfo, - settings, - userCredentials + settings ); await _channel.Writer.WriteAsync(subscriptionMessage, _cts.Token).ConfigureAwait(false); diff --git a/src/KurrentDB.Client/Streams/KurrentDBClient.Append.cs b/src/KurrentDB.Client/Streams/KurrentDBClient.Append.cs index 12e9eea1..b582110c 100644 --- a/src/KurrentDB.Client/Streams/KurrentDBClient.Append.cs +++ b/src/KurrentDB.Client/Streams/KurrentDBClient.Append.cs @@ -68,10 +68,9 @@ ValueTask AppendToStreamInternal( CancellationToken cancellationToken ) { var tags = new ActivityTagsCollection() - .WithRequiredTag(TelemetryTags.KurrentDB.Stream, header.Options.StreamIdentifier.StreamName.ToStringUtf8()) + .WithRequiredTag(TelemetryAttributes.DbCollectionName, header.Options.StreamIdentifier.StreamName.ToStringUtf8()) .WithGrpcChannelServerTags(channelInfo) - .WithClientSettingsServerTags(Settings) - .WithOptionalTag(TelemetryTags.Database.User, userCredentials?.Username ?? Settings.DefaultCredentials?.Username); + .WithClientSettingsServerTags(Settings); return KurrentDBClientDiagnostics.ActivitySource.TraceClientOperation(Operation, TracingConstants.Operations.Append, tags); @@ -228,10 +227,9 @@ ValueTask AppendInternal( CancellationToken cancellationToken ) { var tags = new ActivityTagsCollection() - .WithRequiredTag(TelemetryTags.KurrentDB.Stream, options.StreamIdentifier.StreamName.ToStringUtf8()) + .WithRequiredTag(TelemetryAttributes.DbCollectionName, options.StreamIdentifier.StreamName.ToStringUtf8()) .WithGrpcChannelServerTags(_channelInfo) - .WithClientSettingsServerTags(_settings) - .WithOptionalTag(TelemetryTags.Database.User, _settings.DefaultCredentials?.Username); + .WithClientSettingsServerTags(_settings); return KurrentDBClientDiagnostics.ActivitySource.TraceClientOperation( Operation, diff --git a/src/KurrentDB.Client/Streams/KurrentDBClient.AppendRecords.cs b/src/KurrentDB.Client/Streams/KurrentDBClient.AppendRecords.cs index 6e769a6f..179d5034 100644 --- a/src/KurrentDB.Client/Streams/KurrentDBClient.AppendRecords.cs +++ b/src/KurrentDB.Client/Streams/KurrentDBClient.AppendRecords.cs @@ -51,11 +51,11 @@ public async ValueTask AppendRecordsAsync( var client = new StreamsServiceClient(channelInfo.CallInvoker); var tags = new ActivityTagsCollection() + .WithRequiredTag(TelemetryAttributes.DbOperationBatchSize, recordsList.Count) .WithGrpcChannelServerTags(channelInfo) - .WithClientSettingsServerTags(Settings) - .WithOptionalTag(TelemetryTags.Database.User, Settings.DefaultCredentials?.Username); + .WithClientSettingsServerTags(Settings); - return await KurrentDBClientDiagnostics.ActivitySource.TraceClientOperation(Operation, Operations.MultiAppend, tags).ConfigureAwait(false); + return await KurrentDBClientDiagnostics.ActivitySource.TraceClientOperation(Operation, Operations.BatchAppend, tags).ConfigureAwait(false); async ValueTask Operation() { try { diff --git a/src/KurrentDB.Client/Streams/KurrentDBClient.MultiAppend.cs b/src/KurrentDB.Client/Streams/KurrentDBClient.MultiAppend.cs index a3ff9440..80e5ad56 100644 --- a/src/KurrentDB.Client/Streams/KurrentDBClient.MultiAppend.cs +++ b/src/KurrentDB.Client/Streams/KurrentDBClient.MultiAppend.cs @@ -43,16 +43,20 @@ public async ValueTask MultiStreamAppendAsync( var tags = new ActivityTagsCollection() .WithGrpcChannelServerTags(channelInfo) - .WithClientSettingsServerTags(Settings) - .WithOptionalTag(TelemetryTags.Database.User, Settings.DefaultCredentials?.Username); + .WithClientSettingsServerTags(Settings); - return await KurrentDBClientDiagnostics.ActivitySource.TraceClientOperation(Operation, Operations.MultiAppend, tags).ConfigureAwait(false); + return await KurrentDBClientDiagnostics.ActivitySource.TraceClientOperation(Operation, Operations.BatchAppend, tags).ConfigureAwait(false); + + async ValueTask Operation(Activity? activity) { + activity?.SetTag(TelemetryAttributes.DbOperationBatchSize, 0); - async ValueTask Operation() { try { using var session = client.AppendSession(KurrentDBCallOptions.CreateStreaming(Settings, cancellationToken: cancellationToken)); + var batchSize = 0; await foreach (var request in requests.WithCancellation(cancellationToken)) { + batchSize++; + activity?.SetTag(TelemetryAttributes.DbOperationBatchSize, batchSize); var records = await request.Messages .Map() .ToArrayAsync(cancellationToken) diff --git a/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs b/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs index 52b6998f..639d39b7 100644 --- a/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs +++ b/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs @@ -242,11 +242,10 @@ response.FellBehind.Position is { } position if (subscriptionMessage is StreamMessage.Event evt) KurrentDBClientDiagnostics.ActivitySource.TraceSubscriptionEvent( - SubscriptionId, + null, evt.ResolvedEvent, channelInfo, - _settings, - userCredentials + _settings ); await _channel.Writer diff --git a/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs b/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs index 7fc4bc0a..d4d6063b 100644 --- a/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs +++ b/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs @@ -7,18 +7,22 @@ using KurrentDB.Diagnostics; using KurrentDB.Diagnostics.Telemetry; using KurrentDB.Diagnostics.Tracing; +using OpenTelemetry; +using OpenTelemetry.Context.Propagation; namespace KurrentDB.Client.Tests.Fixtures; public class DiagnosticsFixture : KurrentDBPermanentFixture { readonly ConcurrentDictionary<(string Operation, ActivityTraceId TraceId), List> Activities = []; + readonly TextMapPropagator OriginalPropagator = Propagators.DefaultTextMapPropagator; public DiagnosticsFixture() : base(x => x.RunProjections()) { var diagnosticActivityListener = new ActivityListener { ShouldListenTo = source => source.Name == KurrentDBClientDiagnostics.InstrumentationName, - Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, ActivityStopped = activity => { - var operation = (string?)activity.GetTagItem(TelemetryTags.Database.Operation); + var operation = (string?)activity.GetTagItem(TelemetryAttributes.DbOperationName) + ?? (string?)activity.GetTagItem(TelemetryAttributes.MessagingOperationName); if (operation is null) return; @@ -35,12 +39,17 @@ public DiagnosticsFixture() : base(x => x.RunProjections()) { }; OnSetup += () => { + Sdk.SetDefaultTextMapPropagator(new CompositeTextMapPropagator([ + new TraceContextPropagator(), + new BaggagePropagator() + ])); ActivitySource.AddActivityListener(diagnosticActivityListener); return Task.CompletedTask; }; OnTearDown = () => { diagnosticActivityListener.Dispose(); + Sdk.SetDefaultTextMapPropagator(OriginalPropagator); return Task.CompletedTask; }; } @@ -58,28 +67,35 @@ public List GetActivities(string operation, ActivityTraceId traceId) = public List GetActivities(string operation, ActivityTraceId traceId, string stream) => GetActivities(operation, traceId) - .Where(activity => Equals(activity.GetTagItem(TelemetryTags.KurrentDB.Stream), stream)) + .Where(activity => + Equals(activity.GetTagItem(TelemetryAttributes.DbCollectionName), stream) || + Equals(activity.GetTagItem(TelemetryAttributes.MessagingDestinationName), stream) + ) .ToList(); public void AssertMultiAppendActivityHasExpectedTags(Activity activity) { + activity.DisplayName.ShouldBe(TracingConstants.Operations.BatchAppend); + activity.Kind.ShouldBe(ActivityKind.Client); + var expectedTags = new Dictionary { - { TelemetryTags.Database.System, KurrentDBClientDiagnostics.InstrumentationName }, - { TelemetryTags.Database.Operation, TracingConstants.Operations.MultiAppend }, - { TelemetryTags.Database.User, TestCredentials.Root.Username }, - { TelemetryTags.Otel.StatusCode, ActivityStatusCodeHelper.OkStatusCodeTagValue } + { TelemetryAttributes.DbSystemName, TracingConstants.SystemName }, + { TelemetryAttributes.DbOperationName, TracingConstants.Operations.BatchAppend } }; foreach (var tag in expectedTags) activity.Tags.ShouldContain(tag); + + activity.GetTagItem(TelemetryAttributes.DbOperationBatchSize).ShouldBe(2); } public void AssertAppendActivityHasExpectedTags(Activity activity, string stream) { + activity.DisplayName.ShouldBe($"{TracingConstants.Operations.Append} {stream}"); + activity.Kind.ShouldBe(ActivityKind.Client); + var expectedTags = new Dictionary { - { TelemetryTags.Database.System, KurrentDBClientDiagnostics.InstrumentationName }, - { TelemetryTags.Database.Operation, TracingConstants.Operations.Append }, - { TelemetryTags.KurrentDB.Stream, stream }, - { TelemetryTags.Database.User, TestCredentials.Root.Username }, - { TelemetryTags.Otel.StatusCode, ActivityStatusCodeHelper.OkStatusCodeTagValue } + { TelemetryAttributes.DbSystemName, TracingConstants.SystemName }, + { TelemetryAttributes.DbOperationName, TracingConstants.Operations.Append }, + { TelemetryAttributes.DbCollectionName, stream } }; foreach (var tag in expectedTags) @@ -88,7 +104,7 @@ public void AssertAppendActivityHasExpectedTags(Activity activity, string stream public void AssertErroneousAppendActivityHasExpectedTags(Activity activity, Exception actualException) { var expectedTags = new Dictionary { - { TelemetryTags.Otel.StatusCode, ActivityStatusCodeHelper.ErrorStatusCodeTagValue } + { TelemetryAttributes.ErrorType, actualException.GetType().FullName } }; foreach (var tag in expectedTags) @@ -96,31 +112,34 @@ public void AssertErroneousAppendActivityHasExpectedTags(Activity activity, Exce var actualEvent = activity.Events.ShouldHaveSingleItem(); - actualEvent.Name.ShouldBe(TelemetryTags.Exception.EventName); - actualEvent.Tags.ShouldContain(new KeyValuePair(TelemetryTags.Exception.Type, actualException.GetType().FullName)); + actualEvent.Name.ShouldBe(TracingConstants.ExceptionEventName); + actualEvent.Tags.ShouldContain(new KeyValuePair(TelemetryAttributes.ExceptionType, actualException.GetType().FullName)); - actualEvent.Tags.ShouldContain(new KeyValuePair(TelemetryTags.Exception.Message, actualException.Message)); + actualEvent.Tags.ShouldContain(new KeyValuePair(TelemetryAttributes.ExceptionMessage, actualException.Message)); - actualEvent.Tags.Any(x => x.Key == TelemetryTags.Exception.Stacktrace).ShouldBeTrue(); + actualEvent.Tags.Any(x => x.Key == TelemetryAttributes.ExceptionStacktrace).ShouldBeTrue(); } public void AssertSubscriptionActivityHasExpectedTags( Activity activity, string stream, string eventId, - string? subscriptionId = null + string? consumerGroupName = null ) { + activity.DisplayName.ShouldBe($"{TracingConstants.Operations.Process} {stream}"); + activity.Kind.ShouldBe(ActivityKind.Consumer); + var expectedTags = new Dictionary { - { TelemetryTags.Database.System, KurrentDBClientDiagnostics.InstrumentationName }, - { TelemetryTags.Database.Operation, TracingConstants.Operations.Subscribe }, - { TelemetryTags.KurrentDB.Stream, stream }, - { TelemetryTags.KurrentDB.EventId, eventId }, - { TelemetryTags.KurrentDB.EventType, TestEventType }, - { TelemetryTags.Database.User, TestCredentials.Root.Username } + { TelemetryAttributes.MessagingSystem, TracingConstants.SystemName }, + { TelemetryAttributes.MessagingOperationName, TracingConstants.Operations.Process }, + { TelemetryAttributes.MessagingOperationType, TracingConstants.Operations.Process }, + { TelemetryAttributes.MessagingDestinationName, stream }, + { TelemetryAttributes.MessagingMessageId, eventId }, + { TrogonTelemetryAttributes.EventType, TestEventType } }; - if (subscriptionId != null) - expectedTags[TelemetryTags.KurrentDB.SubscriptionId] = subscriptionId; + if (consumerGroupName != null) + expectedTags[TelemetryAttributes.MessagingConsumerGroupName] = consumerGroupName; foreach (var tag in expectedTags) { activity.Tags.ShouldContain(tag); diff --git a/test/KurrentDB.Client.Tests.Common/Fixtures/KurrentDBPermanentFixture.cs b/test/KurrentDB.Client.Tests.Common/Fixtures/KurrentDBPermanentFixture.cs index a1bb17b5..61787348 100644 --- a/test/KurrentDB.Client.Tests.Common/Fixtures/KurrentDBPermanentFixture.cs +++ b/test/KurrentDB.Client.Tests.Common/Fixtures/KurrentDBPermanentFixture.cs @@ -3,7 +3,6 @@ using System.Net.Http; using KurrentDB.Client.Tests.FluentDocker; using Serilog; -using static System.TimeSpan; namespace KurrentDB.Client.Tests; @@ -11,6 +10,7 @@ namespace KurrentDB.Client.Tests; public partial class KurrentDBPermanentFixture : IAsyncLifetime, IAsyncDisposable { static readonly ILogger Logger; static readonly SemaphoreSlim WarmUpGatekeeper = new(1, 1); + static KurrentDBPermanentTestNode? SharedService; static KurrentDBPermanentFixture() { Logging.Initialize(); @@ -24,14 +24,13 @@ public KurrentDBPermanentFixture() : this(options => options) { } protected KurrentDBPermanentFixture(ConfigureFixture configure) { Options = configure(KurrentDBPermanentTestNode.DefaultOptions()); - Service = new KurrentDBPermanentTestNode(Options); } List TestRuns { get; } = new(); public ILogger Log => Logger; - public KurrentDBPermanentTestNode Service { get; } + public KurrentDBPermanentTestNode Service { get; private set; } = null!; public KurrentDBFixtureOptions Options { get; } public Faker Faker { get; } = new(); @@ -73,7 +72,13 @@ public async Task InitializeAsync() { await WarmUpGatekeeper.WaitAsync(); try { - await Service.Start(); + if (SharedService is null) { + var service = new KurrentDBPermanentTestNode(Options); + await service.Start(); + SharedService = service; + } + + Service = SharedService; DatabaseVersion = TestContainerService.Version; HasLastStreamPosition = (DatabaseVersion?.Major ?? int.MaxValue) >= 21; @@ -121,8 +126,6 @@ public async Task DisposeAsync() { // ignored } - await Service.DisposeAsync().AsTask().WithTimeout(FromMinutes(5)); - foreach (var testRunId in TestRuns) Logging.ReleaseLogs(testRunId); } diff --git a/test/KurrentDB.Client.Tests.Common/KurrentDB.Client.Tests.Common.csproj b/test/KurrentDB.Client.Tests.Common/KurrentDB.Client.Tests.Common.csproj index f2974d22..f98586de 100644 --- a/test/KurrentDB.Client.Tests.Common/KurrentDB.Client.Tests.Common.csproj +++ b/test/KurrentDB.Client.Tests.Common/KurrentDB.Client.Tests.Common.csproj @@ -10,6 +10,7 @@ + diff --git a/test/KurrentDB.Client.Tests/Diagnostics/AppendRecordsTracingTests.cs b/test/KurrentDB.Client.Tests/Diagnostics/AppendRecordsTracingTests.cs index 79017305..76cf97de 100644 --- a/test/KurrentDB.Client.Tests/Diagnostics/AppendRecordsTracingTests.cs +++ b/test/KurrentDB.Client.Tests/Diagnostics/AppendRecordsTracingTests.cs @@ -3,8 +3,8 @@ using System.Diagnostics; using KurrentDB.Client.Diagnostics; using KurrentDB.Client.Tests.Fixtures; -using KurrentDB.Diagnostics.Telemetry; using KurrentDB.Diagnostics; +using KurrentDB.Diagnostics.Telemetry; using KurrentDB.Diagnostics.Tracing; using static KurrentDB.Diagnostics.Tracing.TracingConstants; @@ -12,6 +12,7 @@ namespace KurrentDB.Client.Tests.Diagnostics; [Trait("Category", "Target:Diagnostics")] [Trait("Category", "Operation:AppendRecords")] +[Collection(DiagnosticsCollection.Name)] public class AppendRecordsTracingTests(ITestOutputHelper output, DiagnosticsFixture fixture) : KurrentDBPermanentTests(output, fixture) { [MinimumVersion.Fact(26, 1)] public async Task append_records_creates_trace_activity() { @@ -27,7 +28,7 @@ public async Task append_records_creates_trace_activity() { // Assert result.Position.ShouldBePositive(); - var appendActivities = Fixture.GetActivities(TracingConstants.Operations.MultiAppend, traceId); + var appendActivities = Fixture.GetActivities(TracingConstants.Operations.BatchAppend, traceId); appendActivities.ShouldNotBeEmpty(); appendActivities.Count.ShouldBe(1); @@ -35,14 +36,14 @@ public async Task append_records_creates_trace_activity() { var activity = appendActivities.First(); var expectedTags = new Dictionary { - { TelemetryTags.Database.System, KurrentDBClientDiagnostics.InstrumentationName }, - { TelemetryTags.Database.Operation, TracingConstants.Operations.MultiAppend }, - { TelemetryTags.Database.User, TestCredentials.Root.Username }, - { TelemetryTags.Otel.StatusCode, ActivityStatusCodeHelper.OkStatusCodeTagValue } + { TelemetryAttributes.DbSystemName, TracingConstants.SystemName }, + { TelemetryAttributes.DbOperationName, TracingConstants.Operations.BatchAppend } }; foreach (var tag in expectedTags) activity.Tags.ShouldContain(tag); + + activity.GetTagItem(TelemetryAttributes.DbOperationBatchSize).ShouldBe(3); } [MinimumVersion.Fact(26, 1)] @@ -61,7 +62,7 @@ public async Task append_records_with_exceptions_traces_error() { var rex = await appendTask.ShouldThrowAsync(); // Assert - var appendActivities = Fixture.GetActivities(TracingConstants.Operations.MultiAppend, traceId); + var appendActivities = Fixture.GetActivities(TracingConstants.Operations.BatchAppend, traceId); appendActivities.ShouldNotBeEmpty(); appendActivities.Count.ShouldBe(1); @@ -71,9 +72,9 @@ public async Task append_records_with_exceptions_traces_error() { activity.Events.ShouldHaveSingleItem(); var activityEvent = activity.Events.First(); - activityEvent.Name.ShouldBe(TelemetryTags.Exception.EventName); - activityEvent.Tags.Any(tag => tag.Key == TelemetryTags.Exception.Message).ShouldBeTrue(); - activityEvent.Tags.Any(tag => tag.Key == TelemetryTags.Exception.Stacktrace).ShouldBeTrue(); - activityEvent.Tags.Any(tag => tag.Key == TelemetryTags.Exception.Type && (string?)tag.Value == rex.GetType().FullName).ShouldBeTrue(); + activityEvent.Name.ShouldBe(TracingConstants.ExceptionEventName); + activityEvent.Tags.Any(tag => tag.Key == TelemetryAttributes.ExceptionMessage).ShouldBeTrue(); + activityEvent.Tags.Any(tag => tag.Key == TelemetryAttributes.ExceptionStacktrace).ShouldBeTrue(); + activityEvent.Tags.Any(tag => tag.Key == TelemetryAttributes.ExceptionType && (string?)tag.Value == rex.GetType().FullName).ShouldBeTrue(); } } diff --git a/test/KurrentDB.Client.Tests/Diagnostics/DiagnosticsCollection.cs b/test/KurrentDB.Client.Tests/Diagnostics/DiagnosticsCollection.cs new file mode 100644 index 00000000..b49aa5bd --- /dev/null +++ b/test/KurrentDB.Client.Tests/Diagnostics/DiagnosticsCollection.cs @@ -0,0 +1,6 @@ +namespace KurrentDB.Client.Tests.Diagnostics; + +[CollectionDefinition(Name, DisableParallelization = true)] +public sealed class DiagnosticsCollection { + public const string Name = "Diagnostics"; +} diff --git a/test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs b/test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs new file mode 100644 index 00000000..4d66869e --- /dev/null +++ b/test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs @@ -0,0 +1,293 @@ +using System.Diagnostics; +using System.Text.Json; +using KurrentDB.Client.Diagnostics; +using KurrentDB.Client.Extensions.OpenTelemetry; +using KurrentDB.Diagnostics.Telemetry; +using KurrentDB.Diagnostics.Tracing; +using OpenTelemetry; +using OpenTelemetry.Context.Propagation; +using OpenTelemetry.Trace; + +namespace KurrentDB.Client.Tests.Diagnostics; + +[Trait("Category", "Target:Diagnostics")] +[Collection(DiagnosticsCollection.Name)] +public class OpenTelemetryIntegrationTests { + static readonly TextMapPropagator TestPropagator = new CompositeTextMapPropagator([ + new TraceContextPropagator(), + new BaggagePropagator() + ]); + + [Fact] + public void public_registration_exports_client_activity() { + var exportedActivities = new List(); + using var provider = Sdk + .CreateTracerProviderBuilder() + .AddKurrentDBClientInstrumentation() + .AddInMemoryExporter(exportedActivities) + .Build(); + + using (KurrentDBClientDiagnostics.ActivitySource.StartActivity("test")) { } + + provider.ForceFlush(); + + var activity = exportedActivities.ShouldHaveSingleItem(); + KurrentDBClientDiagnostics.InstrumentationName.ShouldBe("TrogonEventStore.Client"); + activity.Source.Name.ShouldBe(KurrentDBClientDiagnostics.InstrumentationName); + } + + [Fact] + public async Task client_operation_exports_database_semantic_conventions() { + var exportedActivities = new List(); + using var provider = Sdk + .CreateTracerProviderBuilder() + .AddKurrentDBClientInstrumentation() + .AddInMemoryExporter(exportedActivities) + .Build(); + var tags = new ActivityTagsCollection { + { TelemetryAttributes.DbCollectionName, "orders" } + }; + + var result = await KurrentDBClientDiagnostics.ActivitySource.TraceClientOperation( + static () => ValueTask.FromResult(42), + TracingConstants.Operations.Append, + tags + ); + + provider.ForceFlush(); + result.ShouldBe(42); + var activity = exportedActivities.ShouldHaveSingleItem(); + activity.DisplayName.ShouldBe("append orders"); + activity.Kind.ShouldBe(ActivityKind.Client); + activity.GetTagItem(TelemetryAttributes.DbSystemName).ShouldBe(TracingConstants.SystemName); + activity.GetTagItem(TelemetryAttributes.DbOperationName).ShouldBe(TracingConstants.Operations.Append); + activity.GetTagItem(TelemetryAttributes.DbCollectionName).ShouldBe("orders"); + } + + [Fact] + public async Task operation_specific_tags_do_not_mutate_an_unobserved_caller() { + using var source = new ActivitySource($"unobserved-{Guid.NewGuid():N}"); + using var caller = new Activity("caller").Start(); + + await source.TraceClientOperation( + activity => { + activity?.SetTag(TelemetryAttributes.DbOperationBatchSize, 2); + return ValueTask.FromResult(0); + }, + TracingConstants.Operations.BatchAppend + ); + + caller.GetTagItem(TelemetryAttributes.DbOperationBatchSize).ShouldBeNull(); + } + + [Fact] + public async Task operation_specific_tags_are_owned_by_the_client_span() { + using var source = new ActivitySource($"observed-{Guid.NewGuid():N}"); + Activity? completedActivity = null; + using var listener = new ActivityListener { + ShouldListenTo = candidate => candidate == source, + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + ActivityStopped = activity => completedActivity = activity + }; + ActivitySource.AddActivityListener(listener); + using var caller = new Activity("caller").Start(); + + await source.TraceClientOperation( + activity => { + activity?.SetTag(TelemetryAttributes.DbOperationBatchSize, 2); + return ValueTask.FromResult(0); + }, + TracingConstants.Operations.BatchAppend + ); + + completedActivity.ShouldNotBeNull() + .GetTagItem(TelemetryAttributes.DbOperationBatchSize) + .ShouldBe(2); + caller.GetTagItem(TelemetryAttributes.DbOperationBatchSize).ShouldBeNull(); + } + + [Fact] + public void configured_propagator_injects_context_into_json_metadata() { + var originalPropagator = Propagators.DefaultTextMapPropagator; + var originalBaggage = Baggage.Current; + + try { + Sdk.SetDefaultTextMapPropagator(TestPropagator); + Baggage.Current = Baggage.Create(new Dictionary { ["tenant"] = "straw-hat" }); + + using var activity = StartActivity(ActivityTraceFlags.None, "vendor=value"); + ReadOnlyMemory metadata = "{\"custom\":\"value\"}"u8.ToArray(); + + var injected = metadata.InjectTracingContext(activity).ToArray(); + var extracted = Extract(injected); + + extracted.ActivityContext.TraceId.ShouldBe(activity.TraceId); + extracted.ActivityContext.SpanId.ShouldBe(activity.SpanId); + extracted.ActivityContext.TraceFlags.ShouldBe(ActivityTraceFlags.None); + extracted.ActivityContext.TraceState.ShouldBe("vendor=value"); + extracted.Baggage.GetBaggage("tenant").ShouldBe("straw-hat"); + + using var document = JsonDocument.Parse(injected); + document.RootElement.GetProperty("custom").GetString().ShouldBe("value"); + document.RootElement.TryGetProperty("traceparent", out _).ShouldBeTrue(); + document.RootElement.TryGetProperty("tracestate", out _).ShouldBeTrue(); + document.RootElement.TryGetProperty("baggage", out _).ShouldBeTrue(); + } finally { + Baggage.Current = originalBaggage; + Sdk.SetDefaultTextMapPropagator(originalPropagator); + } + } + + [Fact] + public void configured_propagator_injects_context_into_property_metadata() { + var originalPropagator = Propagators.DefaultTextMapPropagator; + var originalBaggage = Baggage.Current; + + try { + Sdk.SetDefaultTextMapPropagator(TestPropagator); + Baggage.Current = Baggage.Create(new Dictionary { ["tenant"] = "straw-hat" }); + + using var activity = StartActivity(ActivityTraceFlags.Recorded, "vendor=value"); + var metadata = new Dictionary { ["custom"] = "value" }; + + metadata.InjectTracingContext(activity); + var extracted = TestPropagator.Extract(default, metadata, Getter); + + extracted.ActivityContext.TraceId.ShouldBe(activity.TraceId); + extracted.ActivityContext.SpanId.ShouldBe(activity.SpanId); + extracted.ActivityContext.TraceFlags.ShouldBe(ActivityTraceFlags.Recorded); + extracted.ActivityContext.TraceState.ShouldBe("vendor=value"); + extracted.Baggage.GetBaggage("tenant").ShouldBe("straw-hat"); + metadata["custom"].ShouldBe("value"); + } finally { + Baggage.Current = originalBaggage; + Sdk.SetDefaultTextMapPropagator(originalPropagator); + } + } + + [Fact] + public void configured_propagator_replaces_mixed_case_fields_in_property_metadata() { + var originalPropagator = Propagators.DefaultTextMapPropagator; + var originalBaggage = Baggage.Current; + + try { + Sdk.SetDefaultTextMapPropagator(TestPropagator); + Baggage.Current = Baggage.Create(new Dictionary { ["tenant"] = "straw-hat" }); + + using var activity = StartActivity(ActivityTraceFlags.Recorded, "vendor=value"); + var metadata = new Dictionary { + ["TraceParent"] = "stale", + ["TraceState"] = "stale", + ["Baggage"] = "stale" + }; + + metadata.InjectTracingContext(activity); + + foreach (var name in new[] { "traceparent", "tracestate", "baggage" }) + metadata.Keys.Count(key => string.Equals(key, name, StringComparison.OrdinalIgnoreCase)).ShouldBe(1); + + var extracted = TestPropagator.Extract(default, metadata, Getter); + extracted.ActivityContext.TraceId.ShouldBe(activity.TraceId); + extracted.ActivityContext.SpanId.ShouldBe(activity.SpanId); + extracted.ActivityContext.TraceState.ShouldBe("vendor=value"); + extracted.Baggage.GetBaggage("tenant").ShouldBe("straw-hat"); + } finally { + Baggage.Current = originalBaggage; + Sdk.SetDefaultTextMapPropagator(originalPropagator); + } + } + + [Fact] + public void configured_propagator_extracts_context_from_json_metadata() { + var originalPropagator = Propagators.DefaultTextMapPropagator; + + try { + Sdk.SetDefaultTextMapPropagator(TestPropagator); + var activityContext = new ActivityContext( + ActivityTraceId.CreateRandom(), + ActivitySpanId.CreateRandom(), + ActivityTraceFlags.None, + "vendor=value" + ); + var expected = new PropagationContext( + activityContext, + Baggage.Create(new Dictionary { ["tenant"] = "straw-hat" }) + ); + var carrier = new Dictionary(); + TestPropagator.Inject(expected, carrier, static (metadata, name, value) => metadata[name] = value); + + ReadOnlyMemory metadata = JsonSerializer.SerializeToUtf8Bytes(carrier); + var extracted = metadata.ExtractPropagationContext(); + + extracted.ActivityContext.TraceId.ShouldBe(activityContext.TraceId); + extracted.ActivityContext.SpanId.ShouldBe(activityContext.SpanId); + extracted.ActivityContext.TraceFlags.ShouldBe(ActivityTraceFlags.None); + extracted.ActivityContext.TraceState.ShouldBe("vendor=value"); + extracted.ActivityContext.IsRemote.ShouldBeTrue(); + extracted.Baggage.GetBaggage("tenant").ShouldBe("straw-hat"); + } finally { + Sdk.SetDefaultTextMapPropagator(originalPropagator); + } + } + + [Fact] + public void configured_propagator_extracts_mixed_case_fields_from_json_metadata() { + var originalPropagator = Propagators.DefaultTextMapPropagator; + + try { + Sdk.SetDefaultTextMapPropagator(TestPropagator); + var activityContext = new ActivityContext( + ActivityTraceId.CreateRandom(), + ActivitySpanId.CreateRandom(), + ActivityTraceFlags.Recorded, + "vendor=value" + ); + var expected = new PropagationContext( + activityContext, + Baggage.Create(new Dictionary { ["tenant"] = "straw-hat" }) + ); + var carrier = new Dictionary(); + TestPropagator.Inject(expected, carrier, static (metadata, name, value) => metadata[name] = value); + + var mixedCaseCarrier = carrier.ToDictionary( + pair => pair.Key switch { + "traceparent" => "TraceParent", + "tracestate" => "TraceState", + "baggage" => "Baggage", + _ => pair.Key + }, + pair => pair.Value + ); + ReadOnlyMemory metadata = JsonSerializer.SerializeToUtf8Bytes(mixedCaseCarrier); + var extracted = metadata.ExtractPropagationContext(); + + extracted.ActivityContext.TraceId.ShouldBe(activityContext.TraceId); + extracted.ActivityContext.SpanId.ShouldBe(activityContext.SpanId); + extracted.ActivityContext.TraceFlags.ShouldBe(ActivityTraceFlags.Recorded); + extracted.ActivityContext.TraceState.ShouldBe("vendor=value"); + extracted.ActivityContext.IsRemote.ShouldBeTrue(); + extracted.Baggage.GetBaggage("tenant").ShouldBe("straw-hat"); + } finally { + Sdk.SetDefaultTextMapPropagator(originalPropagator); + } + } + + static Activity StartActivity(ActivityTraceFlags traceFlags, string traceState) { + var activity = new Activity("parent") + .SetParentId(ActivityTraceId.CreateRandom(), ActivitySpanId.CreateRandom(), traceFlags); + activity.TraceStateString = traceState; + return activity.Start(); + } + + static PropagationContext Extract(ReadOnlyMemory metadata) { + using var document = JsonDocument.Parse(metadata); + var carrier = document.RootElement + .EnumerateObject() + .ToDictionary(property => property.Name, property => property.Value.GetString()!); + + return TestPropagator.Extract(default, carrier, Getter); + } + + static IEnumerable Getter(Dictionary carrier, string name) => + carrier.TryGetValue(name, out var value) ? [value] : []; +} diff --git a/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs b/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs index d1b12632..658703eb 100644 --- a/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs +++ b/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs @@ -6,6 +6,7 @@ namespace KurrentDB.Client.Tests.Diagnostics; [Trait("Category", "Target:Diagnostics")] +[Collection(DiagnosticsCollection.Name)] public class PersistentSubscriptionsTracingInstrumentationTests(ITestOutputHelper output, DiagnosticsFixture fixture) : KurrentDBPermanentTests(output, fixture) { [RetryFact] @@ -36,11 +37,15 @@ await Fixture.Streams.AppendToStreamAsync( .ShouldNotBeNull(); var subscribeActivities = Fixture - .GetActivities(TracingConstants.Operations.Subscribe, traceId, stream) + .GetActivities(TracingConstants.Operations.Process, traceId, stream) + .Where(activity => Equals( + activity.GetTagItem(TelemetryAttributes.MessagingConsumerGroupName), + groupName + )) .ToArray(); var expectedEventIds = events.Select(@event => @event.EventId.ToString()).ToHashSet(); var actualEventIds = subscribeActivities - .Select(activity => Assert.IsType(activity.GetTagItem(TelemetryTags.KurrentDB.EventId))) + .Select(activity => Assert.IsType(activity.GetTagItem(TelemetryAttributes.MessagingMessageId))) .ToHashSet(); subscriptionId.ShouldNotBeNull(); @@ -51,16 +56,13 @@ await Fixture.Streams.AppendToStreamAsync( subscribeActivity.TraceId.ShouldBe(appendActivity.Context.TraceId); subscribeActivity.ParentSpanId.ShouldBe(appendActivity.Context.SpanId); subscribeActivity.HasRemoteParent.ShouldBeTrue(); - Assert.False( - string.IsNullOrWhiteSpace( - Assert.IsType(subscribeActivity.GetTagItem(TelemetryTags.KurrentDB.SubscriptionId)) - ) - ); + subscribeActivity.GetTagItem(TelemetryAttributes.MessagingConsumerGroupName).ShouldBe(groupName); Fixture.AssertSubscriptionActivityHasExpectedTags( subscribeActivity, stream, - Assert.IsType(subscribeActivity.GetTagItem(TelemetryTags.KurrentDB.EventId)) + Assert.IsType(subscribeActivity.GetTagItem(TelemetryAttributes.MessagingMessageId)), + groupName ); } @@ -68,7 +70,7 @@ await Fixture.Streams.AppendToStreamAsync( async Task Subscribe() { await using var subscription = Fixture.Subscriptions.SubscribeToStream(stream, groupName); - await using var enumerator = subscription.Messages.GetAsyncEnumerator(); + await using var enumerator = subscription.Messages.GetAsyncEnumerator(); var remainingEventIds = events.Select(@event => @event.EventId).ToHashSet(); while (await enumerator.MoveNextAsync()) { @@ -113,7 +115,7 @@ await Fixture.Streams.AppendToStreamAsync( async Task Subscribe() { await using var subscription = Fixture.Subscriptions.SubscribeToStream(stream, groupName); - await using var enumerator = subscription.Messages.GetAsyncEnumerator(); + await using var enumerator = subscription.Messages.GetAsyncEnumerator(); var eventsAppeared = 0; while (await enumerator.MoveNextAsync()) { @@ -153,7 +155,7 @@ await Fixture.Streams.AppendToStreamAsync( async Task Subscribe() { await using var subscription = Fixture.Subscriptions.SubscribeToStream(stream, groupName); - await using var enumerator = subscription.Messages.GetAsyncEnumerator(); + await using var enumerator = subscription.Messages.GetAsyncEnumerator(); var eventsAppeared = 0; while (await enumerator.MoveNextAsync()) { diff --git a/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs b/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs index d628bbc8..c51263a4 100644 --- a/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs +++ b/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs @@ -6,15 +6,17 @@ using KurrentDB.Client.Tests.Fixtures; using KurrentDB.Diagnostics.Telemetry; using KurrentDB.Diagnostics.Tracing; -using static KurrentDB.Diagnostics.Tracing.TracingConstants; +using OpenTelemetry; +using OpenTelemetry.Context.Propagation; namespace KurrentDB.Client.Tests.Diagnostics; [Trait("Category", "Target:Diagnostics")] +[Collection(DiagnosticsCollection.Name)] public class StreamsTracingInstrumentationTests(ITestOutputHelper output, DiagnosticsFixture fixture) : KurrentDBPermanentTests(output, fixture) { [Fact] public void trace_contexts_are_independent() { - var first = Fixture.CreateTraceId(); + var first = Fixture.CreateTraceId(); var second = Fixture.CreateTraceId(); Assert.NotEqual(first, second); @@ -69,8 +71,8 @@ public async Task multi_stream_append() { // Assert appendResult.Position.ShouldBePositive(); - var appendActivities = Fixture.GetActivities(TracingConstants.Operations.MultiAppend, traceId); - var subscribeActivities = Fixture.GetActivities(TracingConstants.Operations.Subscribe, traceId); + var appendActivities = Fixture.GetActivities(TracingConstants.Operations.BatchAppend, traceId); + var subscribeActivities = Fixture.GetActivities(TracingConstants.Operations.Process, traceId); appendActivities.ShouldNotBeEmpty(); subscribeActivities.ShouldNotBeEmpty(); @@ -130,7 +132,7 @@ public async Task multi_stream_append_with_exceptions() { var rex = await appendTask.ShouldThrowAsync(); // Assert - var appendActivities = Fixture.GetActivities(TracingConstants.Operations.MultiAppend, traceId); + var appendActivities = Fixture.GetActivities(TracingConstants.Operations.BatchAppend, traceId); appendActivities.ShouldNotBeEmpty(); @@ -142,10 +144,10 @@ public async Task multi_stream_append_with_exceptions() { var activityEvent = activity.Events.First(); - activityEvent.Name.ShouldBe(TelemetryTags.Exception.EventName); - activityEvent.Tags.Any(tag => tag.Key == TelemetryTags.Exception.Message).ShouldBeTrue(); - activityEvent.Tags.Any(tag => tag.Key == TelemetryTags.Exception.Stacktrace).ShouldBeTrue(); - activityEvent.Tags.Any(tag => tag.Key == TelemetryTags.Exception.Type && (string?)tag.Value == rex.GetType().FullName).ShouldBeTrue(); + activityEvent.Name.ShouldBe(TracingConstants.ExceptionEventName); + activityEvent.Tags.Any(tag => tag.Key == TelemetryAttributes.ExceptionMessage).ShouldBeTrue(); + activityEvent.Tags.Any(tag => tag.Key == TelemetryAttributes.ExceptionStacktrace).ShouldBeTrue(); + activityEvent.Tags.Any(tag => tag.Key == TelemetryAttributes.ExceptionType && (string?)tag.Value == rex.GetType().FullName).ShouldBeTrue(); } [Fact] @@ -187,11 +189,59 @@ await Fixture.Streams.AppendToStreamAsync( .ReadStreamAsync(Direction.Forwards, stream, StreamPosition.Start) .ToListAsync(); - var tracingMetadata = readResult[0].OriginalEvent.Metadata.ExtractTracingMetadata(); + var propagationContext = readResult[0].OriginalEvent.Metadata.ExtractPropagationContext(); - tracingMetadata.ShouldNotBe(TracingMetadata.None); - tracingMetadata.TraceId.ShouldBe(activity.TraceId.ToString()); - tracingMetadata.SpanId.ShouldBe(activity.SpanId.ToString()); + propagationContext.ActivityContext.ShouldNotBe(default); + propagationContext.ActivityContext.TraceId.ShouldBe(activity.TraceId); + propagationContext.ActivityContext.SpanId.ShouldBe(activity.SpanId); + } + + [Fact] + public async Task subscription_restores_tracestate_and_baggage() { + var traceId = Fixture.CreateTraceId(); + Activity.Current!.TraceStateString = "vendor=value"; + var originalBaggage = Baggage.Current; + var originalPropagator = Propagators.DefaultTextMapPropagator; + var stream = Fixture.GetStreamName(); + var seedEvent = Fixture.CreateTestEvent(metadata: Fixture.CreateTestJsonMetadata()); + + try { + Sdk.SetDefaultTextMapPropagator(new CompositeTextMapPropagator([ + new TraceContextPropagator(), + new BaggagePropagator() + ])); + Baggage.Current = Baggage.Create(new Dictionary { ["tenant"] = "straw-hat" }); + await Fixture.Streams.AppendToStreamAsync(stream, StreamState.NoStream, [seedEvent]); + + await using var subscription = Fixture.Streams.SubscribeToStream(stream, FromStream.Start); + await using var enumerator = subscription.Messages.GetAsyncEnumerator(); + + Assert.True(await enumerator.MoveNextAsync()); + Assert.IsType(enumerator.Current); + Assert.True(await enumerator.MoveNextAsync()); + Assert.IsType(enumerator.Current); + + var appendActivity = Fixture + .GetActivities(TracingConstants.Operations.Append, traceId) + .ShouldHaveSingleItem(); + var subscriptionActivities = Fixture + .GetActivities(TracingConstants.Operations.Process, traceId, stream) + .Where(activity => Equals(activity.GetTagItem(TelemetryAttributes.MessagingMessageId), seedEvent.EventId.ToString())) + .ToArray(); + + Assert.NotEmpty(subscriptionActivities); + Assert.All( + subscriptionActivities, + subscriptionActivity => { + subscriptionActivity.ParentSpanId.ShouldBe(appendActivity.SpanId); + subscriptionActivity.TraceStateString.ShouldBe("vendor=value"); + subscriptionActivity.Baggage.ShouldContain(new KeyValuePair("tenant", "straw-hat")); + } + ); + } finally { + Baggage.Current = originalBaggage; + Sdk.SetDefaultTextMapPropagator(originalPropagator); + } } [Fact] @@ -214,7 +264,7 @@ await Fixture.Streams.AppendToStreamAsync( } [Fact] - public async Task tracing_context_not_duplicated_when_already_present() { + public async Task tracing_context_replaced_when_already_present() { // Arrange var stream = Fixture.GetStreamName(); @@ -222,8 +272,7 @@ public async Task tracing_context_not_duplicated_when_already_present() { activity.Start(); var metadata = new Dictionary { - [Metadata.TraceId] = activity.TraceId.ToString(), - [Metadata.SpanId] = activity.SpanId.ToString() + ["traceparent"] = $"00-{activity.TraceId}-{activity.SpanId}-00" }; // Act @@ -234,11 +283,19 @@ public async Task tracing_context_not_duplicated_when_already_present() { .ReadStreamAsync(Direction.Forwards, stream, StreamPosition.Start, maxCount: 1) .ToListAsync(); - var tracingMetadata = result.First().OriginalEvent.Metadata.ExtractTracingMetadata(); + var outputMetadata = result.First().OriginalEvent.Metadata; + var propagationContext = outputMetadata.ExtractPropagationContext(); + var appendActivity = Fixture + .GetActivities(TracingConstants.Operations.Append, activity.TraceId) + .ShouldHaveSingleItem(); - tracingMetadata.ShouldNotBe(TracingMetadata.None); - tracingMetadata.TraceId.ShouldBe(activity.TraceId.ToString()); - tracingMetadata.SpanId.ShouldBe(activity.SpanId.ToString()); + propagationContext.ActivityContext.ShouldNotBe(default); + propagationContext.ActivityContext.TraceId.ShouldBe(appendActivity.TraceId); + propagationContext.ActivityContext.SpanId.ShouldBe(appendActivity.SpanId); + propagationContext.ActivityContext.TraceFlags.ShouldBe(appendActivity.ActivityTraceFlags); + + using var document = JsonDocument.Parse(outputMetadata); + document.RootElement.EnumerateObject().Count(property => property.Name == "traceparent").ShouldBe(1); } [Fact] @@ -283,7 +340,7 @@ public async Task json_metadata_traced_non_json_metadata_not_traced() { await Fixture.Streams.AppendToStreamAsync(streamName, StreamState.NoStream, seedEvents); await using var subscription = Fixture.Streams.SubscribeToStream(streamName, FromStream.Start); - await using var enumerator = subscription.Messages.GetAsyncEnumerator(); + await using var enumerator = subscription.Messages.GetAsyncEnumerator(); var appendActivities = Fixture .GetActivities(TracingConstants.Operations.Append, traceId) @@ -296,7 +353,7 @@ public async Task json_metadata_traced_non_json_metadata_not_traced() { await Subscribe(enumerator).WithTimeout(); var subscribeActivities = Fixture - .GetActivities(TracingConstants.Operations.Subscribe, traceId, streamName) + .GetActivities(TracingConstants.Operations.Process, traceId, streamName) .ToArray(); appendActivities.ShouldHaveSingleItem(); @@ -334,7 +391,7 @@ async Task Subscribe(IAsyncEnumerator internalEnumerator) { [Trait("Category", "Special cases")] public async Task no_trace_when_event_is_null() { var traceId = Fixture.CreateTraceId(); - var category = Guid.NewGuid().ToString("N"); + var category = Guid.NewGuid().ToString("N"); var streamName = category + "-123"; var categoryStream = "$ce-" + category; @@ -358,7 +415,7 @@ public async Task no_trace_when_event_is_null() { .ShouldNotBeNull(); var subscribeActivities = Fixture - .GetActivities(TracingConstants.Operations.Subscribe, traceId, categoryStream) + .GetActivities(TracingConstants.Operations.Process, traceId, categoryStream) .ToArray(); appendActivities.ShouldHaveSingleItem(); diff --git a/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj b/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj index cafac282..6f82ad1a 100644 --- a/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj +++ b/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj @@ -7,6 +7,8 @@ + +