Skip to content
Open
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
Original file line number Diff line number Diff line change
@@ -0,0 +1,267 @@
// <copyright file="EventGridCommon.cs" company="Datadog">
// Unless explicitly stated otherwise all files in this repository are licensed under the Apache 2 License.
// This product includes software developed at Datadog (https://www.datadoghq.com/). Copyright 2017 Datadog, Inc.
// </copyright>

#nullable enable

using System;
using System.Collections;
using System.Collections.Generic;
using Datadog.Trace.ClrProfiler.CallTarget;
using Datadog.Trace.Configuration;
using Datadog.Trace.Configuration.Schema;
using Datadog.Trace.DuckTyping;
using Datadog.Trace.Logging;
using Datadog.Trace.Propagators;
using Datadog.Trace.Tagging;

namespace Datadog.Trace.ClrProfiler.AutoInstrumentation.Azure.EventGrid;

internal static class EventGridCommon
{
private static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor(typeof(EventGridCommon));

internal static CallTargetState CreateProducerSpan<TTarget, TEvents>(TTarget instance, ref TEvents events, bool injectContext)
where TTarget : IEventGridPublisherClient
{
var tracer = Tracer.Instance;
if (!IsIntegrationEnabled(tracer))
{
return CallTargetState.GetDefault();
}

var uriBuilder = instance.UriBuilder;
var host = uriBuilder?.Host;
var port = uriBuilder?.Port ?? -1;
return CreateProducerSpan(tracer, host, port, ref events, injectContext);
}

internal static CallTargetState CreateNamespaceProducerSpanForEvent<TTarget, TCloudEvent>(TTarget instance, TCloudEvent cloudEvent)
where TTarget : IEventGridSenderClient
where TCloudEvent : ICloudEvent
{
var tracer = Tracer.Instance;
if (!IsIntegrationEnabled(tracer))
{
return CallTargetState.GetDefault();
}

var endpoint = instance.Endpoint;
object? singleEvent = cloudEvent.Instance is null ? null : cloudEvent;
return CreateProducerSpan(tracer, endpoint?.Host, endpoint?.Port ?? -1, events: null, singleEvent);
}

internal static CallTargetState CreateNamespaceProducerSpanForEvents<TTarget, TEvents>(TTarget instance, ref TEvents cloudEvents)
where TTarget : IEventGridSenderClient
{
var tracer = Tracer.Instance;
if (!IsIntegrationEnabled(tracer))
{
return CallTargetState.GetDefault();
}

var endpoint = instance.Endpoint;
return CreateProducerSpan(tracer, endpoint?.Host, endpoint?.Port ?? -1, ref cloudEvents, injectContext: true);
}

private static bool IsIntegrationEnabled(Tracer tracer) =>
tracer.CurrentTraceSettings.Settings.IsIntegrationEnabled(IntegrationId.AzureEventGrid, defaultValue: false);

private static CallTargetState CreateProducerSpan<TEvents>(Tracer tracer, string? host, int port, ref TEvents events, bool injectContext)
{
var enumerable = events as IEnumerable;
var state = CreateProducerSpan(tracer, host, port, enumerable, singleEvent: null);
Comment thread
pablomartinezbernardo marked this conversation as resolved.
if (state.Scope is not { } scope || enumerable is null)
{
return state;
}

try
{
var observer = new EventGridEnumerableObserver(scope, enumerable is ICollection collection ? collection.Count : null, injectContext);
events = EventGridObservingEnumerable.Wrap(events, observer);
}
catch (Exception ex)
{
Log.Debug(ex, "Error wrapping Azure Event Grid events for context injection");
}

return state;
}

private static CallTargetState CreateProducerSpan(Tracer tracer, string? host, int port, IEnumerable? events, object? singleEvent)

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.

I have not really dug into this but I feel like we could refactor this to pass the single event as an enumerable of one ? Or maybe there is a perf drawback to doing this ?

Because here there are some special tratments that look wrong, like "ProcessEvent" only for single events and not for multiple ones, except if coming from the other CreateProducerSpan where the injection is done... looks a bit scattered

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Creating an enumarable would be an extra unnecessary allocation.

Otherwise, events from ienumerables can't be processed here, because we would actually initiate enumeration before the SDK does. That's what EventGridEnumerableObserver is fixing. I removed the message count check from ProcessEvent because it is only called for single events, do you think that helps avoid confusion? Or do you maybe have another suggestion given these conditions?

{
Scope? scope = null;

try
{
var tags = tracer.CurrentTraceSettings.Schema.Messaging.CreateAzureEventGridTags(SpanKinds.Producer);
tags.MessagingOperation = "send";

tags.NetworkDestinationName = host;

if (port is not -1)
{
tags.NetworkDestinationPort = port.ToString();
}

var messageCount = singleEvent is not null ? 1 : events is ICollection collection ? collection.Count : 0;
if (messageCount > 1)
{
tags.MessagingBatchMessageCount = messageCount.ToString();
}

var (serviceName, serviceNameSource) = tracer.CurrentTraceSettings.Schema.Messaging.GetServiceNameMetadata(MessagingSchema.ServiceType.AzureEventGrid);
scope = tracer.StartActiveInternal("azure_eventgrid.send", tags: tags, serviceName: serviceName, serviceNameSource: serviceNameSource);
var span = scope.Span;

span.Type = SpanTypes.Queue;
span.ResourceName = "eventgrid";

if (singleEvent is not null)
{
ProcessEvent(singleEvent, span, scope);
}

tracer.TracerManager.Telemetry.IntegrationGeneratedSpan(IntegrationId.AzureEventGrid);

return new CallTargetState(scope);
}
catch (Exception ex)
{
Log.Error(ex, "Error creating Azure Event Grid producer span");
scope?.Dispose();
return CallTargetState.GetDefault();
}
}

private static void ProcessEvent(object evt, Span span, Scope scope)
{
SetMessageId(evt, span);
InjectContext(evt, scope);
}

private static void InjectContext(object evt, Scope scope)
{
if (evt.TryDuckCast<ICloudEvent>(out var cloudEvent)
&& cloudEvent.ExtensionAttributes is { } attrs)
{
InjectW3CContext(attrs, scope, Tracer.Instance.Settings.PropagationStyleInject);
}
}

private static void SetMessageId(object evt, Span span)
{
if (evt.TryDuckCast<IEventGridEventId>(out var eventGridEvent) && eventGridEvent.Id is { Length: > 0 } id)
{
span.SetTag(Tags.MessagingMessageId, id);
}
}

/// <summary>
/// Injects the configured W3C traceparent, tracestate, and baggage into CloudEvent ExtensionAttributes.
/// Uses the W3C propagator directly because CloudEvent extension attribute names
/// only allow lowercase letters and digits — Datadog-format headers (x-datadog-*)
/// would throw ArgumentException. Pre-populating the W3C attributes also prevents
/// the Azure SDK from overwriting them with its Activity-based context.
/// </summary>
internal static void InjectW3CContext(IDictionary<string, object> extensionAttributes, Scope scope, string[] propagationStyles)
{
if (scope.Span.Context is not { } spanContext)
{
return;
}

try
{
var context = new PropagationContext(spanContext, Baggage.Current);
var carrier = default(Shared.AzureMessagingCommon.DictionaryContextPropagation);

if (IsPropagationStyleEnabled(propagationStyles, ContextPropagationHeaderStyle.W3CTraceContext) ||
IsPropagationStyleEnabled(propagationStyles, ContextPropagationHeaderStyle.Deprecated.W3CTraceContext))
{
W3CTraceContextPropagator.Instance.Inject(context, extensionAttributes, carrier);
}

if (IsPropagationStyleEnabled(propagationStyles, ContextPropagationHeaderStyle.W3CBaggage))
{
W3CBaggagePropagator.Instance.Inject(context, extensionAttributes, carrier);
}
}
catch (Exception ex)
{
Log.Warning(ex, "Failed to inject W3C trace context into CloudEvent ExtensionAttributes");
}
}

private static bool IsPropagationStyleEnabled(string[] propagationStyles, string expectedStyle)
{
foreach (var propagationStyle in propagationStyles)
{
if (string.Equals(propagationStyle, expectedStyle, StringComparison.OrdinalIgnoreCase))
{
return true;
}
}

return false;
}

private sealed class EventGridEnumerableObserver : EventGridObservingEnumerable.IObserver
{
private readonly Scope _scope;
private readonly int? _knownCount;
private readonly bool _injectContext;

public EventGridEnumerableObserver(Scope scope, int? knownCount, bool injectContext)
{
_scope = scope;
_knownCount = knownCount;
_injectContext = injectContext;
}

public void OnItem(object? item)
{
try
{
if (item is not null)
{
if (_knownCount == 1)
{
SetMessageId(item, _scope.Span);
}

if (_injectContext)
{
InjectContext(item, _scope);
}
}
}
catch (Exception ex)
{
Log.Debug(ex, "Error processing an Azure Event Grid event");
}
}

public void OnEnumerationCompleted(int count, object? firstItem)
{
try
{
if (!_knownCount.HasValue && count > 1)
{
_scope.Span.SetTag(Tags.MessagingBatchMessageCount, count.ToString());
}

if (!_knownCount.HasValue && count == 1 && firstItem is not null)
{
SetMessageId(firstItem, _scope.Span);
}
}
catch (Exception ex)
{
Log.Debug(ex, "Error finalizing Azure Event Grid event processing");
}
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
// <copyright file="EventGridObservingEnumerable.cs" company="Datadog">
// Unless explicitly stated otherwise all files in this repository are licensed under the Apache 2 License.
// This product includes software developed at Datadog (https://www.datadoghq.com/). Copyright 2017 Datadog, Inc.
// </copyright>

#nullable enable

using System;
using System.Collections.Generic;

namespace Datadog.Trace.ClrProfiler.AutoInstrumentation.Azure.EventGrid;

/// <summary>
/// Replaces a declared <see cref="IEnumerable{T}"/> argument with an observing wrapper without
/// taking a compile-time dependency on the enumerable's item type.
/// </summary>
/// <remarks>
/// Event Grid send APIs accept deferred or one-shot sequences. Enumerating them in the CallTarget begin callback
/// could execute customer code early, consume the sequence, or observe different event instances from those the
/// Azure SDK eventually sends. This wrapper observes items only as the SDK enumerates them, preserving the source's
/// enumeration timing and item flow. It is constructed dynamically because the tracer cannot reference the Azure SDK
/// event types.
/// </remarks>
internal static class EventGridObservingEnumerable
{
internal interface IObserver
{
void OnItem(object? item);

void OnEnumerationCompleted(int count, object? firstItem);
}

private interface IEnumerableWrapperFactory
{
object Create(object events, IObserver observer);
}

internal static TEvents Wrap<TEvents>(TEvents events, IObserver observer)
{
if ((object?)events is null || FactoryCache<TEvents>.Factory is not { } factory)
{
return events;
}

return (TEvents)factory.Create(events, observer);
}

private static IEnumerable<TEvent> Observe<TEvent>(IEnumerable<TEvent> events, IObserver observer)
{
object? firstItem = null;
var count = 0;

foreach (var item in events)
{
if (count == 0)
{
firstItem = item;
}

count++;
try
{
observer.OnItem(item);
}
catch
{
// Instrumentation must not affect customer enumeration.
}

yield return item;
}

try
{
observer.OnEnumerationCompleted(count, firstItem);
}
catch
{
// Instrumentation must not affect customer enumeration.
}
}

// CallTarget supplies TEvents as the method's declared parameter type, for example IEnumerable<CloudEvent>.
// The tracer cannot reference CloudEvent directly, so reflection closes EnumerableWrapperFactory<TEvent> once.
// The resulting factory is cached, and subsequent calls perform no reflection.
private static class FactoryCache<TEvents>
{
public static readonly IEnumerableWrapperFactory? Factory = CreateFactory();

private static IEnumerableWrapperFactory? CreateFactory()
{
var eventsType = typeof(TEvents);
if (!eventsType.IsGenericType || eventsType.GetGenericTypeDefinition() != typeof(IEnumerable<>))
{
return null;
}

var eventType = eventsType.GetGenericArguments()[0];
var factoryType = typeof(EnumerableWrapperFactory<>).MakeGenericType(eventType);
return Activator.CreateInstance(factoryType, nonPublic: true) as IEnumerableWrapperFactory;
}
}

private sealed class EnumerableWrapperFactory<TEvent> : IEnumerableWrapperFactory
{
public object Create(object events, IObserver observer)
=> Observe((IEnumerable<TEvent>)events, observer);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
// <copyright file="ICloudEvent.cs" company="Datadog">
// Unless explicitly stated otherwise all files in this repository are licensed under the Apache 2 License.
// This product includes software developed at Datadog (https://www.datadoghq.com/). Copyright 2017 Datadog, Inc.
// </copyright>

#nullable enable

using System.Collections.Generic;
using Datadog.Trace.DuckTyping;

namespace Datadog.Trace.ClrProfiler.AutoInstrumentation.Azure.EventGrid;

/// <summary>
/// Duck type for Azure.Messaging.CloudEvent
/// </summary>
internal interface ICloudEvent : IEventGridEventId, IDuckType
{
IDictionary<string, object> ExtensionAttributes { get; }
}
Loading
Loading