Instruments operations on ServiceBus with DiagnosticsSource and introduces correlation between producer and consumer - #370
Conversation
5a5e263 to
1ccfd22
Compare
Codecov Report
@@ 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
Continue to review full report at Codecov.
|
…duce correlation between producer and consumer
1ccfd22 to
09b4acf
Compare
|
|
||
| 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"; |
There was a problem hiding this comment.
ProcessActivityStartName, ProcessSessionActivityName, and ProcessSessionActivityStartName don't seem to be used.
| null); | ||
| } | ||
|
|
||
| internal void ReceiveStop(Activity activity, int messageCount, TaskStatus? status, IList<Message> messageList) |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Agreed. Therefore MessageCount should not be assigned the value of messageCount but messageList.Length instead.
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
internal void ReportException(Exception ex) [](start = 8, length = 43)
Is there a paradigm of ReportException being associated with an Activity?
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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)
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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)); | ||
| } |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
"{topicPath}/{subscriptionName}" [](start = 68, length = 32)
this should be replaced with this.Path
| } | ||
| finally | ||
| { | ||
| this.diagnosticSource.AddRuleStop(activity, description, addRuleTask?.Status); |
There was a problem hiding this comment.
this.diagnosticSource.AddRuleStop(activity, description, addRuleTask?.Status); [](start = 16, length = 78)
wrap inside isDiagnosticsEnabled?
There was a problem hiding this comment.
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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; | ||
| } | ||
|
|
||
|
|
| MessagingEventSource.Log.MessageSendStop(this.ClientId); | ||
| } | ||
|
|
||
|
|
There was a problem hiding this comment.
nit: Lets not have double empty lines :)
|
|
||
| #region ReceiveDeffered | ||
|
|
||
| internal Activity ReceiveDefferedStart(IEnumerable<long> sequenceNumbers) |
There was a problem hiding this comment.
thanks, will fix!
| this.endpoint = endpoint; | ||
| } | ||
|
|
||
| public static bool IsEnabled() |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
created #372
| if (DiagnosticListener.IsEnabled(activityName, entityPath, tmpActivity)) | ||
| { | ||
| activity = tmpActivity; | ||
| if (DiagnosticListener.IsEnabled(activityName + "Start")) |
There was a problem hiding this comment.
Should "Start" have a leading period?
There was a problem hiding this comment.
right, nice catch!
|
|
||
| MessagingEventSource.Log.MessageReceiveStart(this.ClientId, maxMessageCount); | ||
|
|
||
| bool isDiagnosticsEnabled = ServiceBusDiagnosticSource.IsEnabled(); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
makes sense
| } | ||
| } | ||
|
|
||
| private void SetTags(Activity activity, Message message) |
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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.
This change adds instrumentation with
DiagnosticSourcefor public ServiceBus APIs:MessageSender(i.e.QueueClientandTopicClient):SendAsync,ScueduleMessageandCancelMessageReceiver:Receive,ReceiveDeffered,Peek,Abandon,Defer,Complete,DeadLetter,RenewLockMessageSession:AcceptMessageSessionSubscriptionClient:AddRule,RemoveRule,GetRulesTracing 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 queueCorrelation-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