Sitelet https://github.com/Azure/azure-service-bus-dotnet/pull/370
Skip to content
This repository was archived by the owner on Oct 12, 2023. It is now read-only.

Instruments operations on ServiceBus with DiagnosticsSource and introduces correlation between producer and consumer - #370

Merged
Neeraj Makam (nemakam) merged 9 commits into
Azure:devfrom
lmolkova:dev
Nov 8, 2017
Merged

Neeraj Makam (nemakam) merged 9 commits into
Azure:devfrom
lmolkova:dev

Conversation

@lmolkova

@lmolkova Liudmila Molkova (lmolkova) commented Nov 1, 2017 •

Copy link
Copy Markdown

This change adds instrumentation with DiagnosticSource for public ServiceBus APIs:

  • MessageSender (i.e. QueueClient and TopicClient): SendAsync, ScueduleMessage and Cancel
  • MessageReceiver: Receive, ReceiveDeffered, Peek, Abandon, Defer, Complete, DeadLetter, RenewLock
  • MessageSession: AcceptMessageSession
  • SubscriptionClient : AddRule, RemoveRule, GetRules
  • Message and session handlers

Tracing system (e.g. ApplicationInsights) may subscribe to ServiceBus DiagnosticSource and receive events about these operations including tracing context needed to correlate events and log any important information.

In absence of tracing system, diagnostics is disabled and performance cost of having it is ~zero.

We also introduce here telemetry correlation between producer and consumer: when diagnostics is enabled on a consumer, tracing context is injected into the message. When a producer receives such message, it extracts and restores context making it available for tracing system.

This also adds Message.ExtractActivity() public extension method that allows extracting the context from a message that is useful when handler API is NOT used to process messages.

Note: this change introduces 'standard' fields to pass tracing context through the queues to correlate telemetry:
Diagnostic-Id - uniquely identifies operation that sent message(s) to the queue
Correlation-Context - optional extended context (empty by default)

This is pretty similar to HTTP Correlation protocol implemented in .NET HTTP stack.

Users are encouraged to use these properties on all platforms, not only .NET.

As a result of this change, a tracing system may enable diagnostics information and receive events it's interested in making sure tracing context is propagated through the service bus, without ANY changes in user code (except for non-handler processing).

Test coverage for the changes: 62.05 %


Instrumentation approach backgound:

Neeraj Makam (@nemakam) Tomasz Milos (@TomMilos) brahmnes

@codecov-io

Codecov (codecov-io) commented Nov 2, 2017 •

Copy link
Copy Markdown

Codecov Report

Merging #370 into dev will increase coverage by 2.54%.
The diff coverage is 72.24%.

Impacted file tree graph

@@            Coverage Diff             @@
##              dev     #370      +/-   ##
==========================================
+ Coverage   58.84%   61.38%   +2.54%     
==========================================
  Files          87       89       +2     
  Lines        6150     7013     +863     
  Branches      751     1002     +251     
==========================================
+ Hits         3619     4305     +686     
- Misses       2136     2163      +27     
- Partials      395      545     +150
Impacted Files Coverage Δ
src/Microsoft.Azure.ServiceBus/SessionPumpHost.cs 70.58% <100%> (+1.96%) ⬆️
src/Microsoft.Azure.ServiceBus/SessionClient.cs 71.83% <100%> (-4.92%) ⬇️
src/Microsoft.Azure.ServiceBus/QueueClient.cs 76.11% <100%> (+1.49%) ⬆️
...c/Microsoft.Azure.ServiceBus/SessionReceivePump.cs 64.44% <50%> (+4.34%) ⬆️
src/Microsoft.Azure.ServiceBus/MessageSession.cs 72.86% <58.33%> (-9.28%) ⬇️
...ft.Azure.ServiceBus/ServiceBusDiagnosticsSource.cs 70.81% <70.81%> (ø)
...c/Microsoft.Azure.ServiceBus/Core/MessageSender.cs 76.68% <71.05%> (+0.53%) ⬆️
...Microsoft.Azure.ServiceBus/Core/MessageReceiver.cs 74.11% <74.35%> (+0.96%) ⬆️
...c/Microsoft.Azure.ServiceBus/SubscriptionClient.cs 63.2% <75%> (+2.64%) ⬆️
...viceBus/Extensions/MessageDiagnosticsExtensions.cs 86.15% <86.15%> (ø)
... and 19 more

Continue to review full report at Codecov.

Legend - Click here to learn more
Δ = absolute <relative> (impact), ø = not affected, ? = missing data
Powered by Codecov. Last update 4156114...05303a9. Read the comment docs.


public const string ExceptionEventName = "Microsoft.Azure.ServiceBus.Exception";
public const string ProcessActivityName = "Microsoft.Azure.ServiceBus.Process";
public const string ProcessActivityStartName = "Microsoft.Azure.ServiceBus.Process.Start";

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ProcessActivityStartName, ProcessSessionActivityName, and ProcessSessionActivityStartName don't seem to be used.

