From 9327415a2b21f3f09f95b039a5488da476316b88 Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Wed, 15 Jul 2026 16:30:51 +0200 Subject: [PATCH 01/11] =?UTF-8?q?=E2=9C=A8=20Add=20TraceConnector=20enum?= =?UTF-8?q?=20and=20endpoint-level=20trace=20connector=20defaults=20to=20I?= =?UTF-8?q?nstrumentationOptions?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../InstrumentationOptionsTests.cs | 19 +++++++++++++++++++ .../OpenTelemetry/InstrumentationOptions.cs | 16 ++++++++++++++++ .../OpenTelemetry/TraceConnector.cs | 19 +++++++++++++++++++ 3 files changed, 54 insertions(+) create mode 100644 src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs create mode 100644 src/NServiceBus.Core/OpenTelemetry/TraceConnector.cs diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs new file mode 100644 index 0000000000..13076f0e3f --- /dev/null +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs @@ -0,0 +1,19 @@ +namespace NServiceBus.Core.Tests.OpenTelemetry; + +using NUnit.Framework; + +[TestFixture] +public class InstrumentationOptionsTests +{ + [Test] + public void Should_default_trace_connectors_to_current_behavior() + { + var options = new InstrumentationOptions(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(options.SentMessageTraceConnector, Is.EqualTo(TraceConnector.ChildSpan), "sends continue the trace by default"); + Assert.That(options.PublishedMessageTraceConnector, Is.EqualTo(TraceConnector.SpanLink), "publishes start a new linked trace by default"); + } + } +} diff --git a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs index f920a28b93..28cc9abd65 100644 --- a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs +++ b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs @@ -14,4 +14,20 @@ public class InstrumentationOptions /// Disabled by default for backward compatibility. /// public bool UseMessageDestinationInSpanNames { get; set; } + + /// + /// Controls how the receive-side processing span relates to the send span for messages sent by this endpoint. + /// Defaults to : receivers continue the trace. + /// Can be overridden per message via StartNewTraceOnReceive + /// or ContinueExistingTraceOnReceive. + /// + public TraceConnector SentMessageTraceConnector { get; set; } = TraceConnector.ChildSpan; + + /// + /// Controls how the receive-side processing span relates to the publish span for events published by this endpoint. + /// Defaults to : receivers start a new trace linked back to the publish span. + /// Can be overridden per message via StartNewTraceOnReceive + /// or ContinueExistingTraceOnReceive. + /// + public TraceConnector PublishedMessageTraceConnector { get; set; } = TraceConnector.SpanLink; } diff --git a/src/NServiceBus.Core/OpenTelemetry/TraceConnector.cs b/src/NServiceBus.Core/OpenTelemetry/TraceConnector.cs new file mode 100644 index 0000000000..5854b92101 --- /dev/null +++ b/src/NServiceBus.Core/OpenTelemetry/TraceConnector.cs @@ -0,0 +1,19 @@ +#nullable enable + +namespace NServiceBus; + +/// +/// Controls how the receive-side processing span relates to the outgoing send or publish span. +/// +public enum TraceConnector +{ + /// + /// The receiving endpoint continues the trace: the processing span becomes a child of the outgoing span. + /// + ChildSpan, + + /// + /// The receiving endpoint starts a new trace: the processing span becomes the root of a new trace with a link back to the outgoing span. + /// + SpanLink +} From f4a4131acd64655433ccbd7b793e46976b41726a Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Wed, 15 Jul 2026 16:37:50 +0200 Subject: [PATCH 02/11] =?UTF-8?q?=E2=9C=A8=20Honor=20endpoint-level=20trac?= =?UTF-8?q?e=20connector=20defaults=20in=20send/publish=20behaviors=20with?= =?UTF-8?q?=20per-message=20overrides?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../OpenTelemetryExtensionsTests.cs | 63 +++++++++++++++++++ .../OpenTelemetryPublishBehaviorTests.cs | 55 ++++++++++++++++ .../OpenTelemetrySendBehaviorTests.cs | 55 ++++++++++++++++ .../OpenTelemetry/InstrumentationOptions.cs | 8 +-- .../OpenTelemetry/OpenTelemetryExtensions.cs | 36 +++++++++-- .../OpenTelemetry/OpenTelemetryFeature.cs | 6 +- .../OpenTelemetryPublishBehavior.cs | 17 ++--- .../OpenTelemetrySendBehavior.cs | 17 ++--- 8 files changed, 225 insertions(+), 32 deletions(-) create mode 100644 src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs create mode 100644 src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs create mode 100644 src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs new file mode 100644 index 0000000000..71896cbffe --- /dev/null +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs @@ -0,0 +1,63 @@ +namespace NServiceBus.Core.Tests.OpenTelemetry; + +using NUnit.Framework; + +[TestFixture] +public class OpenTelemetryExtensionsTests +{ + [Test] + public void StartNewTraceOnReceive_should_set_span_link_override_on_send_options() + { + var options = new SendOptions(); + + options.StartNewTraceOnReceive(); + + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceConnector.SpanLink)); + } + + [Test] + public void ContinueExistingTraceOnReceive_should_set_child_span_override_on_send_options() + { + var options = new SendOptions(); + + options.ContinueExistingTraceOnReceive(); + + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceConnector.ChildSpan)); + } + + [Test] + public void StartNewTraceOnReceive_should_set_span_link_override_on_publish_options() + { + var options = new PublishOptions(); + + options.StartNewTraceOnReceive(); + + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceConnector.SpanLink)); + } + + [Test] + public void ContinueExistingTraceOnReceive_should_set_child_span_override_on_publish_options() + { + var options = new PublishOptions(); + + options.ContinueExistingTraceOnReceive(); + + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceConnector.ChildSpan)); + } + + [Test] + public void Last_override_call_wins() + { + var options = new PublishOptions(); + + options.ContinueExistingTraceOnReceive(); + options.StartNewTraceOnReceive(); + + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceConnector.SpanLink)); + } +} diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs new file mode 100644 index 0000000000..f882f8ba2a --- /dev/null +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs @@ -0,0 +1,55 @@ +namespace NServiceBus.Core.Tests.OpenTelemetry; + +using System.Threading.Tasks; +using NUnit.Framework; +using Testing; + +[TestFixture] +public class OpenTelemetryPublishBehaviorTests +{ + [Test] + public async Task Should_start_new_trace_on_receive_by_default() + { + var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions()); + var context = new TestableOutgoingPublishContext(); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); + } + + [Test] + public async Task Should_continue_trace_on_receive_when_endpoint_connector_is_child_span() + { + var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishedMessageTraceConnector = TraceConnector.ChildSpan }); + var context = new TestableOutgoingPublishContext(); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.FalseString)); + } + + [Test] + public async Task Should_prefer_child_span_option_over_endpoint_connector() + { + var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishedMessageTraceConnector = TraceConnector.SpanLink }); + var context = new TestableOutgoingPublishContext(); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.ChildSpan); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.FalseString)); + } + + [Test] + public async Task Should_prefer_span_link_option_over_endpoint_connector() + { + var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishedMessageTraceConnector = TraceConnector.ChildSpan }); + var context = new TestableOutgoingPublishContext(); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.SpanLink); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); + } +} diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs new file mode 100644 index 0000000000..a63cf30e21 --- /dev/null +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs @@ -0,0 +1,55 @@ +namespace NServiceBus.Core.Tests.OpenTelemetry; + +using System.Threading.Tasks; +using NUnit.Framework; +using Testing; + +[TestFixture] +public class OpenTelemetrySendBehaviorTests +{ + [Test] + public async Task Should_continue_trace_on_receive_by_default() + { + var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions()); + var context = new TestableOutgoingSendContext(); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.FalseString)); + } + + [Test] + public async Task Should_start_new_trace_on_receive_when_endpoint_connector_is_span_link() + { + var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SentMessageTraceConnector = TraceConnector.SpanLink }); + var context = new TestableOutgoingSendContext(); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); + } + + [Test] + public async Task Should_prefer_span_link_option_over_endpoint_connector() + { + var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SentMessageTraceConnector = TraceConnector.ChildSpan }); + var context = new TestableOutgoingSendContext(); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.SpanLink); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); + } + + [Test] + public async Task Should_prefer_child_span_option_over_endpoint_connector() + { + var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SentMessageTraceConnector = TraceConnector.SpanLink }); + var context = new TestableOutgoingSendContext(); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.ChildSpan); + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Headers[Headers.StartNewTrace], Is.EqualTo(bool.FalseString)); + } +} diff --git a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs index 28cc9abd65..3dd0eab70c 100644 --- a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs +++ b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs @@ -18,16 +18,16 @@ public class InstrumentationOptions /// /// Controls how the receive-side processing span relates to the send span for messages sent by this endpoint. /// Defaults to : receivers continue the trace. - /// Can be overridden per message via StartNewTraceOnReceive - /// or ContinueExistingTraceOnReceive. + /// Can be overridden per message via + /// or . /// public TraceConnector SentMessageTraceConnector { get; set; } = TraceConnector.ChildSpan; /// /// Controls how the receive-side processing span relates to the publish span for events published by this endpoint. /// Defaults to : receivers start a new trace linked back to the publish span. - /// Can be overridden per message via StartNewTraceOnReceive - /// or ContinueExistingTraceOnReceive. + /// Can be overridden per message via + /// or . /// public TraceConnector PublishedMessageTraceConnector { get; set; } = TraceConnector.SpanLink; } diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs index 9c3ad6e63e..ef6ac1ac25 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs @@ -21,20 +21,48 @@ public static InstrumentationOptions Tracing(this EndpointConfiguration config) } /// - /// Start a new OpenTelemetry trace conversation. + /// Start a new OpenTelemetry trace on receive of this message, linked back to the send span. + /// Overrides for this message. /// /// The option being extended. public static void StartNewTraceOnReceive(this SendOptions sendOptions) { - sendOptions.Context.Set(OpenTelemetrySendBehavior.StartNewTraceOnReceive, true); + ArgumentNullException.ThrowIfNull(sendOptions); + sendOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.SpanLink); } /// - /// Start a new OpenTelemetry trace conversation. + /// Continue the existing OpenTelemetry trace on receive of this message. + /// Overrides for this message. + /// + /// The option being extended. + public static void ContinueExistingTraceOnReceive(this SendOptions sendOptions) + { + ArgumentNullException.ThrowIfNull(sendOptions); + sendOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.ChildSpan); + } + + /// + /// Start a new OpenTelemetry trace on receive of this event, linked back to the publish span. + /// Overrides for this message. + /// + /// The option being extended. + public static void StartNewTraceOnReceive(this PublishOptions publishOptions) + { + ArgumentNullException.ThrowIfNull(publishOptions); + publishOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.SpanLink); + } + + /// + /// Continue the existing OpenTelemetry trace on receive of this event. + /// Overrides for this message. /// /// The option being extended. public static void ContinueExistingTraceOnReceive(this PublishOptions publishOptions) { - publishOptions.Context.Set(OpenTelemetryPublishBehavior.ContinueTraceOnReceive, true); + ArgumentNullException.ThrowIfNull(publishOptions); + publishOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.ChildSpan); } + + internal const string TraceConnectorOverrideKey = "NServiceBus.OpenTelemetry.TraceConnectorOverride"; } \ No newline at end of file diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs index c485129eee..b108b409c9 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs @@ -8,13 +8,15 @@ sealed class OpenTelemetryFeature : Feature { protected override void Setup(FeatureConfigurationContext context) { + var instrumentationOptions = context.Settings.GetOrDefault() ?? new InstrumentationOptions(); + context.Pipeline.Register( - new OpenTelemetryPublishBehavior(), + new OpenTelemetryPublishBehavior(instrumentationOptions), "Manages the depth of the trace for publishes" ); context.Pipeline.Register( - new OpenTelemetrySendBehavior(), + new OpenTelemetrySendBehavior(instrumentationOptions), "Manages the depth of the trace for sends" ); diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs index 55617620d7..630defa1c2 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs @@ -6,22 +6,17 @@ namespace NServiceBus; using System.Threading.Tasks; using Pipeline; -class OpenTelemetryPublishBehavior : IBehavior +class OpenTelemetryPublishBehavior(InstrumentationOptions instrumentationOptions) : IBehavior { public Task Invoke(IOutgoingPublishContext context, Func next) { - // publishes always start a new trace on receive - context.Headers[Headers.StartNewTrace] = bool.TrueString; + // the per-message override wins over the endpoint-level default + var connector = context.Extensions.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector requestedConnector) + ? requestedConnector + : instrumentationOptions.PublishedMessageTraceConnector; - // unless the user explicitly requests to continue the trace - bool continueTraceWasSet = context.Extensions.TryGet(ContinueTraceOnReceive, out var continueTraceRequested); - if (continueTraceWasSet && continueTraceRequested) - { - context.Headers[Headers.StartNewTrace] = bool.FalseString; - } + context.Headers[Headers.StartNewTrace] = connector == TraceConnector.SpanLink ? bool.TrueString : bool.FalseString; return next(context); } - - public const string ContinueTraceOnReceive = "ContinueTraceRequested"; } \ No newline at end of file diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs index b8ca6de993..9424a5eb3b 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs @@ -6,22 +6,17 @@ namespace NServiceBus; using System.Threading.Tasks; using Pipeline; -class OpenTelemetrySendBehavior : IBehavior +class OpenTelemetrySendBehavior(InstrumentationOptions instrumentationOptions) : IBehavior { public Task Invoke(IOutgoingSendContext context, Func next) { - // sends never start a new trace on receive - context.Headers[Headers.StartNewTrace] = bool.FalseString; + // the per-message override wins over the endpoint-level default + var connector = context.Extensions.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector requestedConnector) + ? requestedConnector + : instrumentationOptions.SentMessageTraceConnector; - // unless the user explicitly requests to start a new trace - bool breakTraceWasSet = context.Extensions.TryGet(StartNewTraceOnReceive, out var breakTraceWasRequested); - if (breakTraceWasSet && breakTraceWasRequested) - { - context.Headers[Headers.StartNewTrace] = bool.TrueString; - } + context.Headers[Headers.StartNewTrace] = connector == TraceConnector.SpanLink ? bool.TrueString : bool.FalseString; return next(context); } - - public const string StartNewTraceOnReceive = "BreakTraceRequested"; } \ No newline at end of file From 354c503c3ec8440c14ec5df061884ac62e21bf80 Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Wed, 15 Jul 2026 16:43:30 +0200 Subject: [PATCH 03/11] =?UTF-8?q?=E2=9C=A8=20Add=20acceptance=20tests=20fo?= =?UTF-8?q?r=20endpoint-level=20trace=20connector=20defaults=20and=20per-m?= =?UTF-8?q?essage=20overrides?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Traces/When_publishing_messages.cs | 120 ++++++++++++++++++ .../Traces/When_sending_messages.cs | 80 ++++++++++++ 2 files changed, 200 insertions(+) diff --git a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs index 4281ddc170..737db0d195 100644 --- a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs +++ b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs @@ -242,5 +242,125 @@ public Task Handle(ThisIsAnEvent @event, IMessageHandlerContext context) } } + [Test] + public async Task Should_create_child_on_receive_when_endpoint_defaults_to_child_span() + { + var context = await Scenario.Define() + .WithEndpoint(b => b + .When(ctx => ctx.SomeEventSubscribed, s => s.Publish(new ThisIsAnEvent()))) + .WithEndpoint(b => b.When((session, ctx) => + { + if (ctx.HasNativePubSubSupport) + { + ctx.SomeEventSubscribed = true; + } + + return Task.CompletedTask; + })) + .Run(); + + var publishMessageActivities = NServiceBusActivityListener.CompletedActivities.GetPublishEventActivities(); + var receiveMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities(); + using (Assert.EnterMultipleScope()) + { + Assert.That(publishMessageActivities, Has.Count.EqualTo(1), "1 message is published as part of this test"); + Assert.That(receiveMessageActivities, Has.Count.EqualTo(1), "1 message is received as part of this test"); + } + + var publishRequest = publishMessageActivities[0]; + var receiveRequest = receiveMessageActivities[0]; + + using (Assert.EnterMultipleScope()) + { + Assert.That(receiveRequest.RootId, Is.EqualTo(publishRequest.RootId), "publish and receive operations are part the same root activity"); + Assert.That(receiveRequest.ParentId, Is.Not.Null, "incoming message does have a parent"); + } + + Assert.That(receiveRequest.Links, Is.Empty, "receive does not have links"); + } + + [Test] + public async Task Should_create_new_linked_trace_on_receive_when_option_overrides_endpoint_connector() + { + var context = await Scenario.Define() + .WithEndpoint(b => b + .When(ctx => ctx.SomeEventSubscribed, s => + { + var publishOptions = new PublishOptions(); + publishOptions.StartNewTraceOnReceive(); + return s.Publish(new ThisIsAnEvent(), publishOptions); + })) + .WithEndpoint(b => b.When((session, ctx) => + { + if (ctx.HasNativePubSubSupport) + { + ctx.SomeEventSubscribed = true; + } + + return Task.CompletedTask; + })) + .Run(); + + var publishMessageActivities = NServiceBusActivityListener.CompletedActivities.GetPublishEventActivities(); + var receiveMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities(); + using (Assert.EnterMultipleScope()) + { + Assert.That(publishMessageActivities, Has.Count.EqualTo(1), "1 message is published as part of this test"); + Assert.That(receiveMessageActivities, Has.Count.EqualTo(1), "1 message is received as part of this test"); + } + + var publishRequest = publishMessageActivities[0]; + var receiveRequest = receiveMessageActivities[0]; + + using (Assert.EnterMultipleScope()) + { + Assert.That(receiveRequest.RootId, Is.Not.EqualTo(publishRequest.RootId), "publish and receive operations are part of different root activities"); + Assert.That(receiveRequest.ParentId, Is.Null, "incoming message does not have a parent, it's a root"); + } + + ActivityLink link = receiveRequest.Links.FirstOrDefault(); + Assert.That(link, Is.Not.EqualTo(default(ActivityLink)), "Receive has a link"); + Assert.That(link.Context.TraceId, Is.EqualTo(publishRequest.TraceId), "receive is linked to publish operation"); + } + + public class PublisherWithChildSpanConnector : EndpointConfigurationBuilder + { + public PublisherWithChildSpanConnector() => + EndpointSetup(b => + { + b.Tracing().PublishedMessageTraceConnector = TraceConnector.ChildSpan; + b.OnEndpointSubscribed((s, context) => + { + if (s.SubscriberEndpoint.Contains(Conventions.EndpointNamingConvention(typeof(SubscriberForPublisherWithChildSpanConnector)))) + { + if (s.MessageType == typeof(ThisIsAnEvent).AssemblyQualifiedName) + { + context.SomeEventSubscribed = true; + } + } + }); + }); + } + + public class SubscriberForPublisherWithChildSpanConnector : EndpointConfigurationBuilder + { + public SubscriberForPublisherWithChildSpanConnector() => + EndpointSetup(c => { }, + metadata => + { + metadata.RegisterPublisherFor(); + }); + + [Handler] + public class ThisHandlesSomethingHandler(Context testContext) : IHandleMessages + { + public Task Handle(ThisIsAnEvent @event, IMessageHandlerContext context) + { + testContext.MarkAsCompleted(); + return Task.CompletedTask; + } + } + } + public class ThisIsAnEvent : IEvent; } \ No newline at end of file diff --git a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs index e14df1d05e..e6c9a83c16 100644 --- a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs +++ b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs @@ -163,5 +163,85 @@ public Task Handle(OutgoingMessage message, IMessageHandlerContext context) } } + [Test] + public async Task Should_create_new_linked_trace_on_receive_when_endpoint_defaults_to_span_link() + { + await Scenario.Define() + .WithEndpoint(b => b + .When(s => s.SendLocal(new OutgoingMessage()))) + .Run(); + + var sendMessageActivities = NServiceBusActivityListener.CompletedActivities.GetSendMessageActivities(); + var receiveMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities(); + using (Assert.EnterMultipleScope()) + { + Assert.That(sendMessageActivities, Has.Count.EqualTo(1), "1 message is sent as part of this test"); + Assert.That(receiveMessageActivities, Has.Count.EqualTo(1), "1 message is received as part of this test"); + } + + var sendRequest = sendMessageActivities[0]; + var receiveRequest = receiveMessageActivities[0]; + + using (Assert.EnterMultipleScope()) + { + Assert.That(receiveRequest.RootId, Is.Not.EqualTo(sendRequest.RootId), "send and receive operations are part of different root activities"); + Assert.That(receiveRequest.ParentId, Is.Null, "incoming message does not have a parent, it's a root"); + } + + ActivityLink link = receiveRequest.Links.FirstOrDefault(); + Assert.That(link, Is.Not.EqualTo(default(ActivityLink)), "Receive has a link"); + Assert.That(link.Context.TraceId, Is.EqualTo(sendRequest.TraceId), "receive is linked to send operation"); + } + + [Test] + public async Task Should_create_child_on_receive_when_option_overrides_endpoint_connector() + { + await Scenario.Define() + .WithEndpoint(b => b + .When(s => + { + var sendOptions = new SendOptions(); + sendOptions.RouteToThisEndpoint(); + sendOptions.ContinueExistingTraceOnReceive(); + return s.Send(new OutgoingMessage(), sendOptions); + })) + .Run(); + + var sendMessageActivities = NServiceBusActivityListener.CompletedActivities.GetSendMessageActivities(); + var receiveMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities(); + using (Assert.EnterMultipleScope()) + { + Assert.That(sendMessageActivities, Has.Count.EqualTo(1), "1 message is sent as part of this test"); + Assert.That(receiveMessageActivities, Has.Count.EqualTo(1), "1 message is received as part of this test"); + } + + var sendRequest = sendMessageActivities[0]; + var receiveRequest = receiveMessageActivities[0]; + + using (Assert.EnterMultipleScope()) + { + Assert.That(receiveRequest.RootId, Is.EqualTo(sendRequest.RootId), "send and receive operations are part of the same root activity"); + Assert.That(receiveRequest.ParentId, Is.Not.Null, "incoming message does have a parent"); + } + + Assert.That(receiveRequest.Links, Is.Empty, "receive does not have links"); + } + + public class TestEndpointWithSpanLinkConnector : EndpointConfigurationBuilder + { + public TestEndpointWithSpanLinkConnector() => + EndpointSetup(b => b.Tracing().SentMessageTraceConnector = TraceConnector.SpanLink); + + [Handler] + public class MessageHandler(Context testContext) : IHandleMessages + { + public Task Handle(OutgoingMessage message, IMessageHandlerContext context) + { + testContext.MarkAsCompleted(); + return Task.CompletedTask; + } + } + } + public class OutgoingMessage : IMessage; } \ No newline at end of file From 5a315168be2ac15a694cfea813f4875e978a2930 Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Wed, 15 Jul 2026 16:48:44 +0200 Subject: [PATCH 04/11] =?UTF-8?q?=E2=9C=A8=20Approve=20public=20API=20addi?= =?UTF-8?q?tions=20for=20trace=20connector=20configuration?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../APIApprovals.ApproveNServiceBus.approved.txt | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt b/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt index a5bdc1d389..6849210e76 100644 --- a/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt +++ b/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt @@ -635,6 +635,8 @@ namespace NServiceBus public class InstrumentationOptions { public InstrumentationOptions() { } + public NServiceBus.TraceConnector PublishedMessageTraceConnector { get; set; } + public NServiceBus.TraceConnector SentMessageTraceConnector { get; set; } public bool UseMessageDestinationInSpanNames { get; set; } } public sealed class KeyedServiceKey @@ -776,6 +778,8 @@ namespace NServiceBus public static class OpenTelemetryExtensions { public static void ContinueExistingTraceOnReceive(this NServiceBus.PublishOptions publishOptions) { } + public static void ContinueExistingTraceOnReceive(this NServiceBus.SendOptions sendOptions) { } + public static void StartNewTraceOnReceive(this NServiceBus.PublishOptions publishOptions) { } public static void StartNewTraceOnReceive(this NServiceBus.SendOptions sendOptions) { } public static NServiceBus.InstrumentationOptions Tracing(this NServiceBus.EndpointConfiguration config) { } } @@ -1173,6 +1177,11 @@ namespace NServiceBus public ToSagaExpression(NServiceBus.IConfigureHowToFindSagaWithMessage sagaMessageFindingConfiguration, System.Linq.Expressions.Expression> messageProperty) { } public void ToSaga(System.Linq.Expressions.Expression> sagaEntityProperty) { } } + public enum TraceConnector + { + ChildSpan = 0, + SpanLink = 1, + } public static class TransportConfig { extension(NServiceBus.EndpointConfiguration endpointConfiguration) From b646fb3a55c728d5e7afd93d22d03f9df0d83d0a Mon Sep 17 00:00:00 2001 From: Tomasz Masternak Date: Thu, 16 Jul 2026 13:11:17 +0200 Subject: [PATCH 05/11] Replace `TraceConnector` with `TraceMode` for clearer terminology and updated tracing behavior. --- .../Traces/When_publishing_messages.cs | 2 +- .../Traces/When_sending_messages.cs | 2 +- .../InstrumentationOptionsTests.cs | 4 ++-- .../OpenTelemetryExtensionsTests.cs | 20 +++++++++---------- .../OpenTelemetryPublishBehaviorTests.cs | 10 +++++----- .../OpenTelemetrySendBehaviorTests.cs | 10 +++++----- .../OpenTelemetry/InstrumentationOptions.cs | 8 ++++---- .../OpenTelemetry/OpenTelemetryExtensions.cs | 16 +++++++-------- .../OpenTelemetryPublishBehavior.cs | 6 +++--- .../OpenTelemetrySendBehavior.cs | 6 +++--- .../{TraceConnector.cs => TraceMode.cs} | 6 +++--- 11 files changed, 45 insertions(+), 45 deletions(-) rename src/NServiceBus.Core/OpenTelemetry/{TraceConnector.cs => TraceMode.cs} (89%) diff --git a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs index 737db0d195..7aed06a30a 100644 --- a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs +++ b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_publishing_messages.cs @@ -328,7 +328,7 @@ public class PublisherWithChildSpanConnector : EndpointConfigurationBuilder public PublisherWithChildSpanConnector() => EndpointSetup(b => { - b.Tracing().PublishedMessageTraceConnector = TraceConnector.ChildSpan; + b.Tracing().PublishTraceMode = TraceMode.ContinueExisting; b.OnEndpointSubscribed((s, context) => { if (s.SubscriberEndpoint.Contains(Conventions.EndpointNamingConvention(typeof(SubscriberForPublisherWithChildSpanConnector)))) diff --git a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs index e6c9a83c16..c08e660e88 100644 --- a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs +++ b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_sending_messages.cs @@ -230,7 +230,7 @@ await Scenario.Define() public class TestEndpointWithSpanLinkConnector : EndpointConfigurationBuilder { public TestEndpointWithSpanLinkConnector() => - EndpointSetup(b => b.Tracing().SentMessageTraceConnector = TraceConnector.SpanLink); + EndpointSetup(b => b.Tracing().SendTraceMode = TraceMode.StartNew); [Handler] public class MessageHandler(Context testContext) : IHandleMessages diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs index 13076f0e3f..ded2e6a41c 100644 --- a/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs @@ -12,8 +12,8 @@ public void Should_default_trace_connectors_to_current_behavior() using (Assert.EnterMultipleScope()) { - Assert.That(options.SentMessageTraceConnector, Is.EqualTo(TraceConnector.ChildSpan), "sends continue the trace by default"); - Assert.That(options.PublishedMessageTraceConnector, Is.EqualTo(TraceConnector.SpanLink), "publishes start a new linked trace by default"); + Assert.That(options.SendTraceMode, Is.EqualTo(TraceMode.ContinueExisting), "sends continue the trace by default"); + Assert.That(options.PublishTraceMode, Is.EqualTo(TraceMode.StartNew), "publishes start a new linked trace by default"); } } } diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs index 71896cbffe..1475280284 100644 --- a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryExtensionsTests.cs @@ -12,8 +12,8 @@ public void StartNewTraceOnReceive_should_set_span_link_override_on_send_options options.StartNewTraceOnReceive(); - Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); - Assert.That(connector, Is.EqualTo(TraceConnector.SpanLink)); + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceMode.StartNew)); } [Test] @@ -23,8 +23,8 @@ public void ContinueExistingTraceOnReceive_should_set_child_span_override_on_sen options.ContinueExistingTraceOnReceive(); - Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); - Assert.That(connector, Is.EqualTo(TraceConnector.ChildSpan)); + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceMode.ContinueExisting)); } [Test] @@ -34,8 +34,8 @@ public void StartNewTraceOnReceive_should_set_span_link_override_on_publish_opti options.StartNewTraceOnReceive(); - Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); - Assert.That(connector, Is.EqualTo(TraceConnector.SpanLink)); + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceMode.StartNew)); } [Test] @@ -45,8 +45,8 @@ public void ContinueExistingTraceOnReceive_should_set_child_span_override_on_pub options.ContinueExistingTraceOnReceive(); - Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); - Assert.That(connector, Is.EqualTo(TraceConnector.ChildSpan)); + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceMode.ContinueExisting)); } [Test] @@ -57,7 +57,7 @@ public void Last_override_call_wins() options.ContinueExistingTraceOnReceive(); options.StartNewTraceOnReceive(); - Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector connector), Is.True); - Assert.That(connector, Is.EqualTo(TraceConnector.SpanLink)); + Assert.That(options.Context.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode connector), Is.True); + Assert.That(connector, Is.EqualTo(TraceMode.StartNew)); } } diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs index f882f8ba2a..3cc9e44ee0 100644 --- a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryPublishBehaviorTests.cs @@ -21,7 +21,7 @@ public async Task Should_start_new_trace_on_receive_by_default() [Test] public async Task Should_continue_trace_on_receive_when_endpoint_connector_is_child_span() { - var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishedMessageTraceConnector = TraceConnector.ChildSpan }); + var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishTraceMode = TraceMode.ContinueExisting }); var context = new TestableOutgoingPublishContext(); await behavior.Invoke(context, _ => Task.CompletedTask); @@ -32,9 +32,9 @@ public async Task Should_continue_trace_on_receive_when_endpoint_connector_is_ch [Test] public async Task Should_prefer_child_span_option_over_endpoint_connector() { - var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishedMessageTraceConnector = TraceConnector.SpanLink }); + var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishTraceMode = TraceMode.StartNew }); var context = new TestableOutgoingPublishContext(); - context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.ChildSpan); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceMode.ContinueExisting); await behavior.Invoke(context, _ => Task.CompletedTask); @@ -44,9 +44,9 @@ public async Task Should_prefer_child_span_option_over_endpoint_connector() [Test] public async Task Should_prefer_span_link_option_over_endpoint_connector() { - var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishedMessageTraceConnector = TraceConnector.ChildSpan }); + var behavior = new OpenTelemetryPublishBehavior(new InstrumentationOptions { PublishTraceMode = TraceMode.ContinueExisting }); var context = new TestableOutgoingPublishContext(); - context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.SpanLink); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceMode.StartNew); await behavior.Invoke(context, _ => Task.CompletedTask); diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs index a63cf30e21..8e26d21354 100644 --- a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetrySendBehaviorTests.cs @@ -21,7 +21,7 @@ public async Task Should_continue_trace_on_receive_by_default() [Test] public async Task Should_start_new_trace_on_receive_when_endpoint_connector_is_span_link() { - var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SentMessageTraceConnector = TraceConnector.SpanLink }); + var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SendTraceMode = TraceMode.StartNew }); var context = new TestableOutgoingSendContext(); await behavior.Invoke(context, _ => Task.CompletedTask); @@ -32,9 +32,9 @@ public async Task Should_start_new_trace_on_receive_when_endpoint_connector_is_s [Test] public async Task Should_prefer_span_link_option_over_endpoint_connector() { - var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SentMessageTraceConnector = TraceConnector.ChildSpan }); + var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SendTraceMode = TraceMode.ContinueExisting }); var context = new TestableOutgoingSendContext(); - context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.SpanLink); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceMode.StartNew); await behavior.Invoke(context, _ => Task.CompletedTask); @@ -44,9 +44,9 @@ public async Task Should_prefer_span_link_option_over_endpoint_connector() [Test] public async Task Should_prefer_child_span_option_over_endpoint_connector() { - var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SentMessageTraceConnector = TraceConnector.SpanLink }); + var behavior = new OpenTelemetrySendBehavior(new InstrumentationOptions { SendTraceMode = TraceMode.StartNew }); var context = new TestableOutgoingSendContext(); - context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceConnector.ChildSpan); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceMode.ContinueExisting); await behavior.Invoke(context, _ => Task.CompletedTask); diff --git a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs index 3dd0eab70c..a009e0f488 100644 --- a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs +++ b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs @@ -17,17 +17,17 @@ public class InstrumentationOptions /// /// Controls how the receive-side processing span relates to the send span for messages sent by this endpoint. - /// Defaults to : receivers continue the trace. + /// Defaults to : receivers continue the trace. /// Can be overridden per message via /// or . /// - public TraceConnector SentMessageTraceConnector { get; set; } = TraceConnector.ChildSpan; + public TraceMode SendTraceMode { get; set; } = TraceMode.ContinueExisting; /// /// Controls how the receive-side processing span relates to the publish span for events published by this endpoint. - /// Defaults to : receivers start a new trace linked back to the publish span. + /// Defaults to : receivers start a new trace linked back to the publish span. /// Can be overridden per message via /// or . /// - public TraceConnector PublishedMessageTraceConnector { get; set; } = TraceConnector.SpanLink; + public TraceMode PublishTraceMode { get; set; } = TraceMode.StartNew; } diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs index ef6ac1ac25..6db97f3169 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryExtensions.cs @@ -22,46 +22,46 @@ public static InstrumentationOptions Tracing(this EndpointConfiguration config) /// /// Start a new OpenTelemetry trace on receive of this message, linked back to the send span. - /// Overrides for this message. + /// Overrides for this message. /// /// The option being extended. public static void StartNewTraceOnReceive(this SendOptions sendOptions) { ArgumentNullException.ThrowIfNull(sendOptions); - sendOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.SpanLink); + sendOptions.Context.Set(TraceConnectorOverrideKey, TraceMode.StartNew); } /// /// Continue the existing OpenTelemetry trace on receive of this message. - /// Overrides for this message. + /// Overrides for this message. /// /// The option being extended. public static void ContinueExistingTraceOnReceive(this SendOptions sendOptions) { ArgumentNullException.ThrowIfNull(sendOptions); - sendOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.ChildSpan); + sendOptions.Context.Set(TraceConnectorOverrideKey, TraceMode.ContinueExisting); } /// /// Start a new OpenTelemetry trace on receive of this event, linked back to the publish span. - /// Overrides for this message. + /// Overrides for this message. /// /// The option being extended. public static void StartNewTraceOnReceive(this PublishOptions publishOptions) { ArgumentNullException.ThrowIfNull(publishOptions); - publishOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.SpanLink); + publishOptions.Context.Set(TraceConnectorOverrideKey, TraceMode.StartNew); } /// /// Continue the existing OpenTelemetry trace on receive of this event. - /// Overrides for this message. + /// Overrides for this message. /// /// The option being extended. public static void ContinueExistingTraceOnReceive(this PublishOptions publishOptions) { ArgumentNullException.ThrowIfNull(publishOptions); - publishOptions.Context.Set(TraceConnectorOverrideKey, TraceConnector.ChildSpan); + publishOptions.Context.Set(TraceConnectorOverrideKey, TraceMode.ContinueExisting); } internal const string TraceConnectorOverrideKey = "NServiceBus.OpenTelemetry.TraceConnectorOverride"; diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs index 630defa1c2..367647a014 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryPublishBehavior.cs @@ -11,11 +11,11 @@ class OpenTelemetryPublishBehavior(InstrumentationOptions instrumentationOptions public Task Invoke(IOutgoingPublishContext context, Func next) { // the per-message override wins over the endpoint-level default - var connector = context.Extensions.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector requestedConnector) + var connector = context.Extensions.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode requestedConnector) ? requestedConnector - : instrumentationOptions.PublishedMessageTraceConnector; + : instrumentationOptions.PublishTraceMode; - context.Headers[Headers.StartNewTrace] = connector == TraceConnector.SpanLink ? bool.TrueString : bool.FalseString; + context.Headers[Headers.StartNewTrace] = connector == TraceMode.StartNew ? bool.TrueString : bool.FalseString; return next(context); } diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs index 9424a5eb3b..75338de54e 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetrySendBehavior.cs @@ -11,11 +11,11 @@ class OpenTelemetrySendBehavior(InstrumentationOptions instrumentationOptions) : public Task Invoke(IOutgoingSendContext context, Func next) { // the per-message override wins over the endpoint-level default - var connector = context.Extensions.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceConnector requestedConnector) + var connector = context.Extensions.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode requestedConnector) ? requestedConnector - : instrumentationOptions.SentMessageTraceConnector; + : instrumentationOptions.SendTraceMode; - context.Headers[Headers.StartNewTrace] = connector == TraceConnector.SpanLink ? bool.TrueString : bool.FalseString; + context.Headers[Headers.StartNewTrace] = connector == TraceMode.StartNew ? bool.TrueString : bool.FalseString; return next(context); } diff --git a/src/NServiceBus.Core/OpenTelemetry/TraceConnector.cs b/src/NServiceBus.Core/OpenTelemetry/TraceMode.cs similarity index 89% rename from src/NServiceBus.Core/OpenTelemetry/TraceConnector.cs rename to src/NServiceBus.Core/OpenTelemetry/TraceMode.cs index 5854b92101..590258b3b7 100644 --- a/src/NServiceBus.Core/OpenTelemetry/TraceConnector.cs +++ b/src/NServiceBus.Core/OpenTelemetry/TraceMode.cs @@ -5,15 +5,15 @@ namespace NServiceBus; /// /// Controls how the receive-side processing span relates to the outgoing send or publish span. /// -public enum TraceConnector +public enum TraceMode { /// /// The receiving endpoint continues the trace: the processing span becomes a child of the outgoing span. /// - ChildSpan, + ContinueExisting, /// /// The receiving endpoint starts a new trace: the processing span becomes the root of a new trace with a link back to the outgoing span. /// - SpanLink + StartNew } From 920c4c28fa37c1bdd13423563d73766282b82d2d Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Thu, 16 Jul 2026 12:22:24 +0200 Subject: [PATCH 06/11] =?UTF-8?q?=E2=9C=A8=20Add=20delayed=20send,=20delay?= =?UTF-8?q?ed=20retry,=20and=20error=20message=20trace=20mode=20settings?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../InstrumentationOptionsTests.cs | 13 ++++++++++++ .../OpenTelemetry/InstrumentationOptions.cs | 21 +++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs index ded2e6a41c..641b4b351a 100644 --- a/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/InstrumentationOptionsTests.cs @@ -16,4 +16,17 @@ public void Should_default_trace_connectors_to_current_behavior() Assert.That(options.PublishTraceMode, Is.EqualTo(TraceMode.StartNew), "publishes start a new linked trace by default"); } } + + [Test] + public void Should_default_delayed_and_error_trace_connectors_to_span_link() + { + var options = new InstrumentationOptions(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(options.DelayedSendTraceMode, Is.EqualTo(TraceMode.StartNew), "delayed sends start a new linked trace by default"); + Assert.That(options.DelayedRetryTraceMode, Is.EqualTo(TraceMode.StartNew), "delayed retries start a new linked trace by default"); + Assert.That(options.ErrorMessageTraceMode, Is.EqualTo(TraceMode.StartNew), "messages moved to the error queue start a new linked trace by default"); + } + } } diff --git a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs index a009e0f488..bf6d014892 100644 --- a/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs +++ b/src/NServiceBus.Core/OpenTelemetry/InstrumentationOptions.cs @@ -30,4 +30,25 @@ public class InstrumentationOptions /// or . /// public TraceMode PublishTraceMode { get; set; } = TraceMode.StartNew; + + /// + /// Controls how the receive-side processing span relates to the send span for delayed messages + /// (messages sent with a delivery delay, including saga timeouts). + /// Defaults to : receivers start a new trace linked back to the send span. + /// Can be overridden per message via + /// or . + /// + public TraceMode DelayedSendTraceMode { get; set; } = TraceMode.StartNew; + + /// + /// Controls how the processing span of a delayed retry relates to the trace of the failed attempt. + /// Defaults to : the retry starts a new trace linked back to the failed attempt. + /// + public TraceMode DelayedRetryTraceMode { get; set; } = TraceMode.StartNew; + + /// + /// Controls how the processing span of a message retried from the error queue relates to the trace of the failed attempt. + /// Defaults to : reprocessing starts a new trace linked back to the failed attempt. + /// + public TraceMode ErrorMessageTraceMode { get; set; } = TraceMode.StartNew; } From 2da0a2c769858d14773e5e5f3adea26f8199a9db Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Thu, 16 Jul 2026 12:35:39 +0200 Subject: [PATCH 07/11] =?UTF-8?q?=E2=9C=A8=20Honor=20DelayedSendTraceMode?= =?UTF-8?q?=20and=20per-message=20overrides=20for=20delayed=20sends?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...penTelemetryDelayedMessageBehaviorTests.cs | 98 +++++++++++++++++++ .../AttachSenderRelatedInfoOnMessageTests.cs | 17 ++++ .../OpenTelemetryDelayedMessageBehavior.cs | 42 ++++++++ .../OpenTelemetry/OpenTelemetryFeature.cs | 5 + ...ttachSenderRelatedInfoOnMessageBehavior.cs | 2 - 5 files changed, 162 insertions(+), 2 deletions(-) create mode 100644 src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryDelayedMessageBehaviorTests.cs create mode 100644 src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryDelayedMessageBehaviorTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryDelayedMessageBehaviorTests.cs new file mode 100644 index 0000000000..f9550c83e5 --- /dev/null +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/OpenTelemetryDelayedMessageBehaviorTests.cs @@ -0,0 +1,98 @@ +namespace NServiceBus.Core.Tests.OpenTelemetry; + +using System; +using System.Threading.Tasks; +using DelayedDelivery; +using NServiceBus.Transport; +using NUnit.Framework; +using Testing; + +[TestFixture] +public class OpenTelemetryDelayedMessageBehaviorTests +{ + [Test] + public async Task Should_not_set_context_entry_when_not_delayed() + { + var context = new TestableRoutingContext(); + + await InvokeBehavior(context, new InstrumentationOptions()); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out _), Is.False); + } + + [Test] + public async Task Should_not_set_context_entry_when_dispatch_properties_carry_no_delay() + { + var context = new TestableRoutingContext(); + context.Extensions.Set(new DispatchProperties()); + + await InvokeBehavior(context, new InstrumentationOptions()); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out _), Is.False); + } + + [Test] + public async Task Should_start_new_trace_by_default_when_delayed_with_delay() + { + var context = new TestableRoutingContext(); + context.Extensions.Set(new DispatchProperties { DelayDeliveryWith = new DelayDeliveryWith(TimeSpan.FromSeconds(5)) }); + + await InvokeBehavior(context, new InstrumentationOptions()); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out var startNewTrace), Is.True); + Assert.That(startNewTrace, Is.EqualTo(bool.TrueString)); + } + + [Test] + public async Task Should_start_new_trace_by_default_when_delayed_with_do_not_deliver_before() + { + var context = new TestableRoutingContext(); + context.Extensions.Set(new DispatchProperties { DoNotDeliverBefore = new DoNotDeliverBefore(DateTimeOffset.UtcNow.AddSeconds(5)) }); + + await InvokeBehavior(context, new InstrumentationOptions()); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out var startNewTrace), Is.True); + Assert.That(startNewTrace, Is.EqualTo(bool.TrueString)); + } + + [Test] + public async Task Should_continue_trace_when_delayed_connector_is_child_span() + { + var context = new TestableRoutingContext(); + context.Extensions.Set(new DispatchProperties { DelayDeliveryWith = new DelayDeliveryWith(TimeSpan.FromSeconds(5)) }); + + await InvokeBehavior(context, new InstrumentationOptions { DelayedSendTraceMode = TraceMode.ContinueExisting }); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out var startNewTrace), Is.True); + Assert.That(startNewTrace, Is.EqualTo(bool.FalseString)); + } + + [Test] + public async Task Should_not_set_context_entry_when_per_message_override_present() + { + var context = new TestableRoutingContext(); + context.Extensions.Set(new DispatchProperties { DelayDeliveryWith = new DelayDeliveryWith(TimeSpan.FromSeconds(5)) }); + context.Extensions.Set(OpenTelemetryExtensions.TraceConnectorOverrideKey, TraceMode.ContinueExisting); + + await InvokeBehavior(context, new InstrumentationOptions()); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out _), Is.False, + "the send behavior already resolved the per-message override into the header"); + } + + [Test] + public async Task Should_not_set_context_entry_for_delayed_retries() + { + var context = new TestableRoutingContext(); + context.Message.Headers[Headers.DelayedRetries] = "1"; + context.Extensions.Set(new DispatchProperties { DelayDeliveryWith = new DelayDeliveryWith(TimeSpan.FromSeconds(5)) }); + + await InvokeBehavior(context, new InstrumentationOptions { DelayedSendTraceMode = TraceMode.ContinueExisting }); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out _), Is.False, + "recoverability already decided the trace boundary for delayed retries"); + } + + static Task InvokeBehavior(TestableRoutingContext context, InstrumentationOptions options) => + new OpenTelemetryDelayedMessageBehavior(options).Invoke(context, _ => Task.CompletedTask); +} diff --git a/src/NServiceBus.Core.Tests/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageTests.cs b/src/NServiceBus.Core.Tests/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageTests.cs index d4625aaada..91791254db 100644 --- a/src/NServiceBus.Core.Tests/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageTests.cs +++ b/src/NServiceBus.Core.Tests/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageTests.cs @@ -109,6 +109,23 @@ public async Task Should_not_override_deliver_at_headerAsync() } } + [Test] + public async Task Should_not_set_start_new_trace_context_entry() + { + var message = new OutgoingMessage("id", [], null); + var stash = new ContextBag(); + stash.Set(new DispatchProperties + { + DelayDeliveryWith = new DelayDeliveryWith(TimeSpan.FromSeconds(2)) + }); + var context = new TestableRoutingContext { Message = message, Extensions = stash, RoutingStrategies = new List { new UnicastRoutingStrategy("_") } }; + + await new AttachSenderRelatedInfoOnMessageBehavior().Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Extensions.TryGet(Headers.StartNewTrace, out _), Is.False, + "the trace boundary decision for delayed messages is owned by OpenTelemetryDelayedMessageBehavior"); + } + static async Task InvokeBehaviorAsync(Dictionary headers = null, DispatchProperties dispatchProperties = null, CancellationToken cancellationToken = default) { var message = new OutgoingMessage("id", headers ?? [], null); diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs new file mode 100644 index 0000000000..1b1017a8e6 --- /dev/null +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs @@ -0,0 +1,42 @@ +#nullable enable + +namespace NServiceBus; + +using System; +using System.Threading.Tasks; +using Pipeline; +using Transport; + +class OpenTelemetryDelayedMessageBehavior(InstrumentationOptions instrumentationOptions) : IBehavior +{ + public Task Invoke(IRoutingContext context, Func next) + { + // Recoverability already decided the trace boundary for delayed retries: the retried message + // carries the decision in its headers, copied from recoverability metadata before the routing + // context is created. The DelayedRetries header identifies that dispatch path — the StartNewTrace + // header itself cannot be used because the send behavior stamps it on every outgoing message. + if (context.Message.Headers.ContainsKey(Headers.DelayedRetries)) + { + return next(context); + } + + bool isDelayed = context.Extensions.TryGet(out var dispatchProperties) + && (dispatchProperties.DelayDeliveryWith != null || dispatchProperties.DoNotDeliverBefore != null); + + if (!isDelayed) + { + return next(context); + } + + // A per-message override was already resolved into the header by the send behavior. + if (context.Extensions.TryGet(OpenTelemetryExtensions.TraceConnectorOverrideKey, out TraceMode _)) + { + return next(context); + } + + context.Extensions.Set(Headers.StartNewTrace, + instrumentationOptions.DelayedSendTraceMode == TraceMode.StartNew ? bool.TrueString : bool.FalseString); + + return next(context); + } +} diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs index b108b409c9..f019499484 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs @@ -20,6 +20,11 @@ protected override void Setup(FeatureConfigurationContext context) "Manages the depth of the trace for sends" ); + context.Pipeline.Register( + new OpenTelemetryDelayedMessageBehavior(instrumentationOptions), + "Manages the depth of the trace for delayed messages" + ); + context.Pipeline.Register( new PopulateRecoverabilityTraceMetadataBehavior(), "Populates the recoverability metadata" diff --git a/src/NServiceBus.Core/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageBehavior.cs b/src/NServiceBus.Core/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageBehavior.cs index 3472a9355f..7ba822fbc6 100644 --- a/src/NServiceBus.Core/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageBehavior.cs +++ b/src/NServiceBus.Core/Pipeline/Outgoing/AttachSenderRelatedInfoOnMessageBehavior.cs @@ -35,12 +35,10 @@ public Task Invoke(IRoutingContext context, Func next) { var timeDelay = dispatchProperties.DelayDeliveryWith.Delay; message.Headers[Headers.DeliverAt] = DateTimeOffsetHelper.ToWireFormattedString(utcNow.Add(timeDelay)); - context.Extensions.Set(Headers.StartNewTrace, bool.TrueString); } else if (dispatchProperties.DoNotDeliverBefore != null) { message.Headers[Headers.DeliverAt] = DateTimeOffsetHelper.ToWireFormattedString(dispatchProperties.DoNotDeliverBefore.At); - context.Extensions.Set(Headers.StartNewTrace, bool.TrueString); } } } From eeff00aac86243a40422d92282808e30ac4b81fe Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Thu, 16 Jul 2026 12:42:48 +0200 Subject: [PATCH 08/11] =?UTF-8?q?=E2=9C=A8=20Honor=20DelayedRetryTraceMode?= =?UTF-8?q?=20and=20ErrorMessageTraceMode=20in=20recoverability=20trace=20?= =?UTF-8?q?metadata?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...ecoverabilityTraceMetadataBehaviorTests.cs | 75 ++++++++++++++++++- .../OpenTelemetry/OpenTelemetryFeature.cs | 2 +- ...lateRecoverabilityTraceMetadataBehavior.cs | 14 +++- 3 files changed, 84 insertions(+), 7 deletions(-) diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs index cee304d848..e444890a59 100644 --- a/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs @@ -12,7 +12,7 @@ public class PopulateRecoverabilityTraceMetadataBehaviorTests [Test] public async Task Should_not_write_metadata_when_trace_not_present() { - var behavior = new PopulateRecoverabilityTraceMetadataBehavior(); + var behavior = new PopulateRecoverabilityTraceMetadataBehavior(new InstrumentationOptions()); var context = new TestableRecoverabilityContext(); await behavior.Invoke(context, _ => Task.CompletedTask); @@ -28,7 +28,7 @@ public async Task Should_not_write_metadata_when_trace_not_present() [TestCaseSource(nameof(Actions))] public async Task Should_write_metadata_when_trace_present(RecoverabilityAction recoverabilityAction) { - var behavior = new PopulateRecoverabilityTraceMetadataBehavior(); + var behavior = new PopulateRecoverabilityTraceMetadataBehavior(new InstrumentationOptions()); var context = new TestableRecoverabilityContext { @@ -42,13 +42,82 @@ public async Task Should_write_metadata_when_trace_present(RecoverabilityAction { Assert.That(context.Headers, Does.Not.ContainKey(Headers.StartNewTrace)); Assert.That(context.Metadata, Does.ContainKey(Headers.StartNewTrace)); + Assert.That(context.Metadata[Headers.StartNewTrace], Is.EqualTo(bool.TrueString), "defaults preserve current behavior"); } } + [Test] + public async Task Should_continue_trace_for_delayed_retry_when_connector_is_child_span() + { + var behavior = new PopulateRecoverabilityTraceMetadataBehavior(new InstrumentationOptions { DelayedRetryTraceMode = TraceMode.ContinueExisting }); + + var context = new TestableRecoverabilityContext + { + Headers = { { Headers.DiagnosticsTraceParent, "traceparent" } }, + RecoverabilityAction = new DelayedRetry(TimeSpan.FromSeconds(10)) + }; + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Metadata[Headers.StartNewTrace], Is.EqualTo(bool.FalseString)); + } + + [Test] + public async Task Should_continue_trace_for_move_to_error_when_connector_is_child_span() + { + var behavior = new PopulateRecoverabilityTraceMetadataBehavior(new InstrumentationOptions { ErrorMessageTraceMode = TraceMode.ContinueExisting }); + + var context = new TestableRecoverabilityContext + { + Headers = { { Headers.DiagnosticsTraceParent, "traceparent" } }, + RecoverabilityAction = new MoveToError("errorqueue") + }; + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Metadata[Headers.StartNewTrace], Is.EqualTo(bool.FalseString)); + } + + [Test] + public async Task Should_not_apply_error_connector_to_delayed_retries() + { + var behavior = new PopulateRecoverabilityTraceMetadataBehavior(new InstrumentationOptions { ErrorMessageTraceMode = TraceMode.ContinueExisting }); + + var context = new TestableRecoverabilityContext + { + Headers = { { Headers.DiagnosticsTraceParent, "traceparent" } }, + RecoverabilityAction = new DelayedRetry(TimeSpan.FromSeconds(10)) + }; + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Metadata[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); + } + + [Test] + public async Task Should_start_new_trace_for_other_actions_regardless_of_connectors() + { + var behavior = new PopulateRecoverabilityTraceMetadataBehavior(new InstrumentationOptions + { + DelayedRetryTraceMode = TraceMode.ContinueExisting, + ErrorMessageTraceMode = TraceMode.ContinueExisting + }); + + var context = new TestableRecoverabilityContext + { + Headers = { { Headers.DiagnosticsTraceParent, "traceparent" } }, + RecoverabilityAction = new ImmediateRetry() + }; + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Metadata[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); + } + static IEnumerable Actions() { yield return new ImmediateRetry(); yield return new DelayedRetry(TimeSpan.FromSeconds(10)); yield return new MoveToError("errorqueue"); } -} \ No newline at end of file +} diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs index f019499484..992f942176 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryFeature.cs @@ -26,7 +26,7 @@ protected override void Setup(FeatureConfigurationContext context) ); context.Pipeline.Register( - new PopulateRecoverabilityTraceMetadataBehavior(), + new PopulateRecoverabilityTraceMetadataBehavior(instrumentationOptions), "Populates the recoverability metadata" ); } diff --git a/src/NServiceBus.Core/OpenTelemetry/Tracing/PopulateRecoverabilityTraceMetadataBehavior.cs b/src/NServiceBus.Core/OpenTelemetry/Tracing/PopulateRecoverabilityTraceMetadataBehavior.cs index 0bb256ecef..035dc2ad53 100644 --- a/src/NServiceBus.Core/OpenTelemetry/Tracing/PopulateRecoverabilityTraceMetadataBehavior.cs +++ b/src/NServiceBus.Core/OpenTelemetry/Tracing/PopulateRecoverabilityTraceMetadataBehavior.cs @@ -4,7 +4,7 @@ namespace NServiceBus; using System.Threading.Tasks; using Pipeline; -class PopulateRecoverabilityTraceMetadataBehavior : IBehavior +class PopulateRecoverabilityTraceMetadataBehavior(InstrumentationOptions instrumentationOptions) : IBehavior { public Task Invoke(IRecoverabilityContext context, Func next) { @@ -13,10 +13,18 @@ public Task Invoke(IRecoverabilityContext context, Func instrumentationOptions.DelayedRetryTraceMode, + MoveToError => instrumentationOptions.ErrorMessageTraceMode, + // custom recoverability actions keep the pre-existing behavior of always starting a new trace + _ => TraceMode.StartNew + }; + // Setting it to the metadata makes sure it is propagated to the headers // even in more advanced scenarios like native dead-lettering - context.Metadata[Headers.StartNewTrace] = bool.TrueString; + context.Metadata[Headers.StartNewTrace] = connector == TraceMode.StartNew ? bool.TrueString : bool.FalseString; return next(context); } -} \ No newline at end of file +} From 73ba4475ca2ef5da807bb9db9cefd5740b171a75 Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Thu, 16 Jul 2026 12:48:02 +0200 Subject: [PATCH 09/11] =?UTF-8?q?=E2=9C=A8=20Add=20acceptance=20tests=20fo?= =?UTF-8?q?r=20delayed=20send,=20delayed=20retry,=20and=20error=20message?= =?UTF-8?q?=20trace=20modes?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...n_incoming_message_moved_to_error_queue.cs | 30 ++++ .../When_incoming_message_was_delayed.cs | 155 ++++++++++++++++++ 2 files changed, 185 insertions(+) diff --git a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_moved_to_error_queue.cs b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_moved_to_error_queue.cs index 0c2383481c..ecf0870070 100644 --- a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_moved_to_error_queue.cs +++ b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_moved_to_error_queue.cs @@ -26,6 +26,36 @@ public async Task Should_add_start_new_trace_header() Assert.That(context.ErrorMessageHeaders[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); } + [Test] + public async Task Should_not_add_start_new_trace_header_when_error_connector_is_child_span() + { + var context = await Scenario.Define() + .WithEndpoint(e => e + .When(s => s.SendLocal(new FailingMessage())) + .DoNotFailOnErrorMessages()) + .WithEndpoint() + .Run(); + + Assert.That(context.ErrorMessageHeaders[Headers.StartNewTrace], Is.EqualTo(bool.FalseString)); + } + + public class FailingEndpointWithChildSpanConnector : EndpointConfigurationBuilder + { + static readonly string ErrorQueueAddress = Conventions.EndpointNamingConvention(typeof(ErrorSpy)); + + public FailingEndpointWithChildSpanConnector() => EndpointSetup(c => + { + c.SendFailedMessagesTo(ErrorQueueAddress); + c.Tracing().ErrorMessageTraceMode = TraceMode.ContinueExisting; + }); + + [Handler] + public class FailingMessageHandler() : IHandleMessages + { + public Task Handle(FailingMessage message, IMessageHandlerContext context) => throw new SimulatedException(ErrorMessage); + } + } + public class Context : ScenarioContext { public Dictionary ErrorMessageHeaders { get; set; } diff --git a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs index 253551a392..b48e2ab254 100644 --- a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs +++ b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs @@ -152,6 +152,161 @@ await Scenario.Define() } } + [Test] + public async Task By_sendoptions_Should_continue_trace_when_endpoint_connector_is_child_span() + { + await Scenario.Define() + .WithEndpoint(b => b + .When(s => + { + var sendOptions = new SendOptions(); + sendOptions.DelayDeliveryWith(TimeSpan.FromMilliseconds(5)); + sendOptions.RouteToThisEndpoint(); + return s.Send(new DelayedMessage(), sendOptions); + })) + .Run(); + + var incomingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities(); + var outgoingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetSendMessageActivities(); + using (Assert.EnterMultipleScope()) + { + Assert.That(incomingMessageActivities, Has.Count.EqualTo(1), "1 message is received as part of this test"); + Assert.That(outgoingMessageActivities, Has.Count.EqualTo(1), "1 message is sent as part of this test"); + } + + var sendRequest = outgoingMessageActivities[0]; + var receiveRequest = incomingMessageActivities[0]; + + using (Assert.EnterMultipleScope()) + { + Assert.That(receiveRequest.RootId, Is.EqualTo(sendRequest.RootId), "delayed send and receive operations are part of the same root activity"); + Assert.That(receiveRequest.ParentId, Is.Not.Null, "incoming delayed message does have a parent"); + } + + Assert.That(receiveRequest.Links, Is.Empty, "receive does not have links"); + } + + [Test] + public async Task By_sendoptions_Should_start_new_trace_when_option_overrides_child_span_connector() + { + await Scenario.Define() + .WithEndpoint(b => b + .When(s => + { + var sendOptions = new SendOptions(); + sendOptions.DelayDeliveryWith(TimeSpan.FromMilliseconds(5)); + sendOptions.RouteToThisEndpoint(); + sendOptions.StartNewTraceOnReceive(); + return s.Send(new DelayedMessage(), sendOptions); + })) + .Run(); + + var incomingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities(); + var outgoingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetSendMessageActivities(); + using (Assert.EnterMultipleScope()) + { + Assert.That(incomingMessageActivities, Has.Count.EqualTo(1), "1 message is received as part of this test"); + Assert.That(outgoingMessageActivities, Has.Count.EqualTo(1), "1 message is sent as part of this test"); + } + + var sendRequest = outgoingMessageActivities[0]; + var receiveRequest = incomingMessageActivities[0]; + + using (Assert.EnterMultipleScope()) + { + Assert.That(receiveRequest.RootId, Is.Not.EqualTo(sendRequest.RootId), "per-message override wins: send and receive are different root activities"); + Assert.That(receiveRequest.ParentId, Is.Null, "incoming message does not have a parent, it's a root"); + } + + ActivityLink link = receiveRequest.Links.FirstOrDefault(); + Assert.That(link, Is.Not.Default, "receive has a link"); + Assert.That(link.Context.TraceId, Is.EqualTo(sendRequest.TraceId), "receive is linked to the send operation"); + } + + [Test] + public void By_retry_Should_continue_trace_when_endpoint_connector_is_child_span() + { + _ = Assert.CatchAsync(async () => + { + await Scenario.Define() + .WithEndpoint(b => b + .When(session => session.SendLocal(new MessageToBeRetried()))) + .Run(); + }); + + var incomingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities(); + var outgoingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetSendMessageActivities(); + using (Assert.EnterMultipleScope()) + { + Assert.That(incomingMessageActivities, Has.Count.EqualTo(2), "2 messages are received as part of this test (2 attempts)"); + Assert.That(outgoingMessageActivities, Has.Count.EqualTo(1), "1 message sent as part of this test"); + } + + var sendRequest = outgoingMessageActivities[0]; + var firstAttemptReceiveRequest = incomingMessageActivities[0]; + var secondAttemptReceiveRequest = incomingMessageActivities[1]; + + using (Assert.EnterMultipleScope()) + { + Assert.That(firstAttemptReceiveRequest.RootId, Is.EqualTo(sendRequest.RootId), "first send operation is the root activity"); + Assert.That(secondAttemptReceiveRequest.RootId, Is.EqualTo(sendRequest.RootId), "delayed retry stays in the same trace when connector is child span"); + Assert.That(secondAttemptReceiveRequest.ParentId, Is.Not.Null, "second incoming message does have a parent"); + } + + Assert.That(secondAttemptReceiveRequest.Links, Is.Empty, "second receive does not have links"); + } + + public class ChildSpanConnectorEndpoint : EndpointConfigurationBuilder + { + public ChildSpanConnectorEndpoint() + { + var template = new DefaultServer + { + TransportConfiguration = new ConfigureEndpointAcceptanceTestingTransport(false, true) + }; + EndpointSetup( + template, + (c, _) => c.Tracing().DelayedSendTraceMode = TraceMode.ContinueExisting, + metadata => { }); + } + + [Handler] + public class DelayedMessageHandler(Context testContext) : IHandleMessages + { + public Task Handle(DelayedMessage message, IMessageHandlerContext context) + { + testContext.DelayedMessageReceived = true; + testContext.MaybeCompleted(); + return Task.CompletedTask; + } + } + } + + public class ChildSpanRetryConnectorEndpoint : EndpointConfigurationBuilder + { + public ChildSpanRetryConnectorEndpoint() + { + var template = new DefaultServer + { + TransportConfiguration = new ConfigureEndpointAcceptanceTestingTransport(false, true) + }; + EndpointSetup( + template, + (c, _) => + { + c.Tracing().DelayedRetryTraceMode = TraceMode.ContinueExisting; + var recoverability = c.Recoverability(); + recoverability.Delayed(settings => settings.NumberOfRetries(1).TimeIncrease(TimeSpan.FromMilliseconds(1))); + }, metadata => { }); + } + + [Handler] + public class MessageToBeRetriedHandler : IHandleMessages + { + public Task Handle(MessageToBeRetried message, IMessageHandlerContext context) => throw new SimulatedException(); + } + } + public class Context : ScenarioContext { public bool ReplyMessageReceived { get; set; } From 92ff29e3e9f86d239b98e2e29c1071a074f70303 Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Thu, 16 Jul 2026 12:52:40 +0200 Subject: [PATCH 10/11] =?UTF-8?q?=E2=9C=A8=20Approve=20public=20API=20for?= =?UTF-8?q?=20trace=20mode=20rename=20and=20delayed/error=20message=20addi?= =?UTF-8?q?tions?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../APIApprovals.ApproveNServiceBus.approved.txt | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt b/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt index 6849210e76..ce4e88240f 100644 --- a/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt +++ b/src/NServiceBus.Core.Tests/ApprovalFiles/APIApprovals.ApproveNServiceBus.approved.txt @@ -635,8 +635,11 @@ namespace NServiceBus public class InstrumentationOptions { public InstrumentationOptions() { } - public NServiceBus.TraceConnector PublishedMessageTraceConnector { get; set; } - public NServiceBus.TraceConnector SentMessageTraceConnector { get; set; } + public NServiceBus.TraceMode DelayedRetryTraceMode { get; set; } + public NServiceBus.TraceMode DelayedSendTraceMode { get; set; } + public NServiceBus.TraceMode ErrorMessageTraceMode { get; set; } + public NServiceBus.TraceMode PublishTraceMode { get; set; } + public NServiceBus.TraceMode SendTraceMode { get; set; } public bool UseMessageDestinationInSpanNames { get; set; } } public sealed class KeyedServiceKey @@ -1177,10 +1180,10 @@ namespace NServiceBus public ToSagaExpression(NServiceBus.IConfigureHowToFindSagaWithMessage sagaMessageFindingConfiguration, System.Linq.Expressions.Expression> messageProperty) { } public void ToSaga(System.Linq.Expressions.Expression> sagaEntityProperty) { } } - public enum TraceConnector + public enum TraceMode { - ChildSpan = 0, - SpanLink = 1, + ContinueExisting = 0, + StartNew = 1, } public static class TransportConfig { From 51c7ca7b117196c3ac30e1690f353cc72ac56f00 Mon Sep 17 00:00:00 2001 From: Ramon Smits Date: Thu, 16 Jul 2026 13:01:10 +0200 Subject: [PATCH 11/11] =?UTF-8?q?=E2=9A=9C=EF=B8=8F=20Address=20review=20p?= =?UTF-8?q?olish:=20guard=20comment,=20symmetric=20trace=20mode=20test,=20?= =?UTF-8?q?retry=20assertion?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Traces/When_incoming_message_was_delayed.cs | 1 + ...teRecoverabilityTraceMetadataBehaviorTests.cs | 16 ++++++++++++++++ .../OpenTelemetryDelayedMessageBehavior.cs | 2 ++ 3 files changed, 19 insertions(+) diff --git a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs index b48e2ab254..05f24edbe4 100644 --- a/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs +++ b/src/NServiceBus.AcceptanceTests/Core/OpenTelemetry/Traces/When_incoming_message_was_delayed.cs @@ -249,6 +249,7 @@ await Scenario.Define() using (Assert.EnterMultipleScope()) { Assert.That(firstAttemptReceiveRequest.RootId, Is.EqualTo(sendRequest.RootId), "first send operation is the root activity"); + Assert.That(firstAttemptReceiveRequest.ParentId, Is.EqualTo(sendRequest.Id), "first incoming message is correlated to the first send operation"); Assert.That(secondAttemptReceiveRequest.RootId, Is.EqualTo(sendRequest.RootId), "delayed retry stays in the same trace when connector is child span"); Assert.That(secondAttemptReceiveRequest.ParentId, Is.Not.Null, "second incoming message does have a parent"); } diff --git a/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs b/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs index e444890a59..c881fe7516 100644 --- a/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs +++ b/src/NServiceBus.Core.Tests/OpenTelemetry/PopulateRecoverabilityTraceMetadataBehaviorTests.cs @@ -94,6 +94,22 @@ public async Task Should_not_apply_error_connector_to_delayed_retries() Assert.That(context.Metadata[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); } + [Test] + public async Task Should_not_apply_delayed_retry_connector_to_move_to_error() + { + var behavior = new PopulateRecoverabilityTraceMetadataBehavior(new InstrumentationOptions { DelayedRetryTraceMode = TraceMode.ContinueExisting }); + + var context = new TestableRecoverabilityContext + { + Headers = { { Headers.DiagnosticsTraceParent, "traceparent" } }, + RecoverabilityAction = new MoveToError("errorqueue") + }; + + await behavior.Invoke(context, _ => Task.CompletedTask); + + Assert.That(context.Metadata[Headers.StartNewTrace], Is.EqualTo(bool.TrueString)); + } + [Test] public async Task Should_start_new_trace_for_other_actions_regardless_of_connectors() { diff --git a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs index 1b1017a8e6..5008eb93a9 100644 --- a/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs +++ b/src/NServiceBus.Core/OpenTelemetry/OpenTelemetryDelayedMessageBehavior.cs @@ -15,6 +15,8 @@ public Task Invoke(IRoutingContext context, Func next) // carries the decision in its headers, copied from recoverability metadata before the routing // context is created. The DelayedRetries header identifies that dispatch path — the StartNewTrace // header itself cannot be used because the send behavior stamps it on every outgoing message. + // This assumes the DelayedRetries header only appears on recoverability-generated retries; a user + // manually copying framework-internal incoming headers onto a new delayed send would skip this behavior. if (context.Message.Headers.ContainsKey(Headers.DelayedRetries)) { return next(context);