Skip to content
Open
39 changes: 39 additions & 0 deletions tracer/src/Datadog.Trace/Activity/ActivityListenerHandler.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

#nullable enable

using System;
using System.Collections.Concurrent;
using Datadog.Trace.Activity.DuckTypes;
using Datadog.Trace.Activity.Handlers;
Expand Down Expand Up @@ -89,6 +90,19 @@ public static void OnActivityWithSourceStarted<T>(string sourceName, T activity)
var sName = sourceName ?? "(null)";
if (HandlerBySource.TryGetValue(sName, out var handler))
{
var isSourceNameMissing = StringUtil.IsNullOrEmpty(sourceName);
var isUsingDefaultHandler = handler is DefaultActivityHandler;

// If the source lookup only found the default handler, use the operation name as a fallback.
if (isSourceNameMissing && isUsingDefaultHandler)
{
var integrationHandler = FindIntegrationHandlerByOperationName(activity.OperationName);
if (integrationHandler is not null)
{
handler = integrationHandler;
}
}

handler.ActivityStarted(sName, activity);
}
else
Expand All @@ -110,5 +124,30 @@ public static void OnActivityWithSourceStopped<T>(string sourceName, T activity)
Log.Warning("ActivityListenerHandler: There's no handler to process the ActivityStopped event. [Source={SourceName}]", sName);
}
}

