Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
// <copyright file="FeatureFlagsEvpTransport.cs" company="Datadog">
// Unless explicitly stated otherwise all files in this repository are licensed under the Apache 2 License.
// This product includes software developed at Datadog (https://www.datadoghq.com/). Copyright 2017 Datadog, Inc.
// </copyright>

#nullable enable

using System;
using System.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;

/// <summary>
/// Owns the historical fixed-v2 Feature Flags event sender independently of exposure batching.
/// </summary>
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>(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);
}
}
144 changes: 84 additions & 60 deletions tracer/src/Datadog.Trace/FeatureFlags/Exposure/ExposureApi.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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()
{
Expand All @@ -39,45 +37,38 @@ internal sealed class ExposureApi : IDisposable
}
};

private readonly TaskCompletionSource<bool> _processExit = new();
private readonly TimeSpan _sendInterval = TimeSpan.FromSeconds(10);
private readonly TaskCompletionSource<bool> _processExit = new(TaskCreationOptions.RunContinuationsAsynchronously);
private readonly object _lifecycleLock = new();
private readonly TimeSpan _sendInterval;
private readonly TimeSpan _shutdownTimeout;
private readonly Queue<ExposureEvent> _exposures = new Queue<ExposureEvent>();

private readonly ExposureCache _exposureCache = new ExposureCache(DefaultCapacity);
private IApiRequestFactory _apiRequestFactory;
private Dictionary<string, string> _context;
private int _started;
private readonly FeatureFlagsEvpTransport _transport;
private readonly IDisposable _settingsSubscription;

internal ExposureApi(TracerSettings tracerSettings)
private Dictionary<string, string> _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)
{
Expand All @@ -94,39 +85,37 @@ 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()
{
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
{
Expand All @@ -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");
}
}

Expand All @@ -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<string, string> context, List<ExposureEvent> exposures)
Expand Down
10 changes: 9 additions & 1 deletion tracer/src/Datadog.Trace/FeatureFlags/FeatureFlagsModule.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -65,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;
Expand Down Expand Up @@ -129,6 +131,7 @@ public void Dispose()
ISubscription? subscription;
IFeatureFlagsDeliverySource? agentlessSource;
ExposureApi? exposureApi;
FeatureFlagsEvpTransport? evpTransport;

lock (_stateLock)
{
Expand All @@ -142,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
Expand All @@ -157,6 +162,7 @@ public void Dispose()

agentlessSource?.Dispose();
exposureApi?.Dispose();
evpTransport?.Dispose();
}

/// <summary>
Expand Down Expand Up @@ -501,7 +507,9 @@ private void ReportExposure(in ExposureEvent exposure)
exposureApi = _exposureApi;
if (exposureApi is null)
{
exposureApi = new ExposureApi(_tracerSettings);
// 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);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ internal class TestRequestFactory : IApiRequestFactory
{
private readonly Uri _baseEndpoint;
private readonly Func<Uri, TestApiRequest>[] _requestsToSend;
private readonly object _requestsLock = new();

public TestRequestFactory(params Func<Uri, TestApiRequest>[] requestsToSend)
: this(new Uri("http://localhost"), requestsToSend)
Expand All @@ -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)
Expand Down
Loading
Loading