Skip to content
Draft
Show file tree
Hide file tree
Changes from all 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
@@ -0,0 +1,110 @@
namespace NServiceBus.AcceptanceTests.Core.OpenTelemetry.Metrics;

using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using EndpointTemplates;
using NServiceBus;
using AcceptanceTesting;
using NServiceBus.Pipeline;
using NUnit.Framework;
using global::OpenTelemetry;
using global::OpenTelemetry.Metrics;

public class When_customizing_metric_tags : OpenTelemetryAcceptanceTest
{
const string TotalFetched = "nservicebus.messaging.fetches";
const string MessageDeserializeTime = "nservicebus.messaging.deserialize_time";
const string EndpointDiscriminatorTag = "nservicebus.discriminator";
const string EnclosedMessageTypesTag = "nservicebus.enclosed_message_types";
const string TenantTag = "acceptance.tenant_id";
const string FriendlyMessageTypeName = "Order placed (friendly name)";

[Test]
public async Task Should_allow_adding_removing_and_overriding_tags_per_instrument()
{
using var metricsListener = TestingMetricListener.SetupNServiceBusMetricsListener();

List<Metric> exportedMetrics = [];
using var meterProvider = Sdk.CreateMeterProviderBuilder()
.AddMeter("NServiceBus.Core.Pipeline.Incoming")
.AddView(TotalFetched, new MetricStreamConfiguration
{
TagKeys = ["nservicebus.queue", "nservicebus.message_type", TenantTag]
})
.AddReader(new BaseExportingMetricReader(new CapturingExporter(exportedMetrics)))
.Build();

await Scenario.Define<Context>()
.WithEndpoint<EndpointWithCustomTags>(b => b.CustomConfig(c => c.MakeInstanceUniquelyAddressable("disc"))
.When(async session =>
{
var sendOptions = new SendOptions();
sendOptions.RouteToThisEndpoint();
sendOptions.SetHeader(TenantTag, "acme-corp");
await session.Send(new MyMessage(), sendOptions);
}))
.Run();

meterProvider.ForceFlush();

metricsListener.AssertTags(TotalFetched, new Dictionary<string, object> { [TenantTag] = "acme-corp" });

metricsListener.AssertTagKeyExists(TotalFetched, EndpointDiscriminatorTag);

var overriddenValue = metricsListener.AssertTagKeyExists(MessageDeserializeTime, EnclosedMessageTypesTag);
Assert.That(overriddenValue, Is.EqualTo(FriendlyMessageTypeName));
}

public class Context : ScenarioContext;

public class EndpointWithCustomTags : EndpointConfigurationBuilder
{
public EndpointWithCustomTags() =>
EndpointSetup<DefaultServer>(c => c.Pipeline.Register(
new CustomizeMetricTagsBehavior(), "Adds a tenant tag from a header and overrides the enclosed message type tag"));

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

class CustomizeMetricTagsBehavior : Behavior<IIncomingPhysicalMessageContext>
{
public override Task Invoke(IIncomingPhysicalMessageContext context, Func<Task> next)
{
var tags = context.MetricTags;

if (context.Message.Headers.TryGetValue(TenantTag, out var tenantId))
{
tags.AddOrOverride(TenantTag, tenantId, TotalFetched);
}

tags.AddOrOverride(EnclosedMessageTypesTag, FriendlyMessageTypeName, MessageDeserializeTime);

return next();
}
}

class CapturingExporter(List<Metric> exportedMetrics) : BaseExporter<Metric>
{
public override ExportResult Export(in Batch<Metric> batch)
{
foreach (var metric in batch)
{
exportedMetrics.Add(metric);
}

return ExportResult.Success;
}
}

public class MyMessage : IMessage;
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.ApplicationInsights.Extensibility;
using NServiceBus;
using NServiceBus.AcceptanceTesting;
using NServiceBus.AcceptanceTests.Core.OpenTelemetry;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ public async Task Should_report_failing_message_metrics()
["nservicebus.discriminator"] = "disc",
["nservicebus.message_type"] = typeof(FailingMessage).FullName,
["execution.result"] = "failure",
["error.type"] = typeof(SimulatedException).FullName,
["error.type"] = typeof(SimulatedException).FullName
});
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
namespace NServiceBus.AcceptanceTests.Core.OpenTelemetry.Metrics;

using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using EndpointTemplates;
using NServiceBus;
using AcceptanceTesting;
using NUnit.Framework;
using Conventions = AcceptanceTesting.Customization.Conventions;

