Skip to content
Draft
Show file tree
Hide file tree
Changes from 2 commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
fd48fac
Starting working branch for Otel
irinascurtu Jun 2, 2026
dd3a9db
Added base64 encoded body
irinascurtu Jun 2, 2026
0890965
Apply suggestion from @ramonsmits
ramonsmits Jun 3, 2026
b083629
removed the body serialization
irinascurtu Jun 10, 2026
af64227
access level fixes for acceptance tests
tmasternak Jun 16, 2026
bba80f7
Using DistributedContextPropagator instead of hand-written baggage pr…
tmasternak Jun 16, 2026
ca842e6
🐛 Add regression test for null baggage value (#6983)
ramonsmits Jun 17, 2026
b35ca4b
Merge pull request #7824 from Particular/fix-6983-null-baggage-value
irinascurtu Jun 17, 2026
a33484e
Make DistributedContextPropagator opt-in (keep OTel propagation backw…
ramonsmits Jun 19, 2026
68274cf
Emit handler spans from a dedicated NServiceBus.Core.Handler Activity…
ramonsmits Jul 8, 2026
beeadd0
Optout on dispatching events (#7846)
irinascurtu Jul 16, 2026
338ea5f
Gauge meter for active message processings (#7841)
tmasternak Jul 16, 2026
7d8cc93
Endpoint-level trace connector defaults for sends and publishes (#7867)
ramonsmits Jul 16, 2026
b7eb469
Add meter for total number of messages deduplicated via Outbox (#7864)
irinascurtu Jul 23, 2026
a457e83
Allow changing trace continuation behavior for delayed messages (#7845)
ramonsmits Jul 29, 2026
351d0f9
Added error.type (#7885)
irinascurtu Jul 29, 2026
7e6873b
Spans for recoverability actions (#7890)
tmasternak Jul 31, 2026
9e5b85e
Add option for exception details capturing via logging (#7899)
tmasternak Aug 5, 2026
3345a88
Additional performance-related instruments (#7898)
irinascurtu Aug 10, 2026
a1b0337
Minor tweaks to the OpenTelemetry featue (#7908)
tmasternak Aug 11, 2026
0c135cd
Removed RecordedExceptions tracking in favor of directly using except…
tmasternak Aug 13, 2026
96f8a16
Fix edge case where InstrumentionOptions could not be a shared instan…
ramonsmits Aug 26, 2026
e3eaa5c
Add support for instrument-specific metric tags in the incoming pipel…
tmasternak Sep 1, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -80,5 +80,37 @@ public Task Handle(IncomingMessage message, IMessageHandlerContext context)
}
}

[Test]
public async Task Should_use_receive_address_in_span_name_when_opted_in()
{
await Scenario.Define<Context>()
.WithEndpoint<ReceivingEndpointWithDestinationNaming>(e => e
.When(s => s.SendLocal(new IncomingMessage())))
.Run();

var incomingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetReceiveMessageActivities();
Assert.That(incomingMessageActivities, Has.Count.EqualTo(1));

var incomingActivity = incomingMessageActivities.Single();
Assert.That(incomingActivity.DisplayName, Does.StartWith("process "));
Assert.That(incomingActivity.DisplayName, Is.Not.EqualTo("process message"));
}

public class ReceivingEndpointWithDestinationNaming : EndpointConfigurationBuilder
{
public ReceivingEndpointWithDestinationNaming() =>
EndpointSetup<DefaultServer>(b => b.Tracing().UseMessageDestinationInSpanNames = true);

[Handler]
public class MessageHandler(Context testContext) : IHandleMessages<IncomingMessage>
{
public Task Handle(IncomingMessage message, IMessageHandlerContext context)
{
testContext.MarkAsCompleted();
return Task.CompletedTask;
}
}
}

public class IncomingMessage : IMessage;
}
Original file line number Diff line number Diff line change
Expand Up @@ -179,5 +179,68 @@ public Task Handle(ThisIsAnEvent @event, IMessageHandlerContext context)
}
}

[Test]
public async Task Should_use_event_type_in_span_name_when_opted_in()
{
await Scenario.Define<Context>()
.WithEndpoint<PublisherWithDestinationNaming>(b => b
.When(ctx => ctx.SomeEventSubscribed, s => s.Publish<ThisIsAnEvent>()))
.WithEndpoint<SubscriberForPublisherWithDestinationNaming>(b => b.When((session, ctx) =>
{
if (ctx.HasNativePubSubSupport)
{
ctx.SomeEventSubscribed = true;
}

return Task.CompletedTask;
}))
.Run();

var outgoingEventActivities = NServiceBusActivityListener.CompletedActivities.GetPublishEventActivities();
Assert.That(outgoingEventActivities, Has.Count.EqualTo(1));

var publishedMessage = outgoingEventActivities.Single();
Assert.That(publishedMessage.DisplayName, Is.EqualTo("publish ThisIsAnEvent"));
}

class PublisherWithDestinationNaming : EndpointConfigurationBuilder
{
public PublisherWithDestinationNaming() =>
EndpointSetup<DefaultServer>(b =>
{
b.Tracing().UseMessageDestinationInSpanNames = true;
b.OnEndpointSubscribed<Context>((s, context) =>
{
if (s.SubscriberEndpoint.Contains(Conventions.EndpointNamingConvention(typeof(SubscriberForPublisherWithDestinationNaming))))
{
if (s.MessageType == typeof(ThisIsAnEvent).AssemblyQualifiedName)
{
context.SomeEventSubscribed = true;
}
}
});
});
}

class SubscriberForPublisherWithDestinationNaming : EndpointConfigurationBuilder
{
public SubscriberForPublisherWithDestinationNaming() =>
EndpointSetup<DefaultServer>(c => { },
metadata =>
{
metadata.RegisterPublisherFor<ThisIsAnEvent, PublisherWithDestinationNaming>();
});

[Handler]
public class ThisHandlesSomethingHandler(Context testContext) : IHandleMessages<ThisIsAnEvent>
{
public Task Handle(ThisIsAnEvent @event, IMessageHandlerContext context)
{
testContext.MarkAsCompleted();
return Task.CompletedTask;
}
}
}

public class ThisIsAnEvent : IEvent;
}
Original file line number Diff line number Diff line change
Expand Up @@ -131,5 +131,37 @@ public Task Handle(OutgoingMessage message, IMessageHandlerContext context)
}
}

[Test]
public async Task Should_use_destination_in_send_span_name_when_opted_in()
{
await Scenario.Define<Context>()
.WithEndpoint<TestEndpointWithDestinationNaming>(b => b
.When(s => s.SendLocal(new OutgoingMessage())))
.Run();

var outgoingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetSendMessageActivities();
Assert.That(outgoingMessageActivities, Has.Count.EqualTo(1));

var sentMessage = outgoingMessageActivities.Single();
Assert.That(sentMessage.DisplayName, Does.StartWith("send "));
Assert.That(sentMessage.DisplayName, Is.Not.EqualTo("send message"));
}

public class TestEndpointWithDestinationNaming : EndpointConfigurationBuilder
{
public TestEndpointWithDestinationNaming() =>
EndpointSetup<DefaultServer>(b => b.Tracing().UseMessageDestinationInSpanNames = true);

[Handler]
public class MessageHandler(Context testContext) : IHandleMessages<OutgoingMessage>
{
public Task Handle(OutgoingMessage message, IMessageHandlerContext context)
{
testContext.MarkAsCompleted();
return Task.CompletedTask;
}
}
}

public class OutgoingMessage : IMessage;
}
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,40 @@ public Task Handle(OutgoingReply message, IMessageHandlerContext context)
}
}

[Test]
public async Task Should_use_destination_in_reply_span_name_when_opted_in()
{
await Scenario.Define<Context>()
.WithEndpoint<TestEndpointWithDestinationNaming>(b => b
.When(s => s.SendLocal(new IncomingMessage())))
.Run();

var outgoingMessageActivities = NServiceBusActivityListener.CompletedActivities.GetSendMessageActivities();
Assert.That(outgoingMessageActivities, Has.Count.EqualTo(2), "2 messages are being sent");
var replyMessage = outgoingMessageActivities[1];

Assert.That(replyMessage.DisplayName, Does.StartWith("reply "));
}

public class TestEndpointWithDestinationNaming : EndpointConfigurationBuilder
{
public TestEndpointWithDestinationNaming() =>
EndpointSetup<DefaultServer>(b => b.Tracing().UseMessageDestinationInSpanNames = true);

[Handler]
public class MessageHandler(Context testContext) : IHandleMessages<IncomingMessage>,
IHandleMessages<OutgoingReply>
{
public Task Handle(IncomingMessage message, IMessageHandlerContext context) => context.Reply(new OutgoingReply());

public Task Handle(OutgoingReply message, IMessageHandlerContext context)
{
testContext.MarkAsCompleted();
return Task.CompletedTask;
}
}
}

public class IncomingMessage : IMessage;

public class OutgoingReply : IMessage;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -632,6 +632,12 @@ namespace NServiceBus
StopApplication = 0,
Continue = 1,
}
public class InstrumentationOptions
{
public InstrumentationOptions() { }
public NServiceBus.MessagePayloadAsTag MessagePayloadAsTag { get; set; }
public bool UseMessageDestinationInSpanNames { get; set; }
}
public sealed class KeyedServiceKey
{
public const string Any = "_______<ANY>_______";
Expand Down Expand Up @@ -718,6 +724,13 @@ namespace NServiceBus
Unsubscribe = 4,
Reply = 5,
}
public enum MessagePayloadAsTag
{
None = 0,
IncomingMessage = 1,
OutgoingMessage = 2,
All = 3,
}
public static class MessageProcessingContextExtensions
{
public static System.Threading.Tasks.Task Reply(this NServiceBus.IMessageProcessingContext context, object message) { }
Expand Down Expand Up @@ -772,6 +785,7 @@ namespace NServiceBus
{
public static void ContinueExistingTraceOnReceive(this NServiceBus.PublishOptions publishOptions) { }
public static void StartNewTraceOnReceive(this NServiceBus.SendOptions sendOptions) { }
public static NServiceBus.InstrumentationOptions Tracing(this NServiceBus.EndpointConfiguration config) { }
}
public static class OutboxConfigExtensions
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ namespace NServiceBus.Core.Tests.OpenTelemetry;
[TestFixture]
public class ActivityFactoryTests
{
readonly ActivityFactory activityFactory = new();
readonly ActivityFactory activityFactory = new(new InstrumentationOptions());

TestingActivityListener nsbActivityListener;

Expand All @@ -29,7 +29,7 @@ public class ActivityFactoryTests

class NoDiagnosticListeners
{
readonly ActivityFactory activityFactory = new();
readonly ActivityFactory activityFactory = new(new InstrumentationOptions());

[Test]
public void Should_return_null_incoming_activity_when_no_listeners()
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
#nullable enable

namespace NServiceBus.Core.Tests.OpenTelemetry;

using System;
using System.Collections.Immutable;
using System.Text;
using System.Text.Json;
using System.Threading.Tasks;
using Helpers;
using NServiceBus.Pipeline;
using NUnit.Framework;
using Testing;
using Unicast.Messages;

[TestFixture]
public class MessagePayloadToTagsBehaviorTests
{
TestingActivityListener activityListener;

[SetUp]
public void SetUp() => activityListener = TestingActivityListener.SetupNServiceBusDiagnosticListener();

[TearDown]
public void TearDown() => activityListener.Dispose();

[Test]
public async Task Incoming_should_set_base64_encoded_json_body_tag()
{
using var activity = ActivitySources.Main.StartActivity("test");

var message = new TestMessage { Name = "Hello", Value = 42 };
var context = new TestableIncomingLogicalMessageContext
{
Message = new LogicalMessage(new MessageMetadata(typeof(TestMessage)), message)
};

await new IncomingMessagePayloadToTagsBehavior().Invoke(context, _ => Task.CompletedTask);

var tags = activity!.Tags.ToImmutableDictionary();
Assert.That(tags.ContainsKey("nservicebus.message.body"), Is.True);

var decoded = Encoding.UTF8.GetString(Convert.FromBase64String(tags["nservicebus.message.body"]!));
var deserialized = JsonSerializer.Deserialize<TestMessage>(decoded);

using (Assert.EnterMultipleScope())
{
Assert.That(deserialized!.Name, Is.EqualTo("Hello"));
Assert.That(deserialized.Value, Is.EqualTo(42));
}
}

[Test]
public async Task Incoming_should_not_set_tag_when_no_active_activity()
{
var context = new TestableIncomingLogicalMessageContext
{
Message = new LogicalMessage(new MessageMetadata(typeof(TestMessage)), new TestMessage())
};

Assert.DoesNotThrowAsync(() => new IncomingMessagePayloadToTagsBehavior().Invoke(context, _ => Task.CompletedTask));
}

[Test]
public async Task Outgoing_should_set_base64_encoded_json_body_tag()
{
using var activity = ActivitySources.Main.StartActivity("test");

var message = new TestMessage { Name = "World", Value = 99 };
var context = new TestableOutgoingLogicalMessageContext();
context.UpdateMessage(message);

await new OutgoingMessagePayloadToTagsBehavior().Invoke(context, _ => Task.CompletedTask);

var tags = activity!.Tags.ToImmutableDictionary();
Assert.That(tags.ContainsKey("nservicebus.message.body"), Is.True);

var decoded = Encoding.UTF8.GetString(Convert.FromBase64String(tags["nservicebus.message.body"]!));
var deserialized = JsonSerializer.Deserialize<TestMessage>(decoded);

using (Assert.EnterMultipleScope())
{
Assert.That(deserialized!.Name, Is.EqualTo("World"));
Assert.That(deserialized.Value, Is.EqualTo(99));
}
}

[Test]
public async Task Outgoing_should_not_set_tag_when_no_active_activity()
{
var context = new TestableOutgoingLogicalMessageContext();
context.UpdateMessage(new TestMessage());

Assert.DoesNotThrowAsync(() => new OutgoingMessagePayloadToTagsBehavior().Invoke(context, _ => Task.CompletedTask));
}

class TestMessage
{
public string? Name { get; set; }
public int Value { get; set; }
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -130,7 +130,7 @@ static MainPipelineExecutor CreateMainPipelineExecutor(ServiceProvider servicePr
new TestableMessageOperations(),
new Notification<ReceivePipelineCompleted>(),
receivePipeline,
new ActivityFactory(),
new ActivityFactory(new InstrumentationOptions()),
incomingPipelineMetrics,
new EnvelopeUnwrapper([], incomingPipelineMetrics));

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ class TestableMessageOperations : MessageOperations
public Pipeline<ISubscribeContext> SubscribePipeline => (Pipeline<ISubscribeContext>)subscribePipeline;
public Pipeline<IUnsubscribeContext> UnsubscribePipeline => (Pipeline<IUnsubscribeContext>)unsubscribePipeline;

public TestableMessageOperations() : base(new MessageMapper(), new Pipeline<IOutgoingPublishContext>(), new Pipeline<IOutgoingSendContext>(), new Pipeline<IOutgoingReplyContext>(), new Pipeline<ISubscribeContext>(), new Pipeline<IUnsubscribeContext>(), new ActivityFactory())
public TestableMessageOperations() : base(new MessageMapper(), new Pipeline<IOutgoingPublishContext>(), new Pipeline<IOutgoingSendContext>(), new Pipeline<IOutgoingReplyContext>(), new Pipeline<ISubscribeContext>(), new Pipeline<IUnsubscribeContext>(), new ActivityFactory(new InstrumentationOptions()))
{
}

Expand Down
Loading
Loading