diff --git a/NATS.Net.slnx b/NATS.Net.slnx index 1dca45ac9..187f663d9 100644 --- a/NATS.Net.slnx +++ b/NATS.Net.slnx @@ -64,6 +64,7 @@ + diff --git a/src/NATS.Client.OpenTelemetry/NATS.Client.OpenTelemetry.csproj b/src/NATS.Client.OpenTelemetry/NATS.Client.OpenTelemetry.csproj new file mode 100644 index 000000000..e25ee2dae --- /dev/null +++ b/src/NATS.Client.OpenTelemetry/NATS.Client.OpenTelemetry.csproj @@ -0,0 +1,17 @@ + + + + + opentelemetry;tracing;metrics;observability;nats + OpenTelemetry instrumentation for NATS .NET. Adds the ActivitySource and Meter for the client to a TracerProviderBuilder or MeterProviderBuilder. + + + + + + + + + + + diff --git a/src/NATS.Client.OpenTelemetry/NatsInstrumentationExtensions.cs b/src/NATS.Client.OpenTelemetry/NatsInstrumentationExtensions.cs new file mode 100644 index 000000000..0ac563ce7 --- /dev/null +++ b/src/NATS.Client.OpenTelemetry/NatsInstrumentationExtensions.cs @@ -0,0 +1,38 @@ +using System; +using NATS.Client.Core; +using OpenTelemetry.Metrics; +using OpenTelemetry.Trace; + +namespace NATS.Client.OpenTelemetry; + +public static class NatsInstrumentationExtensions +{ + /// + /// Adds the NATS .NET client to the tracer provider, + /// enabling distributed tracing for publish, subscribe, and request/reply operations. + /// + /// The to add the source to. + /// The supplied for chaining. + public static TracerProviderBuilder AddNatsClientInstrumentation(this TracerProviderBuilder builder) => + builder.AddSource(NatsTelemetry.SourceName); + + /// + /// Adds the NATS .NET client to the tracer provider and + /// configures the shared (filter and enrich callbacks). + /// + /// The to add the source to. + /// Action that mutates the process-wide . + /// The supplied for chaining. + public static TracerProviderBuilder AddNatsClientInstrumentation(this TracerProviderBuilder builder, Action configure) + { + configure?.Invoke(NatsInstrumentationOptions.Default); + return builder.AddSource(NatsTelemetry.SourceName); + } + + /// + /// Adds the NATS .NET client to the meter provider, + /// enabling messaging metrics (published/consumed counters, operation duration, and more). + /// + public static MeterProviderBuilder AddNatsClientInstrumentation(this MeterProviderBuilder builder) => + builder.AddMeter(NatsTelemetry.SourceName); +} diff --git a/src/NATS.Client.OpenTelemetry/NatsInstrumentationOptionsExtensions.cs b/src/NATS.Client.OpenTelemetry/NatsInstrumentationOptionsExtensions.cs new file mode 100644 index 000000000..8a039c8cc --- /dev/null +++ b/src/NATS.Client.OpenTelemetry/NatsInstrumentationOptionsExtensions.cs @@ -0,0 +1,108 @@ +using System; +using NATS.Client.Core; + +namespace NATS.Client.OpenTelemetry; + +/// +/// Extension methods for . +/// +public static class NatsInstrumentationOptionsExtensions +{ + /// + /// Restricts tracing to operations whose subject matches the given NATS subject patterns. + /// + /// The options to configure. + /// + /// Subject patterns to trace. An operation is traced only if its subject matches at least one + /// pattern. When null or empty, every subject is eligible (still subject to ). + /// + /// + /// Subject patterns to skip. An operation matching any of these is not traced, even when it also + /// matches an include pattern. A common use is dropping inbox traffic with _INBOX.>. + /// + /// The same instance for chaining. + /// + /// Patterns use NATS subject wildcards: * matches a single token and > matches one or + /// more trailing tokens. The resulting predicate is combined (logical AND) with any existing + /// , so a previously configured filter still applies. + /// + public static NatsInstrumentationOptions FilterSubjects( + this NatsInstrumentationOptions options, + string[]? include = null, + string[]? exclude = null) + { + if (options is null) + throw new ArgumentNullException(nameof(options)); + + var includeTokens = Tokenize(include); + var excludeTokens = Tokenize(exclude); + + // Nothing to filter on; leave any existing filter untouched. + if (includeTokens is null && excludeTokens is null) + return options; + + var previous = options.Filter; + options.Filter = context => + { + if (previous is not null && !previous(context)) + return false; + + var subject = context.Subject; + + if (excludeTokens is not null) + { + foreach (var pattern in excludeTokens) + { + if (Matches(subject, pattern)) + return false; + } + } + + if (includeTokens is not null) + { + foreach (var pattern in includeTokens) + { + if (Matches(subject, pattern)) + return true; + } + + return false; + } + + return true; + }; + + return options; + } + + private static string[][]? Tokenize(string[]? patterns) + { + if (patterns is null || patterns.Length == 0) + return null; + + var result = new string[patterns.Length][]; + for (var i = 0; i < patterns.Length; i++) + result[i] = patterns[i].Split('.'); + + return result; + } + + // NATS subject match: '*' matches exactly one token, '>' matches one or more trailing tokens. + private static bool Matches(string subject, string[] pattern) + { + var tokens = subject.Split('.'); + for (var i = 0; i < pattern.Length; i++) + { + if (pattern[i] == ">") + return tokens.Length > i; + + if (i >= tokens.Length) + return false; + + if (pattern[i] != "*" && !string.Equals(pattern[i], tokens[i], StringComparison.Ordinal)) + return false; + } + + return tokens.Length == pattern.Length; + } +} 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 8e5dc1bc4..00957156e 100644 --- a/tests/NATS.Net.OpenTelemetry.Tests/NATS.Net.OpenTelemetry.Tests.csproj +++ b/tests/NATS.Net.OpenTelemetry.Tests/NATS.Net.OpenTelemetry.Tests.csproj @@ -11,6 +11,7 @@ + @@ -34,6 +35,7 @@ + diff --git a/tests/NATS.Net.OpenTelemetry.Tests/NatsInstrumentationExtensionsTest.cs b/tests/NATS.Net.OpenTelemetry.Tests/NatsInstrumentationExtensionsTest.cs new file mode 100644 index 000000000..0fce374ca --- /dev/null +++ b/tests/NATS.Net.OpenTelemetry.Tests/NatsInstrumentationExtensionsTest.cs @@ -0,0 +1,125 @@ +using NATS.Client.OpenTelemetry; +using OpenTelemetry; +using OpenTelemetry.Metrics; +using OpenTelemetry.Trace; + +namespace NATS.Client.Core.Tests; + +public class NatsInstrumentationExtensionsTest +{ + [Fact] + public void AddNatsClientInstrumentation_builds_tracer_provider() + { + using var provider = Sdk.CreateTracerProviderBuilder() + .AddNatsClientInstrumentation() + .Build(); + + provider.Should().NotBeNull(); + } + + [Fact] + public void AddNatsClientInstrumentation_builds_meter_provider() + { + using var provider = Sdk.CreateMeterProviderBuilder() + .AddNatsClientInstrumentation() + .Build(); + + provider.Should().NotBeNull(); + } + + [Fact] + public void AddNatsClientInstrumentation_with_configure_sets_options() + { + var configured = false; + try + { + using var provider = Sdk.CreateTracerProviderBuilder() + .AddNatsClientInstrumentation(options => + { + configured = true; + options.Filter = _ => true; + options.Enrich = (_, _) => { }; + }) + .Build(); + + configured.Should().BeTrue(); + NatsInstrumentationOptions.Default.Filter.Should().NotBeNull(); + NatsInstrumentationOptions.Default.Enrich.Should().NotBeNull(); + } + finally + { + NatsInstrumentationOptions.Default.Filter = null; + NatsInstrumentationOptions.Default.Enrich = null; + } + } + + [Fact] + public void SourceName_is_NATS_Net() + { + NatsTelemetry.SourceName.Should().Be("NATS.Net"); + } + + [Theory] + [InlineData("orders.new", true)] + [InlineData("orders.new.eu", true)] + [InlineData("orders", false)] // '>' needs at least one trailing token + [InlineData("payments.new", false)] + public void FilterSubjects_include_matches_only_listed(string subject, bool expected) + { + var options = new NatsInstrumentationOptions().FilterSubjects(include: ["orders.>"]); + + options.Filter!(Context(subject)).Should().Be(expected); + } + + [Theory] + [InlineData("_INBOX.abc.def", false)] + [InlineData("foo.bar", true)] + public void FilterSubjects_exclude_drops_listed(string subject, bool expected) + { + var options = new NatsInstrumentationOptions().FilterSubjects(exclude: ["_INBOX.>"]); + + options.Filter!(Context(subject)).Should().Be(expected); + } + + [Theory] + [InlineData("foo.bar", true)] + [InlineData("foo.bar.baz", false)] // '*' matches a single token only + [InlineData("foo", false)] + public void FilterSubjects_single_token_wildcard(string subject, bool expected) + { + var options = new NatsInstrumentationOptions().FilterSubjects(include: ["foo.*"]); + + options.Filter!(Context(subject)).Should().Be(expected); + } + + [Fact] + public void FilterSubjects_exclude_wins_over_include() + { + var options = new NatsInstrumentationOptions().FilterSubjects(include: ["orders.>"], exclude: ["orders.internal.>"]); + + options.Filter!(Context("orders.new")).Should().BeTrue(); + options.Filter!(Context("orders.internal.audit")).Should().BeFalse(); + } + + [Fact] + public void FilterSubjects_composes_with_existing_filter() + { + var options = new NatsInstrumentationOptions { Filter = ctx => ctx.Subject.StartsWith("orders.", StringComparison.Ordinal) }; + options.FilterSubjects(exclude: ["orders.internal.>"]); + + options.Filter!(Context("orders.new")).Should().BeTrue(); + options.Filter!(Context("orders.internal.audit")).Should().BeFalse(); // dropped by subject exclude + options.Filter!(Context("payments.new")).Should().BeFalse(); // dropped by the pre-existing filter + } + + [Fact] + public void FilterSubjects_without_patterns_leaves_filter_unset() + { + var options = new NatsInstrumentationOptions().FilterSubjects(); + + options.Filter.Should().BeNull(); + } + + private static NatsInstrumentationContext Context(string subject) => + new(subject, Headers: null, ReplyTo: null, QueueGroup: null, BodySize: null, Size: null, Connection: null, ParentContext: default); +} diff --git a/tools/site_src/documentation/advanced/opentelemetry.md b/tools/site_src/documentation/advanced/opentelemetry.md index d88a2f14e..317b32808 100644 --- a/tools/site_src/documentation/advanced/opentelemetry.md +++ b/tools/site_src/documentation/advanced/opentelemetry.md @@ -47,6 +47,19 @@ telemetry for specific requests. When the filter returns `false`, no activity is [!code-csharp[](../../../../tests/NATS.Net.DocsExamples/Advanced/OpenTelemetryPage.cs#filter)] +When the `NATS.Client.OpenTelemetry` package is installed, `FilterSubjects` builds the predicate from +NATS subject patterns (`*` matches one token, `>` matches one or more trailing tokens) instead of writing +the matching by hand. Include patterns allow-list subjects; exclude patterns drop them and win over +include. The predicate is combined (logical AND) with any filter already set: + +```csharp +Sdk.CreateTracerProviderBuilder() + .AddNatsClientInstrumentation(options => options.FilterSubjects( + include: ["orders.>"], + exclude: ["orders.internal.>"])) + .Build(); +``` + ## Enriching Activities Use [`NatsInstrumentationOptions.Default.Enrich`](xref:NATS.Client.Core.NatsInstrumentationOptions) to add @@ -114,3 +127,22 @@ application, per the OTel definition ("messages delivered to the application"). frames consumed internally by the client are excluded: no-responder `503` replies and JetStream heartbeats, flow-control, and protocol notifications. The two counters stay consistent, so `received.bytes / consumed.messages` reflects average delivered message size. + +### Histogram Buckets + +`messaging.client.operation.duration` ships advisory bucket boundaries (`0.005s` to `10s`) through +`InstrumentAdvice`, which the OpenTelemetry SDK applies by default, so no view is required for sensible +latency buckets. To override them, add a view on the meter provider (this needs the `OpenTelemetry` SDK +package, not just `OpenTelemetry.Api`): + +```csharp +Sdk.CreateMeterProviderBuilder() + .AddNatsClientInstrumentation() + .AddView( + "messaging.client.operation.duration", + new ExplicitBucketHistogramConfiguration + { + Boundaries = [0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1, 5], + }) + .Build(); +```