public class When_messages_are_processed_concurrently : OpenTelemetryAcceptanceTest
{
const string ActiveMessagesMetric = "nservicebus.messaging.active_messages";
const int numberOfMessages = 5;

[Test]
public async Task Should_report_active_messages_gauge_that_balances_once_idle()
{
using var metricsListener = TestingMetricListener.SetupNServiceBusMetricsListener();

_ = await Scenario.Define<Context>()
.WithEndpoint<EndpointWithMetrics>(b => b.CustomConfig(c =>
{
c.MakeInstanceUniquelyAddressable("instanceId");
c.LimitMessageProcessingConcurrencyTo(10);
}).When(async (session, ctx) =>
{
for (var x = 0; x < numberOfMessages; x++)
{
await session.SendLocal(new OutgoingMessage());
}
}))
.Run();

Assert.That(metricsListener.ReportedMeters.TryGetValue(ActiveMessagesMetric, out var net), Is.True,
$"'{ActiveMessagesMetric}' gauge should be reported");
Assert.That(net, Is.EqualTo(0),
"increments and decrements should balance once all messages have been processed");

metricsListener.AssertTags(ActiveMessagesMetric,
new Dictionary<string, object>
{
["nservicebus.queue"] = Conventions.EndpointNamingConvention(typeof(EndpointWithMetrics)),
["nservicebus.discriminator"] = "instanceId",
["nservicebus.enclosed_message_types"] = typeof(OutgoingMessage).AssemblyQualifiedName
});
}

public class Context : ScenarioContext
{
public int OutgoingMessagesReceived;
}

public class EndpointWithMetrics : EndpointConfigurationBuilder
{
public EndpointWithMetrics() => EndpointSetup<DefaultServer>();

[Handler]
public class MessageHandler(Context testContext) : IHandleMessages<OutgoingMessage>
{
public Task Handle(OutgoingMessage message, IMessageHandlerContext context)
{
var messagesHandled = Interlocked.Increment(ref testContext.OutgoingMessagesReceived);
testContext.MarkAsCompleted(messagesHandled == numberOfMessages);
return Task.CompletedTask;
}
}
}

public class OutgoingMessage : IMessage;
}
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ public abstract class OpenTelemetryAcceptanceTest : NServiceBusAcceptanceTest
protected TestingActivityListener NServiceBusActivityListener { get; private set; }

[SetUp]
public void Setup() => NServiceBusActivityListener = TestingActivityListener.SetupDiagnosticListener("NServiceBus.Core");
public void Setup() => NServiceBusActivityListener = TestingActivityListener.SetupDiagnosticListener("NServiceBus.Core", "NServiceBus.Core.Recoverability");

[TearDown]
public void Cleanup()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,19 +11,19 @@ public class TestingActivityListener : IDisposable
{
readonly ActivityListener activityListener;

public static TestingActivityListener SetupDiagnosticListener(string sourceName)
public static TestingActivityListener SetupDiagnosticListener(params string[] sourceNames)
{
var testingListener = new TestingActivityListener(sourceName);
var testingListener = new TestingActivityListener(sourceNames);

ActivitySource.AddActivityListener(testingListener.activityListener);
return testingListener;
}

TestingActivityListener(string sourceName = null)
TestingActivityListener(params string[] sourceNames)
{
activityListener = new ActivityListener
{
ShouldListenTo = source => string.IsNullOrEmpty(sourceName) || source.Name == sourceName,
ShouldListenTo = source => sourceNames.Length == 0 || sourceNames.Contains(source.Name),
Sample = (ref ActivityCreationOptions<ActivityContext> _) => ActivitySamplingResult.AllData,
SampleUsingParentId = (ref ActivityCreationOptions<string> options) => ActivitySamplingResult.AllData
};
Expand Down Expand Up @@ -62,4 +62,5 @@ public static List<Activity> GetReceiveMessageActivities(this ConcurrentQueue<Ac
public static List<Activity> GetSendMessageActivities(this ConcurrentQueue<Activity> activities) => activities.Where(a => a.OperationName == "NServiceBus.Diagnostics.SendMessage").ToList();
public static List<Activity> GetPublishEventActivities(this ConcurrentQueue<Activity> activities) => activities.Where(a => a.OperationName == "NServiceBus.Diagnostics.PublishMessage").ToList();
public static List<Activity> GetInvokedHandlerActivities(this ConcurrentQueue<Activity> activities) => activities.Where(a => a.OperationName == "NServiceBus.Diagnostics.InvokeHandler").ToList();
public static List<Activity> GetRecoverabilityActivities(this ConcurrentQueue<Activity> activities) => activities.Where(a => a.OperationName == "NServiceBus.Diagnostics.Recoverability").ToList();
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ public async Task Should_attach_to_ambient_trace()
using var externalActivitySource = new ActivitySource("external trace source");
using var _ = TestingActivityListener.SetupDiagnosticListener(externalActivitySource.Name); // need to have a registered listener for activities to be created

const string wrapperActivityTraceState = "test trace state";
const string wrapperActivityTraceState = "tracekey=traceValue";

var context = await Scenario.Define<Context>()
.WithEndpoint<EndpointWithAmbientActivity>(b => b
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ public async Task Should_propagate_baggage_to_activity()
{
var sendOptions = new SendOptions();
sendOptions.RouteToThisEndpoint();
sendOptions.SetHeader(Headers.DiagnosticsBaggage, "key1=value1,key2=value2,key3=");
sendOptions.SetHeader(Headers.DiagnosticsBaggage, "key1=value1,key2=value2,key3=value3");
await session.Send(new SomeMessage(), sendOptions);
})
)
Expand All @@ -29,7 +29,7 @@ public async Task Should_propagate_baggage_to_activity()

VerifyBaggageItem("key1", "value1");
VerifyBaggageItem("key2", "value2");
VerifyBaggageItem("key3", "");
VerifyBaggageItem("key3", "value3");
return;

void VerifyBaggageItem(string key, string expectedValue)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,9 @@ public async Task Should_propagate_baggage_to_headers()
)
.Run();

// Default (backwards-compatible) propagation produces the legacy comma-separated, percent-encoded format.
// The W3C OWS format ("key3 = , key2 = value2, key1 = value1") is produced only when the
// NServiceBus.Core.OpenTelemetry.UseDistributedContextPropagator AppContext switch is enabled (default in v11).
Assert.That(context.BaggageHeader, Is.EqualTo("key3=,key2=value2,key1=value1"));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,13 @@ public async Task Should_mark_span_as_failed()
handlerActivityTags.VerifyTag("otel.status_code", "ERROR");
handlerActivityTags.VerifyTag("otel.status_description", ErrorMessage);

using (Assert.EnterMultipleScope())
{
Assert.That(failedHandlerActivity.Events, Has.Exactly(1).Items,
"the innermost span (the handler invocation) should record the exception details");
Assert.That(failedPipelineActivity.Events, Is.Empty,
"the outer span should not duplicate the exception details already recorded on the inner span");
}
}

