From 9a8280d92d9bf9fa242721a0954bc8c1d2f2953f Mon Sep 17 00:00:00 2001 From: Leo Romanovsky Date: Thu, 1 Oct 2026 22:14:24 +0000 Subject: [PATCH 1/2] Bound Feature Flags exposure shutdown and own the fixed-v2 sender Preserve historical event routing while draining queued exposures safely during module shutdown. Environment: Datadog workspace --- .../Evp/FeatureFlagsEvpTransport.cs | 84 ++++++++++ .../FeatureFlags/Exposure/ExposureApi.cs | 144 ++++++++++-------- .../FeatureFlags/FeatureFlagsModule.cs | 7 +- .../TransportHelpers/TestRequestFactory.cs | 14 +- .../FeatureFlags/ExposureApiTests.cs | 140 +++++++++++++++++ .../FeatureFlagsFixedEvpTransportTests.cs | 60 ++++++++ 6 files changed, 383 insertions(+), 66 deletions(-) create mode 100644 tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs create mode 100644 tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs create mode 100644 tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs new file mode 100644 index 000000000000..ecceeda61580 --- /dev/null +++ b/tracer/src/Datadog.Trace/FeatureFlags/Evp/FeatureFlagsEvpTransport.cs @@ -0,0 +1,84 @@ +// +// 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.Threading; +using System.Threading.Tasks; +using Datadog.Trace.Agent; +using Datadog.Trace.Agent.Transports; +using Datadog.Trace.Configuration; +using Datadog.Trace.HttpOverStreams; +using Datadog.Trace.SourceGenerators; +using Datadog.Trace.Vendors.Newtonsoft.Json; + +namespace Datadog.Trace.FeatureFlags.Evp; + +/// +/// Owns the historical fixed-v2 Feature Flags event sender independently of exposure batching. +/// +internal sealed class FeatureFlagsEvpTransport : IDisposable +{ + internal const string ExposureIntakePath = "api/v2/exposures"; + internal const string EventPlatformProxyV2 = "evp_proxy/v2"; + + private readonly object _settingsLock = new(); + private readonly IDisposable? _settingsSubscription; + private IApiRequestFactory _localRequestFactory; + private int _disposed; + + internal FeatureFlagsEvpTransport(TracerSettings settings) + { + _localRequestFactory = CreateLocalRequestFactory(settings.Manager.InitialExporterSettings); + _settingsSubscription = settings.Manager.SubscribeToChanges(changes => + { + if (changes.UpdatedExporter is { } exporter) + { + lock (_settingsLock) + { + if (_disposed == 0) + { + Interlocked.Exchange(ref _localRequestFactory, CreateLocalRequestFactory(exporter)); + } + } + } + }); + } + + [TestingOnly] + internal FeatureFlagsEvpTransport(IApiRequestFactory localRequestFactory) + { + _localRequestFactory = localRequestFactory; + } + + [TestingAndPrivateOnly] + internal static IApiRequestFactory CreateLocalRequestFactory(ExporterSettings exporterSettings) + => AgentTransportStrategy.Get( + exporterSettings, + productName: "FeatureFlags exposure", + tcpTimeout: TimeSpan.FromSeconds(5), + httpHeaderHelper: EventPlatformHeaderHelper.Instance); + + public void Dispose() + { + if (Interlocked.Exchange(ref _disposed, 1) == 0) + { + _settingsSubscription?.Dispose(); + } + } + + internal async Task SendAsync(T payload, string intakePath, JsonSerializerSettings serializerSettings) + { + 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); + } +} diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Exposure/ExposureApi.cs b/tracer/src/Datadog.Trace/FeatureFlags/Exposure/ExposureApi.cs index 769ba74503d9..863d431b27b3 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/Exposure/ExposureApi.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/Exposure/ExposureApi.cs @@ -7,15 +7,11 @@ using System; using System.Collections.Generic; using System.Diagnostics.CodeAnalysis; -using System.Linq; -using System.Text; using System.Threading; using System.Threading.Tasks; -using Datadog.Trace.Agent; -using Datadog.Trace.Agent.Transports; using Datadog.Trace.Configuration; +using Datadog.Trace.FeatureFlags.Evp; using Datadog.Trace.FeatureFlags.Exposure.Model; -using Datadog.Trace.HttpOverStreams; using Datadog.Trace.Logging; using Datadog.Trace.SourceGenerators; using Datadog.Trace.Vendors.Newtonsoft.Json; @@ -28,7 +24,9 @@ internal sealed class ExposureApi : IDisposable internal static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor(typeof(ExposureApi)); private const int DefaultCapacity = 1 << 16; // 65536 elements - public const string ExposurePath = "evp_proxy/v2/api/v2/exposures"; + private static readonly TimeSpan DefaultSendInterval = TimeSpan.FromSeconds(10); + private static readonly TimeSpan DefaultShutdownTimeout = TimeSpan.FromSeconds(10); + [TestingAndPrivateOnly] internal static readonly JsonSerializerSettings SerializerSettings = new() { @@ -39,45 +37,38 @@ internal sealed class ExposureApi : IDisposable } }; - private readonly TaskCompletionSource _processExit = new(); - private readonly TimeSpan _sendInterval = TimeSpan.FromSeconds(10); + private readonly TaskCompletionSource _processExit = new(TaskCreationOptions.RunContinuationsAsynchronously); + private readonly object _lifecycleLock = new(); + private readonly TimeSpan _sendInterval; + private readonly TimeSpan _shutdownTimeout; private readonly Queue _exposures = new Queue(); - private readonly ExposureCache _exposureCache = new ExposureCache(DefaultCapacity); - private IApiRequestFactory _apiRequestFactory; - private Dictionary _context; - private int _started; + private readonly FeatureFlagsEvpTransport _transport; + private readonly IDisposable _settingsSubscription; - internal ExposureApi(TracerSettings tracerSettings) + private Dictionary _context; + private Task? _sendLoopTask; + private bool _disposed; + + internal ExposureApi( + TracerSettings tracerSettings, + FeatureFlagsEvpTransport transport, + TimeSpan? sendInterval = null, + TimeSpan? shutdownTimeout = null) { - UpdateApi(tracerSettings.Manager.InitialExporterSettings); + _transport = transport; + _sendInterval = sendInterval ?? DefaultSendInterval; + _shutdownTimeout = shutdownTimeout ?? DefaultShutdownTimeout; UpdateContext(tracerSettings.Manager.InitialMutableSettings); - tracerSettings.Manager.SubscribeToChanges(changes => + _settingsSubscription = tracerSettings.Manager.SubscribeToChanges(changes => { - if (changes.UpdatedExporter is { } exporter) - { - UpdateApi(exporter); - } - if (changes.UpdatedMutable is { } mutable) { UpdateContext(mutable); } }); - [MemberNotNull(nameof(_apiRequestFactory))] - void UpdateApi(ExporterSettings exporterSettings) - { - Log.Debug("ExposureApi::UpdateApi-> Applying settings"); - var apiRequestFactory = AgentTransportStrategy.Get( - exporterSettings, - productName: "FeatureFlags exposure", - tcpTimeout: TimeSpan.FromSeconds(5), - httpHeaderHelper: EventPlatformHeaderHelper.Instance); - Interlocked.Exchange(ref _apiRequestFactory!, apiRequestFactory); - } - [MemberNotNull(nameof(_context))] void UpdateContext(MutableSettings settings) { @@ -94,17 +85,29 @@ void UpdateContext(MutableSettings settings) public void Dispose() { - _processExit.TrySetResult(true); - } + Task? sendLoopTask; + lock (_lifecycleLock) + { + if (_disposed) + { + return; + } - public void TryToStartSendLoopIfNotStarted() - { - if (Interlocked.CompareExchange(ref _started, 1, 0) != 0) + _disposed = true; + _processExit.TrySetResult(true); + sendLoopTask = _sendLoopTask; + } + + if (sendLoopTask is not null) { - return; + var completed = Task.WhenAny(sendLoopTask, Task.Delay(_shutdownTimeout)).GetAwaiter().GetResult(); + if (completed != sendLoopTask) + { + Log.Warning("Could not finish flushing Feature Flags exposures before process end"); + } } - _ = Task.Run(SendLoopAsync).ContinueWith(t => { Log.Error(t.Exception, "FeatureFlags Exposure send loop failed"); }, TaskContinuationOptions.OnlyOnFaulted); + _settingsSubscription.Dispose(); } private async Task SendLoopAsync() @@ -112,21 +115,7 @@ private async Task SendLoopAsync() Log.Debug("ExposureApi::SendLoopAsync -> Enter"); while (!_processExit.Task.IsCompleted) { - try - { - var apiRequestFactory = _apiRequestFactory; - var uri = apiRequestFactory.GetEndpoint(ExposurePath); - var payload = TryGetPayload(); - if (payload is not null) - { - var request = apiRequestFactory.Create(uri); - using var response = await request.PostAsJsonAsync(payload, MultipartCompression.GZip, SerializerSettings).ConfigureAwait(false); - } - } - catch (Exception ex) - { - Log.Error(ex, "Error while sending Feature Flags exposures to the agent"); - } + await FlushAsync().ConfigureAwait(false); try { @@ -136,8 +125,27 @@ private async Task SendLoopAsync() { // We are shutting down, so don't do anything about it } + } - Log.Debug("ExposureApi::SendLoopAsync -> Exit"); + // Dispose signals the loop and then waits for this bounded final flush. This prevents the + // common short-lived-process loss mode without allowing shutdown to hang indefinitely. + await FlushAsync().ConfigureAwait(false); + Log.Debug("ExposureApi::SendLoopAsync -> Exit"); + } + + private async Task FlushAsync() + { + try + { + var payload = TryGetPayload(); + if (payload is not null) + { + await _transport.SendAsync(payload, FeatureFlagsEvpTransport.ExposureIntakePath, SerializerSettings).ConfigureAwait(false); + } + } + catch (Exception ex) + { + Log.Error(ex, "Error while sending Feature Flags exposures"); } } @@ -161,15 +169,31 @@ private async Task SendLoopAsync() public void SendExposure(in ExposureEvent exposure) { - if (_exposureCache.Add(exposure)) + lock (_lifecycleLock) { - lock (_exposures) + if (_disposed) { - _exposures.Enqueue(exposure); + return; } - } - TryToStartSendLoopIfNotStarted(); + if (_exposureCache.Add(exposure)) + { + lock (_exposures) + { + _exposures.Enqueue(exposure); + } + } + + if (_sendLoopTask is null) + { + _sendLoopTask = Task.Run(SendLoopAsync); + _sendLoopTask.ContinueWith( + t => Log.Error(t.Exception, "FeatureFlags Exposure send loop failed"), + CancellationToken.None, + TaskContinuationOptions.OnlyOnFaulted, + TaskScheduler.Default); + } + } } private sealed class ExposuresRequest(Dictionary context, List exposures) diff --git a/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs b/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs index d3acc1a37dbd..70a803e939e1 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs @@ -11,6 +11,7 @@ using System.Threading.Tasks; using Datadog.Trace.Configuration; using Datadog.Trace.FeatureFlags.Agentless; +using Datadog.Trace.FeatureFlags.Evp; using Datadog.Trace.FeatureFlags.Exposure; using Datadog.Trace.FeatureFlags.Exposure.Model; using Datadog.Trace.FeatureFlags.Rcm; @@ -25,6 +26,8 @@ internal sealed class FeatureFlagsModule : IDisposable { internal static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor(typeof(FeatureFlagsModule)); + private readonly FeatureFlagsEvpTransport _evpTransport; + // Activation, disposal and exposure-API creation all mutate the same state from different // threads, so they share one lock rather than individual interlocked flags: a flag set // before its accompanying setup completes lets a concurrent caller observe a half-activated @@ -83,6 +86,7 @@ internal FeatureFlagsModule( _agentlessSourceFactory = agentlessSourceFactory ?? (static module => AgentlessConfigurationSource.Create(module._settings, module._settingsManager, module.ApplyConfiguration)); _rcmSubscriptionManager = rcmSubscriptionManager; + _evpTransport = new FeatureFlagsEvpTransport(settings); Log.Debug("FeatureFlagsModule ENABLED with source {Source}", _settings.Source); } @@ -157,6 +161,7 @@ public void Dispose() agentlessSource?.Dispose(); exposureApi?.Dispose(); + _evpTransport.Dispose(); } /// @@ -501,7 +506,7 @@ private void ReportExposure(in ExposureEvent exposure) exposureApi = _exposureApi; if (exposureApi is null) { - exposureApi = new ExposureApi(_tracerSettings); + exposureApi = new ExposureApi(_tracerSettings, _evpTransport); Volatile.Write(ref _exposureApi, exposureApi); } diff --git a/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs b/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs index d8fd8a76f054..9b84a117f5fa 100644 --- a/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs +++ b/tracer/test/Datadog.Trace.TestHelpers/TransportHelpers/TestRequestFactory.cs @@ -15,6 +15,7 @@ internal class TestRequestFactory : IApiRequestFactory { private readonly Uri _baseEndpoint; private readonly Func[] _requestsToSend; + private readonly object _requestsLock = new(); public TestRequestFactory(params Func[] requestsToSend) : this(new Uri("http://localhost"), requestsToSend) @@ -38,13 +39,16 @@ public Uri GetEndpoint(string relativePath) public IApiRequest Create(Uri endpoint) { - var request = (_requestsToSend is null || RequestsSent.Count >= _requestsToSend.Length) - ? new TestApiRequest(endpoint) - : _requestsToSend[RequestsSent.Count](endpoint); + lock (_requestsLock) + { + var request = (_requestsToSend is null || RequestsSent.Count >= _requestsToSend.Length) + ? new TestApiRequest(endpoint) + : _requestsToSend[RequestsSent.Count](endpoint); - RequestsSent.Add(request); + RequestsSent.Add(request); - return request; + return request; + } } public void SetProxy(WebProxy proxy, NetworkCredential credential) diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs new file mode 100644 index 000000000000..12f79f8db2ba --- /dev/null +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/ExposureApiTests.cs @@ -0,0 +1,140 @@ +// +// 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.Diagnostics; +using System.Threading.Tasks; +using Datadog.Trace.Agent; +using Datadog.Trace.Agent.Transports; +using Datadog.Trace.Configuration; +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.Vendors.Newtonsoft.Json; +using Datadog.Trace.Vendors.Newtonsoft.Json.Linq; +using FluentAssertions; +using Xunit; + +namespace Datadog.Trace.Tests.FeatureFlags; + +public class ExposureApiTests +{ + [Fact] + public async Task DisposeFlushesEventsQueuedWhileAnEarlierBatchIsSending() + { + var firstSendStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseFirstSend = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var local = new TestRequestFactory( + new Uri("http://agent:8126/"), + uri => new BlockingApiRequest(uri, firstSendStarted, releaseFirstSend), + uri => new RecordingJsonApiRequest(uri)); + var transport = CreateLocalTransport(local); + var api = new ExposureApi(CreateSettings(), transport, TimeSpan.FromHours(1), TimeSpan.FromSeconds(2)); + + api.SendExposure(CreateExposure("first")); + (await Task.WhenAny(firstSendStarted.Task, Task.Delay(TimeSpan.FromSeconds(1)))) + .Should().BeSameAs(firstSendStarted.Task); + + api.SendExposure(CreateExposure("second")); + var dispose = Task.Run(api.Dispose); + releaseFirstSend.TrySetResult(true); + + (await Task.WhenAny(dispose, Task.Delay(TimeSpan.FromSeconds(3)))).Should().BeSameAs(dispose); + await dispose; + + 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.Compression.Should().Be(MultipartCompression.GZip); + var payload = JObject.Parse(finalRequest.PayloadJson); + payload["context"]!["service"]!.Value().Should().NotBeNullOrEmpty(); + payload["exposures"]!.Should().ContainSingle(); + payload["exposures"]![0]!["flag"]!["key"]!.Value().Should().Be("second"); + } + + [Fact] + public async Task DisposeHasABoundedWaitWhenTheNetworkSendDoesNotFinish() + { + var sendStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseSend = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var local = new TestRequestFactory( + new Uri("http://agent:8126/"), + uri => new BlockingApiRequest(uri, sendStarted, releaseSend)); + var transport = CreateLocalTransport(local); + var api = new ExposureApi(CreateSettings(), transport, TimeSpan.FromHours(1), TimeSpan.FromMilliseconds(25)); + + api.SendExposure(CreateExposure("first")); + (await Task.WhenAny(sendStarted.Task, Task.Delay(TimeSpan.FromSeconds(1)))) + .Should().BeSameAs(sendStarted.Task); + + var stopwatch = Stopwatch.StartNew(); + api.Dispose(); + stopwatch.Stop(); + + stopwatch.Elapsed.Should().BeLessThan(TimeSpan.FromSeconds(1)); + releaseSend.TrySetResult(true); + } + + [Fact] + public async Task SendExposureAfterDisposeIsIgnored() + { + var local = new TestRequestFactory(new Uri("http://agent:8126/")); + var transport = CreateLocalTransport(local); + var api = new ExposureApi(CreateSettings(), transport, TimeSpan.FromMilliseconds(1), TimeSpan.FromSeconds(1)); + + api.Dispose(); + api.SendExposure(CreateExposure("after-dispose")); + await Task.Delay(TimeSpan.FromMilliseconds(20)); + + local.RequestsSent.Should().BeEmpty(); + } + + private static FeatureFlagsEvpTransport CreateLocalTransport(TestRequestFactory local) + => new(local); + + private static ExposureEvent CreateExposure(string flag) + => new( + DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(), + new Allocation("allocation"), + new Flag(flag), + new Variant("variant"), + new Subject("subject", new Dictionary())); + + private static TracerSettings CreateSettings() + => new(new NameValueConfigurationSource(new NameValueCollection())); + + 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); + } + } + + private sealed class RecordingJsonApiRequest(Uri endpoint) : TestApiRequest(endpoint) + { + public string PayloadJson { get; private set; } = string.Empty; + + public MultipartCompression Compression { get; private set; } + + public override Task PostAsJsonAsync(T payload, MultipartCompression compression, JsonSerializerSettings settings) + { + PayloadJson = JsonConvert.SerializeObject(payload, settings); + Compression = compression; + return base.PostAsJsonAsync(payload, compression, settings); + } + } +} diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs new file mode 100644 index 000000000000..8289794df486 --- /dev/null +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsFixedEvpTransportTests.cs @@ -0,0 +1,60 @@ +// +// 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.Specialized; +using System.IO; +using System.Net; +using System.Threading.Tasks; +using Datadog.Trace.Configuration; +using Datadog.Trace.FeatureFlags.Evp; +using Datadog.Trace.Telemetry; +using Datadog.Trace.TestHelpers; +using Datadog.Trace.Vendors.Newtonsoft.Json; +using FluentAssertions; +using Xunit; + +namespace Datadog.Trace.Tests.FeatureFlags; + +[Collection(nameof(WebRequestCollection))] +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(); + var agentUrl = $"http://127.0.0.1:{TcpPortProvider.GetOpenPort()}/"; + agent.Prefixes.Add(agentUrl); + agent.Start(); + var settings = new TracerSettings(new NameValueConfigurationSource(new NameValueCollection + { + { ConfigurationKeys.FeatureFlags.FeatureFlagsConfigurationSource, source }, + { ConfigurationKeys.AgentUri, agentUrl + "prefix/" }, + { ConfigurationKeys.ApiKey, "must-not-reach-agent" }, + })); + using var transport = new FeatureFlagsEvpTransport(settings); + + foreach (var status in new[] { firstStatus, 200 }) + { + var received = agent.GetContextAsync(); + var send = transport.SendAsync(new { Flag = "test" }, FeatureFlagsEvpTransport.ExposureIntakePath, new JsonSerializerSettings()); + (await Task.WhenAny(received, Task.Delay(TimeSpan.FromSeconds(5)))).Should().BeSameAs(received); + var request = await received; + request.Request.Url!.AbsolutePath.Should().Be("/prefix/evp_proxy/v2/api/v2/exposures"); + request.Request.Headers[TelemetryConstants.ApiKeyHeader].Should().BeNull(); + await request.Request.InputStream.CopyToAsync(Stream.Null); + request.Response.StatusCode = status; + request.Response.Close(); + (await Task.WhenAny(send, Task.Delay(TimeSpan.FromSeconds(5)))).Should().BeSameAs(send); + await send; + } + } +} From 1994bbb10a68d295d5a8a6de2799326267a6392d Mon Sep 17 00:00:00 2001 From: Leo Romanovsky Date: Fri, 2 Oct 2026 13:01:58 +0000 Subject: [PATCH 2/2] Keep Feature Flags EVP transport lazy until first exposure Retain module ownership and disposal without allocating an HTTP client for applications that never use exposures. Environment: Datadog workspace --- .../FeatureFlags/FeatureFlagsModule.cs | 11 +++--- .../FeatureFlags/FeatureFlagsModuleTests.cs | 35 +++++++++++++++++++ 2 files changed, 42 insertions(+), 4 deletions(-) diff --git a/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs b/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs index 70a803e939e1..538657be7beb 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs @@ -26,8 +26,6 @@ internal sealed class FeatureFlagsModule : IDisposable { internal static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor(typeof(FeatureFlagsModule)); - private readonly FeatureFlagsEvpTransport _evpTransport; - // Activation, disposal and exposure-API creation all mutate the same state from different // threads, so they share one lock rather than individual interlocked flags: a flag set // before its accompanying setup completes lets a concurrent caller observe a half-activated @@ -68,6 +66,7 @@ internal sealed class FeatureFlagsModule : IDisposable private FeatureFlagsEvaluator? _evaluator; private IFeatureFlagsDeliverySource? _agentlessSource; private ExposureApi? _exposureApi; + private FeatureFlagsEvpTransport? _evpTransport; private string? _deliveryUnavailableReason; private bool _activated; private bool _disposed; @@ -86,7 +85,6 @@ internal FeatureFlagsModule( _agentlessSourceFactory = agentlessSourceFactory ?? (static module => AgentlessConfigurationSource.Create(module._settings, module._settingsManager, module.ApplyConfiguration)); _rcmSubscriptionManager = rcmSubscriptionManager; - _evpTransport = new FeatureFlagsEvpTransport(settings); Log.Debug("FeatureFlagsModule ENABLED with source {Source}", _settings.Source); } @@ -133,6 +131,7 @@ public void Dispose() ISubscription? subscription; IFeatureFlagsDeliverySource? agentlessSource; ExposureApi? exposureApi; + FeatureFlagsEvpTransport? evpTransport; lock (_stateLock) { @@ -146,10 +145,12 @@ public void Dispose() subscription = _rcmSubscription; agentlessSource = _agentlessSource; exposureApi = _exposureApi; + evpTransport = _evpTransport; _rcmSubscription = null; _agentlessSource = null; Volatile.Write(ref _exposureApi, null); + _evpTransport = null; } // Released the lock first: disposal is not state mutation, and holding it here would @@ -161,7 +162,7 @@ public void Dispose() agentlessSource?.Dispose(); exposureApi?.Dispose(); - _evpTransport.Dispose(); + evpTransport?.Dispose(); } /// @@ -506,6 +507,8 @@ private void ReportExposure(in ExposureEvent exposure) exposureApi = _exposureApi; if (exposureApi is null) { + // Keep the HTTP client and its settings subscription lazy with the first exposure. + _evpTransport = new FeatureFlagsEvpTransport(_tracerSettings); exposureApi = new ExposureApi(_tracerSettings, _evpTransport); Volatile.Write(ref _exposureApi, exposureApi); } diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs index b85d65f68918..d8e0b256f68d 100644 --- a/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/FeatureFlagsModuleTests.cs @@ -9,6 +9,7 @@ using System.Collections.Generic; using System.Collections.Specialized; using System.Diagnostics; +using System.Reflection; using System.Threading; using System.Threading.Tasks; using Datadog.Trace.Configuration; @@ -27,6 +28,40 @@ namespace Datadog.Trace.Tests.FeatureFlags; public class FeatureFlagsModuleTests { + [Fact] + public void DefaultModuleDoesNotCreateExposureTransportUntilFirstUse() + { + var settings = new TracerSettings(new NameValueConfigurationSource(new NameValueCollection())); + using var module = FeatureFlagsModule.Create(settings, new MockRcmSubscriptionManager()); + var transportField = typeof(FeatureFlagsModule).GetField("_evpTransport", BindingFlags.Instance | BindingFlags.NonPublic)!; + + settings.FeatureFlags.Enabled.Should().BeTrue(); + module.Should().NotBeNull(); + transportField.GetValue(module).Should().BeNull("applications that never use exposures need no HTTP client or settings subscription"); + + var exposureApi = module!.GetExposureApi(); + exposureApi.Should().NotBeNull(); + transportField.GetValue(module).Should().NotBeNull(); + module.GetExposureApi().Should().BeSameAs(exposureApi); + + module.Dispose(); + transportField.GetValue(module).Should().BeNull(); + module.GetExposureApi().Should().BeNull(); + } + + [Fact] + public void DisposingUnusedModuleDoesNotCreateExposureTransport() + { + var settings = new TracerSettings(new NameValueConfigurationSource(new NameValueCollection())); + using var module = FeatureFlagsModule.Create(settings, new MockRcmSubscriptionManager()); + + module!.Dispose(); + + module.GetExposureApi().Should().BeNull(); + typeof(FeatureFlagsModule).GetField("_evpTransport", BindingFlags.Instance | BindingFlags.NonPublic)! + .GetValue(module).Should().BeNull(); + } + [Fact] public void UpdateRemoteConfig_WithEmptyList_InvokesCallbackAndReturnsProviderNotReady() {