From 47cf516aa2cb7c3e03f3fd2fec39aa4b8026f44a Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Mon, 1 Jun 2026 09:28:15 +0100 Subject: [PATCH 1/2] otel: fix server.port type and trace tag source server.port and network.peer.port are integers per OTel semconv but the trace activity tags emitted them via ToString(). server.address/server.port also read from ServerInfo.Host/Port, which is the server bind address (often 0.0.0.0), not the endpoint the client dialled. Source the trace host and port from the connect URI to match the metric tags, and cache the host, boxed port and client id once per ServerInfo change so the per-message path no longer re-boxes the port or re-allocates the client id string. --- src/NATS.Client.Core/Internal/Telemetry.cs | 18 +++++----- src/NATS.Client.Core/NatsConnection.cs | 33 +++++++++++++++---- .../OpenTelemetryTest.cs | 13 ++++++++ 3 files changed, 49 insertions(+), 15 deletions(-) diff --git a/src/NATS.Client.Core/Internal/Telemetry.cs b/src/NATS.Client.Core/Internal/Telemetry.cs index 8f3ac6386..4a9ede654 100644 --- a/src/NATS.Client.Core/Internal/Telemetry.cs +++ b/src/NATS.Client.Core/Internal/Telemetry.cs @@ -139,18 +139,19 @@ public static TagList BuildMetricTags(INatsConnection? connection, string operat if (replyTo is not null) len++; - var serverPort = conn.ServerInfo.Port.ToString(); + var serverPort = conn.BoxedServerPort; // boxed once per ServerInfo change, reused for server.port and network.peer.port + var serverHost = conn.ServerHost; // grabbed once per ServerInfo change 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[3] = new KeyValuePair(Constants.ClientId, conn.ServerInfo.ClientId.ToString()); - tags[4] = new KeyValuePair(Constants.ServerAddress, conn.ServerInfo.Host); + tags[3] = new KeyValuePair(Constants.ClientId, conn.ClientId); + tags[4] = new KeyValuePair(Constants.ServerAddress, serverHost); tags[5] = new KeyValuePair(Constants.ServerPort, serverPort); tags[6] = new KeyValuePair(Constants.NetworkProtoName, "nats"); tags[7] = new KeyValuePair(Constants.NetworkTransport, "tcp"); - tags[8] = new KeyValuePair(Constants.NetworkPeerAddress, conn.ServerInfo.Host); + tags[8] = new KeyValuePair(Constants.NetworkPeerAddress, serverHost); tags[9] = new KeyValuePair(Constants.NetworkPeerPort, serverPort); tags[10] = new KeyValuePair(Constants.NetworkLocalAddress, conn.ServerInfo.ClientIp); @@ -238,7 +239,8 @@ public static void AddTraceContextHeaders(Activity? activity, ref NatsHeaders? h KeyValuePair[] tags; if (connection is NatsConnection { ServerInfo: not null } conn) { - var serverPort = conn.ServerInfo.Port.ToString(); + var serverPort = conn.BoxedServerPort; // boxed once per ServerInfo change, reused for server.port and network.peer.port + var serverHost = conn.ServerHost; // grabbed once per ServerInfo change var len = 17; if (replyTo is not null) @@ -256,12 +258,12 @@ public static void AddTraceContextHeaders(Activity? activity, ref NatsHeaders? h tags[6] = new KeyValuePair(Constants.DestPubName, subject); tags[7] = new KeyValuePair(Constants.MsgBodySize, bodySize.ToString()); tags[8] = new KeyValuePair(Constants.MsgTotalSize, size.ToString()); - tags[9] = new KeyValuePair(Constants.ClientId, conn.ServerInfo.ClientId.ToString()); - tags[10] = new KeyValuePair(Constants.ServerAddress, conn.ServerInfo.Host); + tags[9] = new KeyValuePair(Constants.ClientId, conn.ClientId); + tags[10] = new KeyValuePair(Constants.ServerAddress, serverHost); tags[11] = new KeyValuePair(Constants.ServerPort, serverPort); tags[12] = new KeyValuePair(Constants.NetworkProtoName, "nats"); tags[13] = new KeyValuePair(Constants.NetworkTransport, "tcp"); - tags[14] = new KeyValuePair(Constants.NetworkPeerAddress, conn.ServerInfo.Host); + tags[14] = new KeyValuePair(Constants.NetworkPeerAddress, serverHost); tags[15] = new KeyValuePair(Constants.NetworkPeerPort, serverPort); tags[16] = new KeyValuePair(Constants.NetworkLocalAddress, conn.ServerInfo.ClientIp); diff --git a/src/NATS.Client.Core/NatsConnection.cs b/src/NATS.Client.Core/NatsConnection.cs index f80b556c9..70f4219ea 100644 --- a/src/NATS.Client.Core/NatsConnection.cs +++ b/src/NATS.Client.Core/NatsConnection.cs @@ -58,6 +58,9 @@ public partial class NatsConnection : INatsConnection private ServerInfo? _writableServerInfo; private KeyValuePair[]? _metricTagsPrefix; + private object? _boxedServerPort; + private string? _serverHost; + private string? _clientId; private int _pongCount; private int _connectionState; private int _isDisposed; @@ -164,27 +167,37 @@ internal ServerInfo? WritableServerInfo } KeyValuePair[]? prefix = null; + object? port = null; + string? host = null; + string? clientId = 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; + host = connectUri?.Host ?? value.Host; + port = connectUri?.Port ?? value.Port; + clientId = value.ClientId.ToString(); 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.ServerPort, 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. + // Cache the trace-tag host, boxed port and client id once per ServerInfo change so the + // per-message trace path can reuse them instead of re-boxing the port and re-allocating + // the client id string on every message. Host and port come from the connect URI, so the + // trace and metric tags now share the same source. + Volatile.Write(ref _boxedServerPort, port); + Volatile.Write(ref _serverHost, host); + Volatile.Write(ref _clientId, clientId); + + // Publish the cached fields before ServerInfo so any reader observing the new ServerInfo + // is guaranteed to also observe the matching tags. Volatile.Write(ref _metricTagsPrefix, prefix); Interlocked.Exchange(ref _writableServerInfo, value); } @@ -192,6 +205,12 @@ internal ServerInfo? WritableServerInfo internal KeyValuePair[]? MetricTagsPrefix => Volatile.Read(ref _metricTagsPrefix); + internal object? BoxedServerPort => Volatile.Read(ref _boxedServerPort); + + internal string? ServerHost => Volatile.Read(ref _serverHost); + + internal string? ClientId => Volatile.Read(ref _clientId); + internal bool IsDisposed { get => Interlocked.CompareExchange(ref _isDisposed, 0, 0) == 1; diff --git a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs index 21e190bf2..8972f31ec 100644 --- a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs +++ b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs @@ -457,6 +457,12 @@ private void AssertActivityData(string subject, IReadOnlyList activity // Verify network.protocol.version is no longer present Assert.Null(sendActivity.GetTagItem("network.protocol.version")); Assert.Null(receiveActivity.GetTagItem("network.protocol.version")); + + // server.port and network.peer.port are integers per OTel semconv, not strings + AssertIntTag(sendActivity, "server.port"); + AssertIntTag(sendActivity, "network.peer.port"); + AssertIntTag(receiveActivity, "server.port"); + AssertIntTag(receiveActivity, "network.peer.port"); } private void AssertStringTagNotNullOrEmpty(Activity activity, string name) @@ -466,6 +472,13 @@ private void AssertStringTagNotNullOrEmpty(Activity activity, string name) Assert.False(string.IsNullOrEmpty(tag)); } + private void AssertIntTag(Activity activity, string name) + { + var tag = activity.GetTagItem(name); + Assert.NotNull(tag); + Assert.IsType(tag); + } + private sealed class ActivityTracker : IDisposable { private readonly List _started = new(); From d7eb0d97eae93076377e4d2f05e55605bcc0ca8e Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Mon, 1 Jun 2026 09:34:44 +0100 Subject: [PATCH 2/2] otel: assert host and client id trace tags in test --- .../OpenTelemetryTest.cs | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs index 8972f31ec..1046fdc9c 100644 --- a/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs +++ b/tests/NATS.Net.OpenTelemetry.Tests/OpenTelemetryTest.cs @@ -27,7 +27,9 @@ public async Task Publish_subscribe_activities() await sub.Msgs.ReadAsync(cts.Token); - AssertActivityData("foo", tracker.Started); + var expectedHost = new Uri(server.Url).Host; + var expectedClientId = nats.ServerInfo!.ClientId.ToString(); + AssertActivityData("foo", tracker.Started, expectedHost, expectedClientId); tracker.AssertAllStopped(); } @@ -424,7 +426,7 @@ public async Task Direct_request_reply_receive_activity_is_disposed() await reg; } - private void AssertActivityData(string subject, IReadOnlyList activityList) + private void AssertActivityData(string subject, IReadOnlyList activityList, string expectedHost, string expectedClientId) { var activities = activityList.ToArray(); Assert.NotEmpty(activities); @@ -463,6 +465,17 @@ private void AssertActivityData(string subject, IReadOnlyList activity AssertIntTag(sendActivity, "network.peer.port"); AssertIntTag(receiveActivity, "server.port"); AssertIntTag(receiveActivity, "network.peer.port"); + + // server.address/network.peer.address come from the connect URI, not ServerInfo.Host + // (the server bind address, often 0.0.0.0) + Assert.Equal(expectedHost, sendActivity.GetTagItem("server.address")); + Assert.Equal(expectedHost, sendActivity.GetTagItem("network.peer.address")); + Assert.Equal(expectedHost, receiveActivity.GetTagItem("server.address")); + Assert.Equal(expectedHost, receiveActivity.GetTagItem("network.peer.address")); + + // messaging.client_id is the server-assigned client id, cached per ServerInfo change + Assert.Equal(expectedClientId, sendActivity.GetTagItem("messaging.client_id")); + Assert.Equal(expectedClientId, receiveActivity.GetTagItem("messaging.client_id")); } private void AssertStringTagNotNullOrEmpty(Activity activity, string name)