null);
}

internal void ReceiveStop(Activity activity, int messageCount, TaskStatus? status, IList<Message> messageList)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

messageCount is problematic. It represents the maximum number of messages ASB client might fetch, but not necessarily the actual number. The actul number of retrieved messages should be retrieved from messageList.Length to be accurate.

@lmolkova Liudmila Molkova (lmolkova) Nov 6, 2017 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That's a good point. However image you are logging this event:
Requested {messageCount}, received {messageList.Count} in a {Activity.Current.Duration}.
I think it's useful to know how much did you request and how much did you receive in response.

As you notice, every operation has a Start and Stop events and messageCount really makes sense in Start event payload ('Requesting {messageCount} messages'). However Start event is not really interesting, at least is not as interesting as Stop event (when you know the result and duration). So we expect the majority of listeners to care about Stop events only (and reduce verbosity by disabling Start events), so we try to provide full context about every operation in Stop.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. Therefore MessageCount should not be assigned the value of messageCount but messageList.Length instead.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm sorry, I did not understand it.
If I never received Start event (and this is a default case) and on Stop event I want to log
Requested {messageCount}, received {messageList.Count} in a {Activity.Current.Duration}.
How would I do that if messageCount is messageList.Count?

Would you be ok with renaming MessageCount to e.g. RequestedMessageCount?


#endregion

internal void ReportException(Exception ex)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

internal void ReportException(Exception ex) [](start = 8, length = 43)

Is there a paradigm of ReportException being associated with an Activity?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It's useful to track exceptions automatically. E.g. you may even want to track exceptions only.
So yes, we implement events for exceptions in all libraries instrumented with DiagnosticSource. Do you have any objections/concerns about it?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What I meant was, we are just reporting exception here without associating it with an activity. Shouldn't we be passing an activity?
We associate all starts and stops of operations across machines. Will we be able to associate the exception with a correlationId?


In reply to: 149178343 [](ancestors = 149178343)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Liudmila Molkova (@lmolkova) What I meant was, we are just reporting exception here without associating it with an activity. Shouldn't we be passing an activity?
We associate all starts and stops of operations across machines. Will we be able to associate the exception with a correlationId?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Activity has static Current property that flows with the async calls, i.e. users can always use Activity.Current, it does not have to be explicitly passed to the event

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

e.g. in exception event handler you can write

Logger.LogError($"Exception {ex}, Id: {Activity.Current?.Id}")

Endpoint = this.endpoint
},
a => a.AddTag("SessionId", sessionId));
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

How good are these with null params? Couple of methods here have good chance of being passed a null param (like this one). What happens in that case?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nothing would break if there is a null value of sessionId, but I'd still prefer not to add null to Tags (just to optimize amount of data).
I'll add null check then, thanks

}
finally
{
this.diagnosticSource.AcceptMessageSessionStop(activity, sessionId, serverWaitTime, session, acceptMessageSessionTask?.Status);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

session [](start = 100, length = 7)

we don't need to pass the session object. Don't think its very useful.
If you are removing this, make sure instead of passing sessionId, pass session.SessionId

this.ReceiveMode = receiveMode;
this.TokenProvider = this.ServiceBusConnection.CreateTokenProvider();
this.CbsTokenProvider = new TokenProviderAdapter(this.TokenProvider, serviceBusConnection.OperationTimeout);
this.diagnosticSource = new ServiceBusDiagnosticSource($"{topicPath}/{subscriptionName}", serviceBusConnection.Endpoint);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

"{topicPath}/{subscriptionName}" [](start = 68, length = 32)

this should be replaced with this.Path

}
finally
{
this.diagnosticSource.AddRuleStop(activity, description, addRuleTask?.Status);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this.diagnosticSource.AddRuleStop(activity, description, addRuleTask?.Status); [](start = 16, length = 78)

wrap inside isDiagnosticsEnabled?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Similarly at other places


In reply to: 149190095 [](ancestors = 149190095)

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

if IsDiagnosticsEnabled is false, activity is null, all Stop methods do nothing if activity is null

finally
{
return unprocessedMessageList;
this.diagnosticSource.ReceiveStop(activity, maxMessageCount, receiveTask?.Status, processedMessageList);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this.diagnosticSource.ReceiveStop(activity, maxMessageCount, receiveTask?.Status, processedMessageList); [](start = 16, length = 104)

Start-Stop tracking might end up being used as latency counters for simplicity purposes by customers. Lets move the ReceiveStop log before await this.ProcessMessages(..)

return processedMessageList;
}


Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: remove

MessagingEventSource.Log.MessageSendStop(this.ClientId);
}


Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Lets not have double empty lines :)


#region ReceiveDeffered

