From 225b73439c88abce3b0324dc5ac906e4c2de1a8a Mon Sep 17 00:00:00 2001 From: David Erik Jensen Date: Mon, 14 Sep 2026 12:51:08 +0200 Subject: [PATCH] Add TraceCorrelation option to ProcessingOptions Introduce the TraceCorrelation enum (None, Link, Parent) to control how the process activity is correlated with the message's send activity when tracing is enabled: - None (default): no correlation, same as before - Link: the process activity links to the send activity, matching the OpenTelemetry messaging semantic conventions - Parent: the process activity becomes a child of the send activity, so both end up in the same trace, and additionally links to the send activity and to the ambient activity (if any) LinkTraces is kept for backward compatibility and is now a shorthand for TraceCorrelation: true maps to Link and false to None. DotPulsarActivitySource.StartConsumerActivity now takes a TraceCorrelation instead of a bool and passes a parent context to ActivitySource.StartActivity when required. Tests added for both DotPulsarActivitySource and ProcessingOptions, and CHANGELOG updated. --- CHANGELOG.md | 7 + .../Internal/DotPulsarActivitySource.cs | 35 ++-- src/DotPulsar/Internal/MessageProcessor.cs | 6 +- src/DotPulsar/ProcessingOptions.cs | 24 ++- src/DotPulsar/TraceCorrelation.cs | 38 ++++ .../Internal/DotPulsarActivitySourceTests.cs | 177 ++++++++++++++++++ .../DotPulsar.Tests/ProcessingOptionsTests.cs | 85 +++++++++ 7 files changed, 354 insertions(+), 18 deletions(-) create mode 100644 src/DotPulsar/TraceCorrelation.cs create mode 100644 tests/DotPulsar.Tests/Internal/DotPulsarActivitySourceTests.cs create mode 100644 tests/DotPulsar.Tests/ProcessingOptionsTests.cs diff --git a/CHANGELOG.md b/CHANGELOG.md index e476c131e..d180b8d4a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,13 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/) ## [Unreleased] +### Added + +- 'TraceCorrelation' on ProcessingOptions for choosing how the process trace is correlated with the message's send trace, if tracing is enabled + - 'None' (the default): no correlation + - 'Link': the process trace links to the send trace (equivalent to setting 'LinkTraces' to 'true') + - 'Parent': the process trace is a child of the send trace and links to both the send trace and the ambient trace (if any), following the opt-in behavior described in the OpenTelemetry semantic conventions for messaging + ### Changed - Updated the Google.Protobuf dependency from version 3.36.0 to 3.36.1 diff --git a/src/DotPulsar/Internal/DotPulsarActivitySource.cs b/src/DotPulsar/Internal/DotPulsarActivitySource.cs index 5e7f25252..eb8cee2f9 100644 --- a/src/DotPulsar/Internal/DotPulsarActivitySource.cs +++ b/src/DotPulsar/Internal/DotPulsarActivitySource.cs @@ -28,21 +28,33 @@ static DotPulsarActivitySource() public static ActivitySource ActivitySource { get; } - public static Activity? StartConsumerActivity(IMessage message, string operationName, KeyValuePair[] tags, bool linkTraces) + public static Activity? StartConsumerActivity(IMessage message, string operationName, KeyValuePair[] tags, TraceCorrelation traceCorrelation) { if (!ActivitySource.HasListeners()) return null; - IEnumerable? activityLinks = null; + ActivityContext parentContext = default; + List? activityLinks = null; - if (linkTraces) + if (traceCorrelation != TraceCorrelation.None) { - var activityLink = GetActivityLink(message); - if (activityLink is not null) - activityLinks = [activityLink.Value]; + var creationContext = GetCreationContext(message); + if (creationContext is not null) + { + activityLinks = [new ActivityLink(creationContext.Value)]; + + if (traceCorrelation == TraceCorrelation.Parent) + { + parentContext = creationContext.Value; + + var ambientContext = Activity.Current?.Context; + if (ambientContext is not null && ambientContext.Value != default) + activityLinks.Add(new ActivityLink(ambientContext.Value)); + } + } } - return StartActivity(operationName, ActivityKind.Consumer, tags, activityLinks, message.GetConversationId()); + return StartActivity(operationName, ActivityKind.Consumer, tags, activityLinks, message.GetConversationId(), parentContext); } public static Activity? StartProducerActivity(MessageMetadata metadata, string operationName, KeyValuePair[] tags) @@ -53,14 +65,14 @@ static DotPulsarActivitySource() return StartActivity(operationName, ActivityKind.Producer, tags, null, metadata.GetConversationId()); } - private static ActivityLink? GetActivityLink(IMessage message) + private static ActivityContext? GetCreationContext(IMessage message) { if (message.Properties.TryGetValue(Constants.TraceParent, out var traceParent)) { _ = message.Properties.TryGetValue(Constants.TraceState, out var traceState); if (ActivityContext.TryParse(traceParent, traceState, out var context)) - return new ActivityLink(context); + return context; } return null; @@ -71,9 +83,10 @@ static DotPulsarActivitySource() ActivityKind kind, KeyValuePair[] tags, IEnumerable? activityLinks, - string? conversationId) + string? conversationId, + ActivityContext parentContext = default) { - var activity = ActivitySource.StartActivity(kind, name: operationName, tags: tags, links: activityLinks); + var activity = ActivitySource.StartActivity(kind, parentContext: parentContext, name: operationName, tags: tags, links: activityLinks); if (activity is not null && activity.IsAllDataRequested) { diff --git a/src/DotPulsar/Internal/MessageProcessor.cs b/src/DotPulsar/Internal/MessageProcessor.cs index 3ce7492d3..673237f10 100644 --- a/src/DotPulsar/Internal/MessageProcessor.cs +++ b/src/DotPulsar/Internal/MessageProcessor.cs @@ -34,7 +34,7 @@ public sealed class MessageProcessor : IDisposable private readonly SemaphoreSlim _receiveLock; private readonly SemaphoreSlim _acknowledgeLock; private readonly ObjectPool _processInfoPool; - private readonly bool _linkTraces; + private readonly TraceCorrelation _traceCorrelation; private readonly bool _ensureOrderedAcknowledgment; private readonly int _maxDegreeOfParallelism; private readonly int _maxMessagesPerTask; @@ -78,7 +78,7 @@ public MessageProcessor( _acknowledgeLock = new SemaphoreSlim(1, 1); _processInfoPool = new DefaultObjectPool(new DefaultPooledObjectPolicy()); - _linkTraces = options.LinkTraces; + _traceCorrelation = options.TraceCorrelation; _ensureOrderedAcknowledgment = options.EnsureOrderedAcknowledgment; _maxDegreeOfParallelism = options.MaxDegreeOfParallelism; _maxMessagesPerTask = options.MaxMessagesPerTask; @@ -151,7 +151,7 @@ private async ValueTask Processor(CancellationToken cancellationToken) _receiveLock.Release(); } - var activity = DotPulsarActivitySource.StartConsumerActivity(message, _operationName, _activityTags, _linkTraces); + var activity = DotPulsarActivitySource.StartConsumerActivity(message, _operationName, _activityTags, _traceCorrelation); if (activity is not null && activity.IsAllDataRequested) { activity.SetMessageId(message.MessageId); diff --git a/src/DotPulsar/ProcessingOptions.cs b/src/DotPulsar/ProcessingOptions.cs index b1a06cc68..3bf39dea1 100644 --- a/src/DotPulsar/ProcessingOptions.cs +++ b/src/DotPulsar/ProcessingOptions.cs @@ -25,11 +25,11 @@ public sealed class ProcessingOptions public const int Unbounded = -1; private bool _ensureOrderedAcknowledgment; - private bool _linkTraces; private int _maxDegreeOfParallelism; private int _maxMessagesPerTask; private TimeSpan _shutdownGracePeriod; private TaskScheduler _taskScheduler; + private TraceCorrelation _traceCorrelation; /// /// Initializes a new instance with the default values. @@ -37,11 +37,11 @@ public sealed class ProcessingOptions public ProcessingOptions() { _ensureOrderedAcknowledgment = true; - _linkTraces = false; _maxDegreeOfParallelism = 1; _maxMessagesPerTask = Unbounded; _shutdownGracePeriod = TimeSpan.Zero; _taskScheduler = TaskScheduler.Default; + _traceCorrelation = TraceCorrelation.None; } /// @@ -55,11 +55,12 @@ public bool EnsureOrderedAcknowledgment /// /// Whether to link the process trace to the message's send trace, if tracing is enabled. The default is 'false'. + /// This is a shorthand for : 'true' corresponds to and 'false' to . /// public bool LinkTraces { - get => _linkTraces; - set { _linkTraces = value; } + get => _traceCorrelation == TraceCorrelation.Link; + set { _traceCorrelation = value ? TraceCorrelation.Link : TraceCorrelation.None; } } /// @@ -120,4 +121,19 @@ public TaskScheduler TaskScheduler _taskScheduler = value; } } + + /// + /// How the process trace is correlated with the message's send trace, if tracing is enabled. The default is 'None'. + /// + public TraceCorrelation TraceCorrelation + { + get => _traceCorrelation; + set + { + if (!Enum.IsDefined(typeof(TraceCorrelation), value)) + throw new ArgumentOutOfRangeException(nameof(value)); + + _traceCorrelation = value; + } + } } diff --git a/src/DotPulsar/TraceCorrelation.cs b/src/DotPulsar/TraceCorrelation.cs new file mode 100644 index 000000000..14b2524ee --- /dev/null +++ b/src/DotPulsar/TraceCorrelation.cs @@ -0,0 +1,38 @@ +/* + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +namespace DotPulsar; + +/// +/// How the process activity is correlated with the message's send activity, if tracing is enabled. +/// +public enum TraceCorrelation : byte +{ + /// + /// The process activity is not correlated with the message's send activity. + /// + None = 0, + + /// + /// The process activity links to the message's send activity. The process activity is a child of the ambient activity (if any). + /// This is the correlation recommended by the OpenTelemetry semantic conventions for messaging. + /// + Link = 1, + + /// + /// The process activity is a child of the message's send activity, so both end up in the same trace. + /// The process activity also links to the message's send activity and to the ambient activity (if any). + /// + Parent = 2 +} diff --git a/tests/DotPulsar.Tests/Internal/DotPulsarActivitySourceTests.cs b/tests/DotPulsar.Tests/Internal/DotPulsarActivitySourceTests.cs new file mode 100644 index 000000000..290bc34f0 --- /dev/null +++ b/tests/DotPulsar.Tests/Internal/DotPulsarActivitySourceTests.cs @@ -0,0 +1,177 @@ +/* + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +namespace DotPulsar.Tests.Internal; + +using DotPulsar.Abstractions; +using DotPulsar.Internal; +using System.Diagnostics; + +[Trait("Category", "Unit")] +public sealed class DotPulsarActivitySourceTests : IDisposable +{ + private const string OperationName = "test process"; + private static readonly KeyValuePair[] _tags = []; + + private readonly ActivityListener _listener; + private readonly ActivityTraceId _traceId; + private readonly ActivitySpanId _spanId; + private readonly IMessage _message; + + public DotPulsarActivitySourceTests() + { + var activitySource = DotPulsarActivitySource.ActivitySource; + _listener = new ActivityListener + { + ShouldListenTo = source => ReferenceEquals(source, activitySource), + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded + }; + ActivitySource.AddActivityListener(_listener); + + _traceId = ActivityTraceId.CreateRandom(); + _spanId = ActivitySpanId.CreateRandom(); + + _message = Substitute.For(); + _message.Properties.Returns(new Dictionary + { + [Constants.TraceParent] = $"00-{_traceId.ToHexString()}-{_spanId.ToHexString()}-01", + [Constants.TraceState] = "vendor=value" + }); + + Activity.Current = null; + } + + [Fact] + public void StartConsumerActivity_GivenNone_ShouldNotCorrelate() + { + //Act + using var activity = DotPulsarActivitySource.StartConsumerActivity(_message, OperationName, _tags, TraceCorrelation.None); + + //Assert + activity.ShouldNotBeNull(); + activity.TraceId.ShouldNotBe(_traceId); + activity.ParentId.ShouldBeNull(); + activity.Links.ShouldBeEmpty(); + } + + [Fact] + public void StartConsumerActivity_GivenLink_ShouldLinkToCreationContextAndNotUseItAsParent() + { + //Act + using var activity = DotPulsarActivitySource.StartConsumerActivity(_message, OperationName, _tags, TraceCorrelation.Link); + + //Assert + activity.ShouldNotBeNull(); + activity.TraceId.ShouldNotBe(_traceId); + activity.ParentId.ShouldBeNull(); + var link = activity.Links.ShouldHaveSingleItem(); + link.Context.TraceId.ShouldBe(_traceId); + link.Context.SpanId.ShouldBe(_spanId); + link.Context.TraceState.ShouldBe("vendor=value"); + } + + [Fact] + public void StartConsumerActivity_GivenLinkAndAmbientActivity_ShouldBeChildOfAmbientActivity() + { + //Arrange + using var ambient = new Activity("ambient").Start(); + + //Act + using var activity = DotPulsarActivitySource.StartConsumerActivity(_message, OperationName, _tags, TraceCorrelation.Link); + + //Assert + activity.ShouldNotBeNull(); + activity.TraceId.ShouldBe(ambient.TraceId); + activity.ParentSpanId.ShouldBe(ambient.SpanId); + var link = activity.Links.ShouldHaveSingleItem(); + link.Context.TraceId.ShouldBe(_traceId); + link.Context.SpanId.ShouldBe(_spanId); + } + + [Fact] + public void StartConsumerActivity_GivenParent_ShouldUseCreationContextAsParentAndLinkToIt() + { + //Act + using var activity = DotPulsarActivitySource.StartConsumerActivity(_message, OperationName, _tags, TraceCorrelation.Parent); + + //Assert + activity.ShouldNotBeNull(); + activity.TraceId.ShouldBe(_traceId); + activity.ParentSpanId.ShouldBe(_spanId); + activity.TraceStateString.ShouldBe("vendor=value"); + var link = activity.Links.ShouldHaveSingleItem(); + link.Context.TraceId.ShouldBe(_traceId); + link.Context.SpanId.ShouldBe(_spanId); + } + + [Fact] + public void StartConsumerActivity_GivenParentAndAmbientActivity_ShouldLinkToCreationContextAndAmbientActivity() + { + //Arrange + using var ambient = new Activity("ambient").Start(); + + //Act + using var activity = DotPulsarActivitySource.StartConsumerActivity(_message, OperationName, _tags, TraceCorrelation.Parent); + + //Assert + activity.ShouldNotBeNull(); + activity.TraceId.ShouldBe(_traceId); + activity.ParentSpanId.ShouldBe(_spanId); + var links = activity.Links.ToList(); + links.Count.ShouldBe(2); + links.ShouldContain(link => link.Context.TraceId == _traceId && link.Context.SpanId == _spanId); + links.ShouldContain(link => link.Context.TraceId == ambient.TraceId && link.Context.SpanId == ambient.SpanId); + } + + [Theory] + [InlineData(TraceCorrelation.Link)] + [InlineData(TraceCorrelation.Parent)] + public void StartConsumerActivity_GivenNoCreationContextInMessage_ShouldNotCorrelate(TraceCorrelation traceCorrelation) + { + //Arrange + _message.Properties.Returns(new Dictionary()); + + //Act + using var activity = DotPulsarActivitySource.StartConsumerActivity(_message, OperationName, _tags, traceCorrelation); + + //Assert + activity.ShouldNotBeNull(); + activity.TraceId.ShouldNotBe(_traceId); + activity.ParentId.ShouldBeNull(); + activity.Links.ShouldBeEmpty(); + } + + [Theory] + [InlineData(TraceCorrelation.Link)] + [InlineData(TraceCorrelation.Parent)] + public void StartConsumerActivity_GivenInvalidTraceParentInMessage_ShouldNotCorrelate(TraceCorrelation traceCorrelation) + { + //Arrange + _message.Properties.Returns(new Dictionary { [Constants.TraceParent] = "not-a-traceparent" }); + + //Act + using var activity = DotPulsarActivitySource.StartConsumerActivity(_message, OperationName, _tags, traceCorrelation); + + //Assert + activity.ShouldNotBeNull(); + activity.ParentId.ShouldBeNull(); + activity.Links.ShouldBeEmpty(); + } + + public void Dispose() + { + Activity.Current = null; + _listener.Dispose(); + } +} diff --git a/tests/DotPulsar.Tests/ProcessingOptionsTests.cs b/tests/DotPulsar.Tests/ProcessingOptionsTests.cs new file mode 100644 index 000000000..6021fe2f3 --- /dev/null +++ b/tests/DotPulsar.Tests/ProcessingOptionsTests.cs @@ -0,0 +1,85 @@ +/* + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +namespace DotPulsar.Tests; + +[Trait("Category", "Unit")] +public sealed class ProcessingOptionsTests +{ + [Fact] + public void Constructor_GivenDefaults_ShouldNotCorrelateTraces() + { + //Act + var options = new ProcessingOptions(); + + //Assert + options.TraceCorrelation.ShouldBe(TraceCorrelation.None); + options.LinkTraces.ShouldBeFalse(); + } + + [Fact] + public void LinkTraces_GivenTrue_ShouldSetTraceCorrelationToLink() + { + //Arrange + var options = new ProcessingOptions(); + + //Act + options.LinkTraces = true; + + //Assert + options.TraceCorrelation.ShouldBe(TraceCorrelation.Link); + options.LinkTraces.ShouldBeTrue(); + } + + [Fact] + public void LinkTraces_GivenFalse_ShouldSetTraceCorrelationToNone() + { + //Arrange + var options = new ProcessingOptions { TraceCorrelation = TraceCorrelation.Parent }; + + //Act + options.LinkTraces = false; + + //Assert + options.TraceCorrelation.ShouldBe(TraceCorrelation.None); + options.LinkTraces.ShouldBeFalse(); + } + + [Fact] + public void TraceCorrelation_GivenParent_ShouldNotReportLinkTraces() + { + //Arrange + var options = new ProcessingOptions(); + + //Act + options.TraceCorrelation = TraceCorrelation.Parent; + + //Assert + options.TraceCorrelation.ShouldBe(TraceCorrelation.Parent); + options.LinkTraces.ShouldBeFalse(); + } + + [Fact] + public void TraceCorrelation_GivenUndefinedValue_ShouldThrowArgumentOutOfRangeException() + { + //Arrange + var options = new ProcessingOptions(); + + //Act + var exception = Record.Exception(() => options.TraceCorrelation = (TraceCorrelation) 42); + + //Assert + exception.ShouldBeOfType(); + } +}