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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ This changelog will be used to generate documentation on [release notes page](ht
- [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 a bug where TelemetryContext would have missing values on secondary sinks](https://github.com/Microsoft/ApplicationInsights-dotnet/pull/993)
- [Fixed race condition in BroadcastProcessor which caused it to drop TelemetryItems](https://github.com/Microsoft/ApplicationInsights-dotnet/pull/995)
- [Custom Telemetry Item that implements ITelemetry is no longer dropped, bur rather serialized as EventTelemetry and handled by the channels accordingly](https://github.com/Microsoft/ApplicationInsights-dotnet/issues/988)
- [Improved Perf of ITelemetry JsonSerialization](https://github.com/Microsoft/ApplicationInsights-dotnet/issues/997)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,10 @@
using Microsoft.VisualStudio.TestTools.UnitTesting;
using Newtonsoft.Json;
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;

[TestClass]
public class TelemetrySinkTests
Expand Down Expand Up @@ -438,5 +440,114 @@ public void EnsureAllTelemetrySinkItemsAreSimilarAcrossSinks()
new Dictionary<string, string>() { { "Key", "Value" } });
Assert.AreEqual(jsonFromFirstChannel, jsonFromSecondChannel);
}

/// <summary>
/// Ensure broadcast processor does not drop telemetry items.
/// </summary>
[TestMethod]
public void EnsureEventsAreNotDroppedByBroadcastProcessor()
{
var configuration = new TelemetryConfiguration();
var commonChainBuilder = new TelemetryProcessorChainBuilder(configuration);
configuration.TelemetryProcessorChainBuilder = commonChainBuilder;

ConcurrentBag<ITelemetry> itemsReceivedBySink1 = new ConcurrentBag<ITelemetry>();
ConcurrentBag<ITelemetry> itemsReceivedBySink2 = new ConcurrentBag<ITelemetry>();

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<string, string>() { { "Key", "Value" } },
new Dictionary<string, double>() { { "Dimension1", 0.9865 } });

telemetryClient.TrackDependency(
"HTTP",
"Target",
"Test",
"https://azure",
DateTimeOffset.Now,
TimeSpan.FromMilliseconds(100),
"200",
true);

telemetryClient.TrackEvent(
"Event",
new Dictionary<string, string>() { { "Key", "Value" } },
new Dictionary<string, double>() { { "Dimension1", 0.9865 } });

telemetryClient.TrackException(
new Exception("Test"),
new Dictionary<string, string>() { { "Key", "Value" } },
new Dictionary<string, double>() { { "Dimension1", 0.9865 } });

telemetryClient.TrackMetric("Metric", 0.1, new Dictionary<string, string>() { { "Key", "Value" } });

telemetryClient.TrackPageView("PageView");

telemetryClient.TrackRequest(
new RequestTelemetry("GET https://azure.com", DateTimeOffset.Now, TimeSpan.FromMilliseconds(200), "200", true)
{
#pragma warning disable CS0618 // Type or member is obsolete. Using for testing all cases.
HttpMethod = "GET"
#pragma warning restore CS0618 // Type or member is obsolete. Using for testing all cases.
});

telemetryClient.TrackTrace(
"Message",
SeverityLevel.Critical,
new Dictionary<string, string>() { { "Key", "Value" } });

});

Assert.AreEqual(itemsReceivedBySink1.Count, itemsReceivedBySink2.Count);
Assert.AreEqual(8 * 100, itemsReceivedBySink1.Count);
Assert.AreEqual(8 * 100, itemsReceivedBySink2.Count);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -46,14 +46,15 @@ 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++)
{
this.childrenDispatchers[i].Offer(item);
}

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].ProcessOffered();
this.childrenDispatchers[i].SendItemToSink(item);
}
}

Expand All @@ -65,31 +66,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)
/// <summary>
/// Sends the item to sink. If cloning of item is required, clones the item before sending it to <see cref="TelemetrySink"/>
/// </summary>
/// <param name="telemetry">The telemetry item to send to sink.</param>
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");
}
}
}
Expand Down