From c142606dab0d2344893aa288e36dac6538871dc8 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sun, 30 Aug 2026 19:13:35 -0400 Subject: [PATCH] feat(compatibility): protect client interoperability Signed-off-by: Yordis Prieto --- .github/workflows/dotnet.yml | 3 +- .../CompatibilityApplication.cs | 228 ++++++++++++++++++ .../CompatibilityCommand.cs | 49 ++++ .../CompatibilityContract.cs | 11 + .../CompatibilityEvent.cs | 8 + .../CompatibilityOptions.cs | 58 +++++ .../CompatibilityProgram.cs | 51 ++++ .../Program.cs | 3 + ...ogonEventStore.Client.Compatibility.csproj | 19 ++ .../trogon/eventstore/client/spans.yaml | 11 +- .../csharp/subscription-trace-semantics.cs.j2 | 10 + .../templates/registry/csharp/weaver.yaml | 13 + .../Diagnostics/ActivitySourceExtensions.cs | 123 ++++++---- .../Diagnostics/EventMetadataExtensions.cs | 12 +- .../Generated/SubscriptionTraceSemantics.g.cs | 10 + .../Diagnostics/Tracing/TracingConstants.cs | 1 - ...entDBPersistentSubscriptionsClient.Read.cs | 37 ++- .../Streams/KurrentDBClient.Subscriptions.cs | 37 ++- .../Fixtures/DiagnosticsFixture.cs | 8 +- .../CompatibilityApplicationTests.cs | 89 +++++++ .../CompatibilityCommandTests.cs | 51 ++++ .../CompatibilityOptionsTests.cs | 63 +++++ .../OpenTelemetryIntegrationTests.cs | 11 +- ...ubscriptionsTracingInstrumentationTests.cs | 16 +- .../StreamsTracingInstrumentationTests.cs | 213 ++++++++++------ .../KurrentDB.Client.Tests.csproj | 1 + 26 files changed, 978 insertions(+), 158 deletions(-) create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/CompatibilityApplication.cs create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/CompatibilityCommand.cs create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/CompatibilityContract.cs create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/CompatibilityEvent.cs create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/CompatibilityOptions.cs create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/CompatibilityProgram.cs create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/Program.cs create mode 100644 compatibility/TrogonEventStore.Client.Compatibility/TrogonEventStore.Client.Compatibility.csproj create mode 100644 otel/semconv/templates/registry/csharp/subscription-trace-semantics.cs.j2 create mode 100644 src/KurrentDB.Client/Core/Common/Diagnostics/Generated/SubscriptionTraceSemantics.g.cs create mode 100644 test/KurrentDB.Client.Tests/Compatibility/CompatibilityApplicationTests.cs create mode 100644 test/KurrentDB.Client.Tests/Compatibility/CompatibilityCommandTests.cs create mode 100644 test/KurrentDB.Client.Tests/Compatibility/CompatibilityOptionsTests.cs diff --git a/.github/workflows/dotnet.yml b/.github/workflows/dotnet.yml index e407bca..73be8bf 100644 --- a/.github/workflows/dotnet.yml +++ b/.github/workflows/dotnet.yml @@ -64,7 +64,8 @@ jobs: FullyQualifiedName~KurrentDB.Client.Tests.PositionTests.| FullyQualifiedName~KurrentDB.Client.Tests.NodePreferenceComparerTests.| FullyQualifiedName~KurrentDB.Client.Tests.NodeSelectorTests.| - FullyQualifiedName~KurrentDB.Client.Tests.StreamPositionTests. + FullyQualifiedName~KurrentDB.Client.Tests.StreamPositionTests.| + FullyQualifiedName~KurrentDB.Client.Tests.Compatibility. steps: - name: Checkout uses: actions/checkout@v5 diff --git a/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityApplication.cs b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityApplication.cs new file mode 100644 index 0000000..4dd3cb9 --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityApplication.cs @@ -0,0 +1,228 @@ +using System.Text.Json; +using KurrentDB.Client; + +namespace TrogonEventStore.Client.Compatibility; + +internal sealed class CompatibilityApplication(CompatibilityOptions options) { + static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web); + enum AppendRpc { Standard, Batch } + + public async Task ExecuteAsync(CompatibilityCommand command, CancellationToken cancellationToken) { + switch (command) { + case CompatibilityCommand.Write write: + await Write(write.Stream, AppendRpc.Standard, "write", cancellationToken); + break; + case CompatibilityCommand.BatchWrite write: + await Write(write.Stream, AppendRpc.Batch, "batch-write", cancellationToken); + break; + case CompatibilityCommand.Read read: + await Read(read.Stream, cancellationToken); + break; + case CompatibilityCommand.Subscribe subscribe: + await Subscribe(subscribe.Stream, cancellationToken); + break; + case CompatibilityCommand.CreatePersistentSubscription create: + await CreatePersistentSubscription(create.Stream, create.Group, cancellationToken); + break; + case CompatibilityCommand.ConsumePersistentSubscription consume: + await ConsumePersistentSubscription(consume.Stream, consume.Group, cancellationToken); + break; + } + } + + async Task Write( + StreamName stream, + AppendRpc appendRpc, + string command, + CancellationToken cancellationToken + ) { + var settings = CreateClientSettings(); + await using var client = new KurrentDBClient(settings); + var userCredentials = appendRpc switch { + AppendRpc.Standard => settings.DefaultCredentials ?? throw new ArgumentException( + $"Environment variable {CompatibilityContract.ServerUriName} must include credentials for write." + ), + AppendRpc.Batch => null, + _ => throw new ArgumentOutOfRangeException(nameof(appendRpc), appendRpc, null) + }; + var payload = new CompatibilityEvent(CompatibilityContract.Producer, options.RunId); + var eventData = new EventData( + Uuid.NewUuid(), + CompatibilityContract.EventType, + JsonSerializer.SerializeToUtf8Bytes(payload, JsonOptions), + "{}"u8.ToArray() + ); + + await client.AppendToStreamAsync( + stream.Value, + StreamState.NoStream, + [eventData], + userCredentials: userCredentials, + cancellationToken: cancellationToken + ); + + WriteResult(command, stream, null, payload, eventData.EventId.ToString()); + } + + async Task Read(StreamName stream, CancellationToken cancellationToken) { + await using var client = CreateStreamsClient(); + var result = client.ReadStreamAsync( + Direction.Forwards, + stream.Value, + StreamPosition.Start, + cancellationToken: cancellationToken + ); + + await foreach (var resolvedEvent in result.WithCancellation(cancellationToken)) { + var payload = ReadPayload(resolvedEvent); + if (payload is null) + continue; + + WriteResult("read", stream, null, payload, resolvedEvent.Event.EventId.ToString()); + return; + } + + throw MissingEvent(stream); + } + + async Task Subscribe(StreamName stream, CancellationToken cancellationToken) { + await using var client = CreateStreamsClient(); + await using var subscription = client.SubscribeToStream( + stream.Value, + FromStream.Start, + cancellationToken: cancellationToken + ); + + await foreach (var message in subscription.Messages.WithCancellation(cancellationToken)) { + if (message is StreamMessage.SubscriptionConfirmation) { + await SignalReady(cancellationToken); + continue; + } + + if (message is not StreamMessage.Event(var resolvedEvent)) + continue; + + var payload = ReadPayload(resolvedEvent); + if (payload is null) + continue; + + WriteResult("subscribe", stream, null, payload, resolvedEvent.Event.EventId.ToString()); + return; + } + + throw MissingEvent(stream); + } + + async Task CreatePersistentSubscription( + StreamName stream, + GroupName group, + CancellationToken cancellationToken + ) { + await using var client = CreatePersistentSubscriptionsClient(); + await client.CreateToStreamAsync( + stream.Value, + group.Value, + new(startFrom: StreamPosition.Start), + cancellationToken: cancellationToken + ); + + WriteResult( + "create-persistent-subscription", + stream, + group, + new(CompatibilityContract.Producer, options.RunId), + null + ); + } + + async Task ConsumePersistentSubscription( + StreamName stream, + GroupName group, + CancellationToken cancellationToken + ) { + await using var client = CreatePersistentSubscriptionsClient(); + await using var subscription = client.SubscribeToStream( + stream.Value, + group.Value, + cancellationToken: cancellationToken + ); + + await foreach (var message in subscription.Messages.WithCancellation(cancellationToken)) { + if (message is PersistentSubscriptionMessage.SubscriptionConfirmation) { + await SignalReady(cancellationToken); + continue; + } + + if (message is not PersistentSubscriptionMessage.Event(var resolvedEvent, _)) + continue; + + await subscription.Ack(resolvedEvent); + var payload = ReadPayload(resolvedEvent); + if (payload is null) + continue; + + WriteResult( + "consume-persistent-subscription", + stream, + group, + payload, + resolvedEvent.Event.EventId.ToString() + ); + return; + } + + throw MissingEvent(stream); + } + + CompatibilityEvent? ReadPayload(ResolvedEvent resolvedEvent) { + if (resolvedEvent.Event.EventType != CompatibilityContract.EventType) + return null; + + var payload = JsonSerializer.Deserialize(resolvedEvent.Event.Data.Span, JsonOptions) + ?? throw new InvalidDataException("Compatibility event payload is required."); + + if (payload.RunId != options.RunId) + return null; + + if (string.IsNullOrWhiteSpace(payload.Producer)) + throw new InvalidDataException("Compatibility event producer is required."); + + return payload; + } + + KurrentDBClientSettings CreateClientSettings() => + KurrentDBClientSettings.Create(options.ServerUri.OriginalString); + + KurrentDBClient CreateStreamsClient() => new(CreateClientSettings()); + + KurrentDBPersistentSubscriptionsClient CreatePersistentSubscriptionsClient() => + new(CreateClientSettings()); + + async Task SignalReady(CancellationToken cancellationToken) { + if (options.ReadyFile is { } readyFile) + await readyFile.Signal(cancellationToken); + } + + void WriteResult( + string command, + StreamName stream, + GroupName? group, + CompatibilityEvent payload, + string? eventId + ) => Console.WriteLine(JsonSerializer.Serialize( + new CompatibilityResult(command, stream.Value, group?.Value, payload.Producer, payload.RunId, eventId), + JsonOptions + )); + + InvalidDataException MissingEvent(StreamName stream) => + new($"Stream {stream.Value} does not contain a {CompatibilityContract.EventType} event for run {options.RunId}."); +} + +internal sealed record CompatibilityResult( + string Command, + string Stream, + string? Group, + string Producer, + string RunId, + string? EventId +); diff --git a/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityCommand.cs b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityCommand.cs new file mode 100644 index 0000000..3e8a9d7 --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityCommand.cs @@ -0,0 +1,49 @@ +namespace TrogonEventStore.Client.Compatibility; + +internal abstract record CompatibilityCommand(StreamName Stream) { + public sealed record Write(StreamName Stream) : CompatibilityCommand(Stream); + public sealed record BatchWrite(StreamName Stream) : CompatibilityCommand(Stream); + public sealed record Read(StreamName Stream) : CompatibilityCommand(Stream); + public sealed record Subscribe(StreamName Stream) : CompatibilityCommand(Stream); + public sealed record CreatePersistentSubscription(StreamName Stream, GroupName Group) : CompatibilityCommand(Stream); + public sealed record ConsumePersistentSubscription(StreamName Stream, GroupName Group) : CompatibilityCommand(Stream); + + public static CompatibilityCommand Parse(IReadOnlyList arguments) => arguments switch { + ["write", var stream] => new Write(StreamName.Parse(stream)), + ["batch-write", var stream] => new BatchWrite(StreamName.Parse(stream)), + ["read", var stream] => new Read(StreamName.Parse(stream)), + ["subscribe", var stream] => new Subscribe(StreamName.Parse(stream)), + ["create-persistent-subscription", var stream, var group] => + new CreatePersistentSubscription(StreamName.Parse(stream), GroupName.Parse(group)), + ["consume-persistent-subscription", var stream, var group] => + new ConsumePersistentSubscription(StreamName.Parse(stream), GroupName.Parse(group)), + _ => throw new ArgumentException( + "Expected write , batch-write , read , subscribe , " + + "create-persistent-subscription , or " + + "consume-persistent-subscription ." + ) + }; +} + +internal readonly record struct StreamName { + StreamName(string value) => Value = value; + + public string Value { get; } + + public static StreamName Parse(string value) => + new(Required(value, "Stream name")); + + internal static string Required(string value, string label) => + !string.IsNullOrWhiteSpace(value) + ? value.Trim() + : throw new ArgumentException($"{label} is required."); +} + +internal readonly record struct GroupName { + GroupName(string value) => Value = value; + + public string Value { get; } + + public static GroupName Parse(string value) => + new(StreamName.Required(value, "Group name")); +} diff --git a/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityContract.cs b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityContract.cs new file mode 100644 index 0000000..d43b801 --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityContract.cs @@ -0,0 +1,11 @@ +namespace TrogonEventStore.Client.Compatibility; + +internal static class CompatibilityContract { + public const string EventType = "trogon-compatibility"; + public const string Producer = "dotnet"; + public const string ServiceName = "trogon-eventstore-client-dotnet"; + public const string ServerUriName = "TROGON_EVENTSTORE_URI"; + public const string RunIdName = "TROGON_EVENTSTORE_RUN_ID"; + public const string OtlpEndpointName = "OTEL_EXPORTER_OTLP_ENDPOINT"; + public const string ReadyFileName = "TROGON_EVENTSTORE_READY_FILE"; +} diff --git a/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityEvent.cs b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityEvent.cs new file mode 100644 index 0000000..ac2bcab --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityEvent.cs @@ -0,0 +1,8 @@ +using System.Text.Json.Serialization; + +namespace TrogonEventStore.Client.Compatibility; + +internal sealed record CompatibilityEvent( + [property: JsonPropertyName("producer")] string Producer, + [property: JsonPropertyName("runId")] string RunId +); diff --git a/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityOptions.cs b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityOptions.cs new file mode 100644 index 0000000..beca489 --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityOptions.cs @@ -0,0 +1,58 @@ +namespace TrogonEventStore.Client.Compatibility; + +internal sealed record CompatibilityOptions(Uri ServerUri, string RunId, Uri OtlpEndpoint, ReadyFilePath? ReadyFile) { + public static CompatibilityOptions Load(Func readEnvironment) { + ArgumentNullException.ThrowIfNull(readEnvironment); + + var serverUri = ReadAbsoluteUri(CompatibilityContract.ServerUriName); + var runId = ReadRequired(CompatibilityContract.RunIdName); + var otlpEndpoint = ReadAbsoluteUri(CompatibilityContract.OtlpEndpointName); + var readyFileValue = readEnvironment(CompatibilityContract.ReadyFileName); + ReadyFilePath? readyFile = string.IsNullOrWhiteSpace(readyFileValue) + ? null + : ReadyFilePath.Parse(readyFileValue); + + return new(serverUri, runId, otlpEndpoint, readyFile); + + string ReadRequired(string name) { + var value = readEnvironment(name)?.Trim(); + return !string.IsNullOrEmpty(value) + ? value + : throw new ArgumentException($"Environment variable {name} is required."); + } + + Uri ReadAbsoluteUri(string name) { + var value = ReadRequired(name); + return Uri.TryCreate(value, UriKind.Absolute, out var uri) + ? uri + : throw new ArgumentException($"Environment variable {name} must be an absolute URI."); + } + } +} + +internal readonly record struct ReadyFilePath { + static readonly byte[] ReadyContent = "ready\n"u8.ToArray(); + + ReadyFilePath(string fullPath) => FullPath = fullPath; + + public string FullPath { get; } + + public static ReadyFilePath Parse(string value) { + if (string.IsNullOrWhiteSpace(value) || !Path.IsPathFullyQualified(value)) + throw new ArgumentException($"Environment variable {CompatibilityContract.ReadyFileName} must be an absolute path."); + + return new(Path.GetFullPath(value)); + } + + public async Task Signal(CancellationToken cancellationToken) { + await using var file = new FileStream( + FullPath, + FileMode.CreateNew, + FileAccess.Write, + FileShare.Read, + ReadyContent.Length, + FileOptions.Asynchronous + ); + await file.WriteAsync(ReadyContent, cancellationToken); + } +} diff --git a/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityProgram.cs b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityProgram.cs new file mode 100644 index 0000000..4482f5c --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/CompatibilityProgram.cs @@ -0,0 +1,51 @@ +using KurrentDB.Client.Extensions.OpenTelemetry; +using OpenTelemetry; +using OpenTelemetry.Context.Propagation; +using OpenTelemetry.Resources; +using OpenTelemetry.Trace; + +namespace TrogonEventStore.Client.Compatibility; + +internal static class CompatibilityProgram { + const int OperationTimeoutSeconds = 30; + + public static async Task RunAsync( + IReadOnlyList arguments, + Func readEnvironment + ) { + try { + var options = CompatibilityOptions.Load(readEnvironment); + var command = CompatibilityCommand.Parse(arguments); + var originalPropagator = Propagators.DefaultTextMapPropagator; + + try { + Sdk.SetDefaultTextMapPropagator(new CompositeTextMapPropagator([ + new TraceContextPropagator(), + new BaggagePropagator() + ])); + + using var tracerProvider = Sdk + .CreateTracerProviderBuilder() + .SetResourceBuilder(ResourceBuilder.CreateDefault().AddService(CompatibilityContract.ServiceName)) + .AddKurrentDBClientInstrumentation() + .AddOtlpExporter(exporter => exporter.Endpoint = options.OtlpEndpoint) + .Build(); + using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(OperationTimeoutSeconds)); + + await new CompatibilityApplication(options).ExecuteAsync(command, timeout.Token); + if (!tracerProvider.ForceFlush()) + throw new InvalidOperationException("OpenTelemetry export did not flush successfully."); + } finally { + Sdk.SetDefaultTextMapPropagator(originalPropagator); + } + + return 0; + } catch (ArgumentException exception) { + Console.Error.WriteLine(exception.Message); + return 2; + } catch (Exception exception) { + Console.Error.WriteLine(exception); + return 1; + } + } +} diff --git a/compatibility/TrogonEventStore.Client.Compatibility/Program.cs b/compatibility/TrogonEventStore.Client.Compatibility/Program.cs new file mode 100644 index 0000000..653769f --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/Program.cs @@ -0,0 +1,3 @@ +using TrogonEventStore.Client.Compatibility; + +return await CompatibilityProgram.RunAsync(args, Environment.GetEnvironmentVariable); diff --git a/compatibility/TrogonEventStore.Client.Compatibility/TrogonEventStore.Client.Compatibility.csproj b/compatibility/TrogonEventStore.Client.Compatibility/TrogonEventStore.Client.Compatibility.csproj new file mode 100644 index 0000000..35b4e9f --- /dev/null +++ b/compatibility/TrogonEventStore.Client.Compatibility/TrogonEventStore.Client.Compatibility.csproj @@ -0,0 +1,19 @@ + + + Exe + false + + + + + + + + + + + + + + + diff --git a/otel/semconv/registry/trogon/eventstore/client/spans.yaml b/otel/semconv/registry/trogon/eventstore/client/spans.yaml index 0919577..1a35987 100644 --- a/otel/semconv/registry/trogon/eventstore/client/spans.yaml +++ b/otel/semconv/registry/trogon/eventstore/client/spans.yaml @@ -7,7 +7,7 @@ groups: - id: trogon.eventstore.event.type type: string stability: development - brief: Event type processed by the client. + brief: Event type received by the client. examples: [order-created] - id: span.trogon.eventstore.client.database.operation type: span @@ -32,11 +32,14 @@ groups: requirement_level: conditionally_required: If the server port is available. - - id: span.trogon.eventstore.client.process + - id: span.trogon.eventstore.client.receive type: span stability: development - span_kind: consumer - brief: Describes processing an event delivered by a subscription. + span_kind: client + brief: Describes receiving an event delivered by a subscription. + annotations: + code_generation: + operation_name: receive attributes: - ref: messaging.system requirement_level: required diff --git a/otel/semconv/templates/registry/csharp/subscription-trace-semantics.cs.j2 b/otel/semconv/templates/registry/csharp/subscription-trace-semantics.cs.j2 new file mode 100644 index 0000000..955b255 --- /dev/null +++ b/otel/semconv/templates/registry/csharp/subscription-trace-semantics.cs.j2 @@ -0,0 +1,10 @@ +// + +using System.Diagnostics; + +namespace KurrentDB.Diagnostics.Tracing; + +static class SubscriptionTraceSemantics { + public const string Operation = "{{ ctx.receive.annotations.code_generation.operation_name }}"; + public const ActivityKind SpanKind = ActivityKind.{{ ctx.receive.span_kind | pascal_case }}; +}{{- "\n" -}} diff --git a/otel/semconv/templates/registry/csharp/weaver.yaml b/otel/semconv/templates/registry/csharp/weaver.yaml index 8b0b21b..b5e4c95 100644 --- a/otel/semconv/templates/registry/csharp/weaver.yaml +++ b/otel/semconv/templates/registry/csharp/weaver.yaml @@ -52,3 +52,16 @@ templates: end application_mode: single file_name: TrogonTelemetryAttributes.g.cs + - template: subscription-trace-semantics.cs.j2 + filter: > + if $custom_attributes then + { + receive: [semconv_grouped_spans[].spans[] | select( + .id == "span.trogon.eventstore.client.receive" + )][0] + } + else + empty + end + application_mode: single + file_name: SubscriptionTraceSemantics.g.cs diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs index aa79238..ddc5174 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/ActivitySourceExtensions.cs @@ -4,6 +4,7 @@ using System.Diagnostics; using KurrentDB.Diagnostics; using KurrentDB.Diagnostics.Telemetry; +using KurrentDB.Diagnostics.Tracing; using OpenTelemetry; using static KurrentDB.Diagnostics.Tracing.TracingConstants; @@ -44,65 +45,103 @@ public static async ValueTask TraceClientOperation( } } - public static void TraceSubscriptionEvent( + public static SubscriptionReceive StartSubscriptionReceive( this ActivitySource source, - string? consumerGroupName, - ResolvedEvent resolvedEvent, - ChannelInfo channelInfo, - KurrentDBClientSettings settings - ) { - if (source.HasNoActiveListeners() || resolvedEvent.Event is null) - return; - - var propagationContext = resolvedEvent.Event.Metadata.ExtractPropagationContext(); - - if (propagationContext.ActivityContext == default) - return; - - var destination = resolvedEvent.OriginalEvent.EventStreamId; - var tags = new ActivityTagsCollection() - .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); - - using var activity = StartActivity( - source, - $"{Operations.Process} {destination}", - ActivityKind.Consumer, - tags, - propagationContext.ActivityContext - ); + string? consumerGroupName + ) => source.HasNoActiveListeners() + ? default + : new(source, consumerGroupName, Activity.Current?.Context ?? default, DateTimeOffset.UtcNow); + + public readonly struct SubscriptionReceive { + readonly ActivitySource? _source; + readonly string? _consumerGroupName; + readonly ActivityContext _parentContext; + readonly DateTimeOffset _startedAt; + + internal SubscriptionReceive( + ActivitySource source, + string? consumerGroupName, + ActivityContext parentContext, + DateTimeOffset startedAt + ) { + _source = source; + _consumerGroupName = consumerGroupName; + _parentContext = parentContext; + _startedAt = startedAt; + } - if (activity is null) - return; + public void Complete( + ResolvedEvent resolvedEvent, + ChannelInfo channelInfo, + KurrentDBClientSettings settings + ) { + if (_source is null) + return; + + var deliveredEvent = resolvedEvent.Event ?? resolvedEvent.Link; + if (deliveredEvent is null) + return; + + var propagationContext = deliveredEvent.Metadata.ExtractPropagationContext(); + var destination = resolvedEvent.OriginalEvent.EventStreamId; + var tags = new ActivityTagsCollection() + .WithRequiredTag(TelemetryAttributes.MessagingSystem, SystemName) + .WithRequiredTag(TelemetryAttributes.MessagingOperationName, SubscriptionTraceSemantics.Operation) + .WithRequiredTag(TelemetryAttributes.MessagingOperationType, SubscriptionTraceSemantics.Operation) + .WithRequiredTag(TelemetryAttributes.MessagingDestinationName, destination) + .WithOptionalTag(TelemetryAttributes.MessagingConsumerGroupName, _consumerGroupName) + .WithRequiredTag(TelemetryAttributes.MessagingMessageId, resolvedEvent.OriginalEvent.EventId.ToString()) + .WithRequiredTag(TrogonTelemetryAttributes.EventType, deliveredEvent.EventType) + .WithGrpcChannelServerTags(channelInfo) + .WithClientSettingsServerTags(settings); + var links = propagationContext.ActivityContext == default + ? null + : new[] { new ActivityLink(propagationContext.ActivityContext) }; + + using var activity = StartActivity( + _source, + $"{SubscriptionTraceSemantics.Operation} {destination}", + SubscriptionTraceSemantics.SpanKind, + tags, + _parentContext, + links, + _startedAt + ); + + if (activity is null) + return; - foreach (var (name, value) in propagationContext.Baggage.GetBaggage()) - activity.AddBaggage(name, value); + foreach (var (name, value) in propagationContext.Baggage.GetBaggage()) + activity.AddBaggage(name, value); + } } static Activity? StartActivity( this ActivitySource source, string operationName, ActivityKind activityKind, ActivityTagsCollection? tags = null, - ActivityContext? parentContext = null + ActivityContext? parentContext = null, + IEnumerable? links = null, + DateTimeOffset startTime = default ) { if (source.HasNoActiveListeners()) return null; - return source - .CreateActivity( + var activity = source.CreateActivity( operationName, activityKind, parentContext ?? default, tags, + links, idFormat: ActivityIdFormat.W3C - ) - ?.Start(); + ); + + if (activity is null) + return null; + + if (startTime != default) + activity.SetStartTime(startTime.UtcDateTime); + + return activity.Start(); } static bool HasNoActiveListeners(this ActivitySource source) => !source.HasListeners(); diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs index b2391cb..d7c820d 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/EventMetadataExtensions.cs @@ -7,6 +7,9 @@ namespace KurrentDB.Client.Diagnostics; static class EventMetadataExtensions { + const string TraceParent = "traceparent"; + const string TraceState = "tracestate"; + public static void InjectTracingContext(this Dictionary metadata, Activity? activity) { if (activity is null) return; @@ -19,6 +22,9 @@ public static void InjectTracingContext(this Dictionary metadata } static void SetPropagationField(Dictionary metadata, string name, string value) { + if (!IsPersistedTraceField(name)) + return; + while (true) { string? existingName = null; foreach (var key in metadata.Keys) { @@ -38,6 +44,10 @@ static void SetPropagationField(Dictionary metadata, string name metadata[name] = value; } + static bool IsPersistedTraceField(string name) => + string.Equals(name, TraceParent, StringComparison.OrdinalIgnoreCase) || + string.Equals(name, TraceState, StringComparison.OrdinalIgnoreCase); + [MethodImpl(MethodImplOptions.AggressiveInlining)] public static ReadOnlySpan InjectTracingContext( this ReadOnlyMemory eventMetadata, Activity? activity @@ -49,7 +59,7 @@ public static ReadOnlySpan InjectTracingContext( Propagators.DefaultTextMapPropagator.Inject( new PropagationContext(activity.Context, Baggage.Current), propagationMetadata, - static (carrier, name, value) => carrier[name] = value + static (carrier, name, value) => SetPropagationField(carrier, name, value) ); return eventMetadata.InjectPropagationMetadata(propagationMetadata); diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/SubscriptionTraceSemantics.g.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/SubscriptionTraceSemantics.g.cs new file mode 100644 index 0000000..e8844fe --- /dev/null +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/Generated/SubscriptionTraceSemantics.g.cs @@ -0,0 +1,10 @@ +// + +using System.Diagnostics; + +namespace KurrentDB.Diagnostics.Tracing; + +static class SubscriptionTraceSemantics { + public const string Operation = "receive"; + public const ActivityKind SpanKind = ActivityKind.Client; +} diff --git a/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs b/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs index 7dcd65f..e7a653c 100644 --- a/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs +++ b/src/KurrentDB.Client/Core/Common/Diagnostics/Tracing/TracingConstants.cs @@ -9,6 +9,5 @@ static class TracingConstants { 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/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs b/src/KurrentDB.Client/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs index 2aa87ec..0dfaf68 100644 --- a/src/KurrentDB.Client/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs +++ b/src/KurrentDB.Client/PersistentSubscriptions/KurrentDBPersistentSubscriptionsClient.Read.cs @@ -1,7 +1,8 @@ +using System.Runtime.CompilerServices; using System.Threading.Channels; using EventStore.Client; -using KurrentDB.Client.Diagnostics; using Grpc.Core; +using KurrentDB.Client.Diagnostics; using KurrentDB.Protocol.PersistentSubscriptions.V1; using static KurrentDB.Protocol.PersistentSubscriptions.V1.PersistentSubscriptions; using static KurrentDB.Protocol.PersistentSubscriptions.V1.ReadResp.ContentOneofCase; @@ -175,8 +176,10 @@ public class PersistentSubscriptionResult : IAsyncEnumerable, IAs readonly Channel _channel; readonly CancellationTokenSource _cts; readonly CallOptions _callOptions; + readonly KurrentDBClientSettings _settings; AsyncDuplexStreamingCall? _call; + ChannelInfo? _channelInfo; int _messagesEnumerated; /// @@ -206,7 +209,7 @@ public IAsyncEnumerable Messages { async IAsyncEnumerable GetMessages() { try { - await foreach (var message in _channel.Reader.ReadAllAsync(_cts.Token)) { + await foreach (var message in ReadMessages(_cts.Token)) { if (message is PersistentSubscriptionMessage.SubscriptionConfirmation(var subscriptionId)) SubscriptionId = subscriptionId; @@ -230,6 +233,7 @@ CancellationToken cancellationToken GroupName = groupName; _request = request; + _settings = settings; _callOptions = KurrentDBCallOptions.CreateStreaming( settings, @@ -248,6 +252,7 @@ CancellationToken cancellationToken async Task PumpMessages() { try { var channelInfo = await selectChannelInfo(_cts.Token).ConfigureAwait(false); + _channelInfo = channelInfo; var client = new PersistentSubscriptionsClient(channelInfo.CallInvoker); _call = client.Read(_callOptions); @@ -269,14 +274,6 @@ async Task PumpMessages() { _ => PersistentSubscriptionMessage.Unknown.Instance }; - if (subscriptionMessage is PersistentSubscriptionMessage.Event evnt) - KurrentDBClientDiagnostics.ActivitySource.TraceSubscriptionEvent( - GroupName, - evnt.ResolvedEvent, - channelInfo, - settings - ); - await _channel.Writer.WriteAsync(subscriptionMessage, _cts.Token).ConfigureAwait(false); } @@ -456,6 +453,26 @@ public async IAsyncEnumerator GetAsyncEnumerator(CancellationToke yield return resolvedEvent; } } + + async IAsyncEnumerable ReadMessages( + [EnumeratorCancellation] CancellationToken cancellationToken + ) { + await using var messages = _channel.Reader + .ReadAllAsync(cancellationToken) + .GetAsyncEnumerator(cancellationToken); + + while (true) { + var receive = KurrentDBClientDiagnostics.ActivitySource.StartSubscriptionReceive(GroupName); + if (!await messages.MoveNextAsync().ConfigureAwait(false)) + yield break; + + var message = messages.Current; + if (message is PersistentSubscriptionMessage.Event(var resolvedEvent, _)) + receive.Complete(resolvedEvent, _channelInfo!, _settings); + + yield return message; + } + } } } } diff --git a/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs b/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs index 639d39b..b0bb259 100644 --- a/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs +++ b/src/KurrentDB.Client/Streams/KurrentDBClient.Subscriptions.cs @@ -1,6 +1,7 @@ +using System.Runtime.CompilerServices; using System.Threading.Channels; -using KurrentDB.Client.Diagnostics; using Grpc.Core; +using KurrentDB.Client.Diagnostics; using KurrentDB.Protocol.Streams.V1; using static KurrentDB.Protocol.Streams.V1.ReadResp.ContentOneofCase; using static KurrentDB.Protocol.Streams.V1.Streams; @@ -135,6 +136,7 @@ public class StreamSubscriptionResult : IAsyncEnumerable, IAsyncD private readonly CallOptions _callOptions; private readonly KurrentDBClientSettings _settings; private AsyncServerStreamingCall? _call; + private ChannelInfo? _channelInfo; private int _messagesEnumerated; @@ -155,7 +157,7 @@ public IAsyncEnumerable Messages { async IAsyncEnumerable GetMessages() { try { - await foreach (var message in _channel.Reader.ReadAllAsync(_cts.Token)) { + await foreach (var message in ReadMessages(_cts.Token)) { if (message is StreamMessage.SubscriptionConfirmation(var subscriptionId)) SubscriptionId = subscriptionId; @@ -198,6 +200,7 @@ CancellationToken cancellationToken async Task PumpMessages() { try { var channelInfo = await selectChannelInfo(_cts.Token).ConfigureAwait(false); + _channelInfo = channelInfo; var client = new StreamsClient(channelInfo.CallInvoker); _call = client.Read(_request, _callOptions); await foreach (var response in _call.ResponseStream.ReadAllAsync(_cts.Token).ConfigureAwait(false)) { @@ -240,14 +243,6 @@ response.FellBehind.Position is { } position _ => StreamMessage.Unknown.Instance }; - if (subscriptionMessage is StreamMessage.Event evt) - KurrentDBClientDiagnostics.ActivitySource.TraceSubscriptionEvent( - null, - evt.ResolvedEvent, - channelInfo, - _settings - ); - await _channel.Writer .WriteAsync(subscriptionMessage, _cts.Token) .ConfigureAwait(false); @@ -293,7 +288,7 @@ public void Dispose() { /// public async IAsyncEnumerator GetAsyncEnumerator(CancellationToken cancellationToken = default) { try { - await foreach (var message in _channel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { + await foreach (var message in ReadMessages(cancellationToken).ConfigureAwait(false)) { if (message is not StreamMessage.Event e) continue; @@ -304,6 +299,26 @@ public async IAsyncEnumerator GetAsyncEnumerator(CancellationToke await _cts.CancelAsync().ConfigureAwait(false); } } + + async IAsyncEnumerable ReadMessages( + [EnumeratorCancellation] CancellationToken cancellationToken + ) { + await using var messages = _channel.Reader + .ReadAllAsync(cancellationToken) + .GetAsyncEnumerator(cancellationToken); + + while (true) { + var receive = KurrentDBClientDiagnostics.ActivitySource.StartSubscriptionReceive(null); + if (!await messages.MoveNextAsync().ConfigureAwait(false)) + yield break; + + var message = messages.Current; + if (message is StreamMessage.Event(var resolvedEvent)) + receive.Complete(resolvedEvent, _channelInfo!, _settings); + + yield return message; + } + } } } } diff --git a/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs b/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs index d4d6063..79a11ea 100644 --- a/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs +++ b/test/KurrentDB.Client.Tests.Common/Fixtures/DiagnosticsFixture.cs @@ -126,13 +126,13 @@ public void AssertSubscriptionActivityHasExpectedTags( string eventId, string? consumerGroupName = null ) { - activity.DisplayName.ShouldBe($"{TracingConstants.Operations.Process} {stream}"); - activity.Kind.ShouldBe(ActivityKind.Consumer); + activity.DisplayName.ShouldBe($"{SubscriptionTraceSemantics.Operation} {stream}"); + activity.Kind.ShouldBe(SubscriptionTraceSemantics.SpanKind); var expectedTags = new Dictionary { { TelemetryAttributes.MessagingSystem, TracingConstants.SystemName }, - { TelemetryAttributes.MessagingOperationName, TracingConstants.Operations.Process }, - { TelemetryAttributes.MessagingOperationType, TracingConstants.Operations.Process }, + { TelemetryAttributes.MessagingOperationName, SubscriptionTraceSemantics.Operation }, + { TelemetryAttributes.MessagingOperationType, SubscriptionTraceSemantics.Operation }, { TelemetryAttributes.MessagingDestinationName, stream }, { TelemetryAttributes.MessagingMessageId, eventId }, { TrogonTelemetryAttributes.EventType, TestEventType } diff --git a/test/KurrentDB.Client.Tests/Compatibility/CompatibilityApplicationTests.cs b/test/KurrentDB.Client.Tests/Compatibility/CompatibilityApplicationTests.cs new file mode 100644 index 0000000..56d2e43 --- /dev/null +++ b/test/KurrentDB.Client.Tests/Compatibility/CompatibilityApplicationTests.cs @@ -0,0 +1,89 @@ +using System.Buffers.Binary; +using System.Net; +using System.Security.Cryptography; +using System.Security.Cryptography.X509Certificates; +using System.Text; +using Microsoft.AspNetCore.Builder; +using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.Hosting.Server; +using Microsoft.AspNetCore.Hosting.Server.Features; +using Microsoft.AspNetCore.Http; +using Microsoft.AspNetCore.Server.Kestrel.Core; +using Microsoft.Extensions.DependencyInjection; +using TrogonEventStore.Client.Compatibility; + +namespace KurrentDB.Client.Tests.Compatibility; + +public class CompatibilityApplicationTests { + [Fact] + public async Task write_uses_credentials_from_the_server_uri() { + using var certificate = CreateCertificate(); + var builder = WebApplication.CreateSlimBuilder(); + builder.WebHost.ConfigureKestrel(options => options.Listen( + IPAddress.Loopback, + 0, + listen => { + listen.Protocols = HttpProtocols.Http2; + listen.UseHttps(certificate); + } + )); + + var application = builder.Build(); + string? authorizationHeader = null; + + application.MapPost( + "/event_store.client.server_features.ServerFeatures/GetSupportedMethods", + context => CompleteGrpcResponse(context, ReadOnlyMemory.Empty) + ); + application.MapPost( + "/event_store.client.streams.Streams/Append", + async context => { + authorizationHeader = context.Request.Headers.Authorization; + await context.Request.Body.CopyToAsync(Stream.Null); + await CompleteGrpcResponse(context, new byte[] { 0x0a, 0x00 }); + } + ); + + await application.StartAsync(); + + try { + var server = application.Services.GetRequiredService(); + var address = server.Features.Get()!.Addresses.Single(); + var port = new Uri(address).Port; + var options = new CompatibilityOptions( + new($"esdb://writer:secret@127.0.0.1:{port}?tls=true&tlsVerifyCert=false"), + "run-1", + new("http://127.0.0.1:4317"), + null + ); + + await new CompatibilityApplication(options).ExecuteAsync( + CompatibilityCommand.Parse(["write", "compatibility-stream"]), + CancellationToken.None + ); + + var encodedCredentials = Convert.ToBase64String(Encoding.UTF8.GetBytes("writer:secret")); + authorizationHeader.ShouldBe($"Basic {encodedCredentials}"); + } finally { + await application.StopAsync(); + await application.DisposeAsync(); + } + } + + static X509Certificate2 CreateCertificate() { + using var key = RSA.Create(2048); + var request = new CertificateRequest("CN=localhost", key, HashAlgorithmName.SHA256, RSASignaturePadding.Pkcs1); + return request.CreateSelfSigned(DateTimeOffset.UtcNow.AddMinutes(-1), DateTimeOffset.UtcNow.AddMinutes(5)); + } + + static async Task CompleteGrpcResponse(HttpContext context, ReadOnlyMemory payload) { + context.Response.ContentType = "application/grpc"; + context.Response.DeclareTrailer("grpc-status"); + + var header = new byte[5]; + BinaryPrimitives.WriteInt32BigEndian(header.AsSpan(1), payload.Length); + await context.Response.Body.WriteAsync(header); + await context.Response.Body.WriteAsync(payload); + context.Response.AppendTrailer("grpc-status", "0"); + } +} diff --git a/test/KurrentDB.Client.Tests/Compatibility/CompatibilityCommandTests.cs b/test/KurrentDB.Client.Tests/Compatibility/CompatibilityCommandTests.cs new file mode 100644 index 0000000..5dbe1e3 --- /dev/null +++ b/test/KurrentDB.Client.Tests/Compatibility/CompatibilityCommandTests.cs @@ -0,0 +1,51 @@ +using TrogonEventStore.Client.Compatibility; + +namespace KurrentDB.Client.Tests.Compatibility; + +public class CompatibilityCommandTests { + [Theory] + [InlineData("write")] + [InlineData("batch-write")] + [InlineData("read")] + [InlineData("subscribe")] + public void parses_stream_commands(string command) { + var parsed = CompatibilityCommand.Parse([command, "compatibility-stream"]); + + parsed.Stream.Value.ShouldBe("compatibility-stream"); + } + + [Fact] + public void parses_batch_write_as_a_distinct_command() => + CompatibilityCommand.Parse(["batch-write", "compatibility-stream"]) + .ShouldBeOfType(); + + [Theory] + [InlineData("create-persistent-subscription")] + [InlineData("consume-persistent-subscription")] + public void parses_persistent_subscription_commands(string command) { + var parsed = CompatibilityCommand.Parse([command, "compatibility-stream", "compatibility-group"]); + + parsed.Stream.Value.ShouldBe("compatibility-stream"); + var group = parsed switch { + CompatibilityCommand.CreatePersistentSubscription create => create.Group, + CompatibilityCommand.ConsumePersistentSubscription consume => consume.Group, + _ => throw new InvalidOperationException() + }; + group.Value.ShouldBe("compatibility-group"); + } + + [Theory] + [InlineData()] + [InlineData("write")] + [InlineData("write", "")] + [InlineData("batch-write")] + [InlineData("create-persistent-subscription", "stream")] + [InlineData("unknown", "stream")] + public void rejects_invalid_commands(params string[] arguments) => + Should.Throw(() => CompatibilityCommand.Parse(arguments)); + + [Fact] + public void command_error_describes_batch_write() => + Should.Throw(() => CompatibilityCommand.Parse([])) + .Message.ShouldContain("batch-write "); +} diff --git a/test/KurrentDB.Client.Tests/Compatibility/CompatibilityOptionsTests.cs b/test/KurrentDB.Client.Tests/Compatibility/CompatibilityOptionsTests.cs new file mode 100644 index 0000000..4ce617c --- /dev/null +++ b/test/KurrentDB.Client.Tests/Compatibility/CompatibilityOptionsTests.cs @@ -0,0 +1,63 @@ +using TrogonEventStore.Client.Compatibility; + +namespace KurrentDB.Client.Tests.Compatibility; + +public class CompatibilityOptionsTests { + [Fact] + public void loads_required_environment() { + var environment = ValidEnvironment(); + + var options = CompatibilityOptions.Load(environment.GetValueOrDefault); + + options.ServerUri.ShouldBe(new Uri("esdb://admin:changeit@localhost:2113?tls=false")); + options.RunId.ShouldBe("run-1"); + options.OtlpEndpoint.ShouldBe(new Uri("http://localhost:4317")); + options.ReadyFile.ShouldBeNull(); + } + + [Fact] + public void loads_optional_ready_file() { + var environment = ValidEnvironment(); + var path = Path.Combine(Path.GetTempPath(), "trogon-eventstore-ready"); + environment[CompatibilityContract.ReadyFileName] = path; + + var options = CompatibilityOptions.Load(environment.GetValueOrDefault); + + options.ReadyFile.ShouldBe(ReadyFilePath.Parse(path)); + } + + [Theory] + [InlineData(CompatibilityContract.ServerUriName)] + [InlineData(CompatibilityContract.RunIdName)] + [InlineData(CompatibilityContract.OtlpEndpointName)] + public void rejects_missing_environment(string name) { + var environment = ValidEnvironment(); + environment.Remove(name); + + Should.Throw(() => CompatibilityOptions.Load(environment.GetValueOrDefault)); + } + + [Theory] + [InlineData(CompatibilityContract.ServerUriName)] + [InlineData(CompatibilityContract.OtlpEndpointName)] + public void rejects_relative_uris(string name) { + var environment = ValidEnvironment(); + environment[name] = "relative"; + + Should.Throw(() => CompatibilityOptions.Load(environment.GetValueOrDefault)); + } + + [Fact] + public void rejects_relative_ready_file() { + var environment = ValidEnvironment(); + environment[CompatibilityContract.ReadyFileName] = "ready"; + + Should.Throw(() => CompatibilityOptions.Load(environment.GetValueOrDefault)); + } + + static Dictionary ValidEnvironment() => new() { + [CompatibilityContract.ServerUriName] = "esdb://admin:changeit@localhost:2113?tls=false", + [CompatibilityContract.RunIdName] = "run-1", + [CompatibilityContract.OtlpEndpointName] = "http://localhost:4317" + }; +} diff --git a/test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs b/test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs index 4d66869..bf26959 100644 --- a/test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs +++ b/test/KurrentDB.Client.Tests/Diagnostics/OpenTelemetryIntegrationTests.cs @@ -125,13 +125,13 @@ public void configured_propagator_injects_context_into_json_metadata() { 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"); + extracted.Baggage.GetBaggage("tenant").ShouldBeNull(); 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(); + document.RootElement.TryGetProperty("baggage", out _).ShouldBeFalse(); } finally { Baggage.Current = originalBaggage; Sdk.SetDefaultTextMapPropagator(originalPropagator); @@ -157,7 +157,7 @@ public void configured_propagator_injects_context_into_property_metadata() { 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"); + extracted.Baggage.GetBaggage("tenant").ShouldBeNull(); metadata["custom"].ShouldBe("value"); } finally { Baggage.Current = originalBaggage; @@ -183,14 +183,15 @@ public void configured_propagator_replaces_mixed_case_fields_in_property_metadat metadata.InjectTracingContext(activity); - foreach (var name in new[] { "traceparent", "tracestate", "baggage" }) + foreach (var name in new[] { "traceparent", "tracestate" }) metadata.Keys.Count(key => string.Equals(key, name, StringComparison.OrdinalIgnoreCase)).ShouldBe(1); + metadata["Baggage"].ShouldBe("stale"); 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"); + extracted.Baggage.GetBaggage("tenant").ShouldBeNull(); } finally { Baggage.Current = originalBaggage; Sdk.SetDefaultTextMapPropagator(originalPropagator); diff --git a/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs b/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs index 658703e..7c0d104 100644 --- a/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs +++ b/test/KurrentDB.Client.Tests/Diagnostics/PersistentSubscriptionsTracingInstrumentationTests.cs @@ -1,3 +1,4 @@ +using System.Diagnostics; using KurrentDB.Client.Tests.Fixtures; using KurrentDB.Client.Tests.TestNode; using KurrentDB.Diagnostics.Telemetry; @@ -10,8 +11,9 @@ namespace KurrentDB.Client.Tests.Diagnostics; public class PersistentSubscriptionsTracingInstrumentationTests(ITestOutputHelper output, DiagnosticsFixture fixture) : KurrentDBPermanentTests(output, fixture) { [RetryFact] - public async Task persistent_subscription_restores_remote_append_context() { + public async Task persistent_subscription_receive_links_remote_append_context() { var traceId = Fixture.CreateTraceId(); + var subscriber = Activity.Current!; var stream = Fixture.GetStreamName(); var events = Fixture.CreateTestEvents(2, metadata: Fixture.CreateTestJsonMetadata()).ToArray(); @@ -37,7 +39,7 @@ await Fixture.Streams.AppendToStreamAsync( .ShouldNotBeNull(); var subscribeActivities = Fixture - .GetActivities(TracingConstants.Operations.Process, traceId, stream) + .GetActivities(SubscriptionTraceSemantics.Operation, traceId, stream) .Where(activity => Equals( activity.GetTagItem(TelemetryAttributes.MessagingConsumerGroupName), groupName @@ -53,9 +55,13 @@ await Fixture.Streams.AppendToStreamAsync( Assert.True(expectedEventIds.SetEquals(actualEventIds)); foreach (var subscribeActivity in subscribeActivities) { - subscribeActivity.TraceId.ShouldBe(appendActivity.Context.TraceId); - subscribeActivity.ParentSpanId.ShouldBe(appendActivity.Context.SpanId); - subscribeActivity.HasRemoteParent.ShouldBeTrue(); + subscribeActivity.TraceId.ShouldBe(subscriber.TraceId); + subscribeActivity.ParentSpanId.ShouldBe(subscriber.SpanId); + subscribeActivity.HasRemoteParent.ShouldBeFalse(); + var messageLink = subscribeActivity.Links.ShouldHaveSingleItem().Context; + messageLink.TraceId.ShouldBe(appendActivity.TraceId); + messageLink.SpanId.ShouldBe(appendActivity.SpanId); + messageLink.IsRemote.ShouldBeTrue(); subscribeActivity.GetTagItem(TelemetryAttributes.MessagingConsumerGroupName).ShouldBe(groupName); Fixture.AssertSubscriptionActivityHasExpectedTags( diff --git a/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs b/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs index c51263a..7be2924 100644 --- a/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs +++ b/test/KurrentDB.Client.Tests/Diagnostics/StreamsTracingInstrumentationTests.cs @@ -1,13 +1,12 @@ // ReSharper disable AccessToDisposedClosure using System.Diagnostics; +using System.Text; using System.Text.Json; using KurrentDB.Client.Diagnostics; using KurrentDB.Client.Tests.Fixtures; using KurrentDB.Diagnostics.Telemetry; using KurrentDB.Diagnostics.Tracing; -using OpenTelemetry; -using OpenTelemetry.Context.Propagation; namespace KurrentDB.Client.Tests.Diagnostics; @@ -46,6 +45,7 @@ await Fixture.Streams.AppendToStreamAsync( public async Task multi_stream_append() { // Arrange var traceId = Fixture.CreateTraceId(); + var subscriber = Activity.Current!; var seedEvents = Fixture.CreateTestEvents(10).ToList(); @@ -72,7 +72,7 @@ public async Task multi_stream_append() { appendResult.Position.ShouldBePositive(); var appendActivities = Fixture.GetActivities(TracingConstants.Operations.BatchAppend, traceId); - var subscribeActivities = Fixture.GetActivities(TracingConstants.Operations.Process, traceId); + var subscribeActivities = Fixture.GetActivities(SubscriptionTraceSemantics.Operation, traceId); appendActivities.ShouldNotBeEmpty(); subscribeActivities.ShouldNotBeEmpty(); @@ -83,14 +83,15 @@ public async Task multi_stream_append() { // They also have the same duration appendActivities.Select(x => x.Duration).Distinct().Count().ShouldBe(1); - // Check that subscribe activities have the correct parent IDs inherited from append activities - subscribeActivities - .FirstOrDefault(x => x.ParentId == appendActivities.First().Id)?.ParentSpanId - .ShouldBe(appendActivities.First().SpanId); - - subscribeActivities - .FirstOrDefault(x => x.ParentId == appendActivities.Last().Id)?.ParentSpanId - .ShouldBe(appendActivities.Last().SpanId); + Assert.All( + subscribeActivities, + receiveActivity => { + receiveActivity.ParentSpanId.ShouldBe(subscriber.SpanId); + var messageLink = receiveActivity.Links.ShouldHaveSingleItem().Context; + messageLink.TraceId.ShouldBe(appendActivities[0].TraceId); + messageLink.SpanId.ShouldBe(appendActivities[0].SpanId); + } + ); subscribeActivities .All(x => x.StartTimeUtc > appendActivities.First().StartTimeUtc) @@ -197,51 +198,45 @@ await Fixture.Streams.AppendToStreamAsync( } [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; + public async Task subscription_receive_uses_ambient_parent_and_links_message_creation_context() { + var producerTraceId = Fixture.CreateTraceId(); + Activity.Current!.TraceStateString = "vendor=producer"; 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); - } + var metadata = JsonSerializer.SerializeToUtf8Bytes(new Dictionary { + ["baggage"] = "tenant=straw-hat" + }); + var seedEvent = Fixture.CreateTestEvent(metadata: metadata); + + await Fixture.Streams.AppendToStreamAsync(stream, StreamState.NoStream, [seedEvent]); + var appendActivity = Fixture + .GetActivities(TracingConstants.Operations.Append, producerTraceId) + .ShouldHaveSingleItem(); + + Activity.Current = null; + using var subscriber = new Activity("subscriber").SetIdFormat(ActivityIdFormat.W3C).Start(); + subscriber.TraceStateString = "vendor=consumer"; + + 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 receiveActivity = Fixture + .GetActivities(SubscriptionTraceSemantics.Operation, subscriber.TraceId, stream) + .ShouldHaveSingleItem(); + + receiveActivity.Kind.ShouldBe(ActivityKind.Client); + receiveActivity.ParentSpanId.ShouldBe(subscriber.SpanId); + receiveActivity.TraceStateString.ShouldBe("vendor=consumer"); + receiveActivity.Baggage.ShouldContain(new KeyValuePair("tenant", "straw-hat")); + var messageLink = receiveActivity.Links.ShouldHaveSingleItem().Context; + messageLink.TraceId.ShouldBe(appendActivity.TraceId); + messageLink.SpanId.ShouldBe(appendActivity.SpanId); + messageLink.TraceState.ShouldBe("vendor=producer"); + messageLink.IsRemote.ShouldBeTrue(); } [Fact] @@ -326,7 +321,7 @@ await Fixture.Streams.AppendToStreamAsync( } [Fact] - public async Task json_metadata_traced_non_json_metadata_not_traced() { + public async Task subscription_receive_is_emitted_without_propagated_context() { var traceId = Fixture.CreateTraceId(); var streamName = Fixture.GetStreamName(); @@ -353,23 +348,20 @@ public async Task json_metadata_traced_non_json_metadata_not_traced() { await Subscribe(enumerator).WithTimeout(); var subscribeActivities = Fixture - .GetActivities(TracingConstants.Operations.Process, traceId, streamName) + .GetActivities(SubscriptionTraceSemantics.Operation, traceId, streamName) .ToArray(); appendActivities.ShouldHaveSingleItem(); - var jsonMetadataEvent = seedEvents.First(); - - Assert.NotEmpty(subscribeActivities); - Assert.All( - subscribeActivities, - activity => { - Assert.Equal(appendActivities.First().Id, activity.ParentId); - Fixture.AssertSubscriptionActivityHasExpectedTags( - activity, - streamName, - jsonMetadataEvent.EventId.ToString() - ); - } + subscribeActivities.Length.ShouldBe(seedEvents.Length); + var receiveWithoutMessageContext = subscribeActivities.Single(activity => Equals( + activity.GetTagItem(TelemetryAttributes.MessagingMessageId), + seedEvents.Last().EventId.ToString() + )); + receiveWithoutMessageContext.Links.ShouldBeEmpty(); + Fixture.AssertSubscriptionActivityHasExpectedTags( + receiveWithoutMessageContext, + streamName, + seedEvents.Last().EventId.ToString() ); return; @@ -389,13 +381,69 @@ async Task Subscribe(IAsyncEnumerator internalEnumerator) { [RetryFact] [Trait("Category", "Special cases")] - public async Task no_trace_when_event_is_null() { + public async Task unresolved_link_receive_falls_back_to_original_link_semantics() { + var traceId = Fixture.CreateTraceId(); + var targetStream = Fixture.GetStreamName(); + var linkStream = Fixture.GetStreamName(); + + await Fixture.Streams.AppendToStreamAsync( + targetStream, + StreamState.NoStream, + Fixture.CreateTestEvents(1) + ); + var linkEvent = new EventData( + Uuid.NewUuid(), + SystemEventTypes.LinkTo, + Encoding.UTF8.GetBytes($"0@{targetStream}"), + Fixture.CreateTestJsonMetadata(), + Constants.Metadata.ContentTypes.ApplicationOctetStream + ); + await Fixture.Streams.AppendToStreamAsync(linkStream, StreamState.NoStream, [linkEvent]); + var linkAppendActivity = Fixture + .GetActivities(TracingConstants.Operations.Append, traceId, linkStream) + .ShouldHaveSingleItem(); + await Fixture.Streams.DeleteAsync(targetStream, StreamState.StreamExists); + + await using var subscription = Fixture.Streams.SubscribeToStream( + linkStream, + FromStream.Start, + resolveLinkTos: true + ); + await using var enumerator = subscription.Messages.GetAsyncEnumerator(); + + Assert.True(await enumerator.MoveNextAsync()); + Assert.IsType(enumerator.Current); + Assert.True(await enumerator.MoveNextAsync()); + var unresolvedLink = Assert.IsType(enumerator.Current).ResolvedEvent; + + unresolvedLink.Event.ShouldBeNull(); + unresolvedLink.Link.ShouldNotBeNull(); + unresolvedLink.OriginalEvent.EventType.ShouldBe(SystemEventTypes.LinkTo); + + var receiveActivity = Fixture + .GetActivities(SubscriptionTraceSemantics.Operation, traceId, linkStream) + .ShouldHaveSingleItem(); + receiveActivity.GetTagItem(TelemetryAttributes.MessagingDestinationName) + .ShouldBe(unresolvedLink.OriginalEvent.EventStreamId); + receiveActivity.GetTagItem(TelemetryAttributes.MessagingMessageId) + .ShouldBe(unresolvedLink.OriginalEvent.EventId.ToString()); + receiveActivity.GetTagItem(TrogonTelemetryAttributes.EventType) + .ShouldBe(unresolvedLink.OriginalEvent.EventType); + var messageLink = receiveActivity.Links.ShouldHaveSingleItem().Context; + messageLink.TraceId.ShouldBe(linkAppendActivity.TraceId); + messageLink.SpanId.ShouldBe(linkAppendActivity.SpanId); + } + + [RetryFact] + [Trait("Category", "Special cases")] + public async Task resolved_link_receive_uses_target_event_type_and_original_link_identity() { var traceId = Fixture.CreateTraceId(); var category = Guid.NewGuid().ToString("N"); var streamName = category + "-123"; var categoryStream = "$ce-" + category; var seedEvents = Fixture.CreateTestEvents(type: $"{category}-{Fixture.GetStreamName()}").ToArray(); + ResolvedEvent? receivedEvent = null; await Fixture.Streams.AppendToStreamAsync(streamName, StreamState.NoStream, seedEvents); await Fixture.Streams.DeleteAsync(streamName, StreamState.StreamExists); @@ -415,11 +463,26 @@ public async Task no_trace_when_event_is_null() { .ShouldNotBeNull(); var subscribeActivities = Fixture - .GetActivities(TracingConstants.Operations.Process, traceId, categoryStream) + .GetActivities(SubscriptionTraceSemantics.Operation, traceId, categoryStream) .ToArray(); appendActivities.ShouldHaveSingleItem(); - subscribeActivities.ShouldBeEmpty(); + var resolvedLink = receivedEvent!.Value; + var receiveActivity = subscribeActivities.Single(activity => Equals( + activity.GetTagItem(TelemetryAttributes.MessagingMessageId), + resolvedLink.OriginalEvent.EventId.ToString() + )); + resolvedLink.IsResolved.ShouldBeTrue(); + resolvedLink.OriginalEvent.EventType.ShouldBe("$>"); + resolvedLink.Event.EventType.ShouldBe("$metadata"); + receiveActivity.Links.ShouldBeEmpty(); + receiveActivity.Kind.ShouldBe(SubscriptionTraceSemantics.SpanKind); + receiveActivity.GetTagItem(TelemetryAttributes.MessagingDestinationName) + .ShouldBe(resolvedLink.OriginalEvent.EventStreamId); + receiveActivity.GetTagItem(TelemetryAttributes.MessagingMessageId) + .ShouldBe(resolvedLink.OriginalEvent.EventId.ToString()); + receiveActivity.GetTagItem(TrogonTelemetryAttributes.EventType) + .ShouldBe(resolvedLink.Event.EventType); return; @@ -428,8 +491,10 @@ async Task Subscribe() { if (enumerator.Current is not StreamMessage.Event(var resolvedEvent)) continue; - if (resolvedEvent.Event?.EventType is "$metadata") + if (resolvedEvent.Event?.EventType is "$metadata") { + receivedEvent = resolvedEvent; return; + } } } } diff --git a/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj b/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj index 6f82ad1..b68b2bf 100644 --- a/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj +++ b/test/KurrentDB.Client.Tests/KurrentDB.Client.Tests.csproj @@ -1,6 +1,7 @@  +