internal Activity ReceiveDefferedStart(IEnumerable<long> sequenceNumbers)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Typo: Deferred

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks, will fix!

this.endpoint = endpoint;
}

public static bool IsEnabled()

@pharring Paul Harrington (pharring) Nov 7, 2017 •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion: Make this a getter-only property. i.e.

public static bool IsEnabled => DiagnosticListener.IsEnabled();

I think it reads more naturally at the call-site.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

it also creates a false impression on the caller side, that ServiceBusDiagnosticSource.IsEnabled value is static and could be called every time instead of caching the value (like it done now)

            bool isDiagnosticsEnabled = ServiceBusDiagnosticSource.IsEnabled();
            Activity activity = isDiagnosticsEnabled ? this.diagnosticSource.ReceiveDeferredStart(sequenceNumberList) : null;
            try
            {
                ...
            }
            catch (Exception exception)
            {
                if (isDiagnosticsEnabled)
                {
                    this.diagnosticSource.ReportException(exception);
                }
            }

with property it will be more natural to write

            Activity activity = ServiceBusDiagnosticSource.IsEnabled? this.diagnosticSource.ReceiveDeferredStart(sequenceNumberList) : null;
            try
            {
                ...
            }
            catch (Exception exception)
            {
                if (ServiceBusDiagnosticSource.IsEnabled)
                {
                    this.diagnosticSource.ReportException(exception);
                }
            }

While it's quite cheap to call DiagnosticListener.IsEnabled(), I still prefer result to be cached for particular operations not only because of performance, but also consistency

// So, let's wait for timeout and a bit more to make sure all created tasks are completed
sw.Stop();

await Task.Delay((int)(timeout - sw.Elapsed).TotalMilliseconds + 1000);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

await Task.Delay((int)(timeout - sw.Elapsed).TotalMilliseconds + 1000); [](start = 11, length = 72)

Should have caught this before during the code review.
In every test we need to ensure we have a finally block which will do await queueClient.CloseAsync() and similarly for other clients

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Remove the Task.Delay. If the issue still persists, then its a client bug that we need to fix. Once CloseAsync() is invoked, it is expected that every thread it spawns is closed.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I close queueClient and other resources in Dispose method - executed after every xunit test.
It's easier than clean up in every test.

If I remove this, tests become unstable. I've tried to fix it with a small timeout, so AcceptMessageSessionAsync returns faster, but it seems it's not always enough and I've seen tests fail because of it.

I can create an issue for that and add comment that this workaround should be removed once the issue is fixed. In the meantime, tests will be stable. Is it ok?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

created #372

@nemakam Neeraj Makam (nemakam) left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🕐

if (DiagnosticListener.IsEnabled(activityName, entityPath, tmpActivity))
{
activity = tmpActivity;
if (DiagnosticListener.IsEnabled(activityName + "Start"))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should "Start" have a leading period?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

right, nice catch!


MessagingEventSource.Log.MessageReceiveStart(this.ClientId, maxMessageCount);

bool isDiagnosticsEnabled = ServiceBusDiagnosticSource.IsEnabled();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is super nit-picky, I know: Consider either areDiagnosticsEnabled or isDiagnosticsSourceEnabeld. Why? Because "Diagnostics" are plural; whereas "DiagnosticsSource" is singular.

This might be a bit too contorted, but we might be able to avoid the problem altogether by eliminating the local bool and test for activity != null lower down.

@lmolkova Liudmila Molkova (lmolkova) Nov 7, 2017 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, I'm renaming to isDiagnosticSourceEnabled.

Regarding check for activity != null, it works for Stop event, but not for exceptions: we presume some users may want to have Exception events only. With this level of verbosity, there will be no activities created.

MessagingEventSource.Log.MessageReceiveException(this.ClientId, exception);
if (isDiagnosticsEnabled)
{
this.diagnosticSource.ReportException(exception);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Consider putting this before the call to MessagingEventSource.Log.MessageReceiveException. That way, I think, we will have a consistent nesting of EventSource (outer) and DiagnosticSource (inner) messages.

@lmolkova Liudmila Molkova (lmolkova) Nov 7, 2017 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

makes sense

@nemakam
Neeraj Makam (nemakam) merged commit dc40924 into Azure:dev Nov 8, 2017
}
}

private void SetTags(Activity activity, Message message)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Liudmila Molkova (@lmolkova) - shouldn't we set entity path and endpoint as tags here? these are very important for each monitoring solution so they should be made easy to consume

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

as we discussed, we cannot leverage fully 'generic listener' and have to implement the custom listener for ServiceBus. The custom listener will be parsing payloads and extracting endpoint and entity from them.
I wanted to keep instrumentation side as lightweight as possible. Otherwise, I would have to come up with 'standard' tags (as I'm sure similar things exists in almost every other Azure service client) and we would have to stick with them forever on AI side.

Sign up for free to subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants