Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
10 changes: 7 additions & 3 deletions .github/workflows/dotnet.yml
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,17 @@ jobs:

steps:
- name: Checkout
uses: actions/checkout@v2
uses: actions/checkout@v6
- name: Setup .NET 10
uses: actions/setup-dotnet@v5
with:
dotnet-version: '10.0.x'
- name: Setup .NET 7
uses: actions/setup-dotnet@v1
uses: actions/setup-dotnet@v5
with:
dotnet-version: '7.0.x'
- name: Setup .NET 6
uses: actions/setup-dotnet@v1
uses: actions/setup-dotnet@v5
with:
dotnet-version: '6.0.x'
- name: Restore
Expand Down
39 changes: 35 additions & 4 deletions Rebus.Diagnostics.Tests/Incoming/IncomingDiagnosticsStepTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ public async Task StartsActivityWhenTraceStateHeaderIsSet()
using var scope = new RebusTransactionScope();
var context = new IncomingStepContext(transportMessage, scope.TransactionContext);

var step = new IncomingDiagnosticsStep();
var step = new IncomingDiagnosticsStep(new TestLogger());
var callbackInvoked = false;
await step.Process(context, () =>
{
Expand Down Expand Up @@ -87,7 +87,7 @@ public async Task StartsAnEntireNewActivityIfNoActivityIsCurrentlyActive()
using var scope = new RebusTransactionScope();
var context = new IncomingStepContext(transportMessage, scope.TransactionContext);

var step = new IncomingDiagnosticsStep();
var step = new IncomingDiagnosticsStep(new TestLogger());
var callbackInvoked = false;
await step.Process(context, () =>
{
Expand All @@ -100,6 +100,37 @@ await step.Process(context, () =>
Assert.That(callbackInvoked);
}

[Test]
public async Task LogsWarningAndContinuesWhenBaggageHeaderIsInvalid()
{
var logger = new TestLogger();
var headers = new Dictionary<string, string>
{
{Headers.Type, "MyType"},
{Headers.Intent, Headers.IntentOptions.PublishSubscribe},
{Headers.MessageId, "MyMessage"},
{RebusDiagnosticConstants.BaggageHeaderName, "this-is-not-valid-json"}
};

var transportMessage = new TransportMessage(headers, Array.Empty<byte>());

Activity.DefaultIdFormat = ActivityIdFormat.W3C;
using var scope = new RebusTransactionScope();
var context = new IncomingStepContext(transportMessage, scope.TransactionContext);

var step = new IncomingDiagnosticsStep(logger);
var callbackInvoked = false;

await step.Process(context, () =>
{
callbackInvoked = true;
return Task.CompletedTask;
});

Assert.That(callbackInvoked, Is.True);
Assert.That(logger.TraceIds, Has.Some.Matches<(LogLevel Level, string Message, string TraceId)>(t => t.Level == LogLevel.Warn));
}

[Test]
public async Task ActivityIsAvailableInRetryLogging()
{
Expand Down Expand Up @@ -174,7 +205,7 @@ public async Task DefaultsWhenIntentIsNotSet()
using var scope = new RebusTransactionScope();
var context = new IncomingStepContext(transportMessage, scope.TransactionContext);

var step = new IncomingDiagnosticsStep();
var step = new IncomingDiagnosticsStep(new TestLogger());
var callbackInvoked = false;
await step.Process(context, () =>
{
Expand Down Expand Up @@ -203,7 +234,7 @@ public async Task StartsAnEntireNewActivsssityIfNoActivityIsCurrentlyActive()
using var scope = new RebusTransactionScope();
var context = new IncomingStepContext(transportMessage, scope.TransactionContext);

var step = new IncomingDiagnosticsStep();
var step = new IncomingDiagnosticsStep(new TestLogger());
await step.Process(context, () => Task.CompletedTask);

Assert.That(meterObserver.InstrumentCalled("incoming", RebusDiagnosticConstants.MessageCountMeterNameTemplate), Is.True);
Expand Down
2 changes: 1 addition & 1 deletion Rebus.Diagnostics.Tests/Rebus.Diagnostics.Tests.csproj
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<TargetFrameworks>net7.0;net6.0;net48</TargetFrameworks>
<TargetFrameworks>net10.0;net7.0;net6.0;net48</TargetFrameworks>

<IsPackable>false</IsPackable>
<LangVersion>9</LangVersion>
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using System;
using Rebus.Diagnostics.Incoming;
using Rebus.Diagnostics.Outgoing;
using Rebus.Logging;
using Rebus.Pipeline;
using Rebus.Pipeline.Receive;
using Rebus.Pipeline.Send;
Expand All @@ -21,8 +22,11 @@ public static OptionsConfigurer EnableDiagnosticSources(this OptionsConfigurer c
var outgoingStep = new OutgoingDiagnosticsStep();
injector.OnSend(outgoingStep, PipelineRelativePosition.Before,
typeof(SendOutgoingMessageStep));

var incomingStep = new IncomingDiagnosticsStep();


var rebusLoggerFactory = c.Get<IRebusLoggerFactory>();

var incomingStep = new IncomingDiagnosticsStep(rebusLoggerFactory.GetLogger<IncomingDiagnosticsStep>());

var invokerWrapper = new IncomingDiagnosticsHandlerInvokerWrapper();
injector.OnReceive(invokerWrapper, PipelineRelativePosition.After, typeof(ActivateHandlersStep));
Expand Down
107 changes: 66 additions & 41 deletions Rebus.Diagnostics/Diagnostics/Incoming/IncomingDiagnosticsStep.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
using Rebus.Bus;
using Rebus.Diagnostics.Helpers;
using Rebus.Diagnostics.Outgoing;
using Rebus.Logging;
using Rebus.Messages;
using Rebus.Pipeline;

Expand All @@ -14,12 +15,15 @@ namespace Rebus.Diagnostics.Incoming
[StepDocumentation("Extracts trace from the incoming message and starts an activity for it")]
public class IncomingDiagnosticsStep : IIncomingStep
{
private readonly ILog _log;

private static readonly DiagnosticSource DiagnosticListener =
new DiagnosticListener(RebusDiagnosticConstants.ConsumerActivityName);
private readonly StepMeter _stepMeter;

public IncomingDiagnosticsStep()
public IncomingDiagnosticsStep(ILog log)
{
_log = log;
_stepMeter = new StepMeter("incoming");
}

Expand All @@ -41,59 +45,80 @@ public async Task Process(IncomingStepContext context, Func<Task> next)
}
}

private static Activity? StartActivity(IncomingStepContext context, TransportMessage message)
private Activity? StartActivity(IncomingStepContext context, TransportMessage message)
{
Activity? activity = null;
if (RebusDiagnosticConstants.ActivitySource.HasListeners())
try
{
var headers = message.Headers;

var messageType = message.GetMessageType();

var messageWrapper = new TransportMessageWrapper(message);

var initialTags = TagHelper.ExtractInitialTags(messageWrapper);
initialTags.Add("messaging.operation", "receive");

var activityKind = messageWrapper.GetIntentOption() == Headers.IntentOptions.PublishSubscribe
? ActivityKind.Consumer
: ActivityKind.Server;

var activityName = $"{messageType} receive";
if (!headers.TryGetValue(RebusDiagnosticConstants.TraceStateHeaderName, out var traceState))
{
activity = RebusDiagnosticConstants.ActivitySource.StartActivity(activityName, activityKind, default(ActivityContext), initialTags);
}
else
Activity? activity = null;
if (RebusDiagnosticConstants.ActivitySource.HasListeners())
{
activity = RebusDiagnosticConstants.ActivitySource.StartActivity(activityName, activityKind,
traceState, initialTags);
var headers = message.Headers;

var messageType = message.GetMessageType();

var messageWrapper = new TransportMessageWrapper(message);

var initialTags = TagHelper.ExtractInitialTags(messageWrapper);
initialTags.Add("messaging.operation", "receive");

var activityKind = messageWrapper.GetIntentOption() == Headers.IntentOptions.PublishSubscribe
? ActivityKind.Consumer
: ActivityKind.Server;

var activityName = $"{messageType} receive";
if (!headers.TryGetValue(RebusDiagnosticConstants.TraceStateHeaderName, out var traceState))
{
activity = RebusDiagnosticConstants.ActivitySource.StartActivity(activityName, activityKind,
default(ActivityContext), initialTags);
}
else
{
activity = RebusDiagnosticConstants.ActivitySource.StartActivity(activityName, activityKind,
traceState, initialTags);
}

if (activity != null)
{
CopyBaggage(headers, activity);
}

// TODO: Not sure if this is still needed
// DiagnosticListener.OnActivityImport(activity, context);
}

if (activity != null)
{
CopyBaggage(headers, activity);
}
SendBeforeProcessEvent(context, activity);

// TODO: Not sure if this is still needed
// DiagnosticListener.OnActivityImport(activity, context);
return activity;
}
catch (Exception e)
{
_log.Warn(e, "Failed to start message activity. Continuing without");
return null;
}

SendBeforeProcessEvent(context, activity);

return activity;
}

private static void CopyBaggage(Dictionary<string, string> headers, Activity activity)
private void CopyBaggage(Dictionary<string, string> headers, Activity activity)
{
if (headers.TryGetValue(RebusDiagnosticConstants.BaggageHeaderName, out var baggageContent))
{
var baggage =
JsonConvert.DeserializeObject<IEnumerable<KeyValuePair<string, string>>>(baggageContent);

foreach (var keyValuePair in baggage)
try
{
var baggage =
JsonConvert.DeserializeObject<IEnumerable<KeyValuePair<string, string>>>(baggageContent);

if (baggage == null)
{
return;
}

foreach (var keyValuePair in baggage)
{
activity.AddBaggage(keyValuePair.Key, keyValuePair.Value);
}
}
catch (Exception e)
{
activity.AddBaggage(keyValuePair.Key, keyValuePair.Value);
_log.Warn(e, "Failed to process activity baggage: {0}", baggageContent);
}
}
}
Expand Down
Loading