diff --git a/examples/Example.OpenTelemetry/ClientApp.cs b/examples/Example.OpenTelemetry/ClientApp.cs index 45a2ffeb7..d6f8e89bf 100644 --- a/examples/Example.OpenTelemetry/ClientApp.cs +++ b/examples/Example.OpenTelemetry/ClientApp.cs @@ -1,6 +1,9 @@ using System.Diagnostics; +using Microsoft.Extensions.Logging; using NATS.Client.Core; using OpenTelemetry; +using OpenTelemetry.Logs; +using OpenTelemetry.Metrics; using OpenTelemetry.Resources; using OpenTelemetry.Trace; @@ -13,19 +16,41 @@ public static async Task Run() var serviceName = "ClientApp"; var serviceVersion = "1.0.0"; + var resourceBuilder = ResourceBuilder.CreateDefault().AddService(serviceName: serviceName, serviceVersion: serviceVersion); + using var tracerProvider = Sdk.CreateTracerProviderBuilder() .AddOtlpExporter() - .SetResourceBuilder(ResourceBuilder.CreateDefault().AddService(serviceName: serviceName, serviceVersion: serviceVersion)) - .AddSource("NATS.Net") + .SetResourceBuilder(resourceBuilder) + .AddSource(NatsTelemetry.SourceName) .AddSource("MyClientSource") .Build(); + using var meterProvider = Sdk.CreateMeterProviderBuilder() + .AddOtlpExporter() + .SetResourceBuilder(resourceBuilder) + .AddMeter(NatsTelemetry.SourceName) + .Build(); + + using var loggerFactory = LoggerFactory.Create(builder => + { + builder.AddOpenTelemetry(options => + { + options.SetResourceBuilder(resourceBuilder); + options.IncludeFormattedMessage = true; + options.IncludeScopes = true; + options.ParseStateValues = true; + options.AddOtlpExporter(); + }); + }); + var logger = loggerFactory.CreateLogger(serviceName); + ActivitySource activitySource = new("MyClientSource"); - Console.WriteLine("Client App is starting..."); + logger.LogInformation("Client App is starting..."); await using var nats = new NatsConnection(new NatsOpts { + LoggerFactory = loggerFactory, RequestReplyMode = NatsRequestReplyMode.Direct, }); @@ -34,7 +59,7 @@ public static async Task Run() await nats.PublishAsync("greet.presence.client.app", "ClientApp is here!"); var response = await nats.RequestAsync("greet.hi", "Hi, telemetry!"); - Console.WriteLine($"Response: {response}"); + logger.LogInformation("Response: {Response}", response); } } } diff --git a/examples/Example.OpenTelemetry/Program.cs b/examples/Example.OpenTelemetry/Program.cs index fa2b06779..a6c33545d 100644 --- a/examples/Example.OpenTelemetry/Program.cs +++ b/examples/Example.OpenTelemetry/Program.cs @@ -2,26 +2,28 @@ OpenTelemetry Example -(1) Run Jaeger locally and then run the client and server apps. +Both apps export traces and metrics via OTLP. Point them at any backend that +accepts OTLP (Aspire dashboard, Jaeger for traces, Grafana stack, etc.). -https://www.jaegertracing.io/download/ - -https://medium.com/jaegertracing/introducing-native-support-for-opentelemetry-in-jaeger-eb661be8183c +(1) Start Aspire dashboard: ```powershell -> $env:COLLECTOR_OTLP_ENABLED=true -> jaeger-all-in-one.exe +> docker run --rm -it ` + -p 18888:18888 -p 4317:18889 ` + -e DASHBOARD__OTLP__AUTHMODE=Unsecured ` + mcr.microsoft.com/dotnet/aspire-dashboard:latest ``` or ```bash -$ COLLECTOR_OTLP_ENABLED=true jaeger-all-in-one +$ docker run --rm -it \ + -p 18888:18888 -p 4317:18889 \ + -e DASHBOARD__OTLP__AUTHMODE=Unsecured \ + mcr.microsoft.com/dotnet/aspire-dashboard:latest ``` -(2) Jaeger UI default URL http://localhost:16686/search - -(3) In different terminals run: +(2) In different terminals run: ``` nats-server diff --git a/examples/Example.OpenTelemetry/ServiceApp.cs b/examples/Example.OpenTelemetry/ServiceApp.cs index 8b81bf11c..183fc49cb 100644 --- a/examples/Example.OpenTelemetry/ServiceApp.cs +++ b/examples/Example.OpenTelemetry/ServiceApp.cs @@ -1,6 +1,9 @@ using System.Diagnostics; +using Microsoft.Extensions.Logging; using NATS.Client.Core; using OpenTelemetry; +using OpenTelemetry.Logs; +using OpenTelemetry.Metrics; using OpenTelemetry.Resources; using OpenTelemetry.Trace; @@ -13,19 +16,41 @@ public static async Task Run() var serviceName = "ServiceApp"; var serviceVersion = "1.0.0"; + var resourceBuilder = ResourceBuilder.CreateDefault().AddService(serviceName: serviceName, serviceVersion: serviceVersion); + using var tracerProvider = Sdk.CreateTracerProviderBuilder() .AddOtlpExporter() - .SetResourceBuilder(ResourceBuilder.CreateDefault().AddService(serviceName: serviceName, serviceVersion: serviceVersion)) - .AddSource("NATS.Net") + .SetResourceBuilder(resourceBuilder) + .AddSource(NatsTelemetry.SourceName) .AddSource("MyServiceSource") .Build(); + using var meterProvider = Sdk.CreateMeterProviderBuilder() + .AddOtlpExporter() + .SetResourceBuilder(resourceBuilder) + .AddMeter(NatsTelemetry.SourceName) + .Build(); + + using var loggerFactory = LoggerFactory.Create(builder => + { + builder.AddOpenTelemetry(options => + { + options.SetResourceBuilder(resourceBuilder); + options.IncludeFormattedMessage = true; + options.IncludeScopes = true; + options.ParseStateValues = true; + options.AddOtlpExporter(); + }); + }); + var logger = loggerFactory.CreateLogger(serviceName); + ActivitySource activitySource = new("MyServiceSource"); - Console.WriteLine("Service App is starting..."); + logger.LogInformation("Service App is starting..."); await using var nats = new NatsConnection(new NatsOpts { + LoggerFactory = loggerFactory, RequestReplyMode = NatsRequestReplyMode.Direct, }); @@ -35,7 +60,7 @@ public static async Task Run() if (msg.Subject.StartsWith("greet.presence")) { - Console.WriteLine($"{msg.Data} is here!"); + logger.LogInformation("{Data} is here!", msg.Data); activity?.AddEvent(new ActivityEvent("Presence", tags: new() { diff --git a/src/NATS.Client.Core/Commands/CommandWriter.cs b/src/NATS.Client.Core/Commands/CommandWriter.cs index 790238dee..5b75be112 100644 --- a/src/NATS.Client.Core/Commands/CommandWriter.cs +++ b/src/NATS.Client.Core/Commands/CommandWriter.cs @@ -349,6 +349,21 @@ public ValueTask PublishAsync(string subject, T? value, NatsHeaders? headers, { throw new NatsPayloadTooLargeException($"Payload size {size} exceeds server's maximum payload size {info.MaxPayload}"); } + + // Per OTel messaging semconv, messaging.client.published.messages counts publish + // attempts, not successful wire writes -- so we increment here, before the + // _semLock / _disposed / state-machine paths that might fail to enqueue. + // sent.bytes follows the same "attempted" convention for consistency. Operators + // wanting a "successfully published" view subtract operation.duration samples + // that carry an error.type tag. + if (Telemetry.PublishedMessages.Enabled || Telemetry.SentBytes.Enabled) + { + var tags = Telemetry.BuildMetricTags(_connection, Telemetry.Constants.OpPub); + if (Telemetry.PublishedMessages.Enabled) + Telemetry.PublishedMessages.Add(1, tags); + if (Telemetry.SentBytes.Enabled) + Telemetry.SentBytes.Add(size, tags); + } } catch { diff --git a/src/NATS.Client.Core/Internal/ReplyTask.cs b/src/NATS.Client.Core/Internal/ReplyTask.cs index e7c1fbcb3..ae02c85af 100644 --- a/src/NATS.Client.Core/Internal/ReplyTask.cs +++ b/src/NATS.Client.Core/Internal/ReplyTask.cs @@ -19,6 +19,8 @@ internal sealed class ReplyTask : ReplyTaskBase, IDisposable private readonly TimeSpan _requestTimeout; private readonly TaskCompletionSource _tcs; private NatsMsg _msg; + private long _replyBytes; + private bool _isNoResponders; public ReplyTask(ReplyTaskFactory factory, long id, string subject, NatsConnection connection, INatsDeserialize deserializer, TimeSpan requestTimeout) { @@ -51,10 +53,30 @@ await _tcs.Task NatsNoReplyException.Throw(); } + NatsMsg msg; + long bytes; + bool isNoResponders; lock (_gate) { - return _msg; + msg = _msg; + bytes = _replyBytes; + isNoResponders = _isNoResponders; } + + // Count only messages actually delivered to the caller. Late replies that arrive + // after a timeout still hit SetResult, but the user never sees them, so the + // counters belong here on the success path, not in SetResult. 503 NoResponders + // sentinels are also excluded for parity with the SharedInbox path. + if (!isNoResponders && (Telemetry.ConsumedMessages.Enabled || Telemetry.ReceivedBytes.Enabled)) + { + var tags = Telemetry.BuildMetricTags(_connection, Telemetry.Constants.OpRec); + if (Telemetry.ConsumedMessages.Enabled) + Telemetry.ConsumedMessages.Add(1, tags); + if (Telemetry.ReceivedBytes.Enabled) + Telemetry.ReceivedBytes.Add(bytes, tags); + } + + return msg; } public override void SetResult(string? replyTo, ReadOnlySequence payload, ReadOnlySequence? headersBuffer) @@ -62,6 +84,8 @@ public override void SetResult(string? replyTo, ReadOnlySequence payload, lock (_gate) { _msg = NatsMsg.Build(Subject, replyTo, headersBuffer, payload, _connection, _connection.HeaderParser, _deserializer); + _isNoResponders = payload.Length == 0 && NatsSubBase.IsHeader503(headersBuffer); + _replyBytes = payload.Length + (headersBuffer?.Length ?? 0); } _tcs.TrySetResult(); diff --git a/src/NATS.Client.Core/Internal/Telemetry.cs b/src/NATS.Client.Core/Internal/Telemetry.cs index 87fd0bff6..8f3ac6386 100644 --- a/src/NATS.Client.Core/Internal/Telemetry.cs +++ b/src/NATS.Client.Core/Internal/Telemetry.cs @@ -1,17 +1,114 @@ using System.Diagnostics; +using System.Diagnostics.Metrics; namespace NATS.Client.Core.Internal; // https://opentelemetry.io/docs/specs/semconv/attributes-registry/messaging/ // https://opentelemetry.io/docs/specs/semconv/messaging/messaging-spans/#messaging-attributes +// https://opentelemetry.io/docs/specs/semconv/messaging/messaging-metrics/ internal static class Telemetry { - public const string NatsActivitySource = "NATS.Net"; + public const string NatsActivitySource = NatsTelemetry.SourceName; public static readonly ActivitySource NatsActivities = new(name: NatsActivitySource); + + public static readonly Meter NatsMeter = new(name: NatsActivitySource); + + public static readonly Counter PublishedMessages = + NatsMeter.CreateCounter("messaging.client.published.messages", unit: "{message}"); + + public static readonly Counter ConsumedMessages = + NatsMeter.CreateCounter("messaging.client.consumed.messages", unit: "{message}"); + + // OTel messaging semconv recommends these advisory buckets for messaging.client.operation.duration. + // https://opentelemetry.io/docs/specs/semconv/messaging/messaging-metrics/#metric-messagingclientoperationduration + public static readonly Histogram OperationDuration = + NatsMeter.CreateHistogram( + "messaging.client.operation.duration", + unit: "s", + advice: new InstrumentAdvice + { + HistogramBucketBoundaries = new[] { 0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1, 2.5, 5, 7.5, 10 }, + }); + + public static readonly UpDownCounter ActiveSubscriptions = + NatsMeter.CreateUpDownCounter("nats.client.active_subscriptions", unit: "{subscription}"); + + public static readonly Counter Reconnects = + NatsMeter.CreateCounter("nats.client.reconnects", unit: "{reconnect}"); + + public static readonly Counter SentBytes = + NatsMeter.CreateCounter("nats.client.sent.bytes", unit: "By"); + + public static readonly Counter ReceivedBytes = + NatsMeter.CreateCounter("nats.client.received.bytes", unit: "By"); + private static readonly object BoxedTrue = true; + /// + /// Don't use this for metrics. + /// public static bool HasListeners() => NatsActivities.HasListeners(); + public static void RecordOperationDuration(long startTimestamp, INatsConnection? connection, string operation, Exception? error) + { + if (!OperationDuration.Enabled) + return; + + try + { + var elapsed = (Stopwatch.GetTimestamp() - startTimestamp) / (double)Stopwatch.Frequency; + var tags = BuildMetricTags(connection, operation); + if (error is not null) + tags.Add(Constants.ErrorTypeKey, error.GetType().FullName ?? "unknown"); + + OperationDuration.Record(elapsed, tags); + } + catch + { + // Instrumentation must never break the calling operation. A buggy MeterListener + // or tag construction failure here would otherwise replace the in-flight messaging + // exception (catch/finally semantics), hiding the real failure from the caller. + } + } + + public static async ValueTask MeasureOperationAsync(ValueTask task, long startTimestamp, INatsConnection? connection, string operation) + { + try + { + await task.ConfigureAwait(false); + RecordOperationDuration(startTimestamp, connection, operation, null); + } + catch (Exception ex) + { + RecordOperationDuration(startTimestamp, connection, operation, ex); + throw; + } + } + + public static TagList BuildMetricTags(INatsConnection? connection, string operation) + { + var tags = default(TagList); + + // MetricTagsPrefix is read off the concrete NatsConnection by design: exposing it on + // INatsConnection would be a breaking change for external implementers, and default + // interface members aren't available on netstandard2.0. Custom INatsConnection wrappers + // therefore fall through to the minimal tag set below. Normal usage (including via + // NatsJSContext) hits this branch because the runtime type is NatsConnection regardless + // of the static type the caller holds. + if (connection is NatsConnection { MetricTagsPrefix: { } prefix }) + { + for (var i = 0; i < prefix.Length; i++) + tags.Add(prefix[i]); + } + else + { + tags.Add(Constants.SystemKey, Constants.SystemVal); + } + + tags.Add(Constants.OpKey, operation); + return tags; + } + public static Activity? StartSendActivity( string name, INatsConnection? connection, @@ -300,6 +397,10 @@ public class Constants public const string OpKey = "messaging.operation"; public const string OpPub = "publish"; public const string OpRec = "receive"; + public const string OpSub = "subscribe"; + public const string OpReq = "request"; + public const string OpReconnect = "reconnect"; + public const string ErrorTypeKey = "error.type"; public const string MsgBodySize = "messaging.message.body.size"; public const string MsgTotalSize = "messaging.message.envelope.size"; diff --git a/src/NATS.Client.Core/NATS.Client.Core.csproj b/src/NATS.Client.Core/NATS.Client.Core.csproj index b83da55fc..0410eeb58 100644 --- a/src/NATS.Client.Core/NATS.Client.Core.csproj +++ b/src/NATS.Client.Core/NATS.Client.Core.csproj @@ -20,13 +20,18 @@ - all runtime; build; native; contentfiles; analyzers + + + + + + diff --git a/src/NATS.Client.Core/NatsConnection.Publish.cs b/src/NATS.Client.Core/NatsConnection.Publish.cs index 8dd86e629..66b666d33 100644 --- a/src/NATS.Client.Core/NatsConnection.Publish.cs +++ b/src/NATS.Client.Core/NatsConnection.Publish.cs @@ -1,3 +1,4 @@ +using System.Diagnostics; using NATS.Client.Core.Internal; namespace NATS.Client.Core; @@ -13,26 +14,45 @@ public ValueTask PublishAsync(string subject, NatsHeaders? headers = default, st SubjectValidator.ValidateReplyTo(replyTo); } + var measure = Telemetry.OperationDuration.Enabled; + var start = measure ? Stopwatch.GetTimestamp() : 0L; + + ValueTask task; if (Telemetry.HasListeners()) { using var activity = Telemetry.StartSendActivity($"{SpanDestinationName(subject)} {Telemetry.Constants.PublishActivityName}", this, subject, replyTo); Telemetry.AddTraceContextHeaders(activity, ref headers); try { - return ConnectionState != NatsConnectionState.Open + task = ConnectionState != NatsConnectionState.Open ? ConnectAndPublishAsync(subject, default, headers, replyTo, NatsRawSerializer.Default, cancellationToken) : CommandWriter.PublishAsync(subject, default, headers, replyTo, NatsRawSerializer.Default, cancellationToken); } catch (Exception ex) { Telemetry.SetException(activity, ex); + if (measure) + Telemetry.RecordOperationDuration(start, this, Telemetry.Constants.OpPub, ex); + throw; + } + } + else + { + try + { + task = ConnectionState != NatsConnectionState.Open + ? ConnectAndPublishAsync(subject, default, headers, replyTo, NatsRawSerializer.Default, cancellationToken) + : CommandWriter.PublishAsync(subject, default, headers, replyTo, NatsRawSerializer.Default, cancellationToken); + } + catch (Exception ex) + { + if (measure) + Telemetry.RecordOperationDuration(start, this, Telemetry.Constants.OpPub, ex); throw; } } - return ConnectionState != NatsConnectionState.Open - ? ConnectAndPublishAsync(subject, default, headers, replyTo, NatsRawSerializer.Default, cancellationToken) - : CommandWriter.PublishAsync(subject, default, headers, replyTo, NatsRawSerializer.Default, cancellationToken); + return measure ? Telemetry.MeasureOperationAsync(task, start, this, Telemetry.Constants.OpPub) : task; } /// @@ -44,6 +64,10 @@ public ValueTask PublishAsync(string subject, T? data, NatsHeaders? headers = SubjectValidator.ValidateReplyTo(replyTo); } + var measure = Telemetry.OperationDuration.Enabled; + var start = measure ? Stopwatch.GetTimestamp() : 0L; + + ValueTask task; if (Telemetry.HasListeners()) { using var activity = Telemetry.StartSendActivity($"{SpanDestinationName(subject)} {Telemetry.Constants.PublishActivityName}", this, subject, replyTo); @@ -51,21 +75,36 @@ public ValueTask PublishAsync(string subject, T? data, NatsHeaders? headers = try { serializer ??= Opts.SerializerRegistry.GetSerializer(); - return ConnectionState != NatsConnectionState.Open + task = ConnectionState != NatsConnectionState.Open ? ConnectAndPublishAsync(subject, data, headers, replyTo, serializer, cancellationToken) : CommandWriter.PublishAsync(subject, data, headers, replyTo, serializer, cancellationToken); } catch (Exception ex) { Telemetry.SetException(activity, ex); + if (measure) + Telemetry.RecordOperationDuration(start, this, Telemetry.Constants.OpPub, ex); + throw; + } + } + else + { + try + { + serializer ??= Opts.SerializerRegistry.GetSerializer(); + task = ConnectionState != NatsConnectionState.Open + ? ConnectAndPublishAsync(subject, data, headers, replyTo, serializer, cancellationToken) + : CommandWriter.PublishAsync(subject, data, headers, replyTo, serializer, cancellationToken); + } + catch (Exception ex) + { + if (measure) + Telemetry.RecordOperationDuration(start, this, Telemetry.Constants.OpPub, ex); throw; } } - serializer ??= Opts.SerializerRegistry.GetSerializer(); - return ConnectionState != NatsConnectionState.Open - ? ConnectAndPublishAsync(subject, data, headers, replyTo, serializer, cancellationToken) - : CommandWriter.PublishAsync(subject, data, headers, replyTo, serializer, cancellationToken); + return measure ? Telemetry.MeasureOperationAsync(task, start, this, Telemetry.Constants.OpPub) : task; } /// diff --git a/src/NATS.Client.Core/NatsConnection.RequestReply.cs b/src/NATS.Client.Core/NatsConnection.RequestReply.cs index 0dde7c32d..466a429ea 100644 --- a/src/NATS.Client.Core/NatsConnection.RequestReply.cs +++ b/src/NATS.Client.Core/NatsConnection.RequestReply.cs @@ -38,62 +38,79 @@ public async ValueTask> RequestAsync( SubjectValidator.ValidateSubject(subject); } - if (Telemetry.HasListeners()) + var measure = Telemetry.OperationDuration.Enabled; + var start = measure ? Stopwatch.GetTimestamp() : 0L; + Exception? error = null; + + try { - using var activity = Telemetry.StartSendActivity($"{SpanDestinationName(subject)} {Telemetry.Constants.RequestReplyActivityName}", this, subject, null); - try + if (Telemetry.HasListeners()) { - replyOpts = SetReplyOptsDefaults(replyOpts); - - if (Opts.RequestReplyMode == NatsRequestReplyMode.Direct) + using var activity = Telemetry.StartSendActivity($"{SpanDestinationName(subject)} {Telemetry.Constants.RequestReplyActivityName}", this, subject, null); + try { - 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); + replyOpts = SetReplyOptsDefaults(replyOpts); - // Dispose activity from headers to avoid leaking it - msg.Headers?.Activity?.Dispose(); + if (Opts.RequestReplyMode == NatsRequestReplyMode.Direct) + { + 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); - return msg; - } + // Dispose activity from headers to avoid leaking it + msg.Headers?.Activity?.Dispose(); + + return msg; + } + + await using var sub1 = await CreateRequestSubAsync(subject, data, headers, requestSerializer, replySerializer, requestOpts, replyOpts, cancellationToken) + .ConfigureAwait(false); - await using var sub1 = await CreateRequestSubAsync(subject, data, headers, requestSerializer, replySerializer, requestOpts, replyOpts, cancellationToken) - .ConfigureAwait(false); + await foreach (var msg in sub1.Msgs.ReadAllAsync(cancellationToken).ConfigureAwait(false)) + { + return msg; + } - await foreach (var msg in sub1.Msgs.ReadAllAsync(cancellationToken).ConfigureAwait(false)) + throw new NatsNoReplyException(); + } + catch (Exception e) { - return msg; + Telemetry.SetException(activity, e); + throw; } - - throw new NatsNoReplyException(); } - catch (Exception e) + + replyOpts = SetReplyOptsDefaults(replyOpts); + + if (Opts.RequestReplyMode == NatsRequestReplyMode.Direct) { - Telemetry.SetException(activity, e); - throw; + 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); } - } - replyOpts = SetReplyOptsDefaults(replyOpts); + await using var sub = await CreateRequestSubAsync(subject, data, headers, requestSerializer, replySerializer, requestOpts, replyOpts, cancellationToken) + .ConfigureAwait(false); - if (Opts.RequestReplyMode == NatsRequestReplyMode.Direct) + await foreach (var msg in sub.Msgs.ReadAllAsync(cancellationToken).ConfigureAwait(false)) + { + return msg; + } + + throw new NatsNoReplyException(); + } + catch (Exception ex) { - 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); + error = ex; + throw; } - - await using var sub = await CreateRequestSubAsync(subject, data, headers, requestSerializer, replySerializer, requestOpts, replyOpts, cancellationToken) - .ConfigureAwait(false); - - await foreach (var msg in sub.Msgs.ReadAllAsync(cancellationToken).ConfigureAwait(false)) + finally { - return msg; + if (measure) + Telemetry.RecordOperationDuration(start, this, Telemetry.Constants.OpReq, error); } - - throw new NatsNoReplyException(); } /// diff --git a/src/NATS.Client.Core/NatsConnection.Subscribe.cs b/src/NATS.Client.Core/NatsConnection.Subscribe.cs index ede74b4e4..abad357fa 100644 --- a/src/NATS.Client.Core/NatsConnection.Subscribe.cs +++ b/src/NATS.Client.Core/NatsConnection.Subscribe.cs @@ -1,3 +1,4 @@ +using System.Diagnostics; using System.Runtime.CompilerServices; using NATS.Client.Core.Internal; @@ -37,7 +38,24 @@ private async IAsyncEnumerable> SubscribeInternalAsync(string subj serializer ??= Opts.SerializerRegistry.GetDeserializer(); await using var sub = new NatsSub(this, _subscriptionManager.GetManagerFor(subject), subject, queueGroup, opts, serializer, cancellationToken); - await AddSubAsync(sub, cancellationToken: cancellationToken).ConfigureAwait(false); + + var measure = Telemetry.OperationDuration.Enabled; + var start = measure ? Stopwatch.GetTimestamp() : 0L; + Exception? error = null; + try + { + await AddSubAsync(sub, cancellationToken: cancellationToken).ConfigureAwait(false); + } + catch (Exception ex) + { + error = ex; + throw; + } + finally + { + if (measure) + Telemetry.RecordOperationDuration(start, this, Telemetry.Constants.OpSub, error); + } // We don't cancel the channel reader here because we want to keep reading until the subscription // channel writer completes so that messages left in the channel can be consumed before exit the loop. @@ -51,7 +69,25 @@ private async ValueTask> SubscribeCoreInternalAsync(string subjec { serializer ??= Opts.SerializerRegistry.GetDeserializer(); var sub = new NatsSub(this, _subscriptionManager.GetManagerFor(subject), subject, queueGroup, opts, serializer, cancellationToken); - await AddSubAsync(sub, cancellationToken).ConfigureAwait(false); + + var measure = Telemetry.OperationDuration.Enabled; + var start = measure ? Stopwatch.GetTimestamp() : 0L; + Exception? error = null; + try + { + await AddSubAsync(sub, cancellationToken).ConfigureAwait(false); + } + catch (Exception ex) + { + error = ex; + throw; + } + finally + { + if (measure) + Telemetry.RecordOperationDuration(start, this, Telemetry.Constants.OpSub, error); + } + return sub; } } diff --git a/src/NATS.Client.Core/NatsConnection.cs b/src/NATS.Client.Core/NatsConnection.cs index 5e72efbcd..f80b556c9 100644 --- a/src/NATS.Client.Core/NatsConnection.cs +++ b/src/NATS.Client.Core/NatsConnection.cs @@ -57,6 +57,7 @@ public partial class NatsConnection : INatsConnection private readonly HashSet _drainParticipants = new(); private ServerInfo? _writableServerInfo; + private KeyValuePair[]? _metricTagsPrefix; private int _pongCount; private int _connectionState; private int _isDisposed; @@ -162,10 +163,35 @@ internal ServerInfo? WritableServerInfo PushEvent(NatsEvent.LameDuckModeActivated, new NatsLameDuckModeActivatedEventArgs(_currentConnectUri!.Uri)); } + KeyValuePair[]? prefix = null; + if (value is not null) + { + // ServerInfo.Host/Port is the server's bind address (often 0.0.0.0); use the + // address the client actually dialled for OTel server.address/server.port. + var connectUri = _currentConnectUri; + var host = connectUri?.Host ?? value.Host; + var port = connectUri?.Port ?? value.Port; + prefix = new[] + { + new KeyValuePair(Telemetry.Constants.SystemKey, Telemetry.Constants.SystemVal), + new KeyValuePair(Telemetry.Constants.ServerAddress, host), + new KeyValuePair(Telemetry.Constants.ServerPort, (object)port), + new KeyValuePair(Telemetry.Constants.NetworkProtoName, "nats"), + new KeyValuePair(Telemetry.Constants.NetworkTransport, "tcp"), + }; + } + + // Publish prefix before ServerInfo so any reader observing the new ServerInfo + // is guaranteed to also observe the matching prefix. Today no single reader + // reads both, but #1156 will unify the trace and metric tag sources and start + // relying on this ordering. + Volatile.Write(ref _metricTagsPrefix, prefix); Interlocked.Exchange(ref _writableServerInfo, value); } } + internal KeyValuePair[]? MetricTagsPrefix => Volatile.Read(ref _metricTagsPrefix); + internal bool IsDisposed { get => Interlocked.CompareExchange(ref _isDisposed, 0, 0) == 1; @@ -839,18 +865,23 @@ private async void ReconnectLoop() goto CONNECT_AGAIN; } + bool emitReconnect; lock (_gate) { _connectRetry = 0; _backoff = TimeSpan.Zero; _logger.LogInformation(NatsLogEvents.Connection, "Connection succeeded {Name}, NATS {Url} [{ReconnectCount}]", _name, url, reconnectCount); ConnectionState = NatsConnectionState.Open; + emitReconnect = Telemetry.Reconnects.Enabled; _pingTimerCancellationTokenSource = new CancellationTokenSource(); StartPingTimer(_pingTimerCancellationTokenSource.Token); _waitForOpenConnection.TrySetResult(); _reconnectLoopTask = Task.Run(ReconnectLoop); PushEvent(NatsEvent.ConnectionOpened, new NatsEventArgs(url.ToString())); } + + if (emitReconnect) + Telemetry.Reconnects.Add(1, Telemetry.BuildMetricTags(this, Telemetry.Constants.OpReconnect)); } catch (Exception ex) { diff --git a/src/NATS.Client.Core/NatsSubBase.cs b/src/NATS.Client.Core/NatsSubBase.cs index bd2271035..26893e799 100644 --- a/src/NATS.Client.Core/NatsSubBase.cs +++ b/src/NATS.Client.Core/NatsSubBase.cs @@ -54,6 +54,7 @@ public abstract class NatsSubBase private int _pendingMsgs; private Exception? _exception; private int _isSlowConsumer; + private int _telemetryActive; private TaskCompletionSource? _readerExited; /// @@ -208,6 +209,25 @@ protected NatsSubBase( /// A that represents the asynchronous operation. public virtual ValueTask ReadyAsync() { + // nats.client.active_subscriptions counts every NatsSubBase that reaches Ready, + // regardless of whether it corresponds to a wire-level SUB. Under SharedInbox + // request/reply mode each in-flight RequestAsync registers a transient reply + // NatsSub with the shared inbox muxer (no extra wire SUB) and that registration + // still bumps the gauge. Operator dashboards should expect oscillation under + // request load. Direct mode uses ReplyTask and does not increment. + if (Telemetry.ActiveSubscriptions.Enabled) + { + var emit = false; + lock (_gate) + { + if (!_unsubscribed && Interlocked.Exchange(ref _telemetryActive, 1) == 0) + emit = true; + } + + if (emit) + Telemetry.ActiveSubscriptions.Add(1, Telemetry.BuildMetricTags(Connection, Telemetry.Constants.OpSub)); + } + // Let idle timer start with the first message, in case // we're allowed to wait longer for the first message. if (_startUpTimeoutTimer == null) @@ -233,6 +253,8 @@ public ValueTask UnsubscribeAsync() _unsubscribed = true; } + DecrementActiveSubscription(); + _timeoutTimer?.Change(Timeout.Infinite, Timeout.Infinite); _idleTimeoutTimer?.Change(Timeout.Infinite, Timeout.Infinite); _startUpTimeoutTimer?.Change(Timeout.Infinite, Timeout.Infinite); @@ -321,6 +343,23 @@ public virtual async ValueTask ReceiveAsync(string subject, string? replyTo, Rea { // Need to await to handle any exceptions await ReceiveInternalAsync(subject, replyTo, headersBuffer, payloadBuffer).ConfigureAwait(false); + + // Per OTel messaging semconv, consumed.messages counts messages "delivered to + // the application", so we only record after ReceiveInternalAsync has completed + // the hand-off to the subscription channel. received.bytes follows the same + // convention. Messages dropped by cancellation, channel closure, or serializer + // exceptions are not counted. + if (Telemetry.ConsumedMessages.Enabled || Telemetry.ReceivedBytes.Enabled) + { + var tags = Telemetry.BuildMetricTags(Connection, Telemetry.Constants.OpRec); + if (Telemetry.ConsumedMessages.Enabled) + Telemetry.ConsumedMessages.Add(1, tags); + if (Telemetry.ReceivedBytes.Enabled) + { + var bytes = payloadBuffer.Length + (headersBuffer?.Length ?? 0); + Telemetry.ReceivedBytes.Add(bytes, tags); + } + } } catch (TaskCanceledException) { @@ -393,6 +432,8 @@ internal async ValueTask DrainAsync() if (needsUnsub) { + DecrementActiveSubscription(); + _timeoutTimer?.Change(Timeout.Infinite, Timeout.Infinite); _idleTimeoutTimer?.Change(Timeout.Infinite, Timeout.Infinite); _startUpTimeoutTimer?.Change(Timeout.Infinite, Timeout.Infinite); @@ -608,6 +649,12 @@ protected void EndSubscription(NatsSubEndReason reason) #pragma warning restore CA2012 } + private void DecrementActiveSubscription() + { + if (Interlocked.Exchange(ref _telemetryActive, 0) == 1) + Telemetry.ActiveSubscriptions.Add(-1, Telemetry.BuildMetricTags(Connection, Telemetry.Constants.OpSub)); + } + private async ValueTask WaitForReaderDrainCoreAsync(Task task, TimeSpan timeout) { try diff --git a/src/NATS.Client.Core/NatsTelemetry.cs b/src/NATS.Client.Core/NatsTelemetry.cs new file mode 100644 index 000000000..6c8583617 --- /dev/null +++ b/src/NATS.Client.Core/NatsTelemetry.cs @@ -0,0 +1,12 @@ +namespace NATS.Client.Core; + +/// +/// Telemetry identifiers for NATS .NET. Use these when configuring OpenTelemetry +/// or any other listener directly. The same name is used for both the +/// and the +/// . +/// +public static class NatsTelemetry +{ + public const string SourceName = "NATS.Net"; +} diff --git a/tests/NATS.Client.TestUtilities/NATS.Client.TestUtilities.csproj b/tests/NATS.Client.TestUtilities/NATS.Client.TestUtilities.csproj index a9c661236..4944afde5 100644 --- a/tests/NATS.Client.TestUtilities/NATS.Client.TestUtilities.csproj +++ b/tests/NATS.Client.TestUtilities/NATS.Client.TestUtilities.csproj @@ -25,7 +25,7 @@ - + all runtime; build; native; contentfiles; analyzers diff --git a/tests/NATS.Net.DocsExamples/Advanced/OpenTelemetryPage.cs b/tests/NATS.Net.DocsExamples/Advanced/OpenTelemetryPage.cs index 7bb117d36..a1d08efee 100644 --- a/tests/NATS.Net.DocsExamples/Advanced/OpenTelemetryPage.cs +++ b/tests/NATS.Net.DocsExamples/Advanced/OpenTelemetryPage.cs @@ -48,6 +48,36 @@ public async Task Run() #endregion } + { + #region metrics-setup + // The NATS.Net client emits metrics through System.Diagnostics.Metrics.Meter + // under the same "NATS.Net" name used for activities. Metrics are opt-in: + // nothing is recorded until something subscribes to the meter. + + // Using the OpenTelemetry SDK and the NATS.Client.OpenTelemetry package: + // + // using var meterProvider = Sdk.CreateMeterProviderBuilder() + // .AddNatsClientInstrumentation() // or .AddMeter("NATS.Net") + // .AddOtlpExporter() + // .Build(); + + // Or using a plain MeterListener (no extra packages): + using System.Diagnostics.Metrics.MeterListener meterListener = new() + { + InstrumentPublished = (instrument, listener) => + { + if (instrument.Meter.Name == "NATS.Net") + listener.EnableMeasurementEvents(instrument); + }, + }; + meterListener.SetMeasurementEventCallback((inst, value, tags, _) => + Console.WriteLine($"{inst.Name}: {value}")); + meterListener.SetMeasurementEventCallback((inst, value, tags, _) => + Console.WriteLine($"{inst.Name}: {value}")); + meterListener.Start(); + #endregion + } + { #region publish-subscribe await using NatsConnection nats = new NatsConnection(); diff --git a/tests/NATS.Net.OpenTelemetry.Tests/NATS.Net.OpenTelemetry.Tests.csproj b/tests/NATS.Net.OpenTelemetry.Tests/NATS.Net.OpenTelemetry.Tests.csproj index 55caf094c..8e5dc1bc4 100644 --- a/tests/NATS.Net.OpenTelemetry.Tests/NATS.Net.OpenTelemetry.Tests.csproj +++ b/tests/NATS.Net.OpenTelemetry.Tests/NATS.Net.OpenTelemetry.Tests.csproj @@ -49,5 +49,4 @@ - diff --git a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs index 42bac7fa5..21e190bf2 100644 --- a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs +++ b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs @@ -1,4 +1,5 @@ using System.Diagnostics; +using System.Diagnostics.Metrics; using NATS.Client.JetStream; using NATS.Client.JetStream.Models; using Synadia.Orbit.Testing.NatsServerProcessManager; @@ -74,6 +75,326 @@ public async Task JetStream_consume_start_activity_with_interface() tracker.AssertAllStopped(); } + [Fact] + public async Task Publish_counter() + { + using var meter = new MeterTracker(); + 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)); + + // Intentionally not calling ConnectAsync first; the first publish triggers connect. + // Counter must not record measurements with (Unset) server tags. + for (var i = 0; i < 5; i++) + { + await nats.PublishAsync("foo", i, cancellationToken: cts.Token); + } + + var published = meter.LongMeasurements + .Where(m => m.Name == "messaging.client.published.messages") + .ToList(); + + published.Sum(m => m.Value).Should().Be(5); + + // Every measurement must carry the full server tag set, including the first publish + // that triggered the connect. + foreach (var m in published) + { + var tags = m.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("publish"); + tags.Should().ContainKey("server.address"); + tags.Should().ContainKey("server.port"); + tags.Should().NotContainKey("messaging.destination.name"); + tags.Should().NotContainKey("messaging.nats.message.subject"); + } + } + + [Fact] + public async Task Consume_counter() + { + using var meter = new MeterTracker(); + 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)); + + await using var sub = await nats.SubscribeCoreAsync("foo.consume", cancellationToken: cts.Token); + + for (var i = 0; i < 5; i++) + { + await nats.PublishAsync("foo.consume", i, cancellationToken: cts.Token); + } + + for (var i = 0; i < 5; i++) + { + await sub.Msgs.ReadAsync(cts.Token); + } + + var consumed = meter.LongMeasurements + .Where(m => m.Name == "messaging.client.consumed.messages") + .ToList(); + + consumed.Sum(m => m.Value).Should().Be(5); + + var tags = consumed[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"); + tags.Should().ContainKey("server.address"); + tags.Should().ContainKey("server.port"); + tags.Should().NotContainKey("messaging.destination.name"); + tags.Should().NotContainKey("messaging.nats.message.subject"); + } + + [Fact] + public async Task Consume_counter_direct_request_reply() + { + using var meter = new MeterTracker(); + 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.consume", 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.consume", 21, cancellationToken: cts.Token); + reply.Data.Should().Be(42); + + var consumed = meter.LongMeasurements + .Where(m => m.Name == "messaging.client.consumed.messages") + .Sum(m => m.Value); + + consumed.Should().Be(2); + + await sub.DisposeAsync(); + await reg; + } + + [Fact] + public async Task Active_subscriptions_updown_counter() + { + using var meter = new MeterTracker(); + 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 sub1 = await nats.SubscribeCoreAsync("foo.active.1", cancellationToken: cts.Token); + var sub2 = await nats.SubscribeCoreAsync("foo.active.2", cancellationToken: cts.Token); + + var active = meter.LongMeasurements + .Where(m => m.Name == "nats.client.active_subscriptions") + .ToList(); + + active.Sum(m => m.Value).Should().Be(2); + + var tags = active[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("subscribe"); + tags.Should().ContainKey("server.address"); + tags.Should().ContainKey("server.port"); + + await sub1.DisposeAsync(); + + meter.LongMeasurements + .Where(m => m.Name == "nats.client.active_subscriptions") + .Sum(m => m.Value).Should().Be(1); + + await sub2.DisposeAsync(); + + meter.LongMeasurements + .Where(m => m.Name == "nats.client.active_subscriptions") + .Sum(m => m.Value).Should().Be(0); + } + + [Fact] + public async Task Publish_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 cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + await nats.ConnectAsync(); + + await nats.PublishAsync("foo.duration", 42, cancellationToken: cts.Token); + + var durations = meter.DoubleMeasurements + .Where(m => m.Name == "messaging.client.operation.duration") + .ToList(); + + var publish = durations.Where(m => m.Tags.Any(t => t.Key == "messaging.operation" && (string?)t.Value == "publish")).ToList(); + publish.Should().NotBeEmpty(); + publish[0].Value.Should().BeGreaterThan(0); + + var tags = publish[0].Tags.ToDictionary(t => t.Key, t => t.Value); + tags.Should().ContainKey("messaging.system").WhoseValue.Should().Be("nats"); + tags.Should().ContainKey("server.address"); + tags.Should().ContainKey("server.port"); + tags.Should().NotContainKey("error.type"); + } + + [Fact] + public async Task Request_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 cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + + var sub = await nats.SubscribeCoreAsync("foo.req", cancellationToken: cts.Token); + var reg = sub.Register(async msg => await msg.ReplyAsync(msg.Data * 2, cancellationToken: cts.Token)); + + var reply = await nats.RequestAsync("foo.req", 21, cancellationToken: cts.Token); + reply.Data.Should().Be(42); + + var request = meter.DoubleMeasurements + .Where(m => m.Name == "messaging.client.operation.duration") + .Where(m => m.Tags.Any(t => t.Key == "messaging.operation" && (string?)t.Value == "request")) + .ToList(); + + request.Should().HaveCount(1); + request[0].Value.Should().BeGreaterThan(0); + + var tags = request[0].Tags.ToDictionary(t => t.Key, t => t.Value); + tags.Should().ContainKey("messaging.system").WhoseValue.Should().Be("nats"); + tags.Should().NotContainKey("error.type"); + + await sub.DisposeAsync(); + await reg; + } + + [Fact] + public async Task Reconnect_counter() + { + using var meter = new MeterTracker(); + await using var server = await NatsServerProcess.StartAsync(); + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + + // Attach before ConnectAsync so we observe both opens and can wait for the second + // (the reconnect) without racing the async event-channel delivery of the first. + var opened = 0; + var openedTwice = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + nats.ConnectionOpened += (_, _) => + { + if (Interlocked.Increment(ref opened) == 2) + openedTwice.TrySetResult(); + return default; + }; + + await nats.ConnectAsync(); + await nats.ReconnectAsync(); + await openedTwice.Task.WaitAsync(TimeSpan.FromSeconds(30)); + + var reconnects = meter.LongMeasurements + .Where(m => m.Name == "nats.client.reconnects") + .ToList(); + + reconnects.Sum(m => m.Value).Should().Be(1); + + var tags = reconnects[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("reconnect"); + tags.Should().ContainKey("server.address"); + tags.Should().ContainKey("server.port"); + } + + [Fact] + public async Task Sent_and_received_bytes_counters() + { + using var meter = new MeterTracker(); + 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)); + + await using var sub = await nats.SubscribeCoreAsync("foo.bytes", cancellationToken: cts.Token); + + var payload = new byte[1024]; + await nats.PublishAsync("foo.bytes", payload, cancellationToken: cts.Token); + + var msg = await sub.Msgs.ReadAsync(cts.Token); + msg.Data!.Length.Should().Be(1024); + + var sent = meter.LongMeasurements + .Where(m => m.Name == "nats.client.sent.bytes") + .Sum(m => m.Value); + var received = meter.LongMeasurements + .Where(m => m.Name == "nats.client.received.bytes") + .Sum(m => m.Value); + + sent.Should().Be(1024); + received.Should().Be(1024); + + var sentTags = meter.LongMeasurements.First(m => m.Name == "nats.client.sent.bytes").Tags.ToDictionary(t => t.Key, t => t.Value); + sentTags.Should().ContainKey("messaging.system").WhoseValue.Should().Be("nats"); + sentTags.Should().ContainKey("messaging.operation").WhoseValue.Should().Be("publish"); + + var recvTags = meter.LongMeasurements.First(m => m.Name == "nats.client.received.bytes").Tags.ToDictionary(t => t.Key, t => t.Value); + recvTags.Should().ContainKey("messaging.operation").WhoseValue.Should().Be("receive"); + } + + [Fact] + public async Task Subscribe_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 cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + + await using var sub = await nats.SubscribeCoreAsync("foo.sub.duration", cancellationToken: cts.Token); + + var subscribe = meter.DoubleMeasurements + .Where(m => m.Name == "messaging.client.operation.duration") + .Where(m => m.Tags.Any(t => t.Key == "messaging.operation" && (string?)t.Value == "subscribe")) + .ToList(); + + subscribe.Should().HaveCount(1); + subscribe[0].Value.Should().BeGreaterThan(0); + + var tags = subscribe[0].Tags.ToDictionary(t => t.Key, t => t.Value); + tags.Should().ContainKey("messaging.system").WhoseValue.Should().Be("nats"); + tags.Should().ContainKey("server.address"); + tags.Should().ContainKey("server.port"); + tags.Should().NotContainKey("error.type"); + } + + [Fact] + public async Task Request_operation_duration_records_error_type_on_failure() + { + using var meter = new MeterTracker(); + 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 act = async () => await nats.RequestAsync( + "foo.no.responders", + 1, + replyOpts: new NatsSubOpts { Timeout = TimeSpan.FromSeconds(1) }, + cancellationToken: cts.Token); + + await act.Should().ThrowAsync(); + + var request = meter.DoubleMeasurements + .Where(m => m.Name == "messaging.client.operation.duration") + .Where(m => m.Tags.Any(t => t.Key == "messaging.operation" && (string?)t.Value == "request")) + .ToList(); + + request.Should().HaveCount(1); + var tags = request[0].Tags.ToDictionary(t => t.Key, t => t.Value); + tags.Should().ContainKey("error.type"); + ((string?)tags["error.type"]).Should().Contain("NatsNoRespondersException"); + } + [Fact] public async Task Direct_request_reply_receive_activity_is_disposed() { @@ -185,4 +506,34 @@ public void AssertAllStopped() public void Dispose() => _listener.Dispose(); } + + private sealed class MeterTracker : IDisposable + { + private readonly MeterListener _listener; + + public MeterTracker() + { + _listener = new MeterListener + { + InstrumentPublished = (instrument, listener) => + { + if (instrument.Meter.Name == "NATS.Net") + listener.EnableMeasurementEvents(instrument); + }, + }; + + _listener.SetMeasurementEventCallback((inst, val, tags, _) => + LongMeasurements.Add((inst.Name, val, tags.ToArray()))); + _listener.SetMeasurementEventCallback((inst, val, tags, _) => + DoubleMeasurements.Add((inst.Name, val, tags.ToArray()))); + + _listener.Start(); + } + + public List<(string Name, long Value, KeyValuePair[] Tags)> LongMeasurements { get; } = new(); + + public List<(string Name, double Value, KeyValuePair[] Tags)> DoubleMeasurements { get; } = new(); + + public void Dispose() => _listener.Dispose(); + } } diff --git a/tools/site_src/documentation/advanced/opentelemetry.md b/tools/site_src/documentation/advanced/opentelemetry.md index 84ce55ac0..c062ca001 100644 --- a/tools/site_src/documentation/advanced/opentelemetry.md +++ b/tools/site_src/documentation/advanced/opentelemetry.md @@ -1,10 +1,12 @@ # OpenTelemetry -NATS.Net has built-in distributed tracing support using [`System.Diagnostics.Activity`](https://learn.microsoft.com/dotnet/api/system.diagnostics.activity), -the standard .NET API for OpenTelemetry. Activities are created automatically for publish and subscribe operations, -and trace context is propagated through message headers so that send and receive spans are linked across services. +NATS.Net has built-in distributed tracing and metrics support through `System.Diagnostics.Activity` and +`System.Diagnostics.Metrics.Meter`, the standard .NET APIs for OpenTelemetry. Activities are created +automatically for publish and subscribe operations, trace context is propagated through message headers +so send and receive spans are linked across services, and a set of standard messaging metrics is emitted +when a meter listener is attached. -The activity source name is `NATS.Net`. +The activity source name and the meter name are both `NATS.Net`. ## Setting Up Tracing @@ -13,6 +15,13 @@ with an exporter (Jaeger, Zipkin, OTLP, etc.) or a plain `ActivityListener` for [!code-csharp[](../../../../tests/NATS.Net.DocsExamples/Advanced/OpenTelemetryPage.cs#setup)] +## Setting Up Metrics + +Metrics are emitted through the same `NATS.Net` name. No measurements are recorded until a listener +subscribes; the runtime cost is a single boolean check per operation when no listener is attached: + +[!code-csharp[](../../../../tests/NATS.Net.DocsExamples/Advanced/OpenTelemetryPage.cs#metrics-setup)] + ## Automatic Trace Context Propagation When you publish a message, the client injects the current trace context into the message headers. @@ -72,3 +81,30 @@ Receive activities include additional attributes: | `messaging.message.body.size` | `1024` | Message body size in bytes | | `messaging.message.envelope.size` | `1280` | Total message size in bytes | | `messaging.consumer.group.name` | `workers` | Queue group (if used) | + +## Metrics + +The following instruments are exposed on the `NATS.Net` meter: + +| Name | Type | Unit | Description | +|---|---|---|---| +| `messaging.client.published.messages` | Counter | `{message}` | Messages published by the client | +| `messaging.client.consumed.messages` | Counter | `{message}` | Messages received by the client | +| `messaging.client.operation.duration` | Histogram | `s` | Duration of publish, request, and subscribe operations | +| `nats.client.active_subscriptions` | UpDownCounter | `{subscription}` | Active `NatsSubBase` instances. Under `SharedInbox` request/reply mode each in-flight `RequestAsync` registers a transient reply subscription with the shared inbox muxer and is included here; `Direct` mode uses a reply task and is not counted | +| `nats.client.reconnects` | Counter | `{reconnect}` | Successful reconnects since process start | +| `nats.client.sent.bytes` | Counter | `By` | Bytes sent in published messages (body + headers) | +| `nats.client.received.bytes` | Counter | `By` | Bytes received in consumed messages (body + headers) | + +All instruments carry these tags: + +| Tag | Example | Description | +|---|---|---| +| `messaging.system` | `nats` | Always `nats` | +| `messaging.operation` | `publish` / `receive` / `subscribe` / `request` / `reconnect` | Operation type | +| `server.address` | `localhost` | Server host | +| `server.port` | `4222` | Server port | +| `network.protocol.name` | `nats` | Protocol name | +| `network.transport` | `tcp` | Transport protocol | + +`messaging.client.operation.duration` adds `error.type` (full exception type name) when the operation fails.