diff --git a/src/NATS.Client.Core/Internal/Telemetry.cs b/src/NATS.Client.Core/Internal/Telemetry.cs index 4a9ede654..ad53e13f9 100644 --- a/src/NATS.Client.Core/Internal/Telemetry.cs +++ b/src/NATS.Client.Core/Internal/Telemetry.cs @@ -42,6 +42,13 @@ internal static class Telemetry public static readonly Counter ReceivedBytes = NatsMeter.CreateCounter("nats.client.received.bytes", unit: "By"); + // No messaging semantic convention covers drops, so this is a deliberately NATS-specific + // metric. Shares the consumed.messages tag set (messaging.operation=receive) so it correlates + // with the rest of the receive-path signals. Pending channel depth at drop time is not added + // as a tag because its value is unbounded; it stays available on the MessageDropped event. + public static readonly Counter DroppedMessages = + NatsMeter.CreateCounter("nats.client.messages.dropped", unit: "{message}"); + private static readonly object BoxedTrue = true; /// @@ -144,7 +151,7 @@ public static TagList BuildMetricTags(INatsConnection? connection, string operat tags = new KeyValuePair[len]; tags[0] = new KeyValuePair(Constants.SystemKey, Constants.SystemVal); tags[1] = new KeyValuePair(Constants.OpKey, Constants.OpPub); - tags[2] = new KeyValuePair(Constants.DestName, subject); + tags[2] = new KeyValuePair(Constants.DestName, LowCardinalitySubject(conn, subject)); tags[3] = new KeyValuePair(Constants.ClientId, conn.ClientId); tags[4] = new KeyValuePair(Constants.ServerAddress, serverHost); @@ -156,7 +163,7 @@ public static TagList BuildMetricTags(INatsConnection? connection, string operat tags[10] = new KeyValuePair(Constants.NetworkLocalAddress, conn.ServerInfo.ClientIp); if (replyTo is not null) - tags[11] = new KeyValuePair(Constants.ReplyTo, replyTo); + tags[11] = new KeyValuePair(Constants.ReplyTo, LowCardinalitySubject(conn, replyTo)); } else { @@ -253,9 +260,10 @@ public static void AddTraceContextHeaders(Activity? activity, ref NatsHeaders? h tags[1] = new KeyValuePair(Constants.OpKey, Constants.OpRec); tags[2] = new KeyValuePair(Constants.DestTemplate, subscriptionSubject); tags[3] = new KeyValuePair(Constants.DestIsTemporary, subscriptionSubject.StartsWith(conn.InboxPrefix, StringComparison.Ordinal) ? Constants.True : Constants.False); - tags[4] = new KeyValuePair(Constants.Subject, subject); - tags[5] = new KeyValuePair(Constants.DestName, subject); - tags[6] = new KeyValuePair(Constants.DestPubName, subject); + var lowCardinalitySubject = LowCardinalitySubject(conn, subject); + tags[4] = new KeyValuePair(Constants.Subject, lowCardinalitySubject); + tags[5] = new KeyValuePair(Constants.DestName, lowCardinalitySubject); + tags[6] = new KeyValuePair(Constants.DestPubName, lowCardinalitySubject); tags[7] = new KeyValuePair(Constants.MsgBodySize, bodySize.ToString()); tags[8] = new KeyValuePair(Constants.MsgTotalSize, size.ToString()); tags[9] = new KeyValuePair(Constants.ClientId, conn.ClientId); @@ -269,7 +277,7 @@ public static void AddTraceContextHeaders(Activity? activity, ref NatsHeaders? h var index = 17; if (replyTo is not null) - tags[index++] = new KeyValuePair(Constants.ReplyTo, replyTo); + tags[index++] = new KeyValuePair(Constants.ReplyTo, LowCardinalitySubject(conn, replyTo)); if (queueGroup is not null) tags[index] = new KeyValuePair(Constants.QueueGroup, queueGroup); } @@ -346,6 +354,16 @@ static string GetStackTrace(Exception? exception) } } + // Inbox subjects (_INBOX.[.]) are unique per request, so emitting them as + // indexed tag values pushes unbounded cardinality into tracing backends (Tempo, Jaeger). + // Collapse them to a constant, matching SpanDestinationName. The raw subject stays + // available to the Enrich callback via NatsInstrumentationContext. + // Uses Opts.InboxPrefix (the bare prefix, e.g. "_INBOX") rather than conn.InboxPrefix + // (this connection's "_INBOX."): a reply-to address can belong to any connection, + // so the match must be broad. Do not "normalise" this to conn.InboxPrefix. + private static string LowCardinalitySubject(NatsConnection conn, string subject) + => subject.StartsWith(conn.Opts.InboxPrefix, StringComparison.Ordinal) ? Constants.InboxName : subject; + private static bool TryParseTraceContext(NatsHeaders headers, out ActivityContext context) { DistributedContextPropagator.Current.ExtractTraceIdAndState( @@ -388,6 +406,7 @@ public class Constants { public const string True = "true"; public const string False = "false"; + public const string InboxName = "inbox"; public const string RequestReplyActivityName = "request"; public const string PublishActivityName = "publish"; public const string SubscribeActivityName = "subscribe"; @@ -401,6 +420,7 @@ public class Constants public const string OpRec = "receive"; public const string OpSub = "subscribe"; public const string OpReq = "request"; + public const string OpAck = "ack"; public const string OpReconnect = "reconnect"; public const string ErrorTypeKey = "error.type"; public const string MsgBodySize = "messaging.message.body.size"; diff --git a/src/NATS.Client.Core/NatsConnection.cs b/src/NATS.Client.Core/NatsConnection.cs index 743287a31..38f986c03 100644 --- a/src/NATS.Client.Core/NatsConnection.cs +++ b/src/NATS.Client.Core/NatsConnection.cs @@ -291,6 +291,9 @@ public async ValueTask ConnectAsync() /// public void OnMessageDropped(NatsSubBase natsSub, int pending, NatsMsg msg) { + if (Telemetry.DroppedMessages.Enabled) + Telemetry.DroppedMessages.Add(1, Telemetry.BuildMetricTags(this, Telemetry.Constants.OpRec)); + var subject = msg.Subject; PushEvent(NatsEvent.MessageDropped, new NatsMessageDroppedEventArgs(natsSub, pending, subject, msg.ReplyTo, msg.Headers, msg.Data)); diff --git a/src/NATS.Client.JetStream/NatsJSMsg.cs b/src/NATS.Client.JetStream/NatsJSMsg.cs index 503afbb8d..797d34b26 100644 --- a/src/NATS.Client.JetStream/NatsJSMsg.cs +++ b/src/NATS.Client.JetStream/NatsJSMsg.cs @@ -1,7 +1,9 @@ using System.Buffers; +using System.Diagnostics; using System.Diagnostics.CodeAnalysis; using System.Text; using NATS.Client.Core; +using NATS.Client.Core.Internal; using NATS.Client.JetStream.Internal; namespace NATS.Client.JetStream; @@ -231,21 +233,41 @@ private async ValueTask SendAckAsync(ReadOnlySequence payload, AckOpts? op if (_msg == default) throw new NatsJSException("No user message, can't acknowledge"); - if (opts?.DoubleAck ?? _context.Opts.DoubleAck) + // All ack-protocol messages (Ack, Nak, AckProgress, AckTerminate) route through here and + // record under a single OpAck operation, so their durations share one histogram. This is + // intentional: the operation tag tracks "sent an ack-protocol reply", not the ack kind. + // Do not split per kind without weighing the added metric cardinality. + var measure = Telemetry.OperationDuration.Enabled; + var start = measure ? Stopwatch.GetTimestamp() : 0L; + Exception? error = null; + try + { + if (opts?.DoubleAck ?? _context.Opts.DoubleAck) + { + await Connection.RequestAsync, object?>( + subject: ReplyTo, + data: payload, + requestSerializer: NatsRawSerializer>.Default, + replySerializer: NatsRawSerializer.Default, + cancellationToken: cancellationToken); + } + else + { + await _msg.ReplyAsync( + data: payload, + serializer: NatsRawSerializer>.Default, + cancellationToken: cancellationToken); + } + } + catch (Exception ex) { - await Connection.RequestAsync, object?>( - subject: ReplyTo, - data: payload, - requestSerializer: NatsRawSerializer>.Default, - replySerializer: NatsRawSerializer.Default, - cancellationToken: cancellationToken); + error = ex; + throw; } - else + finally { - await _msg.ReplyAsync( - data: payload, - serializer: NatsRawSerializer>.Default, - cancellationToken: cancellationToken); + if (measure) + Telemetry.RecordOperationDuration(start, Connection, Telemetry.Constants.OpAck, error); } } diff --git a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs index 97aa0fd84..a6917b9f6 100644 --- a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs +++ b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs @@ -442,6 +442,165 @@ public async Task Direct_request_reply_receive_activity_is_disposed() await reg; } + [Fact] + public async Task Ack_operation_duration_histogram() + { + using var meter = new MeterTracker(); + await using var server = await NatsServerProcess.StartAsync(); + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + + var js = new NatsJSContext(nats); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + + await js.CreateStreamAsync(new StreamConfig { Name = "ack-stream", Subjects = ["ack.>"] }, cts.Token); + await js.CreateOrUpdateConsumerAsync("ack-stream", new ConsumerConfig("ack-consumer"), cts.Token); + + await js.PublishAsync("ack.subject", "test-message", cancellationToken: cts.Token); + + var consumer = await js.GetConsumerAsync("ack-stream", "ack-consumer", cts.Token); + + await foreach (var msg in consumer.ConsumeAsync(cancellationToken: cts.Token)) + { + await msg.AckAsync(cancellationToken: cts.Token); + break; + } + + var ack = meter.DoubleMeasurements + .Where(m => m.Name == "messaging.client.operation.duration") + .Where(m => m.Tags.Any(t => t.Key == "messaging.operation" && (string?)t.Value == "ack")) + .ToList(); + + ack.Should().NotBeEmpty(); + ack[0].Value.Should().BeGreaterThan(0); + + var tags = ack[0].Tags.ToDictionary(t => t.Key, t => t.Value); + tags.Should().ContainKey("messaging.system").WhoseValue.Should().Be("nats"); + tags.Should().NotContainKey("error.type"); + } + + [Fact] + public async Task Dropped_messages_counter() + { + using var meter = new MeterTracker(); + await using var server = await NatsServerProcess.StartAsync(); + await using var nats = new NatsConnection(new NatsOpts + { + Url = server.Url, + SubPendingChannelCapacity = 3, + }); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + var cancellationToken = cts.Token; + + var dropped = 0; + nats.MessageDropped += (_, _) => + { + Interlocked.Increment(ref dropped); + return default; + }; + + var sync = 0; + var signal = new WaitSignal(); + + // Block the consumer after the sync message so the bounded channel overflows. + var subscription = Task.Run( + async () => + { + await foreach (var msg in nats.SubscribeAsync("drop.>", cancellationToken: cancellationToken)) + { + if (msg.Subject == "drop.sync") + { + Interlocked.Increment(ref sync); + await signal; + continue; + } + + if (msg.Subject == "drop.end") + { + break; + } + } + }, + cancellationToken); + + await Retry.Until( + "subscription is ready", + () => Volatile.Read(ref sync) > 0, + async () => await nats.PublishAsync("drop.sync", cancellationToken: cancellationToken)); + + for (var i = 0; i < 20; i++) + { + await nats.PublishAsync("drop.data", $"msg{i}", cancellationToken: cancellationToken); + } + + await Retry.Until("messages are dropped", () => Volatile.Read(ref dropped) > 0); + + signal.Pulse(); + await Retry.Until( + "subscription ended", + () => subscription.IsCompleted, + async () => await nats.PublishAsync("drop.end", cancellationToken: cancellationToken)); + await subscription; + + var droppedMeasurements = meter.LongMeasurements + .Where(m => m.Name == "nats.client.messages.dropped") + .ToList(); + + droppedMeasurements.Sum(m => m.Value).Should().BeGreaterThan(0); + + var tags = droppedMeasurements[0].Tags.ToDictionary(t => t.Key, t => t.Value); + tags.Should().ContainKey("messaging.system").WhoseValue.Should().Be("nats"); + tags.Should().ContainKey("messaging.operation").WhoseValue.Should().Be("receive"); + } + + [Fact] + public async Task Inbox_subjects_collapsed_in_trace_tags() + { + using var tracker = new ActivityTracker(); + await using var server = await NatsServerProcess.StartAsync(); + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + + var sub = await nats.SubscribeCoreAsync("foo.inbox", cancellationToken: cts.Token); + var reg = sub.Register(async msg => await msg.ReplyAsync(msg.Data * 2, cancellationToken: cts.Token)); + + var reply = await nats.RequestAsync("foo.inbox", 21, cancellationToken: cts.Token); + reply.Data.Should().Be(42); + + // The request's reply-to is an inbox; it must be collapsed to "inbox" rather than + // emitting the unique _INBOX. value that would blow up backend tag cardinality. + var replyToTags = tracker.Started + .Select(a => a.GetTagItem("messaging.nats.message.reply_to") as string) + .Where(v => v is not null) + .ToList(); + + replyToTags.Should().NotBeEmpty(); + replyToTags.Should().AllSatisfy(v => v.Should().Be("inbox")); + + // No high-cardinality tag should leak a raw inbox subject. + string[] cardinalityTags = + [ + "messaging.destination.name", + "messaging.destination_publish.name", + "messaging.nats.message.subject", + "messaging.nats.message.reply_to", + ]; + + foreach (var activity in tracker.Started) + { + foreach (var tag in cardinalityTags) + { + if (activity.GetTagItem(tag) is string value) + value.Should().NotStartWith("_INBOX", $"{tag} should not carry a raw inbox subject"); + } + } + + await sub.DisposeAsync(); + await reg; + } + [Fact] public async Task Consumed_counter_excludes_jetstream_control_messages() {