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..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); @@ -457,6 +459,23 @@ 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"); + + // 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) @@ -466,6 +485,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();