Skip to content
Merged
Show file tree
Hide file tree
Changes from 5 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions NATS.Net.slnx
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@
<Project Path="src/NATS.Client.JetStream/NATS.Client.JetStream.csproj" />
<Project Path="src/NATS.Client.KeyValueStore/NATS.Client.KeyValueStore.csproj" />
<Project Path="src/NATS.Client.ObjectStore/NATS.Client.ObjectStore.csproj" />
<Project Path="src/NATS.Client.OpenTelemetry/NATS.Client.OpenTelemetry.csproj" />
<Project Path="src/NATS.Client.Serializers.Json/NATS.Client.Serializers.Json.csproj" />
<Project Path="src/NATS.Client.Services/NATS.Client.Services.csproj" />
<Project Path="src/NATS.Client.Simplified/NATS.Client.Simplified.csproj" />
Expand Down
17 changes: 17 additions & 0 deletions src/NATS.Client.OpenTelemetry/NATS.Client.OpenTelemetry.csproj
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<!-- NuGet Packaging -->
<PackageTags>opentelemetry;tracing;metrics;observability;nats</PackageTags>
<Description>OpenTelemetry instrumentation for NATS .NET. Adds the ActivitySource and Meter for the client to a TracerProviderBuilder or MeterProviderBuilder.</Description>
</PropertyGroup>

<ItemGroup>
<PackageReference Include="OpenTelemetry.Api" Version="1.15.3" />
</ItemGroup>

<ItemGroup>
<ProjectReference Include="..\NATS.Client.Core\NATS.Client.Core.csproj" />
</ItemGroup>

</Project>
38 changes: 38 additions & 0 deletions src/NATS.Client.OpenTelemetry/NatsInstrumentationExtensions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
using System;
using NATS.Client.Core;
using OpenTelemetry.Metrics;
using OpenTelemetry.Trace;

namespace NATS.Client.OpenTelemetry;

public static class NatsInstrumentationExtensions
{
/// <summary>
/// Adds the NATS .NET client <see cref="System.Diagnostics.ActivitySource"/> to the tracer provider,
/// enabling distributed tracing for publish, subscribe, and request/reply operations.
/// </summary>
/// <param name="builder">The <see cref="TracerProviderBuilder"/> to add the source to.</param>
/// <returns>The supplied <paramref name="builder"/> for chaining.</returns>
public static TracerProviderBuilder AddNatsClientInstrumentation(this TracerProviderBuilder builder) =>
builder.AddSource(NatsTelemetry.SourceName);

/// <summary>
/// Adds the NATS .NET client <see cref="System.Diagnostics.ActivitySource"/> to the tracer provider and
/// configures the shared <see cref="NatsInstrumentationOptions"/> (filter and enrich callbacks).
/// </summary>
/// <param name="builder">The <see cref="TracerProviderBuilder"/> to add the source to.</param>
/// <param name="configure">Action that mutates the process-wide <see cref="NatsInstrumentationOptions.Default"/>.</param>
/// <returns>The supplied <paramref name="builder"/> for chaining.</returns>
public static TracerProviderBuilder AddNatsClientInstrumentation(this TracerProviderBuilder builder, Action<NatsInstrumentationOptions> configure)
{
configure?.Invoke(NatsInstrumentationOptions.Default);
return builder.AddSource(NatsTelemetry.SourceName);
}

/// <summary>
/// Adds the NATS .NET client <see cref="System.Diagnostics.Metrics.Meter"/> to the meter provider,
/// enabling messaging metrics (published/consumed counters, operation duration, and more).
/// </summary>
public static MeterProviderBuilder AddNatsClientInstrumentation(this MeterProviderBuilder builder) =>
builder.AddMeter(NatsTelemetry.SourceName);
Comment thread
mtmk marked this conversation as resolved.
}
108 changes: 108 additions & 0 deletions src/NATS.Client.OpenTelemetry/NatsInstrumentationOptionsExtensions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
using System;
using NATS.Client.Core;

namespace NATS.Client.OpenTelemetry;

/// <summary>
/// Extension methods for <see cref="NatsInstrumentationOptions"/>.
/// </summary>
public static class NatsInstrumentationOptionsExtensions
{
/// <summary>
/// Restricts tracing to operations whose subject matches the given NATS subject patterns.
/// </summary>
/// <param name="options">The options to configure.</param>
/// <param name="include">
/// Subject patterns to trace. An operation is traced only if its subject matches at least one
/// pattern. When <c>null</c> or empty, every subject is eligible (still subject to <paramref name="exclude"/>).
/// </param>
/// <param name="exclude">
/// 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 <c>_INBOX.&gt;</c>.
/// </param>
/// <returns>The same <paramref name="options"/> instance for chaining.</returns>
/// <remarks>
/// Patterns use NATS subject wildcards: <c>*</c> matches a single token and <c>&gt;</c> matches one or
/// more trailing tokens. The resulting predicate is combined (logical AND) with any existing
/// <see cref="NatsInstrumentationOptions.Filter"/>, so a previously configured filter still applies.
/// </remarks>
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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@

<ItemGroup>
<PackageReference Include="FluentAssertions" Version="[7,8)" />
<PackageReference Include="OpenTelemetry" Version="1.15.3" />
<PackageReference Include="ProcessX" Version="1.5.6" />
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="17.12.0" />
<PackageReference Include="xunit.v3" Version="1.0.1" />
Expand All @@ -34,6 +35,7 @@

<ItemGroup>
<ProjectReference Include="..\..\src\NATS.Client.JetStream\NATS.Client.JetStream.csproj" />
<ProjectReference Include="..\..\src\NATS.Client.OpenTelemetry\NATS.Client.OpenTelemetry.csproj" />
<ProjectReference Include="..\..\src\NATS.Client.Serializers.Json\NATS.Client.Serializers.Json.csproj" />
<ProjectReference Include="..\NATS.Client.TestUtilities\NATS.Client.TestUtilities.csproj" />
<ProjectReference Include="..\..\src\NATS.Client.Core\NATS.Client.Core.csproj" />
Expand Down
Original file line number Diff line number Diff line change
@@ -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);
}
32 changes: 32 additions & 0 deletions tools/site_src/documentation/advanced/opentelemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -108,3 +121,22 @@ All instruments carry these tags:
| `network.transport` | `tcp` | Transport protocol |

`messaging.client.operation.duration` adds `error.type` (full exception type name) when the operation fails.

### 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();
```
Loading