From 8d551242443631950e99904502c7e0e8210c1b29 Mon Sep 17 00:00:00 2001 From: Leo Romanovsky Date: Thu, 1 Oct 2026 22:23:55 +0000 Subject: [PATCH 1/2] Deliver Agentless Feature Flags events through safe EVP fallback Prefer identity-capable local relays and keep a credential-isolated direct route with conservative replay and sticky selection. Environment: Datadog workspace --- .../build/Datadog.Trace.Trimming.xml | 1 + .../Agent/AgentTransportStrategy.cs | 13 +- .../Agent/Transports/ApiWebRequestFactory.cs | 7 +- .../Transports/HttpClientRequestFactory.cs | 15 +- .../Transports/SocketHandlerRequestFactory.cs | 5 +- .../Agentless/AgentlessEndpoint.cs | 98 +- .../Evp/FeatureFlagsEvpHeaderHelper.cs | 45 + .../Evp/FeatureFlagsEvpTransport.cs | 573 +++++++++++- .../FeatureFlags/FeatureFlagsModule.cs | 11 +- .../src/Datadog.Trace/TracerManagerFactory.cs | 2 +- .../TransportHelpers/TestRequestFactory.cs | 2 +- .../FeatureFlags/AgentlessEndpointTests.cs | 6 + .../FeatureFlags/ExposureApiTests.cs | 11 +- .../FeatureFlagsEvpTransportTests.cs | 872 ++++++++++++++++++ .../FeatureFlagsFixedEvpTransportTests.cs | 5 +- .../FeatureFlags/FeatureFlagsModuleTests.cs | 33 + 16 files changed, 1636 insertions(+), 63 deletions(-) create mode 100644 tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpHeaderHelper.cs create mode 100644 tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs diff --git a/tracer/src/Datadog.Trace.Trimming/build/Datadog.Trace.Trimming.xml b/tracer/src/Datadog.Trace.Trimming/build/Datadog.Trace.Trimming.xml index 8925d38839bb..0c6e862295df 100644 --- a/tracer/src/Datadog.Trace.Trimming/build/Datadog.Trace.Trimming.xml +++ b/tracer/src/Datadog.Trace.Trimming/build/Datadog.Trace.Trimming.xml @@ -930,6 +930,7 @@ + diff --git a/tracer/src/Datadog.Trace/Agent/AgentTransportStrategy.cs b/tracer/src/Datadog.Trace/Agent/AgentTransportStrategy.cs index 9bff5f2cbd17..55a1be9ee015 100644 --- a/tracer/src/Datadog.Trace/Agent/AgentTransportStrategy.cs +++ b/tracer/src/Datadog.Trace/Agent/AgentTransportStrategy.cs @@ -31,12 +31,14 @@ internal static class AgentTransportStrategy /// A func that returns the endpoint to send requests to for a given "base" endpoint. /// The base endpoint will be for TCP requests and /// http://localhost/ for named pipes/UDS if null, the default base endpoint is used + /// Whether HTTP handlers may follow redirects. Stream transports never follow them. public static IApiRequestFactory Get( ExporterSettings settings, string productName, TimeSpan? tcpTimeout, HttpHeaderHelperBase httpHeaderHelper, - Func? getBaseEndpoint = null) + Func? getBaseEndpoint = null, + bool allowAutoRedirect = true) { var strategy = settings.TracesTransport; @@ -56,7 +58,8 @@ public static IApiRequestFactory Get( return new SocketHandlerRequestFactory( new UnixDomainSocketStreamFactory(settings.TracesUnixDomainSocketPath), httpHeaderHelper.DefaultHeaders, - getBaseEndpoint?.Invoke(Localhost) ?? Localhost); + getBaseEndpoint?.Invoke(Localhost) ?? Localhost, + allowAutoRedirect: allowAutoRedirect); #elif NETCOREAPP3_1_OR_GREATER Log.Information("Using " + nameof(UnixDomainSocketStreamFactory) + " for {ProductName} transport, with Unix Domain Sockets path {TracesUnixDomainSocketPath} and timeout {TracesPipeTimeoutMs}ms.", productName, settings.TracesUnixDomainSocketPath, settings.TracesPipeTimeoutMs); return new HttpStreamRequestFactory( @@ -74,13 +77,15 @@ public static IApiRequestFactory Get( return new HttpClientRequestFactory( getBaseEndpoint?.Invoke(settings.AgentUri) ?? settings.AgentUri, httpHeaderHelper.DefaultHeaders, - timeout: tcpTimeout); + timeout: tcpTimeout, + allowAutoRedirect: allowAutoRedirect); #else Log.Information("Using " + nameof(ApiWebRequestFactory) + " for {ProductName} transport.", productName); return new ApiWebRequestFactory( getBaseEndpoint?.Invoke(settings.AgentUri) ?? settings.AgentUri, httpHeaderHelper.DefaultHeaders, - timeout: tcpTimeout); + timeout: tcpTimeout, + allowAutoRedirect: allowAutoRedirect); #endif } } diff --git a/tracer/src/Datadog.Trace/Agent/Transports/ApiWebRequestFactory.cs b/tracer/src/Datadog.Trace/Agent/Transports/ApiWebRequestFactory.cs index 46160bba2f3e..dab42dd8b2f4 100644 --- a/tracer/src/Datadog.Trace/Agent/Transports/ApiWebRequestFactory.cs +++ b/tracer/src/Datadog.Trace/Agent/Transports/ApiWebRequestFactory.cs @@ -17,17 +17,21 @@ internal sealed class ApiWebRequestFactory : IApiRequestFactory { private readonly KeyValuePair[] _defaultHeaders; private readonly Uri _baseEndpoint; + private readonly bool _allowAutoRedirect; private WebProxy _proxy; private NetworkCredential _credential; private TimeSpan? _timeout; - public ApiWebRequestFactory(Uri baseEndpoint, KeyValuePair[] defaultHeaders, TimeSpan? timeout = null) + public ApiWebRequestFactory(Uri baseEndpoint, KeyValuePair[] defaultHeaders, TimeSpan? timeout = null, bool allowAutoRedirect = true) { _baseEndpoint = baseEndpoint; _defaultHeaders = defaultHeaders; _timeout = timeout; + _allowAutoRedirect = allowAutoRedirect; } + internal bool AllowAutoRedirect => _allowAutoRedirect; + public string Info(Uri endpoint) { return endpoint.ToString(); @@ -38,6 +42,7 @@ public string Info(Uri endpoint) public IApiRequest Create(Uri endpoint) { var request = WebRequest.CreateHttp(endpoint); + request.AllowAutoRedirect = _allowAutoRedirect; if (_proxy is not null) { request.Proxy = _proxy; diff --git a/tracer/src/Datadog.Trace/Agent/Transports/HttpClientRequestFactory.cs b/tracer/src/Datadog.Trace/Agent/Transports/HttpClientRequestFactory.cs index a6b5b8360e89..26f1b409ce0f 100644 --- a/tracer/src/Datadog.Trace/Agent/Transports/HttpClientRequestFactory.cs +++ b/tracer/src/Datadog.Trace/Agent/Transports/HttpClientRequestFactory.cs @@ -1,4 +1,4 @@ -// +// // 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. // @@ -22,9 +22,9 @@ internal sealed class HttpClientRequestFactory : IApiRequestFactory private readonly HttpMessageHandler _handler; private readonly Uri _baseEndpoint; - public HttpClientRequestFactory(Uri baseEndpoint, KeyValuePair[] defaultHeaders, HttpMessageHandler handler = null, TimeSpan? timeout = null) + public HttpClientRequestFactory(Uri baseEndpoint, KeyValuePair[] defaultHeaders, HttpMessageHandler handler = null, TimeSpan? timeout = null, bool allowAutoRedirect = true) { - _handler = handler ?? new HttpClientHandler(); + _handler = handler ?? new HttpClientHandler { AllowAutoRedirect = allowAutoRedirect }; _client = new HttpClient(_handler); _baseEndpoint = baseEndpoint; if (timeout.HasValue) @@ -41,6 +41,15 @@ public HttpClientRequestFactory(Uri baseEndpoint, KeyValuePair[] _client.DefaultRequestHeaders.ConnectionClose = true; } + internal bool AllowAutoRedirect => _handler switch + { + HttpClientHandler handler => handler.AllowAutoRedirect, +#if NET5_0_OR_GREATER + SocketsHttpHandler handler => handler.AllowAutoRedirect, +#endif + _ => true, + }; + public Uri GetEndpoint(string relativePath) => relativePath is null ? _baseEndpoint : UriHelpers.Combine(_baseEndpoint, relativePath); #if NET5_0_OR_GREATER // in .NET 6 we derive a SocketHandlerRequestFactory diff --git a/tracer/src/Datadog.Trace/Agent/Transports/SocketHandlerRequestFactory.cs b/tracer/src/Datadog.Trace/Agent/Transports/SocketHandlerRequestFactory.cs index 6bcbd1427dd4..b0cd23a73628 100644 --- a/tracer/src/Datadog.Trace/Agent/Transports/SocketHandlerRequestFactory.cs +++ b/tracer/src/Datadog.Trace/Agent/Transports/SocketHandlerRequestFactory.cs @@ -1,4 +1,4 @@ -// +// // 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. // @@ -15,7 +15,7 @@ internal sealed class SocketHandlerRequestFactory : HttpClientRequestFactory { private readonly IStreamFactory _streamFactory; - public SocketHandlerRequestFactory(IStreamFactory streamFactory, KeyValuePair[] defaultHeaders, Uri baseEndpoint, TimeSpan? timeout = null) + public SocketHandlerRequestFactory(IStreamFactory streamFactory, KeyValuePair[] defaultHeaders, Uri baseEndpoint, TimeSpan? timeout = null, bool allowAutoRedirect = true) : base( // HttpClient requires a "valid" host header, and will only accept http:// or https:// schemes // The host part of the endpoint is irrelevant, as we're using the UDS socket/named pipe @@ -25,6 +25,7 @@ public SocketHandlerRequestFactory(IStreamFactory streamFactory, KeyValuePair await streamFactory.GetBidirectionalStreamAsync(token).ConfigureAwait(false) }) { diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Agentless/AgentlessEndpoint.cs b/tracer/src/Datadog.Trace/FeatureFlags/Agentless/AgentlessEndpoint.cs index 64bda187b754..8dbf9f8777a8 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/Agentless/AgentlessEndpoint.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/Agentless/AgentlessEndpoint.cs @@ -78,28 +78,12 @@ public static bool TryCreate(string? site, string? baseUrl, [NotNullWhen(true)] var configured = baseUrl?.Trim(); if (StringUtil.IsNullOrEmpty(configured)) { - var trimmedSite = site?.Trim(); - if (StringUtil.IsNullOrEmpty(trimmedSite)) + if (!TryNormalizeSite(site, out var normalizedSite, out error)) { - error = "No Datadog site is configured"; return false; } - // The site is concatenated into a host, so every character that can change what a URL means - // has to be rejected before that happens. "@" is the dangerous one: it would make the rest of - // the value the real host, and the API key would be sent there. "/", "?" and "#" would start a - // path, query or fragment, and ":" a port or a scheme. Uri.TryCreate accepts several of these, - // so it cannot be relied on to catch them. The other tracers reject the same set. - foreach (var character in trimmedSite) - { - if (char.IsWhiteSpace(character) || character is '/' or '?' or '#' or '@' or ':') - { - error = "The configured Datadog site is not valid"; - return false; - } - } - - var managedHost = ManagedHostPrefix + StringUtil.ToLowerInvariant(trimmedSite); + var managedHost = ManagedHostPrefix + normalizedSite; if (!Uri.TryCreate($"https://{managedHost}{DefaultPath}", UriKind.Absolute, out var managedUri)) { @@ -144,6 +128,84 @@ public static bool TryCreate(string? site, string? baseUrl, [NotNullWhen(true)] return true; } + /// + /// Validates and normalizes a Datadog site before it is appended to a managed hostname. + /// Configuration and event delivery use this same method so credentials cannot be routed by + /// two subtly different parsers. + /// + internal static bool TryNormalizeSite(string? site, out string normalizedSite, out string? error) + { + normalizedSite = string.Empty; + error = null; + + var trimmedSite = site?.Trim(); + if (StringUtil.IsNullOrEmpty(trimmedSite)) + { + error = "No Datadog site is configured"; + return false; + } + + // The complete managed host must remain below the DNS limit once either the configuration + // or event-intake prefix is applied. + if (trimmedSite.Length > 230) + { + error = "The configured Datadog site is not valid"; + return false; + } + + // Accept only DNS label characters before concatenating the site into a managed host. + // Uri parsing alone accepts delimiters such as '@' that can redirect the API key to a + // different host, as well as '/', '?', '#', and ':' that change the URL's meaning. + var labelLength = 0; + var previousWasHyphen = false; + foreach (var character in trimmedSite) + { + if (character == '.') + { + if (labelLength == 0 || previousWasHyphen) + { + error = "The configured Datadog site is not valid"; + return false; + } + + labelLength = 0; + previousWasHyphen = false; + continue; + } + + if (character is >= 'A' and <= 'Z' or >= 'a' and <= 'z' or >= '0' and <= '9') + { + labelLength++; + previousWasHyphen = false; + } + else if (character == '-' && labelLength > 0) + { + labelLength++; + previousWasHyphen = true; + } + else + { + error = "The configured Datadog site is not valid"; + return false; + } + + if (labelLength > 63) + { + error = "The configured Datadog site is not valid"; + return false; + } + } + + if (labelLength == 0 || previousWasHyphen) + { + error = "The configured Datadog site is not valid"; + return false; + } + + normalizedSite = StringUtil.ToLowerInvariant(trimmedSite); + return true; + } + /// /// Returns the URI to request configuration for . The environment is /// added as a query parameter rather than baked into the endpoint, because it can change while diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpHeaderHelper.cs b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpHeaderHelper.cs new file mode 100644 index 000000000000..24c6726ce1c4 --- /dev/null +++ b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpHeaderHelper.cs @@ -0,0 +1,45 @@ +// +// 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. +// + +#nullable enable + +using System.Collections.Generic; +using Datadog.Trace.HttpOverStreams; + +namespace Datadog.Trace.FeatureFlags.Evp; + +/// +/// Headers used when a Feature Flags event is delivered through a local EVP relay. +/// +internal sealed class FeatureFlagsEvpHeaderHelper : HttpHeaderHelperBase +{ + internal const string EvpSubdomainHeader = "X-Datadog-EVP-Subdomain"; + internal const string EvpSubdomain = "event-platform-intake"; + internal const string EvpOriginHeader = "DD-EVP-ORIGIN"; + internal const string EvpOrigin = "dd-trace-dotnet"; + internal const string EvpOriginVersionHeader = "DD-EVP-ORIGIN-VERSION"; + + public static readonly FeatureFlagsEvpHeaderHelper Instance = new(); + + private FeatureFlagsEvpHeaderHelper() + { + DefaultHeaders = + [ + .. AgentHttpHeaderNames.MinimalHeaders, + new(EvpSubdomainHeader, EvpSubdomain), + new(EvpOriginHeader, EvpOrigin), + new(EvpOriginVersionHeader, TracerConstants.ThreePartVersion), + ]; + HttpSerializedDefaultHeaders = + $"{AgentHttpHeaderNames.HttpSerializedMinimalHeaders}" + + $"{EvpSubdomainHeader}: {EvpSubdomain}{DatadogHttpValues.CrLf}" + + $"{EvpOriginHeader}: {EvpOrigin}{DatadogHttpValues.CrLf}" + + $"{EvpOriginVersionHeader}: {TracerConstants.ThreePartVersion}{DatadogHttpValues.CrLf}"; + } + + public override KeyValuePair[] DefaultHeaders { get; } + + protected override string HttpSerializedDefaultHeaders { get; } +} diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs index ecceeda61580..be25f5fb3ffc 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs @@ -6,79 +6,604 @@ #nullable enable using System; +using System.Collections.Generic; +using System.Diagnostics.CodeAnalysis; +using System.IO; +using System.Net; +#if NETCOREAPP +using System.Net.Http; +#endif +using System.Net.Sockets; using System.Threading; using System.Threading.Tasks; using Datadog.Trace.Agent; +using Datadog.Trace.Agent.DiscoveryService; using Datadog.Trace.Agent.Transports; using Datadog.Trace.Configuration; +using Datadog.Trace.FeatureFlags.Agentless; using Datadog.Trace.HttpOverStreams; +using Datadog.Trace.Logging; using Datadog.Trace.SourceGenerators; +using Datadog.Trace.Telemetry; using Datadog.Trace.Vendors.Newtonsoft.Json; namespace Datadog.Trace.FeatureFlags.Evp; /// -/// Owns the historical fixed-v2 Feature Flags event sender independently of exposure batching. +/// Selects and sends to a Feature Flags EVP route without changing the product payload. /// internal sealed class FeatureFlagsEvpTransport : IDisposable { internal const string ExposureIntakePath = "api/v2/exposures"; + internal const string FlagEvaluationIntakePath = "api/v2/flagevaluation"; + internal const string EventPlatformProxyV4 = "evp_proxy/v4"; internal const string EventPlatformProxyV2 = "evp_proxy/v2"; - private readonly object _settingsLock = new(); + private static readonly TimeSpan InitialDiscoveryWait = TimeSpan.FromSeconds(5); + private static readonly TimeSpan RouteRecoveryCooldown = TimeSpan.FromSeconds(30); + private static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor(typeof(FeatureFlagsEvpTransport)); + + private readonly FeatureFlagsSource _source; + private readonly IApiRequestFactory? _directRequestFactory; + private readonly IDiscoveryService _discoveryService; + private readonly Action _discoveryCallback; + private readonly Action? _warningSink; private readonly IDisposable? _settingsSubscription; - private IApiRequestFactory _localRequestFactory; + private readonly TimeSpan _initialDiscoveryWait; + private readonly TimeSpan _routeRecoveryCooldown; + private readonly Func _utcNow; + private readonly bool _discoverySubscribed; + private readonly object _settingsLock = new(); + private readonly IDatadogLogger _log = Log; + + private LocalEndpoint _localEndpoint; + private int _directIsSticky; private int _disposed; + private int _localRecoveryProbeInProgress; + private int _unavailableWarningLogged; + private long _localUnavailableUntilUtcTicks; - internal FeatureFlagsEvpTransport(TracerSettings settings) + internal FeatureFlagsEvpTransport(TracerSettings settings, IDiscoveryService discoveryService, Action? warningSink = null) { - _localRequestFactory = CreateLocalRequestFactory(settings.Manager.InitialExporterSettings); + _source = settings.FeatureFlags.Source; + var exporter = settings.Manager.InitialExporterSettings; + _localEndpoint = new(CreateLocalRequestFactory(exporter, _source == FeatureFlagsSource.RemoteConfig), exporter); + _discoveryService = discoveryService; + _discoveryCallback = UpdateAgentConfiguration; + _warningSink = warningSink; + _initialDiscoveryWait = InitialDiscoveryWait; + _routeRecoveryCooldown = RouteRecoveryCooldown; + _utcNow = static () => DateTimeOffset.UtcNow; + + if (_source == FeatureFlagsSource.Agentless) + { + _directRequestFactory = CreateDirectRequestFactory(settings.FeatureFlags, out var invalidSite); + if (invalidSite) + { + WarnInvalidSite(); + } + + // The shared discovery service owns its own bounded retry/backoff loop. Event flushes + // wait for its first result once and then consume later callbacks; they never start an + // independent /info request on every flush. + if (discoveryService is NullDiscoveryService) + { + _localEndpoint.Discovery.TrySetResult(false); + } + else + { + _discoveryService.SubscribeToChanges(_discoveryCallback); + _discoverySubscribed = true; + } + } + else if (_source == FeatureFlagsSource.RemoteConfig) + { + // Preserve the historical Remote Config transport: fixed EVP v2, no /info discovery, + // and no direct credentials. + _localEndpoint.ProxyEndpoint = EventPlatformProxyV2; + _localEndpoint.Discovery.TrySetResult(true); + } + _settingsSubscription = settings.Manager.SubscribeToChanges(changes => { - if (changes.UpdatedExporter is { } exporter) + if (changes.UpdatedExporter is { } updatedExporter) { - lock (_settingsLock) - { - if (_disposed == 0) - { - Interlocked.Exchange(ref _localRequestFactory, CreateLocalRequestFactory(exporter)); - } - } + UpdateExporter(updatedExporter); } }); } [TestingOnly] - internal FeatureFlagsEvpTransport(IApiRequestFactory localRequestFactory) + internal FeatureFlagsEvpTransport( + FeatureFlagsSource source, + IApiRequestFactory localRequestFactory, + IApiRequestFactory? directRequestFactory, + IDiscoveryService discoveryService, + string? initialLocalProxyEndpoint = null, + bool initialDiscoveryKnown = true, + TimeSpan? initialDiscoveryWait = null, + TimeSpan? routeRecoveryCooldown = null, + Func? utcNow = null, + Action? warningSink = null, + ExporterSettings? exporterSettings = null, + IDatadogLogger? logger = null) + { + _source = source; + _localEndpoint = new(localRequestFactory, exporterSettings); + _log = logger ?? Log; + _directRequestFactory = source == FeatureFlagsSource.Agentless ? directRequestFactory : null; + _discoveryService = discoveryService; + _discoveryCallback = UpdateAgentConfiguration; + _warningSink = warningSink; + _initialDiscoveryWait = initialDiscoveryWait ?? InitialDiscoveryWait; + _routeRecoveryCooldown = routeRecoveryCooldown ?? RouteRecoveryCooldown; + _utcNow = utcNow ?? (static () => DateTimeOffset.UtcNow); + + if (source == FeatureFlagsSource.RemoteConfig) + { + _localEndpoint.ProxyEndpoint = EventPlatformProxyV2; + _localEndpoint.Discovery.TrySetResult(true); + } + else + { + _localEndpoint.ProxyEndpoint = initialLocalProxyEndpoint; + if (initialDiscoveryKnown || discoveryService is NullDiscoveryService) + { + _localEndpoint.Discovery.TrySetResult(true); + } + + if (discoveryService is not NullDiscoveryService) + { + _discoveryService.SubscribeToChanges(_discoveryCallback); + _discoverySubscribed = true; + } + } + } + + private enum NetworkFailure { - _localRequestFactory = localRequestFactory; + None, + DefinitivePreSend, + Ambiguous, + } + + internal static KeyValuePair[] GetDirectHeaders(string apiKey) => + [ + new(TelemetryConstants.ApiKeyHeader, apiKey), + new(FeatureFlagsEvpHeaderHelper.EvpOriginHeader, FeatureFlagsEvpHeaderHelper.EvpOrigin), + new(FeatureFlagsEvpHeaderHelper.EvpOriginVersionHeader, TracerConstants.ThreePartVersion), + new(HttpHeaderNames.TracingEnabled, "false"), + ]; + + [SuppressMessage("Performance", "CA1859:Use concrete types when possible for improved performance", Justification = "Different implementation types are returned for different TFMs")] + internal static IApiRequestFactory? CreateDirectRequestFactory(FeatureFlagsSettings settings) + => CreateDirectRequestFactory(settings, out _); + + [SuppressMessage("Performance", "CA1859:Use concrete types when possible for improved performance", Justification = "Different implementation types are returned for different TFMs")] + private static IApiRequestFactory? CreateDirectRequestFactory(FeatureFlagsSettings settings, out bool invalidSite) + { + invalidSite = false; + if (string.IsNullOrWhiteSpace(settings.ApiKey)) + { + return null; + } + + if (!AgentlessEndpoint.TryNormalizeSite(settings.Site, out var site, out _)) + { + invalidSite = true; + return null; + } + + var expectedHost = $"event-platform-intake.{site}"; + if (!Uri.TryCreate($"https://{expectedHost}", UriKind.Absolute, out var endpoint) + || endpoint.Scheme != Uri.UriSchemeHttps + || !string.Equals(endpoint.IdnHost, expectedHost, StringComparison.Ordinal)) + { + invalidSite = true; + return null; + } + + var headers = GetDirectHeaders(settings.ApiKey!); +#if NETCOREAPP + // The default HttpClientHandler honours the runtime's HTTPS_PROXY and NO_PROXY settings. + // Redirects stay disabled so the API key is sent only to the configured intake host. + return new HttpClientRequestFactory(endpoint, headers, timeout: TimeSpan.FromSeconds(5), allowAutoRedirect: false); +#else + // HttpWebRequest uses the platform default proxy and bypass list. Redirects stay disabled + // so the API key is sent only to the configured intake host. + return new ApiWebRequestFactory(endpoint, headers, timeout: TimeSpan.FromSeconds(5), allowAutoRedirect: false); +#endif } [TestingAndPrivateOnly] - internal static IApiRequestFactory CreateLocalRequestFactory(ExporterSettings exporterSettings) + internal static IApiRequestFactory CreateLocalRequestFactory(ExporterSettings exporterSettings, bool allowAutoRedirect = false) => AgentTransportStrategy.Get( exporterSettings, - productName: "FeatureFlags exposure", + productName: "Feature Flags EVP", tcpTimeout: TimeSpan.FromSeconds(5), - httpHeaderHelper: EventPlatformHeaderHelper.Instance); + httpHeaderHelper: FeatureFlagsEvpHeaderHelper.Instance, + allowAutoRedirect: allowAutoRedirect); + + [TestingAndPrivateOnly] + internal static bool IsDefinitivePreSendSocketFailure(SocketException exception) + => exception.SocketErrorCode is SocketError.HostNotFound + or SocketError.TryAgain + or SocketError.ConnectionRefused + or SocketError.NetworkUnreachable + or SocketError.HostUnreachable + or SocketError.AddressNotAvailable + || exception.ErrorCode == 2 // ENOENT for a missing Unix domain socket + || exception.ErrorCode == 10061; // WSAECONNREFUSED + + private static NetworkFailure ClassifyNetworkFailure(Exception exception) + { + var isNetworkFailure = false; + for (Exception? current = exception; current is not null; current = current.InnerException) + { + switch (current) + { + case SocketException socketException: + return IsDefinitivePreSendSocketFailure(socketException) + ? NetworkFailure.DefinitivePreSend + : NetworkFailure.Ambiguous; + case FileNotFoundException: + return NetworkFailure.DefinitivePreSend; + case WebException webException: + switch (webException.Status) + { + case WebExceptionStatus.ConnectFailure: + case WebExceptionStatus.NameResolutionFailure: + case WebExceptionStatus.ProxyNameResolutionFailure: + return NetworkFailure.DefinitivePreSend; + case WebExceptionStatus.ConnectionClosed: + case WebExceptionStatus.KeepAliveFailure: + case WebExceptionStatus.PipelineFailure: + case WebExceptionStatus.ReceiveFailure: + case WebExceptionStatus.RequestCanceled: + case WebExceptionStatus.SendFailure: + case WebExceptionStatus.Timeout: + return NetworkFailure.Ambiguous; + } + + isNetworkFailure = true; + break; +#if NETCOREAPP + case HttpRequestException: +#endif + case IOException: + case OperationCanceledException: + case TimeoutException: + isNetworkFailure = true; + break; + } + } + + return isNetworkFailure ? NetworkFailure.Ambiguous : NetworkFailure.None; + } public void Dispose() { - if (Interlocked.Exchange(ref _disposed, 1) == 0) + if (Interlocked.Exchange(ref _disposed, 1) != 0) { - _settingsSubscription?.Dispose(); + return; + } + + _settingsSubscription?.Dispose(); + lock (_settingsLock) + { + _localEndpoint.Discovery.TrySetResult(false); + if (_discoverySubscribed) + { + _discoveryService.RemoveSubscription(_discoveryCallback); + } } } - internal async Task SendAsync(T payload, string intakePath, JsonSerializerSettings serializerSettings) + internal Task SendAsync(T payload, string intakePath, JsonSerializerSettings serializerSettings) + => SendAsync(intakePath, request => request.PostAsJsonAsync(payload, MultipartCompression.GZip, serializerSettings)); + + private void UpdateExporter(ExporterSettings exporter) + { + lock (_settingsLock) + { + if (_disposed != 0) + { + return; + } + + var replacement = new LocalEndpoint(CreateLocalRequestFactory(exporter, _source == FeatureFlagsSource.RemoteConfig), exporter); + if (_source == FeatureFlagsSource.RemoteConfig) + { + replacement.ProxyEndpoint = EventPlatformProxyV2; + replacement.Discovery.TrySetResult(true); + } + else if (_discoveryService is NullDiscoveryService) + { + replacement.Discovery.TrySetResult(false); + } + + var previous = Interlocked.Exchange(ref _localEndpoint, replacement); + previous.Discovery.TrySetResult(false); + if (_discoverySubscribed) + { + // The shared discovery callback may precede this settings callback. Subscribe + // again to consume any result already obtained for the new immutable settings. + _discoveryService.RemoveSubscription(_discoveryCallback); + _discoveryService.SubscribeToChanges(_discoveryCallback); + } + } + } + + private async Task SendAsync(string intakePath, Func> sendAsync) { if (Volatile.Read(ref _disposed) != 0) { return; } - var factory = Volatile.Read(ref _localRequestFactory); - var request = factory.Create(factory.GetEndpoint($"{EventPlatformProxyV2}/{intakePath}")); - using var response = await request.PostAsJsonAsync(payload, MultipartCompression.GZip, serializerSettings).ConfigureAwait(false); + if (Volatile.Read(ref _directIsSticky) != 0) + { + await SendDirectAsync(intakePath, sendAsync).ConfigureAwait(false); + return; + } + + var local = Volatile.Read(ref _localEndpoint); + var localProxyEndpoint = local.ProxyEndpoint; + if (localProxyEndpoint is not null) + { + await TrySendLocalAsync(intakePath, local, localProxyEndpoint, sendAsync).ConfigureAwait(false); + return; + } + + if (_source == FeatureFlagsSource.Agentless) + { + // Endpoint replacement wakes waiters on the previous generation. Keep waiting + // for the replacement's capabilities without extending this batch's deadline. + var discoveryDeadline = Task.Delay(_initialDiscoveryWait); + do + { + await Task.WhenAny(local.Discovery.Task, discoveryDeadline).ConfigureAwait(false); + if (Volatile.Read(ref _disposed) != 0) + { + return; + } + + local = Volatile.Read(ref _localEndpoint); + if (discoveryDeadline.IsCompleted) + { + local.Discovery.TrySetResult(false); + break; + } + } + while (!local.Discovery.Task.IsCompleted); + } + + // Use the checked factory snapshot, not a newer unvalidated generation. + localProxyEndpoint = local.ProxyEndpoint; + if (localProxyEndpoint is not null) + { + await TrySendLocalAsync(intakePath, local, localProxyEndpoint, sendAsync).ConfigureAwait(false); + return; + } + + if (_source == FeatureFlagsSource.Agentless && _directRequestFactory is not null) + { + Interlocked.Exchange(ref _directIsSticky, 1); + await SendDirectAsync(intakePath, sendAsync).ConfigureAwait(false); + return; + } + + if (Interlocked.Exchange(ref _unavailableWarningLogged, 1) == 0) + { + WarnUnavailable(); + } + } + + private void WarnInvalidSite() + { + const string Message = "Feature Flags direct event delivery is disabled because DD_SITE is not a valid DNS site suffix"; + if (_warningSink is { } warningSink) + { + warningSink(Message); + } + else + { + Log.Warning("Feature Flags direct event delivery is disabled because DD_SITE is not a valid DNS site suffix"); + } + } + + private void WarnUnavailable() + { + const string Message = "Feature Flags event delivery is unavailable because no compatible local EVP route or direct intake credentials are available"; + if (_warningSink is { } warningSink) + { + warningSink(Message); + } + else + { + Log.Warning("Feature Flags event delivery is unavailable because no compatible local EVP route or direct intake credentials are available"); + } + } + + private void UpdateAgentConfiguration(AgentConfiguration configuration) + { + var local = Volatile.Read(ref _localEndpoint); + if (local.ExporterSettings is not null && !ReferenceEquals(local.ExporterSettings, configuration.DiscoverySettings)) + { + return; + } + + // Agentless event delivery requires the Agent to preserve the logical producer identity. + // Older Agents can advertise EVP while silently dropping these headers, so use direct + // intake instead unless /info explicitly advertises both forwarding capabilities. + var advertisedEndpoint = configuration.EventPlatformProxyEndpoint; + var endpoint = configuration.EventPlatformProxySupportsEvpOriginHeaders + && (string.Equals(advertisedEndpoint, EventPlatformProxyV4, StringComparison.OrdinalIgnoreCase) + || string.Equals(advertisedEndpoint, EventPlatformProxyV2, StringComparison.OrdinalIgnoreCase)) + ? advertisedEndpoint + : null; + + local.ProxyEndpoint = endpoint; + if (endpoint is not null) + { + Interlocked.Exchange(ref _localUnavailableUntilUtcTicks, 0); + } + + local.Discovery.TrySetResult(true); + } + + private async Task TrySendLocalAsync(string intakePath, LocalEndpoint local, string localProxyEndpoint, Func> sendAsync) + { + var unavailableUntilUtcTicks = Interlocked.Read(ref _localUnavailableUntilUtcTicks); + var isRecoveryProbe = unavailableUntilUtcTicks != 0; + if (isRecoveryProbe + && (_utcNow().UtcTicks < unavailableUntilUtcTicks + || Interlocked.CompareExchange(ref _localRecoveryProbeInProgress, 1, 0) != 0)) + { + if (Interlocked.Exchange(ref _unavailableWarningLogged, 1) == 0) + { + WarnUnavailable(); + } + + return; + } + + try + { + await SendLocalAsync(intakePath, local.Factory, localProxyEndpoint, sendAsync).ConfigureAwait(false); + } + finally + { + if (isRecoveryProbe) + { + Interlocked.Exchange(ref _localRecoveryProbeInProgress, 0); + } + } + } + + private async Task SendLocalAsync(string intakePath, IApiRequestFactory localFactory, string localProxyEndpoint, Func> sendAsync) + { + var endpoint = localFactory.GetEndpoint($"{localProxyEndpoint}/{intakePath}"); + + try + { + var request = localFactory.Create(endpoint); + using var response = await sendAsync(request).ConfigureAwait(false); + if (response.StatusCode is >= 200 and < 300) + { + Interlocked.Exchange(ref _localUnavailableUntilUtcTicks, 0); + return; + } + + // The Agent contract proves these statuses mean the proxy route did not accept the + // payload. An upstream 403 is not safe to replay because it may have been forwarded. + if (response.StatusCode is 404 or 405) + { + if (LeaveLocalRoute()) + { + await SendDirectAsync(intakePath, sendAsync).ConfigureAwait(false); + } + else + { + Log.Warning("Feature Flags local EVP request failed with HTTP status code {StatusCode}", response.StatusCode); + } + + return; + } + + // Other responses may have come from upstream after the Agent accepted the payload. + // Never replay this batch, but leave the failed route for future Agentless batches. + LeaveLocalRoute(); + Log.Warning("Feature Flags local EVP request failed with HTTP status code {StatusCode}", response.StatusCode); + } + catch (Exception ex) when (ClassifyNetworkFailure(ex) is NetworkFailure.DefinitivePreSend) + { + if (LeaveLocalRoute()) + { + await SendDirectAsync(intakePath, sendAsync).ConfigureAwait(false); + } + else + { + Log.ErrorSkipTelemetry(ex, "Feature Flags local EVP request failed before the payload was sent"); + } + } + catch (Exception ex) when (ClassifyNetworkFailure(ex) is NetworkFailure.Ambiguous) + { + // The local relay may have received this payload. Switch only future payloads so the + // current one can never be duplicated across the local and direct routes. + LeaveLocalRoute(); + + Log.ErrorSkipTelemetry(ex, "Feature Flags local EVP request failed ambiguously; the current event batch will not be replayed"); + } + } + + private bool LeaveLocalRoute() + { + if (_source != FeatureFlagsSource.Agentless) + { + return false; + } + + if (_directRequestFactory is not null) + { + Interlocked.Exchange(ref _directIsSticky, 1); + return true; + } + + var unavailableUntil = _utcNow().Add(_routeRecoveryCooldown).UtcTicks; + Interlocked.Exchange(ref _localUnavailableUntilUtcTicks, unavailableUntil); + return false; + } + + private async Task SendDirectAsync(string intakePath, Func> sendAsync) + { + var directFactory = _directRequestFactory; + if (directFactory is null) + { + return; + } + + var endpoint = directFactory.GetEndpoint(intakePath); + try + { + var request = directFactory.Create(endpoint); + using var response = await sendAsync(request).ConfigureAwait(false); + if (response.StatusCode is < 200 or >= 300) + { + const string Message = "Feature Flags direct EVP request to {Endpoint} failed with HTTP status code {StatusCode} after 1 attempt; the batch will not be replayed. See https://docs.datadoghq.com/feature_flags/"; + if (response.StatusCode == 400) + { + _log.Error(Message, intakePath, response.StatusCode); + } + else + { + _log.ErrorSkipTelemetry(Message, intakePath, response.StatusCode); + } + } + } + catch (Exception ex) + { + // Direct intake is terminal: never loop a failed direct request back through the Agent. + Log.ErrorSkipTelemetry(ex, "Feature Flags direct EVP request failed"); + } + } + + // Snapshot the factory together with its capabilities: a send can finish on the old Agent, + // but must never use its capabilities with the replacement Agent's factory. + private sealed class LocalEndpoint(IApiRequestFactory factory, ExporterSettings? exporterSettings) + { + private string? _proxyEndpoint; + + public IApiRequestFactory Factory { get; } = factory; + + public ExporterSettings? ExporterSettings { get; } = exporterSettings; + + public TaskCompletionSource Discovery { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + + public string? ProxyEndpoint + { + get => Volatile.Read(ref _proxyEndpoint); + set => Volatile.Write(ref _proxyEndpoint, value); + } } } diff --git a/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs b/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs index 70a803e939e1..2eff2f9792d7 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs @@ -9,6 +9,7 @@ using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; +using Datadog.Trace.Agent.DiscoveryService; using Datadog.Trace.Configuration; using Datadog.Trace.FeatureFlags.Agentless; using Datadog.Trace.FeatureFlags.Evp; @@ -76,7 +77,8 @@ internal sealed class FeatureFlagsModule : IDisposable internal FeatureFlagsModule( TracerSettings settings, IRcmSubscriptionManager rcmSubscriptionManager, - Func? agentlessSourceFactory = null) + Func? agentlessSourceFactory = null, + IDiscoveryService? discoveryService = null) { _settings = settings.FeatureFlags; _settingsManager = settings.Manager; @@ -86,7 +88,7 @@ internal FeatureFlagsModule( _agentlessSourceFactory = agentlessSourceFactory ?? (static module => AgentlessConfigurationSource.Create(module._settings, module._settingsManager, module.ApplyConfiguration)); _rcmSubscriptionManager = rcmSubscriptionManager; - _evpTransport = new FeatureFlagsEvpTransport(settings); + _evpTransport = new FeatureFlagsEvpTransport(settings, discoveryService ?? NullDiscoveryService.Instance); Log.Debug("FeatureFlagsModule ENABLED with source {Source}", _settings.Source); } @@ -109,14 +111,15 @@ internal FeatureFlagsModule( public static FeatureFlagsModule? Create( TracerSettings settings, IRcmSubscriptionManager rcmSubscriptionManager, - Func? agentlessSourceFactory = null) + Func? agentlessSourceFactory = null, + IDiscoveryService? discoveryService = null) { if (!settings.FeatureFlags.Enabled) { return null; } - var module = new FeatureFlagsModule(settings, rcmSubscriptionManager, agentlessSourceFactory); + var module = new FeatureFlagsModule(settings, rcmSubscriptionManager, agentlessSourceFactory, discoveryService); // Subscribing from here rather than the constructor, so the callback can only ever reach // a fully constructed module. diff --git a/tracer/src/Datadog.Trace/TracerManagerFactory.cs b/tracer/src/Datadog.Trace/TracerManagerFactory.cs index f984352592cc..d634f00002e5 100644 --- a/tracer/src/Datadog.Trace/TracerManagerFactory.cs +++ b/tracer/src/Datadog.Trace/TracerManagerFactory.cs @@ -191,7 +191,7 @@ internal TracerManager CreateTracerManager( } } - featureFlags = FeatureFlagsModule.Create(settings, RcmSubscriptionManager.Instance); + featureFlags = FeatureFlagsModule.Create(settings, RcmSubscriptionManager.Instance, discoveryService: discoveryService); return CreateTracerManagerFrom( settings, diff --git a/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs b/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs index 9b84a117f5fa..4276217b8aa4 100644 --- a/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs +++ b/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs @@ -1,4 +1,4 @@ -// +// // 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. // diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/AgentlessEndpointTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/AgentlessEndpointTests.cs index 6e024bfe7456..7f1252252ff1 100644 --- a/tracer/test/Datadog.Trace.Tests/FeatureFlags/AgentlessEndpointTests.cs +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/AgentlessEndpointTests.cs @@ -19,6 +19,7 @@ public class AgentlessEndpointTests [Theory] [InlineData("datadoghq.com", "https://ufc-server.ff-cdn.datadoghq.com" + DefaultPath)] [InlineData("DATADOGHQ.COM", "https://ufc-server.ff-cdn.datadoghq.com" + DefaultPath)] // site is lowercased + [InlineData(" datadoghq.com ", "https://ufc-server.ff-cdn.datadoghq.com" + DefaultPath)] // surrounding whitespace is trimmed [InlineData("datad0g.com", "https://ufc-server.ff-cdn.datad0g.com" + DefaultPath)] // staging [InlineData("ddog-gov.com", "https://ufc-server.ff-cdn.ddog-gov.com" + DefaultPath)] // govcloud public void DerivesManagedEndpointFromSite(string site, string expected) @@ -140,6 +141,11 @@ public void RejectsWhitespaceOnlySiteWithoutBaseUrl() [InlineData("datadoghq.com/../evil")] // a path escapes the host [InlineData("datadoghq.com?x=1")] // a query escapes the host [InlineData("datadoghq.com#f")] // a fragment escapes the host + [InlineData("datadoghq.com\\attacker.example")] // a backslash can be normalized as a URL separator + [InlineData("dátadoghq.com")] // managed endpoints do not implicitly convert Unicode to IDN + [InlineData("-datadoghq.com")] + [InlineData("datadoghq.com-")] + [InlineData("datadoghq..com")] public void RejectsMalformedSiteWithoutThrowing(string site) { AgentlessEndpoint.TryCreate(site, baseUrl: null, out var endpoint, out var error) diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs index 12f79f8db2ba..5fe8f3d10b70 100644 --- a/tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs @@ -13,11 +13,13 @@ using Datadog.Trace.Agent; using Datadog.Trace.Agent.Transports; using Datadog.Trace.Configuration; +using Datadog.Trace.FeatureFlags; using Datadog.Trace.FeatureFlags.Evp; using Datadog.Trace.FeatureFlags.Exposure; using Datadog.Trace.FeatureFlags.Exposure.Model; using Datadog.Trace.HttpOverStreams; using Datadog.Trace.TestHelpers.TransportHelpers; +using Datadog.Trace.Tests.Agent; using Datadog.Trace.Vendors.Newtonsoft.Json; using Datadog.Trace.Vendors.Newtonsoft.Json.Linq; using FluentAssertions; @@ -52,7 +54,7 @@ public async Task DisposeFlushesEventsQueuedWhileAnEarlierBatchIsSending() local.RequestsSent.Should().HaveCount(2, "the second exposure is sent by the shutdown flush"); var finalRequest = local.RequestsSent[1].Should().BeOfType().Subject; - finalRequest.Endpoint.AbsolutePath.Should().Be("/evp_proxy/v2/api/v2/exposures"); + finalRequest.Endpoint.AbsolutePath.Should().Be("/evp_proxy/v4/api/v2/exposures"); finalRequest.Compression.Should().Be(MultipartCompression.GZip); var payload = JObject.Parse(finalRequest.PayloadJson); payload["context"]!["service"]!.Value().Should().NotBeNullOrEmpty(); @@ -98,7 +100,12 @@ public async Task SendExposureAfterDisposeIsIgnored() } private static FeatureFlagsEvpTransport CreateLocalTransport(TestRequestFactory local) - => new(local); + => new( + FeatureFlagsSource.Agentless, + local, + directRequestFactory: null, + discoveryService: new DiscoveryServiceMock(), + initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); private static ExposureEvent CreateExposure(string flag) => new( diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs new file mode 100644 index 000000000000..757d6bfefb55 --- /dev/null +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs @@ -0,0 +1,872 @@ +// +// 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. +// + +#nullable enable + +using System; +using System.Collections.Generic; +using System.Collections.Specialized; +using System.IO; +using System.Linq; +using System.Net; +using System.Net.Sockets; +using System.Threading.Tasks; +using Datadog.Trace.Agent; +using Datadog.Trace.Agent.DiscoveryService; +using Datadog.Trace.Agent.Transports; +using Datadog.Trace.ClrProfiler.AutoInstrumentation.ManualInstrumentation; +using Datadog.Trace.Configuration; +using Datadog.Trace.Configuration.ConfigurationSources; +using Datadog.Trace.Configuration.Telemetry; +using Datadog.Trace.FeatureFlags; +using Datadog.Trace.FeatureFlags.Evp; +using Datadog.Trace.HttpOverStreams; +using Datadog.Trace.Logging; +using Datadog.Trace.PlatformHelpers; +using Datadog.Trace.Telemetry; +using Datadog.Trace.TestHelpers; +using Datadog.Trace.TestHelpers.TransportHelpers; +using Datadog.Trace.Tests.Agent; +using Datadog.Trace.Vendors.Newtonsoft.Json; +using FluentAssertions; +using Moq; +using Xunit; + +namespace Datadog.Trace.Tests.FeatureFlags; + +[Collection(nameof(WebRequestCollection))] +public class FeatureFlagsEvpTransportTests +{ + private static readonly JsonSerializerSettings SerializerSettings = new(); + + public static IEnumerable AmbiguousFailures() + { + yield return [new IOException("broken pipe")]; + yield return [new TimeoutException("timeout")]; + yield return [new TaskCanceledException("HTTP timeout")]; + yield return [new SocketException((int)SocketError.ConnectionReset)]; + yield return [new WebException("send failed", WebExceptionStatus.SendFailure)]; + } + + public static IEnumerable DefinitivePreSendFailures() + { + yield return [new SocketException((int)SocketError.HostNotFound)]; + yield return [new SocketException((int)SocketError.TryAgain)]; + yield return [new SocketException((int)SocketError.ConnectionRefused)]; + yield return [new SocketException((int)SocketError.NetworkUnreachable)]; + yield return [new SocketException((int)SocketError.HostUnreachable)]; + yield return [new SocketException((int)SocketError.AddressNotAvailable)]; + yield return [new SocketException(10061)]; // Windows WSAECONNREFUSED + yield return [new FileNotFoundException("missing Unix domain socket")]; + yield return [new WebException("refused", WebExceptionStatus.ConnectFailure)]; + } + + [Theory] + [InlineData(false, true, false)] + [InlineData(true, true, false)] + [InlineData(false, false, false)] + [InlineData(false, true, true)] + public async Task RuntimeEndpointChangeRequiresCapabilitiesFromReplacementAgent(bool discoverBeforeTransportUpdate, bool replacementSupportsIdentity, bool sendBeforeUpdate) + { + const string CapableInfo = "{\"endpoints\":[\"evp_proxy/v4\"],\"evp_proxy_allowed_headers\":[\"DD-EVP-ORIGIN\",\"DD-EVP-ORIGIN-VERSION\"]}"; + using var replacement = new HttpListener(); + var replacementUrl = $"http://127.0.0.1:{TcpPortProvider.GetOpenPort()}/"; + replacement.Prefixes.Add(replacementUrl); + replacement.Start(); + var received = replacement.GetContextAsync(); + var settings = CreateSettings( + (ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "agentless"), + (ConfigurationKeys.AgentUri, "http://old-agent:8126/")); + var factoryA = CreateFactory("http://old-agent:8126/", uri => new TestApiRequest(uri, responseContent: CapableInfo)); + var factoryB = CreateFactory(replacementUrl, uri => new TestApiRequest(uri, responseContent: replacementSupportsIdentity ? CapableInfo : "{\"endpoints\":[\"evp_proxy/v4\"]}")); + await using var discovery = new DiscoveryService(factoryA, new ServiceRemappingHash(null), 1, 1, 30_000, autoStartLoop: false, exporterSettings: settings.Manager.InitialExporterSettings); + // Mirrors the production subscription order: shared discovery subscribes before EVP. + using var discoverySettings = settings.Manager.SubscribeToChanges(changes => + { + if (changes.UpdatedExporter is { } exporter) + { + discovery.UpdateRequestFactory(factoryB, exporter); + if (discoverBeforeTransportUpdate) + { + discovery.RunOneIterationAsync(null).GetAwaiter().GetResult(); + } + } + }); + using var transport = new FeatureFlagsEvpTransport(settings, discovery); + Task? send = null; + if (sendBeforeUpdate) + { + send = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + send.IsCompleted.Should().BeFalse(); + } + else + { + await discovery.RunOneIterationAsync(null); + } + + try + { + settings.Manager.UpdateManualConfigurationSettings( + new ManualInstrumentationConfigurationSource(new Dictionary { { TracerSettingKeyConstants.AgentUriKey, new Uri(replacementUrl) } }, useDefaultSources: true), + NullConfigurationTelemetry.Instance).Should().BeTrue(); + send ??= transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + if (!discoverBeforeTransportUpdate) + { + if (sendBeforeUpdate) + { + // Give the waiter released from A a chance to run before publishing B. + var stillWaiting = Task.Delay(TimeSpan.FromMilliseconds(100)); + (await Task.WhenAny(send, stillWaiting)).Should().BeSameAs(stillWaiting, "replacing A must not abandon the current batch while B discovery is pending"); + } + + send.IsCompleted.Should().BeFalse("new endpoint discovery is still pending"); + await discovery.RunOneIterationAsync(null); + } + + if (replacementSupportsIdentity) + { + (await Task.WhenAny(received, Task.Delay(TimeSpan.FromSeconds(5)))).Should().BeSameAs(received); + var request = await received; + request.Request.Url!.AbsolutePath.Should().Be("/evp_proxy/v4/api/v2/exposures"); + request.Request.Headers[TelemetryConstants.ApiKeyHeader].Should().BeNull(); + await request.Request.InputStream.CopyToAsync(Stream.Null); + request.Response.StatusCode = 200; + request.Response.Close(); + } + + (await Task.WhenAny(send, Task.Delay(TimeSpan.FromSeconds(6)))).Should().BeSameAs(send); + await send; + received.IsCompleted.Should().Be(replacementSupportsIdentity, "the old Agent's capabilities cannot authorize an upload to the new Agent"); + factoryB.RequestsSent.Should().ContainSingle("an equal capability body still validates the new endpoint"); + } + finally + { + replacement.Close(); + try + { + await received; + } + catch (Exception ex) when (ex is HttpListenerException or ObjectDisposedException) + { + } + } + } + + [Theory] + [InlineData(FeatureFlagsEvpTransport.EventPlatformProxyV4, FeatureFlagsEvpTransport.ExposureIntakePath)] + [InlineData(FeatureFlagsEvpTransport.EventPlatformProxyV4, FeatureFlagsEvpTransport.FlagEvaluationIntakePath)] + [InlineData(FeatureFlagsEvpTransport.EventPlatformProxyV2, FeatureFlagsEvpTransport.ExposureIntakePath)] + [InlineData(FeatureFlagsEvpTransport.EventPlatformProxyV2, FeatureFlagsEvpTransport.FlagEvaluationIntakePath)] + [InlineData("EVP_PROXY/V4", FeatureFlagsEvpTransport.ExposureIntakePath)] + [InlineData("Evp_Proxy/V2", FeatureFlagsEvpTransport.FlagEvaluationIntakePath)] + public async Task DiscoverySelectsAdvertisedLocalRoute(string proxyEndpoint, string intakePath) + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + var discovery = new DiscoveryServiceMock(); + using var transport = CreateTransport(local, direct, discovery); + + discovery.TriggerChange(eventPlatformProxyEndpoint: proxyEndpoint); + await transport.SendAsync(new object(), intakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle() + .Which.Endpoint.AbsolutePath.Should().Be($"/{proxyEndpoint}/{intakePath}"); + direct.RequestsSent.Should().BeEmpty(); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task AgentlessLocalSelectionRequiresIdentityHeaderForwarding(bool supportsBothIdentityHeaders) + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + var discovery = new DiscoveryServiceMock(); + using var transport = CreateTransport(local, direct, discovery); + + discovery.TriggerChange( + eventPlatformProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4, + eventPlatformProxySupportsEvpOriginHeaders: supportsBothIdentityHeaders); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + if (supportsBothIdentityHeaders) + { + local.RequestsSent.Should().ContainSingle(); + direct.RequestsSent.Should().BeEmpty(); + } + else + { + local.RequestsSent.Should().BeEmpty("the Agent cannot preserve the logical SDK identity"); + direct.RequestsSent.Should().ContainSingle(); + } + } + + [Fact] + public async Task ConcurrentInitialSendsShareOneDiscoverySubscription() + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + var discovery = new DiscoveryServiceMock(); + using var transport = CreateTransport(local, direct, discovery, initialDiscoveryKnown: false); + + var exposure = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + var evaluation = transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + discovery.Callbacks.Should().ContainSingle("the transport consumes the tracer's shared /info discovery stream"); + discovery.TriggerChange(eventPlatformProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + await Task.WhenAll(exposure, evaluation); + + local.RequestsSent.Should().HaveCount(2); + direct.RequestsSent.Should().BeEmpty(); + } + + [Fact] + public async Task DirectRouteIsStickyAfterDiscoveryReportsNoCompatibleProxy() + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + var discovery = new DiscoveryServiceMock(); + using var transport = CreateTransport(local, direct, discovery); + + discovery.TriggerChange(eventPlatformProxyEndpoint: "v0.4/traces"); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + discovery.TriggerChange(eventPlatformProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().BeEmpty(); + direct.RequestsSent.Select(r => r.Endpoint.AbsolutePath).Should().Equal( + $"/{FeatureFlagsEvpTransport.ExposureIntakePath}", + $"/{FeatureFlagsEvpTransport.FlagEvaluationIntakePath}"); + } + + [Fact] + public async Task InitialDiscoveryTimeoutSelectsDirectAndStaysDirect() + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + using var transport = CreateTransport( + local, + direct, + initialDiscoveryKnown: false, + initialDiscoveryWait: TimeSpan.FromMilliseconds(10)); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().BeEmpty(); + direct.RequestsSent.Should().HaveCount(2); + } + + [Fact] + public async Task UnavailableRouteWarnsOnceAndRecoversWhenAgentAppears() + { + var warnings = new List(); + var local = CreateFactory("http://agent:8126/"); + var discovery = new DiscoveryServiceMock(); + using var transport = CreateTransport(local, direct: null, discovery, warningSink: warnings.Add); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + warnings.Should().ContainSingle(); + local.RequestsSent.Should().BeEmpty(); + + discovery.TriggerChange(eventPlatformProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle(); + } + + [Fact] + public async Task FailedLocalRouteWithoutDirectCredentialsRecoversAfterCooldown() + { + var now = new DateTimeOffset(2026, 9, 12, 0, 0, 0, TimeSpan.Zero); + var warnings = new List(); + var local = CreateFactory( + "http://agent:8126/", + uri => new ThrowingApiRequest(uri, new IOException("ambiguous reset")), + uri => new TestApiRequest(uri)); + using var transport = CreateTransport( + local, + direct: null, + initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4, + routeRecoveryCooldown: TimeSpan.FromMinutes(1), + utcNow: () => now, + warningSink: warnings.Add); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle("flushes must not probe a failed route during cooldown"); + warnings.Should().ContainSingle(); + + now = now.AddMinutes(1); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + local.RequestsSent.Should().HaveCount(3, "one post-cooldown probe restores normal local delivery"); + } + + [Fact] + public async Task ConcurrentFlushesShareOnePostCooldownRecoveryProbe() + { + var now = new DateTimeOffset(2026, 9, 12, 0, 0, 0, TimeSpan.Zero); + var probeStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseProbe = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var local = CreateFactory( + "http://agent:8126/", + uri => new ThrowingApiRequest(uri, new IOException("ambiguous reset")), + uri => new BlockingApiRequest(uri, probeStarted, releaseProbe)); + using var transport = CreateTransport( + local, + direct: null, + initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4, + routeRecoveryCooldown: TimeSpan.FromMinutes(1), + utcNow: () => now); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + now = now.AddMinutes(1); + + var probe = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + (await Task.WhenAny(probeStarted.Task, Task.Delay(TimeSpan.FromSeconds(1)))) + .Should().BeSameAs(probeStarted.Task); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + local.RequestsSent.Should().HaveCount(2, "a concurrent flush must not start a second recovery probe"); + + releaseProbe.TrySetResult(true); + await probe; + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().HaveCount(3, "a successful probe restores normal local delivery"); + } + + [Theory] + [InlineData(404)] + [InlineData(405)] + public async Task UnsupportedLocalRouteWithoutDirectCredentialsEntersCooldown(int statusCode) + { + var now = new DateTimeOffset(2026, 9, 12, 0, 0, 0, TimeSpan.Zero); + var local = CreateFactory( + "http://agent:8126/", + uri => new TestApiRequest(uri, statusCode), + uri => new TestApiRequest(uri)); + using var transport = CreateTransport( + local, + direct: null, + initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4, + routeRecoveryCooldown: TimeSpan.FromMinutes(1), + utcNow: () => now); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle(); + + now = now.AddMinutes(1); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + local.RequestsSent.Should().HaveCount(2); + } + + [Fact] + public async Task RemoteConfigurationAlwaysUsesHistoricalV2WithoutDiscoveryOrDirectCredentials() + { + var local = CreateFactory( + "http://agent:8126/", + uri => new TestApiRequest(uri, statusCode: 404), + uri => new ThrowingApiRequest(uri, new SocketException((int)SocketError.ConnectionRefused)), + uri => new TestApiRequest(uri)); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + var discovery = new DiscoveryServiceMock(); + using var transport = new FeatureFlagsEvpTransport(FeatureFlagsSource.RemoteConfig, local, direct, discovery); + + discovery.Callbacks.Should().BeEmpty("Remote Config must not change its historical transport contract"); + discovery.TriggerChange(eventPlatformProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + local.RequestsSent.Should().HaveCount(3); + local.RequestsSent.Should().OnlyContain( + request => request.Endpoint.AbsolutePath == $"/{FeatureFlagsEvpTransport.EventPlatformProxyV2}/{FeatureFlagsEvpTransport.ExposureIntakePath}"); + direct.RequestsSent.Should().BeEmpty(); + } + + [Fact] + public void RemoteConfigurationProductionConstructionDoesNotSubscribeToDiscovery() + { + var settings = CreateSettings( + (ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "remote_config"), + (ConfigurationKeys.ApiKey, "must-not-be-used")); + var discovery = new DiscoveryServiceMock(); + + using var transport = new FeatureFlagsEvpTransport(settings, discovery); + + discovery.Callbacks.Should().BeEmpty(); + } + + [Theory] + [InlineData(404)] + [InlineData(405)] + public async Task UnsupportedLocalRouteReplaysCurrentBatchDirectAndStaysDirect(int statusCode) + { + var local = CreateFactory("http://agent:8126/", uri => new TestApiRequest(uri, statusCode)); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + using var transport = CreateTransport(local, direct, initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle(); + direct.RequestsSent.Should().HaveCount(2, "the safe-to-replay batch and later batches use direct intake"); + } + + [Theory] + [InlineData(403)] + [InlineData(429)] + [InlineData(500)] + [InlineData(503)] + public async Task PotentiallyForwardedOrRetryableHttpFailureChangesOnlyFutureBatchesToDirect(int statusCode) + { + var local = CreateFactory( + "http://agent:8126/", + uri => new TestApiRequest(uri, statusCode)); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + using var transport = CreateTransport(local, direct, initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle(); + direct.RequestsSent.Should().ContainSingle("the failed batch is not replayed, but the next batch uses direct intake"); + } + + [Theory] + [InlineData(403)] + [InlineData(429)] + [InlineData(500)] + [InlineData(503)] + public async Task PotentiallyForwardedOrRetryableHttpFailureWithoutDirectCredentialsEntersCooldown(int statusCode) + { + var now = new DateTimeOffset(2026, 9, 12, 0, 0, 0, TimeSpan.Zero); + var local = CreateFactory( + "http://agent:8126/", + uri => new TestApiRequest(uri, statusCode), + uri => new TestApiRequest(uri)); + using var transport = CreateTransport( + local, + direct: null, + initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4, + routeRecoveryCooldown: TimeSpan.FromMinutes(1), + utcNow: () => now); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle("the failed route remains unavailable during cooldown"); + + now = now.AddMinutes(1); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().HaveCount(2, "one post-cooldown local probe restores delivery"); + } + + [Theory] + [MemberData(nameof(DefinitivePreSendFailures))] + public async Task DefinitivePreSendFailureReplaysCurrentBatchDirectAndStaysDirect(Exception failure) + { + var local = CreateFactory("http://agent:8126/", uri => new ThrowingApiRequest(uri, failure)); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + using var transport = CreateTransport(local, direct, initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle(); + direct.RequestsSent.Should().HaveCount(2); + } + + [Theory] + [MemberData(nameof(AmbiguousFailures))] + public async Task AmbiguousLocalFailureChangesOnlyFutureBatchesToDirect(Exception failure) + { + var local = CreateFactory("http://agent:8126/", uri => new ThrowingApiRequest(uri, failure)); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + using var transport = CreateTransport(local, direct, initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + direct.RequestsSent.Should().BeEmpty("an ambiguously failed batch may already have reached the relay"); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + local.RequestsSent.Should().ContainSingle(); + direct.RequestsSent.Should().ContainSingle("only a later batch is safe to send direct"); + } + + [Fact] + public async Task DirectFailureNeverLoopsBackToAgent() + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory( + "https://event-platform-intake.datadoghq.com/", + uri => new ThrowingApiRequest(uri, new IOException("direct reset"))); + using var transport = CreateTransport(local, direct); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + + direct.RequestsSent.Should().ContainSingle(); + local.RequestsSent.Should().BeEmpty(); + } + + [Fact] + public void DirectAndLocalHeadersHaveExactIdentityAndCredentialIsolation() + { + var directHeaders = FeatureFlagsEvpTransport.GetDirectHeaders("test-api-key").ToDictionary(x => x.Key, x => x.Value); + var localHeaders = FeatureFlagsEvpHeaderHelper.Instance.DefaultHeaders.ToDictionary(x => x.Key, x => x.Value); + + directHeaders.Should().HaveCount(4); + directHeaders.Should().Contain(HttpHeaderNames.TracingEnabled, "false"); + directHeaders.Should().Contain(TelemetryConstants.ApiKeyHeader, "test-api-key"); + directHeaders.Should().Contain(FeatureFlagsEvpHeaderHelper.EvpOriginHeader, FeatureFlagsEvpHeaderHelper.EvpOrigin); + directHeaders.Should().Contain(FeatureFlagsEvpHeaderHelper.EvpOriginVersionHeader, TracerConstants.ThreePartVersion); + directHeaders.Should().NotContainKey(FeatureFlagsEvpHeaderHelper.EvpSubdomainHeader); + + localHeaders.Should().Contain(FeatureFlagsEvpHeaderHelper.EvpSubdomainHeader, FeatureFlagsEvpHeaderHelper.EvpSubdomain); + localHeaders.Should().Contain(FeatureFlagsEvpHeaderHelper.EvpOriginHeader, FeatureFlagsEvpHeaderHelper.EvpOrigin); + localHeaders.Should().Contain(FeatureFlagsEvpHeaderHelper.EvpOriginVersionHeader, TracerConstants.ThreePartVersion); + localHeaders.Should().NotContainKey(TelemetryConstants.ApiKeyHeader); + localHeaders.Should().Contain(HttpHeaderNames.TracingEnabled, "false"); + } + + [Fact] + public async Task DisabledDiscoveryDoesNotDelayDirectDelivery() + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + using var transport = new FeatureFlagsEvpTransport( + FeatureFlagsSource.Agentless, + local, + direct, + NullDiscoveryService.Instance, + initialDiscoveryKnown: false, + initialDiscoveryWait: TimeSpan.FromMinutes(1)); + + var send = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + send.IsCompleted.Should().BeTrue("disabled discovery cannot produce a callback, and the test sender completes synchronously"); + await send; + direct.RequestsSent.Should().ContainSingle(); + local.RequestsSent.Should().BeEmpty(); + } + + [Fact] + public async Task ProductionConstructionDoesNotWaitForDisabledDiscovery() + { + var settings = CreateSettings((ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "agentless")); + using var transport = new FeatureFlagsEvpTransport(settings, NullDiscoveryService.Instance); + + var send = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + send.IsCompleted.Should().BeTrue("there is no discovery loop or direct credential to wait for"); + await send; + } + + [Theory] + [InlineData(400)] + [InlineData(403)] + [InlineData(429)] + [InlineData(500)] + public async Task DirectTerminalFailureLogsSafeErrorContextWithoutRetry(int status) + { + var logger = new Mock(); + var direct = CreateFactory("https://secret@event-platform-intake.datadoghq.com/", uri => new TestApiRequest(uri, statusCode: status, responseContent: "secret-response")); + var local = CreateFactory("http://agent:8126/"); + using var transport = new FeatureFlagsEvpTransport( + FeatureFlagsSource.Agentless, local, direct, NullDiscoveryService.Instance, logger: logger.Object); + + await transport.SendAsync(new { Secret = "secret-payload" }, FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + + var log = logger.Invocations.Should().ContainSingle().Which; + log.Method.Name.Should().Be(status == 400 ? "Error" : "ErrorSkipTelemetry"); + log.Arguments[0].Should().BeOfType().Which.Should().Contain("after 1 attempt").And.Contain("will not be replayed").And.NotContain("secret"); + log.Arguments[1].Should().Be(FeatureFlagsEvpTransport.ExposureIntakePath); + log.Arguments[2].Should().Be(status); + direct.RequestsSent.Should().ContainSingle(); + local.RequestsSent.Should().BeEmpty(); + } + + [Theory] + [InlineData("http://localhost:8126/")] +#if NET5_0_OR_GREATER + [InlineData("unix:///tmp/dd-evp-redirect-test.socket")] +#endif + public void LocalFactoriesDisableRedirectsWithoutChangingHistoricalDefaults(string agentUrl) + { + var exporter = CreateSettings((ConfigurationKeys.AgentUri, agentUrl)).Manager.InitialExporterSettings; + var agentless = FeatureFlagsEvpTransport.CreateLocalRequestFactory(exporter); + var historical = FeatureFlagsEvpTransport.CreateLocalRequestFactory(exporter, allowAutoRedirect: true); +#if NETCOREAPP + agentless.Should().BeAssignableTo().Which.AllowAutoRedirect.Should().BeFalse(); + historical.Should().BeAssignableTo().Which.AllowAutoRedirect.Should().BeTrue(); +#else + agentless.Should().BeOfType().Which.AllowAutoRedirect.Should().BeFalse(); + historical.Should().BeOfType().Which.AllowAutoRedirect.Should().BeTrue(); +#endif + } + + [Fact] + public async Task LocalRedirectIsNotFollowedOrReplayedAndOnlyFutureBatchUsesDirect() + { + using var relay = new HttpListener(); + using var destination = new HttpListener(); + var relayUrl = $"http://127.0.0.1:{TcpPortProvider.GetOpenPort()}/"; + relay.Prefixes.Add(relayUrl); + relay.Start(); + var destinationUrl = $"http://127.0.0.1:{TcpPortProvider.GetOpenPort()}/"; + destination.Prefixes.Add(destinationUrl); + destination.Start(); + var receivedAtDestination = destination.GetContextAsync(); + var receivedAtRelay = relay.GetContextAsync(); + var local = FeatureFlagsEvpTransport.CreateLocalRequestFactory(CreateSettings((ConfigurationKeys.AgentUri, relayUrl)).Manager.InitialExporterSettings); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + using var transport = new FeatureFlagsEvpTransport( + FeatureFlagsSource.Agentless, + local, + direct, + new DiscoveryServiceMock(), + initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + + try + { + var send = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + (await Task.WhenAny(receivedAtRelay, Task.Delay(TimeSpan.FromSeconds(5)))).Should().BeSameAs(receivedAtRelay); + var request = await receivedAtRelay; + await request.Request.InputStream.CopyToAsync(Stream.Null); + request.Request.Headers[TelemetryConstants.ApiKeyHeader].Should().BeNull(); + request.Request.Headers[FeatureFlagsEvpHeaderHelper.EvpOriginHeader].Should().Be("dd-trace-dotnet"); + request.Response.StatusCode = 307; + request.Response.RedirectLocation = destinationUrl; + request.Response.Close(); + + await send; + receivedAtDestination.IsCompleted.Should().BeFalse("the selected local route must never follow an HTTP redirect"); + direct.RequestsSent.Should().BeEmpty("a redirect does not prove that the current batch was not accepted"); + + await transport.SendAsync(new object(), FeatureFlagsEvpTransport.FlagEvaluationIntakePath, SerializerSettings); + direct.RequestsSent.Should().ContainSingle(); + } + finally + { + destination.Close(); + try + { + await receivedAtDestination; + } + catch (Exception ex) when (ex is HttpListenerException or ObjectDisposedException) + { + } + } + } + + [Theory] + [InlineData("http://agent:8126/base/path/", "http://agent:8126/base/path/evp_proxy/v4/api/v2/exposures")] + [InlineData("https://agent:8126/base/path/", "https://agent:8126/base/path/evp_proxy/v4/api/v2/exposures")] + public void LocalHttpFactoryPreservesSchemeAndBasePath(string agentUrl, string expectedEndpoint) + { + var settings = CreateSettings((ConfigurationKeys.AgentUri, agentUrl)); + + var factory = FeatureFlagsEvpTransport.CreateLocalRequestFactory(settings.Manager.InitialExporterSettings); + + factory.GetEndpoint($"{FeatureFlagsEvpTransport.EventPlatformProxyV4}/{FeatureFlagsEvpTransport.ExposureIntakePath}") + .Should().Be(new Uri(expectedEndpoint)); + } + +#if NETCOREAPP3_1_OR_GREATER + [Fact] + public void LocalUnixDomainSocketFactoryBuildsEvpRouteOnStreamEndpoint() + { + var settings = CreateSettings((ConfigurationKeys.AgentUri, "unix:///tmp/dd-apm-test.socket")); + + var factory = FeatureFlagsEvpTransport.CreateLocalRequestFactory(settings.Manager.InitialExporterSettings); + + factory.GetEndpoint($"{FeatureFlagsEvpTransport.EventPlatformProxyV4}/{FeatureFlagsEvpTransport.ExposureIntakePath}") + .Should().Be(new Uri("http://localhost/evp_proxy/v4/api/v2/exposures")); + } +#endif + + [Fact] + public void LocalNamedPipeFactoryBuildsEvpRouteOnStreamEndpoint() + { + var settings = CreateSettings((ConfigurationKeys.TracesPipeName, "dd-apm-test-pipe")); + + var factory = FeatureFlagsEvpTransport.CreateLocalRequestFactory(settings.Manager.InitialExporterSettings); + + factory.GetEndpoint($"{FeatureFlagsEvpTransport.EventPlatformProxyV2}/{FeatureFlagsEvpTransport.FlagEvaluationIntakePath}") + .Should().Be(new Uri("http://localhost/evp_proxy/v2/api/v2/flagevaluation")); + } + + [Theory] + [InlineData("datadoghq.eu", "https://event-platform-intake.datadoghq.eu/")] + [InlineData(" DATADOGHQ.COM ", "https://event-platform-intake.datadoghq.com/")] + public void DirectFactoryUsesNormalizedHttpsIntakeAndDisablesRedirects(string site, string expectedBase) + { + var settings = CreateSettings( + (ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "agentless"), + (ConfigurationKeys.ApiKey, "test-api-key"), + (ConfigurationKeys.Site, site)); + + var factory = FeatureFlagsEvpTransport.CreateDirectRequestFactory(settings.FeatureFlags); + + factory.Should().NotBeNull(); + factory!.GetEndpoint(FeatureFlagsEvpTransport.ExposureIntakePath) + .Should().Be(new Uri(expectedBase + FeatureFlagsEvpTransport.ExposureIntakePath)); + factory.GetEndpoint(FeatureFlagsEvpTransport.FlagEvaluationIntakePath) + .Should().Be(new Uri(expectedBase + FeatureFlagsEvpTransport.FlagEvaluationIntakePath)); +#if NETCOREAPP3_1_OR_GREATER + factory.Should().BeOfType() + .Which.AllowAutoRedirect.Should().BeFalse("DD-API-KEY must never follow an intake redirect"); +#else + factory.Should().BeOfType() + .Which.AllowAutoRedirect.Should().BeFalse("DD-API-KEY must never follow an intake redirect"); +#endif + } + + [Fact] + public void DirectFactoryUsesDefaultSiteWhenDdSiteIsUnset() + { + var settings = CreateSettings( + (ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "agentless"), + (ConfigurationKeys.ApiKey, "test-api-key")); + + var factory = FeatureFlagsEvpTransport.CreateDirectRequestFactory(settings.FeatureFlags); + + factory.Should().NotBeNull(); + factory!.GetEndpoint(FeatureFlagsEvpTransport.ExposureIntakePath) + .Should().Be(new Uri("https://event-platform-intake.datadoghq.com/api/v2/exposures")); + } + + [Theory] + [InlineData("datadoghq.com@attacker.example")] + [InlineData("https://attacker.example")] + [InlineData("datadoghq.com/path")] + [InlineData("datadoghq.com?redirect=attacker.example")] + [InlineData("datadoghq.com#attacker.example")] + [InlineData("datadoghq.com:443")] + [InlineData("datadoghq.com\\attacker.example")] + [InlineData("data doghq.com")] + [InlineData("-datadoghq.com")] + [InlineData("datadoghq.com-")] + [InlineData("datadoghq..com")] + [InlineData("dátadoghq.com")] + public void DirectFactoryRejectsSiteThatIsNotAnAsciiDnsSuffix(string site) + { + var settings = CreateSettings( + (ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "agentless"), + (ConfigurationKeys.ApiKey, "test-api-key"), + (ConfigurationKeys.Site, site)); + + FeatureFlagsEvpTransport.CreateDirectRequestFactory(settings.FeatureFlags).Should().BeNull(); + } + + [Fact] + public void InvalidSiteWarningIsGenericAndEmittedOnce() + { + const string Site = "secret@attacker.example"; + const string ApiKey = "secret-api-key"; + var warnings = new List(); + var settings = CreateSettings( + (ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "agentless"), + (ConfigurationKeys.ApiKey, ApiKey), + (ConfigurationKeys.Site, Site)); + + using var transport = new FeatureFlagsEvpTransport(settings, new DiscoveryServiceMock(), warnings.Add); + + warnings.Should().ContainSingle(); + warnings[0].Should().NotContain(Site).And.NotContain(ApiKey); + } + + [Fact] + public async Task DisposeUnsubscribesAndUnblocksInitialDiscoveryWait() + { + var local = CreateFactory("http://agent:8126/"); + var direct = CreateFactory("https://event-platform-intake.datadoghq.com/"); + var discovery = new DiscoveryServiceMock(); + var transport = CreateTransport( + local, + direct, + discovery, + initialDiscoveryKnown: false, + initialDiscoveryWait: TimeSpan.FromMinutes(1)); + + var send = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + discovery.Callbacks.Should().ContainSingle(); + + transport.Dispose(); + + discovery.Callbacks.Should().BeEmpty(); + (await Task.WhenAny(send, Task.Delay(TimeSpan.FromSeconds(1)))).Should().BeSameAs(send); + await send; + local.RequestsSent.Should().BeEmpty(); + direct.RequestsSent.Should().BeEmpty(); + } + + private static FeatureFlagsEvpTransport CreateTransport( + TestRequestFactory local, + TestRequestFactory? direct, + DiscoveryServiceMock? discovery = null, + string? initialLocalProxyEndpoint = null, + bool initialDiscoveryKnown = true, + TimeSpan? initialDiscoveryWait = null, + TimeSpan? routeRecoveryCooldown = null, + Func? utcNow = null, + Action? warningSink = null) + => new( + FeatureFlagsSource.Agentless, + local, + direct, + discovery ?? new DiscoveryServiceMock(), + initialLocalProxyEndpoint, + initialDiscoveryKnown, + initialDiscoveryWait, + routeRecoveryCooldown, + utcNow, + warningSink); + + private static TestRequestFactory CreateFactory(string baseEndpoint, params Func[] requests) + => new(new Uri(baseEndpoint), requests); + + private static TracerSettings CreateSettings(params (string Key, string Value)[] values) + { + var source = new NameValueCollection(); + foreach (var (key, value) in values) + { + source[key] = value; + } + + return new TracerSettings(new NameValueConfigurationSource(source)); + } + + private sealed class ThrowingApiRequest(Uri endpoint, Exception exception) : TestApiRequest(endpoint) + { + public override Task PostAsJsonAsync(T payload, MultipartCompression compression, JsonSerializerSettings settings) + => Task.FromException(exception); + } + + private sealed class BlockingApiRequest( + Uri endpoint, + TaskCompletionSource sendStarted, + TaskCompletionSource releaseSend) : TestApiRequest(endpoint) + { + public override async Task PostAsJsonAsync(T payload, MultipartCompression compression, JsonSerializerSettings settings) + { + sendStarted.TrySetResult(true); + await releaseSend.Task.ConfigureAwait(false); + return await base.PostAsJsonAsync(payload, compression, settings).ConfigureAwait(false); + } + } +} diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs index 8289794df486..e48ca6089bb3 100644 --- a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs @@ -10,6 +10,7 @@ using System.IO; using System.Net; using System.Threading.Tasks; +using Datadog.Trace.Agent.DiscoveryService; using Datadog.Trace.Configuration; using Datadog.Trace.FeatureFlags.Evp; using Datadog.Trace.Telemetry; @@ -26,8 +27,6 @@ public class FeatureFlagsFixedEvpTransportTests [Theory] [InlineData("remote_config", 403)] [InlineData("remote_config", 500)] - [InlineData("agentless", 403)] - [InlineData("agentless", 500)] public async Task ProductionSenderKeepsFixedV2AfterHttpFailure(string source, int firstStatus) { using var agent = new HttpListener(); @@ -40,7 +39,7 @@ public async Task ProductionSenderKeepsFixedV2AfterHttpFailure(string source, in { ConfigurationKeys.AgentUri, agentUrl + "prefix/" }, { ConfigurationKeys.ApiKey, "must-not-reach-agent" }, })); - using var transport = new FeatureFlagsEvpTransport(settings); + using var transport = new FeatureFlagsEvpTransport(settings, NullDiscoveryService.Instance); foreach (var status in new[] { firstStatus, 200 }) { diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs index b85d65f68918..3ca735523a6f 100644 --- a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs @@ -17,6 +17,7 @@ using Datadog.Trace.RemoteConfigurationManagement; using Datadog.Trace.RemoteConfigurationManagement.Protocol; using Datadog.Trace.TestHelpers; +using Datadog.Trace.Tests.Agent; using Datadog.Trace.Vendors.Newtonsoft.Json; using FluentAssertions; using Xunit; @@ -27,6 +28,38 @@ namespace Datadog.Trace.Tests.FeatureFlags; public class FeatureFlagsModuleTests { + [Fact] + public void AgentlessModuleUsesSharedDiscoveryForEventDelivery() + { + var rcmManager = new MockRcmSubscriptionManager(); + var discovery = new DiscoveryServiceMock(); + var source = new FakeDeliverySource(); + var collection = new NameValueCollection + { + { ConfigurationKeys.FeatureFlags.FlaggingProviderEnabled, "true" }, + { ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, "agentless" }, + { ConfigurationKeys.ApiKey, "test-api-key" }, + { ConfigurationKeys.Site, "datadoghq.com" }, + }; + var settings = new TracerSettings(new NameValueConfigurationSource(collection)); + + using (var module = FeatureFlagsModule.Create(settings, rcmManager, _ => source, discovery)) + { + module.Should().NotBeNull(); + source.Started.Should().Be(0); + rcmManager.HasAnySubscription.Should().BeFalse(); + module!.GetExposureApi().Should().NotBeNull(); + discovery.Callbacks.Should().ContainSingle(); + + module.Activate(); + module.Activate(); + source.Started.Should().Be(1); + } + + discovery.Callbacks.Should().BeEmpty(); + source.Disposed.Should().Be(1); + } + [Fact] public void UpdateRemoteConfig_WithEmptyList_InvokesCallbackAndReturnsProviderNotReady() { From 3cc8a5aa979bf7a6de6f9629b69ebe802509d686 Mon Sep 17 00:00:00 2001 From: Leo Romanovsky Date: Fri, 2 Oct 2026 13:10:53 +0000 Subject: [PATCH 2/2] Handle EVP stream EOF and stop fallback after disposal Preserve ambiguous-delivery semantics, cover production header wiring, and correct old-runtime test compilation. Environment: Datadog workspace --- .../Evp/FeatureFlagsEvpTransport.cs | 4 +- .../HttpOverStreams/DatadogHttpClient.cs | 6 +- .../FeatureFlagsEvpTransportTests.cs | 139 +++++++++++++++++- .../HttpOverStreams/DatadogHttpClientTests.cs | 19 ++- 4 files changed, 162 insertions(+), 6 deletions(-) diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs index be25f5fb3ffc..807db77c6d78 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs @@ -558,7 +558,9 @@ private bool LeaveLocalRoute() private async Task SendDirectAsync(string intakePath, Func> sendAsync) { var directFactory = _directRequestFactory; - if (directFactory is null) + // The local send can outlive the bounded shutdown wait. Its continuation must not + // start a new fallback request after the transport has been disposed. + if (Volatile.Read(ref _disposed) != 0 || directFactory is null) { return; } diff --git a/tracer/src/Datadog.Trace/HttpOverStreams/DatadogHttpClient.cs b/tracer/src/Datadog.Trace/HttpOverStreams/DatadogHttpClient.cs index c6fac8b8c2b8..0c24ee070227 100644 --- a/tracer/src/Datadog.Trace/HttpOverStreams/DatadogHttpClient.cs +++ b/tracer/src/Datadog.Trace/HttpOverStreams/DatadogHttpClient.cs @@ -1,4 +1,4 @@ -// +// // 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. // @@ -89,7 +89,7 @@ async Task GoNextChar() var bytesRead = await responseStream.ReadAsync(chArray, offset: 0, count: 1).ConfigureAwait(false); if (bytesRead == 0) { - ThrowHelper.ThrowInvalidOperationException($"Unexpected end of stream at position {streamPosition}"); + throw new EndOfStreamException($"Unexpected end of stream at position {streamPosition}"); } currentChar = Encoding.ASCII.GetChars(chArray)[0]; @@ -112,7 +112,7 @@ async Task SkipUntil(int requiredStreamPosition) lastBytesRead = await responseStream.ReadAsync(chArray, offset: 0, count: bytesToRead).ConfigureAwait(false); if (lastBytesRead == 0) { - ThrowHelper.ThrowInvalidOperationException($"Unexpected end of stream at position {streamPosition}"); + throw new EndOfStreamException($"Unexpected end of stream at position {streamPosition}"); } bytesRemaining -= lastBytesRead; diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs index 757d6bfefb55..7bf229d5a5ba 100644 --- a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsEvpTransportTests.cs @@ -9,9 +9,13 @@ using System.Collections.Generic; using System.Collections.Specialized; using System.IO; +using System.IO.Pipes; using System.Linq; using System.Net; +using System.Net.Http; using System.Net.Sockets; +using System.Reflection; +using System.Text; using System.Threading.Tasks; using Datadog.Trace.Agent; using Datadog.Trace.Agent.DiscoveryService; @@ -610,7 +614,7 @@ public void LocalFactoriesDisableRedirectsWithoutChangingHistoricalDefaults(stri var exporter = CreateSettings((ConfigurationKeys.AgentUri, agentUrl)).Manager.InitialExporterSettings; var agentless = FeatureFlagsEvpTransport.CreateLocalRequestFactory(exporter); var historical = FeatureFlagsEvpTransport.CreateLocalRequestFactory(exporter, allowAutoRedirect: true); -#if NETCOREAPP +#if NETCOREAPP3_1_OR_GREATER agentless.Should().BeAssignableTo().Which.AllowAutoRedirect.Should().BeFalse(); historical.Should().BeAssignableTo().Which.AllowAutoRedirect.Should().BeTrue(); #else @@ -815,6 +819,118 @@ public async Task DisposeUnsubscribesAndUnblocksInitialDiscoveryWait() direct.RequestsSent.Should().BeEmpty(); } + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task AcceptedNamedPipeEofDoesNotReplayAndChangesFutureRouting(bool hasDirectCredentials) + { + var pipeName = "evp" + Guid.NewGuid().ToString("N").Substring(0, 8); + var settings = CreateSettings((ConfigurationKeys.TracesPipeName, pipeName)); + var local = FeatureFlagsEvpTransport.CreateLocalRequestFactory(settings.Manager.InitialExporterSettings); + var direct = CreateFactory("https://event-platform-intake.mock-intake.invalid/"); + var now = new DateTimeOffset(2026, 10, 2, 0, 0, 0, TimeSpan.Zero); + using var transport = new FeatureFlagsEvpTransport( + FeatureFlagsSource.Agentless, + local, + hasDirectCredentials ? direct : null, + new DiscoveryServiceMock(), + initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4, + routeRecoveryCooldown: TimeSpan.FromMinutes(1), + utcNow: () => now); + using var server = new NamedPipeServerStream(pipeName, PipeDirection.InOut, 1, PipeTransmissionMode.Byte, PipeOptions.Asynchronous); + var received = ReadPipeRequest(server); + var firstSend = transport.SendAsync(new { Batch = 1 }, FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + (await Task.WhenAny(received, Task.Delay(TimeSpan.FromSeconds(5)))).Should().BeSameAs(received); + var request = await received; + request.PathAndQuery.Should().Be("/evp_proxy/v4/api/v2/exposures"); + request.Headers.Contains("DD-API-KEY").Should().BeFalse(); + server.Disconnect(); + + await firstSend; + direct.RequestsSent.Should().BeEmpty("an accepted batch must not be replayed after EOF"); + + var nextConnection = server.WaitForConnectionAsync(); + await transport.SendAsync(new { Batch = 2 }, FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + nextConnection.IsCompleted.Should().BeFalse("future batches use direct or wait for the unavailable cooldown"); + direct.RequestsSent.Should().HaveCount(hasDirectCredentials ? 1 : 0); + + if (!hasDirectCredentials) + { + now = now.AddMinutes(1); + var recovery = transport.SendAsync(new { Batch = 3 }, FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + (await Task.WhenAny(nextConnection, Task.Delay(TimeSpan.FromSeconds(5)))).Should().BeSameAs(nextConnection); + await nextConnection; + var recoveredRequest = await MockHttpParser.ReadRequest(server); + recoveredRequest.ReadStreamBody().Should().NotBeEmpty(); + var response = Encoding.ASCII.GetBytes("HTTP/1.1 200 OK\r\nContent-Length: 0\r\n\r\n"); + await server.WriteAsync(response, 0, response.Length); + await recovery; + direct.RequestsSent.Should().BeEmpty(); + } + else + { + server.Dispose(); + await Record.ExceptionAsync(() => nextConnection); + } + } + + [Theory] + [InlineData(404)] + [InlineData(405)] + [InlineData(0)] + public async Task DisposingTransportPreventsFallbackAfterLateLocalCompletion(int statusCode) + { + var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var local = CreateFactory("http://agent:8126/", uri => new DelayedFailureApiRequest(uri, statusCode, started, release)); + var direct = CreateFactory("https://event-platform-intake.mock-intake.invalid/"); + using var transport = CreateTransport(local, direct, initialLocalProxyEndpoint: FeatureFlagsEvpTransport.EventPlatformProxyV4); + var send = transport.SendAsync(new object(), FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings); + await started.Task; + + transport.Dispose(); + release.TrySetResult(true); + await send; + + local.RequestsSent.Should().ContainSingle(); + direct.RequestsSent.Should().BeEmpty("a late response or pre-send failure must not start new egress after disposal"); + } + + [Fact] + public void ProductionDirectFactoryAttachesHeadersToItsRequest() + { + var settings = CreateSettings( + (ConfigurationKeys.ApiKey, "test-api-key"), + (ConfigurationKeys.Site, "mock-intake.invalid")); + var factory = FeatureFlagsEvpTransport.CreateDirectRequestFactory(settings.FeatureFlags)!; + var endpoint = factory.GetEndpoint(FeatureFlagsEvpTransport.ExposureIntakePath); + var request = factory.Create(endpoint); + + // Inspect the real request/client, not GetDirectHeaders: omitting factory wiring must fail. + // No TLS bypass or outbound request is needed for this construction regression. +#if NETCOREAPP3_1_OR_GREATER + using var client = (HttpClient)typeof(HttpClientRequest).GetField("_client", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(request)!; + var headers = client.DefaultRequestHeaders.ToDictionary(header => header.Key, header => header.Value.Single(), StringComparer.OrdinalIgnoreCase); +#else + var webRequest = (HttpWebRequest)typeof(ApiWebRequest).GetField("_request", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(request)!; + var headers = webRequest.Headers.AllKeys.ToDictionary(name => name!, name => webRequest.Headers[name]!, StringComparer.OrdinalIgnoreCase); +#endif + endpoint.Should().Be(new Uri("https://event-platform-intake.mock-intake.invalid/api/v2/exposures")); + headers.Should().Contain("DD-API-KEY", "test-api-key"); + headers.Should().Contain("DD-EVP-ORIGIN", "dd-trace-dotnet"); + headers.Should().Contain("DD-EVP-ORIGIN-VERSION", TracerConstants.ThreePartVersion); + headers.Should().Contain("x-datadog-tracing-enabled", "false"); + headers.Should().NotContainKey("X-Datadog-EVP-Subdomain"); + } + + private static async Task ReadPipeRequest(NamedPipeServerStream server) + { + await server.WaitForConnectionAsync(); + var request = await MockHttpParser.ReadRequest(server); + request.ReadStreamBody().Should().NotBeEmpty("the relay has accepted the full serialized batch before disconnecting"); + return request; + } + private static FeatureFlagsEvpTransport CreateTransport( TestRequestFactory local, TestRequestFactory? direct, @@ -857,6 +973,27 @@ public override Task PostAsJsonAsync(T payload, MultipartCompre => Task.FromException(exception); } + private sealed class DelayedFailureApiRequest( + Uri endpoint, + int statusCode, + TaskCompletionSource started, + TaskCompletionSource release) : TestApiRequest(endpoint, statusCode) + { + private readonly bool _failBeforeSend = statusCode == 0; + + public override async Task PostAsJsonAsync(T payload, MultipartCompression compression, JsonSerializerSettings settings) + { + started.TrySetResult(true); + await release.Task.ConfigureAwait(false); + if (_failBeforeSend) + { + throw new SocketException((int)SocketError.ConnectionRefused); + } + + return await base.PostAsJsonAsync(payload, compression, settings).ConfigureAwait(false); + } + } + private sealed class BlockingApiRequest( Uri endpoint, TaskCompletionSource sendStarted, diff --git a/tracer/test/Datadog.Trace.Tests/HttpOverStreams/DatadogHttpClientTests.cs b/tracer/test/Datadog.Trace.Tests/HttpOverStreams/DatadogHttpClientTests.cs index f9376a434d43..2446ce8b88c5 100644 --- a/tracer/test/Datadog.Trace.Tests/HttpOverStreams/DatadogHttpClientTests.cs +++ b/tracer/test/Datadog.Trace.Tests/HttpOverStreams/DatadogHttpClientTests.cs @@ -1,4 +1,4 @@ -// +// // 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. // @@ -18,6 +18,23 @@ namespace Datadog.Trace.Tests.HttpOverStreams { public class DatadogHttpClientTests { + [Theory] + [InlineData("")] + [InlineData("HTTP/1.1 ")] + [InlineData("HTTP/1.1 200 OK\r\nContent-Length:")] + public async Task TruncatedResponseThrowsEndOfStreamAfterSendingRequest(string response) + { + var client = new DatadogHttpClient(TraceAgentHttpHeaderHelper.Instance); + var request = new HttpRequest("POST", "localhost", "/test", new HttpHeaders(), null); + using var requestStream = new MemoryStream(); + using var responseStream = new MemoryStream(Encoding.ASCII.GetBytes(response)); + + Func send = () => client.SendAsync(request, requestStream, responseStream); + + await send.Should().ThrowAsync(); + requestStream.Length.Should().BeGreaterThan(0, "response EOF does not prove a zero-byte send"); + } + [Fact] public async Task DatadogHttpClient_CanParseResponse() {