public class Context : ScenarioContext;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
namespace NServiceBus.AcceptanceTests.Core.OpenTelemetry.Traces;

using System.Linq;
using System.Threading.Tasks;
using AcceptanceTesting;
using Configuration.AdvancedExtensibility;
using EndpointTemplates;
using NServiceBus;
using NUnit.Framework;

// The OTEL_SEMCONV_EXCEPTION_SIGNAL_OPT_IN override is applied while the OpenTelemetryFeature defaults
// run, which is AFTER the endpoint's activity factory has already been built from the instrumentation
// options. The opt-in only takes effect when the activity factory and the settings share a single
// InstrumentationOptions instance.
public class When_processing_fails_with_exception_logs_opt_in : OpenTelemetryAcceptanceTest
{
[Test]
public async Task Should_record_the_exception_as_a_log_instead_of_a_span_event()
{
var context = await Scenario.Define<Context>()
.WithEndpoint<FailingEndpoint>(e => e
.DoNotFailOnErrorMessages()
.When(s => s.SendLocal(new FailingMessage())))
.Run();

Assert.That(context.FailedMessages, Has.Count.EqualTo(1), "the message should have failed");

var handlerActivity = NServiceBusActivityListener.CompletedActivities.GetInvokedHandlerActivities().Single();

using (Assert.EnterMultipleScope())
{
Assert.That(handlerActivity.Events, Is.Empty, "the exception should not be recorded as a span event when the endpoint opted in to exceptions as logs");
Assert.That(context.Logs.Any(l => l.LoggerName == "NServiceBus.ActivityFactory" && l.Level == Logging.LogLevel.Error && l.Message.Contains(ErrorMessage)), Is.True, "the exception should be recorded as an error log instead");
}
}

public class Context : ScenarioContext;

public class FailingEndpoint : EndpointConfigurationBuilder
{
// Does not call endpointConfiguration.Tracing(): the instrumentation options only come into existence while the endpoint is being created.
public FailingEndpoint() => EndpointSetup<DefaultServer>(c => c.GetSettings().Set($"ACCEPTANCETEST_ENV:{OptInEnvironmentVariable}", "logs"));

[Handler]
public class FailingMessageHandler(Context testContext) : IHandleMessages<FailingMessage>
{
public Task Handle(FailingMessage message, IMessageHandlerContext context)
{
testContext.MarkAsCompleted();
throw new SimulatedException(ErrorMessage);
}
}
}

public class FailingMessage : IMessage;

const string OptInEnvironmentVariable = "OTEL_SEMCONV_EXCEPTION_SIGNAL_OPT_IN";
const string ErrorMessage = "boom!";
}
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;
}
Loading