From 2638cf7d8a85149ca1eccd2d6ffba8ebda55dd81 Mon Sep 17 00:00:00 2001 From: Ramjot Singh Date: Sun, 11 Nov 2018 17:53:44 -0800 Subject: [PATCH 1/2] Fixing Broadcast processor dropping telemetry items --- .../Extensibility/TelemetrySinkTests.cs | 107 ++++++++++++++++++ .../Implementation/BroadcastProcessor.cs | 27 ++--- 2 files changed, 117 insertions(+), 17 deletions(-) diff --git a/Test/Microsoft.ApplicationInsights.Test/Shared/Extensibility/TelemetrySinkTests.cs b/Test/Microsoft.ApplicationInsights.Test/Shared/Extensibility/TelemetrySinkTests.cs index 0ee3125514..38b3f3257c 100644 --- a/Test/Microsoft.ApplicationInsights.Test/Shared/Extensibility/TelemetrySinkTests.cs +++ b/Test/Microsoft.ApplicationInsights.Test/Shared/Extensibility/TelemetrySinkTests.cs @@ -295,5 +295,112 @@ public void ConfigurationDisposesAllSinks() Assert.Fail(ex.ToString()); } } + + /// + /// Ensure broadcast processor does not drop telemetry items. + /// + [TestMethod] + public void EnsureEventsAreNotDroppedByBroadcastProcessor() + { + var configuration = new TelemetryConfiguration(); + var commonChainBuilder = new TelemetryProcessorChainBuilder(configuration); + configuration.TelemetryProcessorChainBuilder = commonChainBuilder; + + ConcurrentBag itemsReceivedBySink1 = new ConcurrentBag(); + ConcurrentBag itemsReceivedBySink2 = new ConcurrentBag(); + + ITelemetryChannel firstTelemetryChannel = new StubTelemetryChannel + { + OnSend = telemetry => + { + itemsReceivedBySink1.Add(telemetry); + } + }; + + ITelemetryChannel secondTelemetryChannel = new StubTelemetryChannel + { + OnSend = telemetry => + { + itemsReceivedBySink2.Add(telemetry); + } + }; + + configuration.DefaultTelemetrySink.TelemetryChannel = firstTelemetryChannel; + configuration.TelemetrySinks.Add(new TelemetrySink(configuration, secondTelemetryChannel)); + + configuration.TelemetryProcessorChainBuilder.Build(); + + TelemetryClient telemetryClient = new TelemetryClient(configuration); + + // Setup TelemetryContext in a way that it is filledup. + telemetryClient.Context.Operation.Id = "OpId"; + telemetryClient.Context.Cloud.RoleName = "UnitTest"; + telemetryClient.Context.Component.Version = "TestVersion"; + telemetryClient.Context.Device.Id = "TestDeviceId"; + telemetryClient.Context.Flags = 1234; + telemetryClient.Context.InstrumentationKey = Guid.Empty.ToString(); + telemetryClient.Context.Location.Ip = "127.0.0.1"; + telemetryClient.Context.Session.Id = "SessionId"; + telemetryClient.Context.User.Id = "userId"; + + Parallel.ForEach( + new int[100], + new ParallelOptions + { + MaxDegreeOfParallelism = 100 + }, + (value) => + { + telemetryClient.TrackAvailability( + "Availability", + DateTimeOffset.Now, + TimeSpan.FromMilliseconds(200), + "Local", + true, + "Message", + new Dictionary() { { "Key", "Value" } }, + new Dictionary() { { "Dimension1", 0.9865 } }); + + telemetryClient.TrackDependency( + "HTTP", + "Target", + "Test", + "https://azure", + DateTimeOffset.Now, + TimeSpan.FromMilliseconds(100), + "200", + true); + + telemetryClient.TrackEvent( + "Event", + new Dictionary() { { "Key", "Value" } }, + new Dictionary() { { "Dimension1", 0.9865 } }); + + telemetryClient.TrackException( + new Exception("Test"), + new Dictionary() { { "Key", "Value" } }, + new Dictionary() { { "Dimension1", 0.9865 } }); + + telemetryClient.TrackMetric("Metric", 0.1, new Dictionary() { { "Key", "Value" } }); + + telemetryClient.TrackPageView("PageView"); + + telemetryClient.TrackRequest( + new RequestTelemetry("GET https://azure.com", DateTimeOffset.Now, TimeSpan.FromMilliseconds(200), "200", true) + { + HttpMethod = "GET" + }); + + telemetryClient.TrackTrace( + "Message", + SeverityLevel.Critical, + new Dictionary() { { "Key", "Value" } }); + + }); + + Assert.AreEqual(itemsReceivedBySink1.Count, itemsReceivedBySink2.Count); + Assert.AreEqual(8 * 100, itemsReceivedBySink1.Count); + Assert.AreEqual(8 * 100, itemsReceivedBySink2.Count); + } } } diff --git a/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs b/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs index 8a4447712f..62eb1f3b75 100644 --- a/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs +++ b/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs @@ -48,12 +48,7 @@ public void Process(ITelemetry item) // But it is also true that a misbehaving telemetry processor or channel in one of the sinks will affect other sinks. for (int i = 0; i < this.childrenDispatchers.Length; i++) { - this.childrenDispatchers[i].Offer(item); - } - - for (int i = 0; i < this.childrenDispatchers.Length; i++) - { - this.childrenDispatchers[i].ProcessOffered(); + this.childrenDispatchers[i].SendItemToSink(item); } } @@ -65,31 +60,29 @@ private class TelemetryDispatcher { private bool cloneBeforeDispatch; private TelemetrySink sink; - private ITelemetry nextTelemetryToProcess; public TelemetryDispatcher(TelemetrySink sink, bool cloneBeforeDispatch) { Debug.Assert(sink != null, "Telemetry sink should not be null"); this.sink = sink; this.cloneBeforeDispatch = cloneBeforeDispatch; - this.nextTelemetryToProcess = null; } - public void Offer(ITelemetry telemetry) + /// + /// Sends the item to sink. If cloning of item is required, clones the item before sending it to + /// + /// The telemetry item to send to sink. + public void SendItemToSink(ITelemetry telemetry) { - this.nextTelemetryToProcess = this.cloneBeforeDispatch ? telemetry.DeepClone() : telemetry; - } + ITelemetry itemToSendToSink = this.cloneBeforeDispatch ? telemetry.DeepClone() : telemetry; - public void ProcessOffered() - { - if (this.nextTelemetryToProcess != null) + if (itemToSendToSink != null) { - this.sink.Process(this.nextTelemetryToProcess); - this.nextTelemetryToProcess = null; + this.sink.Process(itemToSendToSink); } else { - Debug.Fail("We should not be asked to process a telemetry item if none was offered"); + Debug.Fail("Telemetry item should not be null"); } } } From b86595a2ae37fdd9046d832791bb266b3f23bb81 Mon Sep 17 00:00:00 2001 From: Ramjot Singh Date: Sun, 11 Nov 2018 18:11:50 -0800 Subject: [PATCH 2/2] Fixing a bug in code and updating the Changelog --- CHANGELOG.md | 1 + .../Extensibility/Implementation/BroadcastProcessor.cs | 8 +++++++- 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index df350f0ad4..56c9c00546 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ This changelog will be used to generate documentation on [release notes page](ht - [RequestTelemetry modified to lazily instantiate ConcurrentDictionary for Properties](https://github.com/Microsoft/ApplicationInsights-dotnet/issues/969) - [RequestTelemetry modified to not service public fields with data class to avoid converting between types.](https://github.com/Microsoft/ApplicationInsights-dotnet/issues/965) - [Fixed a bug in TelemetryContext which prevented rawobject store to be not available in all sinks.](https://github.com/Microsoft/ApplicationInsights-dotnet/issues/974) +- [Fixed BroadcastProcessor dropping TelemetryItems](https://github.com/Microsoft/ApplicationInsights-dotnet/pull/995) ## Version 2.8.1 [Patch release addressing perf regression.](https://github.com/Microsoft/ApplicationInsights-dotnet/issues/952) diff --git a/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs b/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs index 62eb1f3b75..d19b9eb08e 100644 --- a/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs +++ b/src/Microsoft.ApplicationInsights/Extensibility/Implementation/BroadcastProcessor.cs @@ -46,7 +46,13 @@ public void Process(ITelemetry item) // 2. Channels are reliable and process data asynchronously (ITelemetryChannel.Send() just queues up the telemetry and returns quickly). // As a result of these assumptions we can just let each sink process the data synchronously, with acceptable performance. // But it is also true that a misbehaving telemetry processor or channel in one of the sinks will affect other sinks. - for (int i = 0; i < this.childrenDispatchers.Length; i++) + + // Why the reverse traversal? As a perf optimization we want to avoid unecessary .DeepClone(). So we send the + // original item to the very first TelemetrySink, however first telemetry sink can choose to modify this object. + // In this case all the telemetry sinks will get the modified object. Hence as a protection against this, we are going to + // send the object through the first telemetry sink at the very last. At this point the first telemetry sink is free to + // modify the object as we have no further use of it. + for (int i = this.childrenDispatchers.Length - 1; i >= 0; i--) { this.childrenDispatchers[i].SendItemToSink(item); }