-
Notifications
You must be signed in to change notification settings - Fork 171
[Tracer] (Event Grid 1/5) Add contract and core behavior #8911
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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); | ||
| 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) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 |
||
| { | ||
| 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; } | ||
| } |
Uh oh!
There was an error while loading. Please reload this page.