From 1cc5671dc2bfd36d187c64fa7d1d363878424903 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 22 Jan 2026 11:47:36 +0000 Subject: [PATCH 1/2] Fix OTel activity leak in Direct request mode When using NatsRequestReplyMode.Direct, the OpenTelemetry activity stored in message headers was never disposed. SharedInbox mode handles this via ActivityEndingMsgReader, but Direct mode bypasses that path. Add explicit activity disposal after receiving the response message. --- src/NATS.Client.Core/NatsConnection.RequestReply.cs | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/src/NATS.Client.Core/NatsConnection.RequestReply.cs b/src/NATS.Client.Core/NatsConnection.RequestReply.cs index a3ddc8c86..aa6885bce 100644 --- a/src/NATS.Client.Core/NatsConnection.RequestReply.cs +++ b/src/NATS.Client.Core/NatsConnection.RequestReply.cs @@ -50,7 +50,9 @@ public async ValueTask> RequestAsync( using var rt = _replyTaskFactory.CreateReplyTask(replySerializer, replyOpts.Timeout); requestSerializer ??= Opts.SerializerRegistry.GetSerializer(); await PublishAsync(subject, data, headers, rt.Subject, requestSerializer, requestOpts, cancellationToken).ConfigureAwait(false); - return await rt.GetResultAsync(cancellationToken).ConfigureAwait(false); + var msg = await rt.GetResultAsync(cancellationToken).ConfigureAwait(false); + msg.Headers?.Activity?.Dispose(); + return msg; } await using var sub1 = await CreateRequestSubAsync(subject, data, headers, requestSerializer, replySerializer, requestOpts, replyOpts, cancellationToken) @@ -77,7 +79,9 @@ public async ValueTask> RequestAsync( using var rt = _replyTaskFactory.CreateReplyTask(replySerializer, replyOpts.Timeout); requestSerializer ??= Opts.SerializerRegistry.GetSerializer(); await PublishAsync(subject, data, headers, rt.Subject, requestSerializer, requestOpts, cancellationToken).ConfigureAwait(false); - return await rt.GetResultAsync(cancellationToken).ConfigureAwait(false); + var msg = await rt.GetResultAsync(cancellationToken).ConfigureAwait(false); + msg.Headers?.Activity?.Dispose(); + return msg; } await using var sub = await CreateRequestSubAsync(subject, data, headers, requestSerializer, replySerializer, requestOpts, replyOpts, cancellationToken) From 9743343ad19df2182f6f5ae317676103fc78a3d8 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 29 Jan 2026 12:51:55 +0000 Subject: [PATCH 2/2] Add tests Fix activity disposal for Direct request mode Ensure OpenTelemetry activities are explicitly disposed in Direct request/reply workflows. Refactor tests with `ActivityTracker` to detect leaks. --- examples/Example.OpenTelemetry/ClientApp.cs | 5 +- examples/Example.OpenTelemetry/ServiceApp.cs | 5 +- .../NatsConnection.RequestReply.cs | 7 +- .../OpenTelemetryTest.cs | 92 +++++++++++++++---- 4 files changed, 88 insertions(+), 21 deletions(-) diff --git a/examples/Example.OpenTelemetry/ClientApp.cs b/examples/Example.OpenTelemetry/ClientApp.cs index d2dba19ca..45a2ffeb7 100644 --- a/examples/Example.OpenTelemetry/ClientApp.cs +++ b/examples/Example.OpenTelemetry/ClientApp.cs @@ -24,7 +24,10 @@ public static async Task Run() Console.WriteLine("Client App is starting..."); - await using var nats = new NatsConnection(); + await using var nats = new NatsConnection(new NatsOpts + { + RequestReplyMode = NatsRequestReplyMode.Direct, + }); using (var activity = activitySource.StartActivity("SayHi")) { diff --git a/examples/Example.OpenTelemetry/ServiceApp.cs b/examples/Example.OpenTelemetry/ServiceApp.cs index 25b55bbc7..8b81bf11c 100644 --- a/examples/Example.OpenTelemetry/ServiceApp.cs +++ b/examples/Example.OpenTelemetry/ServiceApp.cs @@ -24,7 +24,10 @@ public static async Task Run() Console.WriteLine("Service App is starting..."); - await using var nats = new NatsConnection(); + await using var nats = new NatsConnection(new NatsOpts + { + RequestReplyMode = NatsRequestReplyMode.Direct, + }); await foreach (var msg in nats.SubscribeAsync("greet.>")) { diff --git a/src/NATS.Client.Core/NatsConnection.RequestReply.cs b/src/NATS.Client.Core/NatsConnection.RequestReply.cs index aa6885bce..fd70eb082 100644 --- a/src/NATS.Client.Core/NatsConnection.RequestReply.cs +++ b/src/NATS.Client.Core/NatsConnection.RequestReply.cs @@ -51,7 +51,10 @@ public async ValueTask> RequestAsync( requestSerializer ??= Opts.SerializerRegistry.GetSerializer(); await PublishAsync(subject, data, headers, rt.Subject, requestSerializer, requestOpts, cancellationToken).ConfigureAwait(false); var msg = await rt.GetResultAsync(cancellationToken).ConfigureAwait(false); + + // Dispose activity from headers to avoid leaking it msg.Headers?.Activity?.Dispose(); + return msg; } @@ -79,9 +82,7 @@ public async ValueTask> RequestAsync( using var rt = _replyTaskFactory.CreateReplyTask(replySerializer, replyOpts.Timeout); requestSerializer ??= Opts.SerializerRegistry.GetSerializer(); await PublishAsync(subject, data, headers, rt.Subject, requestSerializer, requestOpts, cancellationToken).ConfigureAwait(false); - var msg = await rt.GetResultAsync(cancellationToken).ConfigureAwait(false); - msg.Headers?.Activity?.Dispose(); - return msg; + return await rt.GetResultAsync(cancellationToken).ConfigureAwait(false); } await using var sub = await CreateRequestSubAsync(subject, data, headers, requestSerializer, replySerializer, requestOpts, replyOpts, cancellationToken) diff --git a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs index c7d05e924..7e41edfaa 100644 --- a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs +++ b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs @@ -14,8 +14,7 @@ public class OpenTelemetryTest [Fact] public async Task Publish_subscribe_activities() { - var activities = new List(); - using var activityListener = StartActivityListener(activities); + using var tracker = new ActivityTracker(); await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); @@ -27,14 +26,14 @@ public async Task Publish_subscribe_activities() await sub.Msgs.ReadAsync(cts.Token); - AssertActivityData("foo", activities); + AssertActivityData("foo", tracker.Started); + tracker.AssertAllStopped(); } [Fact] public async Task JetStream_consume_start_activity_with_interface() { - var activities = new List(); - using var activityListener = StartActivityListener(activities); + using var tracker = new ActivityTracker(); await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); @@ -61,30 +60,50 @@ public async Task JetStream_consume_start_activity_with_interface() } // Verify the publish activity was recorded - var publishActivity = activities.FirstOrDefault(x => x.OperationName == "test.subject publish"); + var publishActivity = tracker.Started.FirstOrDefault(x => x.OperationName == "test.subject publish"); Assert.NotNull(publishActivity); // Verify our custom activity was recorded - var consumeActivity = activities.FirstOrDefault(x => x.OperationName == "test.consume"); + var consumeActivity = tracker.Started.FirstOrDefault(x => x.OperationName == "test.consume"); Assert.NotNull(consumeActivity); Assert.Equal(ActivityKind.Internal, consumeActivity.Kind); // Verify the parent relationship (consume activity should have publish as parent via trace context) Assert.Equal(publishActivity.TraceId, consumeActivity.TraceId); + + tracker.AssertAllStopped(); } - private static ActivityListener StartActivityListener(List activities) + [Fact] + public async Task Direct_request_reply_receive_activity_is_disposed() { - var activityListener = new ActivityListener(); - activityListener.Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded; - activityListener.SampleUsingParentId = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded; - activityListener.ShouldListenTo = activitySource => activitySource.Name.StartsWith("NATS.Net"); - activityListener.ActivityStarted = activities.Add; - ActivitySource.AddActivityListener(activityListener); - return activityListener; + using var tracker = new ActivityTracker(); + + await using var server = await NatsServerProcess.StartAsync(); + await using var nats = new NatsConnection(new NatsOpts + { + Url = server.Url, + RequestReplyMode = NatsRequestReplyMode.Direct, + }); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + + var sub = await nats.SubscribeCoreAsync("foo.direct", cancellationToken: cts.Token); + var reg = sub.Register(async msg => + { + await msg.ReplyAsync(msg.Data * 2, cancellationToken: cts.Token); + }); + + var reply = await nats.RequestAsync("foo.direct", 21, cancellationToken: cts.Token); + Assert.Equal(42, reply.Data); + + tracker.AssertAllStopped(); + + await sub.DisposeAsync(); + await reg; } - private void AssertActivityData(string subject, List activityList) + private void AssertActivityData(string subject, IReadOnlyList activityList) { var activities = activityList.ToArray(); Assert.NotEmpty(activities); @@ -117,4 +136,45 @@ private void AssertStringTagNotNullOrEmpty(Activity activity, string name) Assert.NotNull(tag); Assert.False(string.IsNullOrEmpty(tag)); } + + private sealed class ActivityTracker : IDisposable + { + private readonly List _started = new(); + private readonly List _stopped = new(); + private readonly ActivityListener _listener; + + public ActivityTracker() + { + _listener = new ActivityListener + { + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded, + SampleUsingParentId = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded, + ShouldListenTo = source => source.Name.StartsWith("NATS.Net"), + ActivityStarted = _started.Add, + ActivityStopped = _stopped.Add, + }; + ActivitySource.AddActivityListener(_listener); + } + + public IReadOnlyList Started => _started; + + public IReadOnlyList Stopped => _stopped; + + public void AssertAllStopped() + { + Assert.NotEmpty(_started); + + var leaked = _started + .Where(started => !_stopped.Any(stopped => stopped.Id == started.Id)) + .ToList(); + + if (leaked.Count > 0) + { + var details = string.Join("\n", leaked.Select(a => $" [{a.Kind}] {a.OperationName} id={a.Id}")); + Assert.Fail($"Activity leak detected. {leaked.Count} activity(s) started but never stopped:\n{details}"); + } + } + + public void Dispose() => _listener.Dispose(); + } }