private static IActivityHandler? FindIntegrationHandlerByOperationName(string? operationName)
{
if (StringUtil.IsNullOrEmpty(operationName))
{
return null;
}

foreach (var handler in ActivityHandlersRegister.Handlers)
{
if (handler is DefaultActivityHandler)
{
return null;
}

if (handler is not DisableActivityHandler
&& handler is not IgnoreActivityHandler
&& handler.ShouldListenTo(operationName, version: null))
{
return handler;
}
}

return null;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@ internal sealed class ActivityHandlerCommon
{
private static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor(typeof(ActivityHandlerCommon));
internal static readonly ConcurrentDictionary<ActivityKey, ActivityMapping> ActivityMappingById = new();
private static readonly IntegrationId IntegrationId = IntegrationId.OpenTelemetry;

// Periodic sweep that handles two reliability cases that the hot-path Stop callback
// cannot: 1. Activities the customer abandoned without calling Stop (their WeakReference
Expand All @@ -48,6 +47,7 @@ internal sealed class ActivityHandlerCommon
/// <summary>
/// Handles when a new Activity is started to map it to a new <see cref="Span"/>/<see cref="Scope"/>.
/// </summary>
/// <param name="integrationId">The integration that generated the Activity.</param>
/// <param name="sourceName">The name of the Activity source</param>
/// <param name="activity">The Activity object</param>
/// <param name="tags">
Expand All @@ -56,10 +56,10 @@ internal sealed class ActivityHandlerCommon
/// </param>
/// <param name="activityMapping">The mapping of Activity to its <see cref="Scope"/>.</param>
/// <typeparam name="T">The <see cref="IActivity"/>.</typeparam>
public static void ActivityStarted<T>(string sourceName, T activity, OpenTelemetryTags? tags, out ActivityMapping activityMapping)
public static void ActivityStarted<T>(IntegrationId integrationId, string sourceName, T activity, OpenTelemetryTags? tags, out ActivityMapping activityMapping)
where T : IActivity
{
Tracer.Instance.TracerManager.Telemetry.IntegrationRunning(IntegrationId);
Tracer.Instance.TracerManager.Telemetry.IntegrationRunning(integrationId);

// Propagate Trace and Parent Span ids
SpanContext? parent = null;
Expand Down Expand Up @@ -221,10 +221,10 @@ public static void ActivityStarted<T>(string sourceName, T activity, OpenTelemet
// Avoid closure allocation if we can
activityMapping = ActivityMappingById.GetOrAdd(
activityKey.Value,
static (_, details) => new(new(details.activity.Instance!), CreateScopeFromActivity(details.activity, details.tags, details.parent, details.traceId, details.spanId, details.rawTraceId, details.rawSpanId)),
(activity, tags, parent, traceId, spanId, rawTraceId, rawSpanId));
static (_, details) => new(new(details.activity.Instance!), CreateScopeFromActivity(details.integrationId, details.activity, details.tags, details.parent, details.traceId, details.spanId, details.rawTraceId, details.rawSpanId)),
(integrationId, activity, tags, parent, traceId, spanId, rawTraceId, rawSpanId));
#else
activityMapping = ActivityMappingById.GetOrAdd(activityKey.Value, _ => new(new(activity.Instance!), CreateScopeFromActivity(activity, tags, parent, traceId, spanId, rawTraceId, rawSpanId)));
activityMapping = ActivityMappingById.GetOrAdd(activityKey.Value, _ => new(new(activity.Instance!), CreateScopeFromActivity(integrationId, activity, tags, parent, traceId, spanId, rawTraceId, rawSpanId)));
#endif
}
catch (Exception ex)
Expand All @@ -233,7 +233,7 @@ public static void ActivityStarted<T>(string sourceName, T activity, OpenTelemet
activityMapping = default;
}

static Scope CreateScopeFromActivity(T activity, OpenTelemetryTags? tags, SpanContext? parent, TraceId traceId, ulong spanId, string? rawTraceId, string? rawSpanId)
static Scope CreateScopeFromActivity(IntegrationId integrationId, T activity, OpenTelemetryTags? tags, SpanContext? parent, TraceId traceId, ulong spanId, string? rawTraceId, string? rawSpanId)
{
var span = Tracer.Instance.StartSpan(
activity.OperationName,
Expand All @@ -245,7 +245,7 @@ static Scope CreateScopeFromActivity(T activity, OpenTelemetryTags? tags, SpanCo
rawTraceId: rawTraceId,
rawSpanId: rawSpanId);

Tracer.Instance.TracerManager.Telemetry.IntegrationGeneratedSpan(IntegrationId);
Tracer.Instance.TracerManager.Telemetry.IntegrationGeneratedSpan(integrationId);
return Tracer.Instance.ActivateSpan(span, finishOnClose: false);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,20 +6,22 @@
#nullable enable

using Datadog.Trace.Activity.DuckTypes;
using Datadog.Trace.Configuration;
using Datadog.Trace.Tagging;

namespace Datadog.Trace.Activity.Handlers
{
internal sealed class AzureServiceBusActivityHandler : IActivityHandler
{
public bool ShouldListenTo(string sourceName, string? version)
=> sourceName.StartsWith("Azure.Messaging.ServiceBus");
=> sourceName.StartsWith("Azure.Messaging.ServiceBus")
&& Tracer.Instance.CurrentTraceSettings.Settings.IsIntegrationEnabled(IntegrationId.AzureServiceBus);

public void ActivityStarted<T>(string sourceName, T activity)
where T : IActivity
{
var tags = Tracer.Instance.CurrentTraceSettings.Schema.Client.CreateAzureServiceBusTags();
ActivityHandlerCommon.ActivityStarted(sourceName, activity, tags: tags, out var activityMapping);
ActivityHandlerCommon.ActivityStarted(IntegrationId.AzureServiceBus, sourceName, activity, tags: tags, out var activityMapping);
}

public void ActivityStopped<T>(string sourceName, T activity)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
#nullable enable

using Datadog.Trace.Activity.DuckTypes;
using Datadog.Trace.Configuration;
using Datadog.Trace.Tagging;

namespace Datadog.Trace.Activity.Handlers
Expand All @@ -22,7 +23,7 @@ public bool ShouldListenTo(string sourceName, string? version)

public void ActivityStarted<T>(string sourceName, T activity)
where T : IActivity
=> ActivityHandlerCommon.ActivityStarted(sourceName, activity, tags: null, out _);
=> ActivityHandlerCommon.ActivityStarted(IntegrationId.OpenTelemetry, sourceName, activity, tags: null, out _);

public void ActivityStopped<T>(string sourceName, T activity)
where T : IActivity
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,14 @@ internal sealed class QuartzActivityHandler : IActivityHandler
public bool ShouldListenTo(string sourceName, string? version)
{
// Listen to Quartz diagnostic source
return sourceName.StartsWith("Quartz");
return sourceName.StartsWith("Quartz")
&& Tracer.Instance.CurrentTraceSettings.Settings.IsIntegrationEnabled(IntegrationId.Quartz);
}

public void ActivityStarted<T>(string sourceName, T activity)
where T : IActivity
{
ActivityHandlerCommon.ActivityStarted(sourceName, activity, tags: new OpenTelemetryTags(), out var activityMapping);
ActivityHandlerCommon.ActivityStarted(IntegrationId.Quartz, sourceName, activity, tags: new OpenTelemetryTags(), out var activityMapping);
Comment thread
chojomok marked this conversation as resolved.
}

public void ActivityStopped<T>(string sourceName, T activity)
Expand Down
1 change: 1 addition & 0 deletions tracer/src/Datadog.Trace/Configuration/IntegrationId.cs
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,7 @@ internal enum IntegrationId
DatadogTraceVersionConflict,
Hangfire,
OpenFeature,
Quartz,
ServerlessCompat
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
using Datadog.Trace.Activity;
using Datadog.Trace.Activity.DuckTypes;
using Datadog.Trace.ClrProfiler.AutoInstrumentation.Quartz;
using Datadog.Trace.Configuration;
using Datadog.Trace.Logging;

namespace Datadog.Trace.DiagnosticListeners;
Expand All @@ -28,6 +29,11 @@ internal sealed class QuartzDiagnosticObserver : DiagnosticObserver

protected override void OnNext(string eventName, object arg)
{
if (!Tracer.Instance.CurrentTraceSettings.Settings.IsIntegrationEnabled(IntegrationId.Quartz))
{
return;
}

switch (eventName)
{
case "Quartz.Job.Execute.Start":
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ internal static partial class IntegrationIdExtensions
/// The number of members in the enum.
/// This is a non-distinct count of defined names.
/// </summary>
public const int Length = 80;
public const int Length = 81;

/// <summary>
/// Returns the string representation of the <see cref="Datadog.Trace.Configuration.IntegrationId"/> value.
Expand Down Expand Up @@ -109,6 +109,7 @@ public static string ToStringFast(this Datadog.Trace.Configuration.IntegrationId
Datadog.Trace.Configuration.IntegrationId.DatadogTraceVersionConflict => nameof(Datadog.Trace.Configuration.IntegrationId.DatadogTraceVersionConflict),
Datadog.Trace.Configuration.IntegrationId.Hangfire => nameof(Datadog.Trace.Configuration.IntegrationId.Hangfire),
Datadog.Trace.Configuration.IntegrationId.OpenFeature => nameof(Datadog.Trace.Configuration.IntegrationId.OpenFeature),
Datadog.Trace.Configuration.IntegrationId.Quartz => nameof(Datadog.Trace.Configuration.IntegrationId.Quartz),
Datadog.Trace.Configuration.IntegrationId.ServerlessCompat => nameof(Datadog.Trace.Configuration.IntegrationId.ServerlessCompat),
_ => value.ToString(),
};
Expand Down Expand Up @@ -202,6 +203,7 @@ public static Datadog.Trace.Configuration.IntegrationId[] GetValues()
Datadog.Trace.Configuration.IntegrationId.DatadogTraceVersionConflict,
Datadog.Trace.Configuration.IntegrationId.Hangfire,
Datadog.Trace.Configuration.IntegrationId.OpenFeature,
Datadog.Trace.Configuration.IntegrationId.Quartz,
Datadog.Trace.Configuration.IntegrationId.ServerlessCompat,
};

Expand Down Expand Up @@ -295,6 +297,7 @@ public static string[] GetNames()
nameof(Datadog.Trace.Configuration.IntegrationId.DatadogTraceVersionConflict),
nameof(Datadog.Trace.Configuration.IntegrationId.Hangfire),
nameof(Datadog.Trace.Configuration.IntegrationId.OpenFeature),
nameof(Datadog.Trace.Configuration.IntegrationId.Quartz),
nameof(Datadog.Trace.Configuration.IntegrationId.ServerlessCompat),
};
}
Original file line number Diff line number Diff line change
Expand Up @@ -259,6 +259,9 @@ public static string[] GetAllIntegrationEnabledKeys() =>
"DD_TRACE_OPENFEATURE_ENABLED", "DD_TRACE_OpenFeature_ENABLED", "DD_OpenFeature_ENABLED",
"DD_TRACE_OPENFEATURE_ANALYTICS_ENABLED", "DD_TRACE_OpenFeature_ANALYTICS_ENABLED", "DD_OpenFeature_ANALYTICS_ENABLED",
"DD_TRACE_OPENFEATURE_ANALYTICS_SAMPLE_RATE", "DD_TRACE_OpenFeature_ANALYTICS_SAMPLE_RATE", "DD_OpenFeature_ANALYTICS_SAMPLE_RATE",
"DD_TRACE_QUARTZ_ENABLED", "DD_TRACE_Quartz_ENABLED", "DD_Quartz_ENABLED",
"DD_TRACE_QUARTZ_ANALYTICS_ENABLED", "DD_TRACE_Quartz_ANALYTICS_ENABLED", "DD_Quartz_ANALYTICS_ENABLED",
"DD_TRACE_QUARTZ_ANALYTICS_SAMPLE_RATE", "DD_TRACE_Quartz_ANALYTICS_SAMPLE_RATE", "DD_Quartz_ANALYTICS_SAMPLE_RATE",
"DD_TRACE_SERVERLESSCOMPAT_ENABLED", "DD_TRACE_ServerlessCompat_ENABLED", "DD_ServerlessCompat_ENABLED",
"DD_TRACE_SERVERLESSCOMPAT_ANALYTICS_ENABLED", "DD_TRACE_ServerlessCompat_ANALYTICS_ENABLED", "DD_ServerlessCompat_ANALYTICS_ENABLED",
"DD_TRACE_SERVERLESSCOMPAT_ANALYTICS_SAMPLE_RATE", "DD_TRACE_ServerlessCompat_ANALYTICS_SAMPLE_RATE", "DD_ServerlessCompat_ANALYTICS_SAMPLE_RATE",
Expand Down Expand Up @@ -351,6 +354,7 @@ public static System.Collections.Generic.KeyValuePair<string, string[]> GetInteg
"DatadogTraceVersionConflict" => new("DD_TRACE_DATADOGTRACEVERSIONCONFLICT_ENABLED", ["DD_TRACE_DatadogTraceVersionConflict_ENABLED", "DD_DatadogTraceVersionConflict_ENABLED"]),
"Hangfire" => new("DD_TRACE_HANGFIRE_ENABLED", ["DD_TRACE_Hangfire_ENABLED", "DD_Hangfire_ENABLED"]),
"OpenFeature" => new("DD_TRACE_OPENFEATURE_ENABLED", ["DD_TRACE_OpenFeature_ENABLED", "DD_OpenFeature_ENABLED"]),
"Quartz" => new("DD_TRACE_QUARTZ_ENABLED", ["DD_TRACE_Quartz_ENABLED", "DD_Quartz_ENABLED"]),
"ServerlessCompat" => new("DD_TRACE_SERVERLESSCOMPAT_ENABLED", ["DD_TRACE_ServerlessCompat_ENABLED", "DD_ServerlessCompat_ENABLED"]),
_ => GetIntegrationEnabledKeysFallback(integrationName) // we should never get here
};
Expand Down Expand Up @@ -443,6 +447,7 @@ public static System.Collections.Generic.KeyValuePair<string, string[]> GetInteg
"DatadogTraceVersionConflict" => new("DD_TRACE_DATADOGTRACEVERSIONCONFLICT_ANALYTICS_ENABLED", ["DD_TRACE_DatadogTraceVersionConflict_ANALYTICS_ENABLED", "DD_DatadogTraceVersionConflict_ANALYTICS_ENABLED"]),
"Hangfire" => new("DD_TRACE_HANGFIRE_ANALYTICS_ENABLED", ["DD_TRACE_Hangfire_ANALYTICS_ENABLED", "DD_Hangfire_ANALYTICS_ENABLED"]),
"OpenFeature" => new("DD_TRACE_OPENFEATURE_ANALYTICS_ENABLED", ["DD_TRACE_OpenFeature_ANALYTICS_ENABLED", "DD_OpenFeature_ANALYTICS_ENABLED"]),
"Quartz" => new("DD_TRACE_QUARTZ_ANALYTICS_ENABLED", ["DD_TRACE_Quartz_ANALYTICS_ENABLED", "DD_Quartz_ANALYTICS_ENABLED"]),
"ServerlessCompat" => new("DD_TRACE_SERVERLESSCOMPAT_ANALYTICS_ENABLED", ["DD_TRACE_ServerlessCompat_ANALYTICS_ENABLED", "DD_ServerlessCompat_ANALYTICS_ENABLED"]),
_ => GetIntegrationAnalyticsEnabledKeysFallback(integrationName) // we should never get here
};
Expand Down Expand Up @@ -535,6 +540,7 @@ public static System.Collections.Generic.KeyValuePair<string, string[]> GetInteg
"DatadogTraceVersionConflict" => new("DD_TRACE_DATADOGTRACEVERSIONCONFLICT_ANALYTICS_SAMPLE_RATE", ["DD_TRACE_DatadogTraceVersionConflict_ANALYTICS_SAMPLE_RATE", "DD_DatadogTraceVersionConflict_ANALYTICS_SAMPLE_RATE"]),
"Hangfire" => new("DD_TRACE_HANGFIRE_ANALYTICS_SAMPLE_RATE", ["DD_TRACE_Hangfire_ANALYTICS_SAMPLE_RATE", "DD_Hangfire_ANALYTICS_SAMPLE_RATE"]),
"OpenFeature" => new("DD_TRACE_OPENFEATURE_ANALYTICS_SAMPLE_RATE", ["DD_TRACE_OpenFeature_ANALYTICS_SAMPLE_RATE", "DD_OpenFeature_ANALYTICS_SAMPLE_RATE"]),
"Quartz" => new("DD_TRACE_QUARTZ_ANALYTICS_SAMPLE_RATE", ["DD_TRACE_Quartz_ANALYTICS_SAMPLE_RATE", "DD_Quartz_ANALYTICS_SAMPLE_RATE"]),
"ServerlessCompat" => new("DD_TRACE_SERVERLESSCOMPAT_ANALYTICS_SAMPLE_RATE", ["DD_TRACE_ServerlessCompat_ANALYTICS_SAMPLE_RATE", "DD_ServerlessCompat_ANALYTICS_SAMPLE_RATE"]),
_ => GetIntegrationAnalyticsSampleRateKeysFallback(integrationName) // we should never get here
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
namespace Datadog.Trace.Telemetry;
internal sealed partial class CiVisibilityMetricsTelemetryCollector
{
private const int CountSharedLength = 340;
private const int CountSharedLength = 344;

/// <summary>
/// Creates the buffer for the <see cref="Datadog.Trace.Telemetry.Metrics.CountShared" /> values.
Expand Down Expand Up @@ -356,6 +356,10 @@ private static AggregatedMetric[] GetCountSharedBuffer()
new(new[] { "integration_name:hangfire", "error_type:invoker" }),
new(new[] { "integration_name:hangfire", "error_type:execution" }),
new(new[] { "integration_name:hangfire", "error_type:missing_member" }),
new(new[] { "integration_name:quartz", "error_type:duck_typing" }),
new(new[] { "integration_name:quartz", "error_type:invoker" }),
new(new[] { "integration_name:quartz", "error_type:execution" }),
new(new[] { "integration_name:quartz", "error_type:missing_member" }),
new(new[] { "integration_name:serverlesscompat", "error_type:duck_typing" }),
new(new[] { "integration_name:serverlesscompat", "error_type:invoker" }),
new(new[] { "integration_name:serverlesscompat", "error_type:execution" }),
Expand All @@ -368,7 +372,7 @@ private static AggregatedMetric[] GetCountSharedBuffer()
/// It is equal to the cardinality of the tag combinations (or 1 if there are no tags)
/// </summary>
private static int[] CountSharedEntryCounts { get; }
= new int[]{ 340, };
= new int[]{ 344, };

public void RecordCountSharedIntegrationsError(Datadog.Trace.Telemetry.Metrics.MetricTags.IntegrationName tag1, Datadog.Trace.Telemetry.Metrics.MetricTags.InstrumentationError tag2, int increment = 1)
{
Expand Down
Loading
Loading