Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -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<Context>()
.WithEndpoint<FailingEndpointWithChildSpanConnector>(e => e
.When(s => s.SendLocal(new FailingMessage()))
.DoNotFailOnErrorMessages())
.WithEndpoint<ErrorSpy>()
.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<DefaultServer>(c =>
{
c.SendFailedMessagesTo(ErrorQueueAddress);
c.Tracing().ErrorMessageTraceMode = TraceMode.ContinueExisting;
});

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

public class Context : ScenarioContext
{
public Dictionary<string, string> ErrorMessageHeaders { get; set; }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,162 @@ await Scenario.Define<SagaContext>()
}
}

[Test]
public async Task By_sendoptions_Should_continue_trace_when_endpoint_connector_is_child_span()
{
await Scenario.Define<Context>()
.WithEndpoint<ChildSpanConnectorEndpoint>(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<Context>()
.WithEndpoint<ChildSpanConnectorEndpoint>(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<Context>()
.WithEndpoint<ChildSpanRetryConnectorEndpoint>(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(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");
}

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<DelayedMessage>
{
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<MessageToBeRetried>
{
public Task Handle(MessageToBeRetried message, IMessageHandlerContext context) => throw new SimulatedException();
}
}

public class Context : ScenarioContext
{
public bool ReplyMessageReceived { get; set; }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Context>()
.WithEndpoint<PublisherWithChildSpanConnector>(b => b
.When(ctx => ctx.SomeEventSubscribed, s => s.Publish(new ThisIsAnEvent())))
.WithEndpoint<SubscriberForPublisherWithChildSpanConnector>(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<Context>()
.WithEndpoint<PublisherWithChildSpanConnector>(b => b
.When(ctx => ctx.SomeEventSubscribed, s =>
{
var publishOptions = new PublishOptions();
publishOptions.StartNewTraceOnReceive();
return s.Publish(new ThisIsAnEvent(), publishOptions);
}))
.WithEndpoint<SubscriberForPublisherWithChildSpanConnector>(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<DefaultServer>(b =>
{
b.Tracing().PublishTraceMode = TraceMode.ContinueExisting;
b.OnEndpointSubscribed<Context>((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<DefaultServer>(c => { },
metadata =>
{
metadata.RegisterPublisherFor<ThisIsAnEvent, PublisherWithChildSpanConnector>();
});

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

public class ThisIsAnEvent : IEvent;
}
Loading