diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Agentless/AgentlessConfigurationSource.cs b/tracer/src/Datadog.Trace/FeatureFlags/Agentless/AgentlessConfigurationSource.cs new file mode 100644 index 000000000000..d42449987fdb --- /dev/null +++ b/tracer/src/Datadog.Trace/FeatureFlags/Agentless/AgentlessConfigurationSource.cs @@ -0,0 +1,451 @@ +// +// 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.IO; +using System.IO.Compression; +using System.Threading; +using System.Threading.Tasks; +using Datadog.Trace.Agent; +using Datadog.Trace.Agent.Transports; +using Datadog.Trace.Configuration; +using Datadog.Trace.FeatureFlags.Rcm.Model; +using Datadog.Trace.Headers; +using Datadog.Trace.Logging; +using Datadog.Trace.Telemetry; +using Datadog.Trace.Util; + +namespace Datadog.Trace.FeatureFlags.Agentless; + +/// +/// Polls the agentless endpoint for flag configuration. Polling is billable, so it is only +/// started once application code has activated the provider. +/// +internal sealed class AgentlessConfigurationSource : IDisposable +{ + private const int MaxAttempts = 3; + private const double RetryJitter = 0.2; + + private static readonly TimeSpan FirstRetryMin = TimeSpan.FromSeconds(2); + private static readonly TimeSpan FirstRetryMax = TimeSpan.FromSeconds(10); + private static readonly TimeSpan SecondRetryMin = TimeSpan.FromSeconds(5); + private static readonly TimeSpan SecondRetryMax = TimeSpan.FromSeconds(30); + + // A jittered retry delay never drops below this, so a short poll interval cannot turn + // retries into a burst against the endpoint. + private static readonly TimeSpan MinRetryDelay = TimeSpan.FromSeconds(1); + + private static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor(typeof(AgentlessConfigurationSource)); + + private readonly IApiRequestFactory _requestFactory; + private readonly AgentlessEndpoint _endpoint; + private readonly TimeSpan _pollInterval; + private readonly Func _applyConfiguration; + private readonly Func _waitAsync; + + // Not a CancellationTokenSource: cancellation throws, and an exception on the shutdown path can + // crash the runtime, so shutdown is signalled by completing a task instead. + private readonly TaskCompletionSource _shutdown = new(TaskCreationOptions.RunContinuationsAsynchronously); + + // Only ever touched from the poll loop. + private readonly HashSet _loggedFailureCategories = new(); + + private readonly IDisposable? _environmentSubscription; + + // Written by the settings-change callback, read by the poll loop. The URI is stored rather than + // the environment it carries, so it is only built when the environment changes. + private Uri _requestUri; + + // Only ever touched from the poll loop. + private bool _malformedPayloadLogged; + private bool _applyFailureLogged; + private string? _etag; + private Uri? _etagUri; + + private int _started; + + internal AgentlessConfigurationSource( + AgentlessEndpoint endpoint, + IApiRequestFactory requestFactory, + TimeSpan pollInterval, + Func applyConfiguration, + string? environment = null, + TracerSettings.SettingsManager? settingsManager = null, + Func? waitAsync = null) + { + _endpoint = endpoint; + _requestFactory = requestFactory; + _pollInterval = pollInterval; + _applyConfiguration = applyConfiguration; + _requestUri = endpoint.BuildRequestUri(environment); + _waitAsync = waitAsync ?? Task.Delay; + + // Subscribed last, so every field the callback touches is already set. The environment is + // tracked rather than captured: customers can change it in code while the application runs, + // and flags are targeted per environment. + _environmentSubscription = settingsManager?.SubscribeToChanges(changes => + { + if (changes.UpdatedMutable is { } mutable) + { + UpdateEnvironment(mutable.Environment); + } + }); + } + + /// + /// Creates the source, or returns null when it cannot be operated: a base URL that is + /// not a URL, or the managed endpoint without an API key. Polling anyway would only produce + /// failures every interval. + /// + public static AgentlessConfigurationSource? Create( + FeatureFlagsSettings settings, + TracerSettings.SettingsManager manager, + Func applyConfiguration) + { + if (!AgentlessEndpoint.TryCreate(settings.Site, settings.AgentlessBaseUrl, out var endpoint, out var error)) + { + Log.Error("Feature Flags agentless source is unavailable: {Error}", error); + return null; + } + + if (endpoint.IsManaged && StringUtil.IsNullOrEmpty(settings.ApiKey)) + { + Log.Error("Feature Flags agentless source requires an API key. Set DD_API_KEY, or point DD_FEATURE_FLAGS_CONFIGURATION_SOURCE_AGENTLESS_BASE_URL at an endpoint of your own."); + return null; + } + + return new AgentlessConfigurationSource( + endpoint, + CreateRequestFactory(endpoint, settings), + settings.PollInterval, + applyConfiguration, + manager.InitialMutableSettings.Environment, + manager); + } + + /// + /// Records the environment to request configuration for, as the URI that carries it. Picked up + /// by the poll loop on its next request, so a change never disturbs a request already in flight. + /// + internal void UpdateEnvironment(string? environment) + => Volatile.Write(ref _requestUri, _endpoint.BuildRequestUri(environment)); + + /// + /// Starts polling. Idempotent. + /// + public void Start() + { + if (Interlocked.CompareExchange(ref _started, 1, 0) != 0) + { + return; + } + + // The loop runs on the thread pool, so nothing of it happens on the caller's thread. Nothing + // ever awaits it either, so a fault is observed here or not at all. + _ = Task.Run(RunAsync).ContinueWith(t => Log.Error(t.Exception, "Feature Flags agentless poll loop failed"), TaskContinuationOptions.OnlyOnFaulted); + } + + /// + /// Runs a single poll, including its in-tick retries. + /// + internal async Task PollAsync() + { + var result = default(PollResult); + + for (var attempt = 1; attempt <= MaxAttempts; attempt++) + { + result = await RequestAsync().ConfigureAwait(false); + + if (_shutdown.Task.IsCompleted) + { + // A shutdown mid-poll leaves the response unusable for state transitions: keep + // last-known-good and the current ETag. + return; + } + + if (!IsRetryable(in result)) + { + break; + } + + if (attempt == MaxAttempts) + { + // Every attempt failed in a retryable way. Last-known-good stays in place. + WarnFailure(in result, MaxAttempts); + return; + } + + await WaitAsync(RetryDelay(attempt)).ConfigureAwait(false); + + if (_shutdown.Task.IsCompleted) + { + return; + } + } + + if (_shutdown.Task.IsCompleted) + { + // A shutdown during the final attempt leaves the response unusable for state + // transitions: keep last-known-good and the current ETag. + return; + } + + Apply(in result); + + // A failure with no status code never reached the endpoint, so it is worth another attempt. + static bool IsRetryable(in PollResult result) + => result.StatusCode is not { } status || status is 408 or 429 or (>= 500 and <= 599); + } + + public void Dispose() + { + // The request in flight is bounded by the request timeout, and the loop is never joined, + // so a shutdown does not wait for it. A poll that completes after disposal is prevented + // from applying its result by the shutdown check in PollAsync. + _shutdown.TrySetResult(true); + + try + { + _environmentSubscription?.Dispose(); + } + catch (Exception ex) + { + Log.Debug(ex, "Error unsubscribing the Feature Flags agentless poll loop from settings changes"); + } + } + + // The concrete type is returned rather than the interface because CA1859 asks for it on a + // private member, which is also why the signature varies by target framework. +#if NETCOREAPP + private static HttpClientRequestFactory CreateRequestFactory(AgentlessEndpoint endpoint, FeatureFlagsSettings settings) +#else + private static ApiWebRequestFactory CreateRequestFactory(AgentlessEndpoint endpoint, FeatureFlagsSettings settings) +#endif + { + var headers = new List> + { + // The endpoint serves gzip, and neither transport decompresses for us. + new("Accept-Encoding", "gzip"), + new(TelemetryConstants.ClientLibraryLanguageHeader, TracerConstants.Language), + new(TelemetryConstants.ClientLibraryVersionHeader, TracerConstants.ThreePartVersion), + + // Without this the poll is itself instrumented, producing a span per poll and letting + // auto-instrumentation recurse through the poller's own client. + new(HttpHeaderNames.TracingEnabled, "false"), + }; + + if (endpoint.IsManaged) + { + // A custom endpoint is left to report its own authentication failure rather than + // having the Datadog credential sent to it. + headers.Add(new(TelemetryConstants.ApiKeyHeader, settings.ApiKey!)); + } + + // The endpoint is only the factory's default: both transports honour the URI passed to + // Create, which is what carries the current environment. +#if NETCOREAPP + return new HttpClientRequestFactory(endpoint.Uri, headers.ToArray(), timeout: settings.RequestTimeout); +#else + return new ApiWebRequestFactory(endpoint.Uri, headers.ToArray(), timeout: settings.RequestTimeout); +#endif + } + + private async Task RunAsync() + { + Log.Debug("AgentlessConfigurationSource::RunAsync -> Enter"); + + while (!_shutdown.Task.IsCompleted) + { + try + { + await PollAsync().ConfigureAwait(false); + } + catch (Exception ex) + { + Log.Error(ex, "Feature Flags agentless poll failed unexpectedly"); + } + + // Fixed delay after completion, so polls never overlap. + await WaitAsync(_pollInterval).ConfigureAwait(false); + } + + Log.Debug("AgentlessConfigurationSource::RunAsync -> Exit"); + } + + // A shutdown ends the wait early. The delay itself is left to expire on its own: it holds no + // thread, and the loop has already exited by the time it does. + private async Task WaitAsync(TimeSpan delay) + => await Task.WhenAny(_waitAsync(delay), _shutdown.Task).ConfigureAwait(false); + + private TimeSpan RetryDelay(int attempt) + { + var seconds = attempt == 1 + ? Clamp(_pollInterval.TotalSeconds / 6, FirstRetryMin, FirstRetryMax) + : Clamp(_pollInterval.TotalSeconds / 3, SecondRetryMin, SecondRetryMax); + + var jitter = 1 - RetryJitter + (ThreadSafeRandom.Shared.NextDouble() * RetryJitter * 2); + + return TimeSpan.FromSeconds(Math.Max(MinRetryDelay.TotalSeconds, seconds * jitter)); + + static double Clamp(double value, TimeSpan minimum, TimeSpan maximum) + => Math.Max(minimum.TotalSeconds, Math.Min(maximum.TotalSeconds, value)); + } + + private async Task RequestAsync() + { + try + { + // Read per request rather than captured, because the environment it carries can be + // changed in code after startup. + var uri = Volatile.Read(ref _requestUri); + + // An ETag only identifies the configuration served for the URI it came from. Sending it + // against a different environment would earn a 304 and pin the process to the previous + // environment's flags, with no way back. + if (_etagUri is not null && uri != _etagUri) + { + _etag = null; + } + + _etagUri = uri; + + var request = _requestFactory.Create(uri); + if (_etag is { } etag) + { + request.AddHeader("If-None-Match", etag); + } + + using var response = await request.GetAsync().ConfigureAwait(false); + + // Only a 200 carries configuration; other bodies are never decoded as one. The payload + // is parsed here, while the response is still open, so it never has to be held as a + // string. A payload that does not parse is reported, not thrown: it is not retryable. + ServerConfiguration? configuration = null; + string? parseError = null; + + if (response.StatusCode == 200) + { + using var stream = await response.GetStreamAsync().ConfigureAwait(false); + + using var decompressed = + response.GetContentEncodingType() == ContentEncodingType.GZip + ? new GZipStream(stream, CompressionMode.Decompress, leaveOpen: true) + : null; + + // Every parameter has to be given to reach leaveOpen. A byte order mark is not + // expected, and letting one be detected would override the declared encoding. + using var reader = new StreamReader( + decompressed ?? stream, + response.GetCharsetEncoding(), + detectEncodingFromByteOrderMarks: false, + bufferSize: 1024, // the default + leaveOpen: true); + + if (UfcConfigurationParser.TryParse(reader, out var parsed, out parseError)) + { + configuration = parsed; + } + } + + return new PollResult(response.StatusCode, response.GetHeader("ETag"), configuration, parseError, error: null); + } + catch (Exception ex) + { + return new PollResult(statusCode: null, etag: null, configuration: null, parseError: null, error: ex); + } + } + + private void Apply(in PollResult result) + { + switch (result.StatusCode) + { + case 304: + // Nothing changed, and the ETag stays as it is. + return; + case 401 or 403: + WarnFailure(in result, attempts: 1); + return; + case not 200: + WarnFailure(in result, attempts: 1); + return; + } + + if (result.Configuration is not { } configuration) + { + if (!_malformedPayloadLogged) + { + _malformedPayloadLogged = true; + Log.Error("Feature Flags agentless endpoint returned an unusable payload: {Error}", result.ParseError); + } + + return; + } + + if (!_applyConfiguration(configuration)) + { + if (!_applyFailureLogged) + { + _applyFailureLogged = true; + Log.Warning("Feature Flags agentless configuration could not be applied"); + } + + return; + } + + // The ETag advances only once parsing and applying have both succeeded. Advancing on + // receipt would acknowledge a payload that was never applied, and every later poll would + // answer 304, pinning the process to stale configuration with no way back. + var newEtag = result.ETag?.Trim(); + _etag = StringUtil.IsNullOrEmpty(newEtag) ? null : newEtag; + } + + /// + /// Warns once per failure category. A dead endpoint would otherwise produce a warning every + /// poll interval, indefinitely. + /// + private void WarnFailure(in PollResult result, int attempts) + { + var category = result.StatusCode switch + { + 401 or 403 => "authentication", + not null => "http", + _ => "request", + }; + + if (!_loggedFailureCategories.Add(category)) + { + return; + } + + switch (result.StatusCode) + { + case 401 or 403: + Log.Warning("Feature Flags agentless endpoint returned HTTP {StatusCode}; verify endpoint authentication", result.StatusCode!.Value); + break; + case not null: + Log.Warning("Feature Flags agentless endpoint returned HTTP {StatusCode} after {Attempts} attempts", result.StatusCode.Value, attempts); + break; + default: + Log.Warning(result.Error, "Feature Flags agentless request failed after {Attempts} attempts", attempts); + break; + } + } + + internal readonly struct PollResult(int? statusCode, string? etag, ServerConfiguration? configuration, string? parseError, Exception? error) + { + public int? StatusCode { get; } = statusCode; + + public string? ETag { get; } = etag; + + public ServerConfiguration? Configuration { get; } = configuration; + + public string? ParseError { get; } = parseError; + + public Exception? Error { get; } = error; + } +} diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Agentless/UfcConfigurationParser.cs b/tracer/src/Datadog.Trace/FeatureFlags/Agentless/UfcConfigurationParser.cs new file mode 100644 index 000000000000..1dddb52db0bc --- /dev/null +++ b/tracer/src/Datadog.Trace/FeatureFlags/Agentless/UfcConfigurationParser.cs @@ -0,0 +1,176 @@ +// +// 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.Diagnostics.CodeAnalysis; +using System.IO; +using Datadog.Trace.FeatureFlags.Rcm.Model; +using Datadog.Trace.Util.Json; +using Datadog.Trace.Vendors.Newtonsoft.Json; + +namespace Datadog.Trace.FeatureFlags.Agentless; + +/// +/// Reads the JSON:API envelope returned by the agentless endpoint. +/// +internal static class UfcConfigurationParser +{ + private const string ResourceType = "universal-flag-configuration"; + + private const string MalformedError = "Malformed UFC payload"; + private const string ResourceError = "Expected a JSON:API Universal Flag Configuration resource"; + private const string AttributesError = "Expected a Universal Flag Configuration v1 object"; + + /// + /// Validates a JSON:API Universal Flag Configuration response and returns data.attributes, + /// which is the document the evaluator consumes. A raw UFC document is rejected, including from + /// a custom endpoint, so that every source agrees on one wire format. + /// + /// The envelope is walked with the reader rather than loaded into a JSON tree, and + /// data.attributes is deserialized in place, so the payload is read exactly once and no + /// copy of it is ever held. + /// + /// + /// The response body. Read straight from the response, so it never has to be held as a string. + /// The parsed configuration. + /// Why the payload was rejected. + /// true when the payload matches the contract. + public static bool TryParse(TextReader body, [NotNullWhen(true)] out ServerConfiguration? configuration, out string? error) + { + configuration = null; + error = null; + + var sawData = false; + string? resourceType = null; + ServerConfiguration? attributes = null; + + try + { + // Timestamps stay strings: the model carries createdAt verbatim, and letting Newtonsoft + // turn it into a date would also make the type check below fail. The reader belongs to + // the caller, which owns the response it came from. + using var reader = new JsonTextReader(body) { DateParseHandling = DateParseHandling.None, CloseInput = false, ArrayPool = JsonArrayPool.Shared }; + var serializer = new JsonSerializer { DateParseHandling = DateParseHandling.None }; + + if (!reader.Read()) + { + // Nothing at all, so there is no document to judge against the contract. + error = MalformedError; + return false; + } + + if (reader.TokenType != JsonToken.StartObject) + { + error = ResourceError; + return false; + } + + while (reader.Read() && reader.TokenType == JsonToken.PropertyName) + { + if ((string?)reader.Value != "data") + { + reader.Skip(); + continue; + } + + sawData = true; + + if (!reader.Read()) + { + // The document ended where the resource should have been. + error = MalformedError; + return false; + } + + if (reader.TokenType != JsonToken.StartObject) + { + error = ResourceError; + return false; + } + + while (reader.Read() && reader.TokenType == JsonToken.PropertyName) + { + switch ((string?)reader.Value) + { + case "type": + if (!reader.Read()) + { + error = MalformedError; + return false; + } + + // A type that is not a string cannot identify the resource. Checked on + // the token, because a number would otherwise be read as its digits. + if (reader.TokenType != JsonToken.String) + { + error = ResourceError; + return false; + } + + resourceType = (string?)reader.Value; + break; + + case "attributes": + if (!reader.Read()) + { + error = MalformedError; + return false; + } + + if (reader.TokenType != JsonToken.StartObject) + { + error = AttributesError; + return false; + } + + attributes = serializer.Deserialize(reader); + break; + + default: + reader.Skip(); + break; + } + } + } + + // A document that ends before the root object closes was truncated in transit, whatever + // was found in it up to that point. + if (reader.TokenType != JsonToken.EndObject) + { + error = MalformedError; + return false; + } + } + catch (Exception) + { + error = MalformedError; + return false; + } + + if (!sawData || resourceType != ResourceType) + { + error = ResourceError; + return false; + } + + // Every member of the v1 contract has to be there. A member of the wrong shape arrives as + // null, because the flag collection rejects anything that is not an object and Newtonsoft + // leaves a member it cannot convert unset. + if (attributes is null + || attributes.Format is null + || attributes.CreatedAt is null + || attributes.Environment?.Name is null + || attributes.Flags is null) + { + error = AttributesError; + return false; + } + + configuration = attributes; + return true; + } +} diff --git a/tracer/src/Datadog.Trace/FeatureFlags/Rcm/Model/FlagCollectionJsonConverter.cs b/tracer/src/Datadog.Trace/FeatureFlags/Rcm/Model/FlagCollectionJsonConverter.cs index 87a4a8937986..101818093e56 100644 --- a/tracer/src/Datadog.Trace/FeatureFlags/Rcm/Model/FlagCollectionJsonConverter.cs +++ b/tracer/src/Datadog.Trace/FeatureFlags/Rcm/Model/FlagCollectionJsonConverter.cs @@ -25,6 +25,15 @@ internal sealed class FlagCollectionJsonConverter : JsonConverter +// 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.IO; +using System.IO.Compression; +using System.Linq; +using System.Text; +using System.Threading.Tasks; +using Datadog.Trace.Agent; +using Datadog.Trace.FeatureFlags.Agentless; +using Datadog.Trace.FeatureFlags.Rcm.Model; +using Datadog.Trace.TestHelpers.TransportHelpers; +using FluentAssertions; +using Xunit; + +namespace Datadog.Trace.Tests.FeatureFlags; + +public class AgentlessConfigurationSourceTests +{ + private const string Body = """ + { "data": { "type": "universal-flag-configuration", + "attributes": { "format": "SERVER", "createdAt": "2025-01-01T00:00:00Z", + "environment": { "name": "production" }, "flags": {} } } } + """; + + private const string EndpointUrl = "https://ufc-server.ff-cdn.datadoghq.com/api/v2/feature-flagging/config/rules-based/server"; + + [Fact] + public async Task AppliesConfigurationFromA200() + { + var applied = new List(); + var factory = new TestRequestFactory(uri => new TestApiRequest(uri, responseContent: Body)); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().ContainSingle(); + applied[0].Environment!.Name.Should().Be("production"); + factory.RequestsSent.Should().ContainSingle(); + } + + [Fact] + public async Task SendsTheEtagOfTheLastAppliedConfiguration() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, responseContent: Body, responseHeaders: new() { { "ETag", "\"ufc-v1\"" } })); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + applied.Should().ContainSingle(); + + // Second poll should send If-None-Match + await source.PollAsync(); + factory.RequestsSent.Should().HaveCount(2); + factory.RequestsSent[1].ExtraHeaders.Should().ContainKey("If-None-Match"); + factory.RequestsSent[1].ExtraHeaders["If-None-Match"].Should().Be("\"ufc-v1\""); + } + + [Fact] + public async Task RequestsTheConfiguredEnvironment() + { + var applied = new List(); + var factory = new TestRequestFactory(uri => new TestApiRequest(uri, responseContent: Body)); + using var source = CreateSource(factory, applied, environment: "production"); + + await source.PollAsync(); + + factory.RequestsSent[0].Endpoint.Should().Be(new Uri(EndpointUrl + "?dd_env=production")); + } + + [Fact] + public async Task DropsTheEtagWhenTheEnvironmentChanges() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, responseContent: Body, responseHeaders: new() { { "ETag", "\"ufc-v1\"" } })); + using var source = CreateSource(factory, applied, environment: "production"); + + await source.PollAsync(); + + // The ETag identifies production's configuration, so it must not be sent against staging: + // a 304 would pin the process to production's flags with no way back. + source.UpdateEnvironment("staging"); + await source.PollAsync(); + + factory.RequestsSent[1].Endpoint.Should().Be(new Uri(EndpointUrl + "?dd_env=staging")); + factory.RequestsSent[1].ExtraHeaders.Should().NotContainKey("If-None-Match"); + } + + [Fact] + public async Task DoesNotApplyOn304() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, statusCode: 304, responseContent: "{}")); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().BeEmpty(); + } + + [Fact] + public async Task DoesNotApplyOn401() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, statusCode: 401, responseContent: "Unauthorized")); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().BeEmpty(); + } + + [Fact] + public async Task DoesNotApplyOnMalformedPayload() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, statusCode: 200, responseContent: "not json")); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().BeEmpty(); + } + + [Fact] + public async Task DoesNotApplyAfterDisposal() + { + var applied = new List(); + var factory = new TestRequestFactory(uri => new TestApiRequest(uri, responseContent: Body)); + var source = CreateSource(factory, applied); + + // A shutdown mid-poll leaves the response unusable for a state transition. + source.Dispose(); + await source.PollAsync(); + + factory.RequestsSent.Should().ContainSingle(); + applied.Should().BeEmpty(); + } + + [Fact] + public async Task DoesNotApplyWhenDisposedAfterRequestSucceeds() + { + var applied = new List(); + AgentlessConfigurationSource? sourceRef = null; + var factory = new TestRequestFactory(uri => + { + var request = new DisposingApiRequest(uri, Body); + request.Source = sourceRef; + return request; + }); + using var source = CreateSource(factory, applied); + sourceRef = source; + + await source.PollAsync(); + + applied.Should().BeEmpty(); + } + + [Fact] + public async Task RetriesOn500ThenAppliesOnSuccess() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, statusCode: 500, responseContent: "error"), + uri => new TestApiRequest(uri, responseContent: Body)); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().ContainSingle(); + factory.RequestsSent.Should().HaveCount(2); + } + + [Fact] + public async Task RetriesUpToMaxAttemptsOn500() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, statusCode: 500, responseContent: "error"), + uri => new TestApiRequest(uri, statusCode: 500, responseContent: "error"), + uri => new TestApiRequest(uri, statusCode: 500, responseContent: "error")); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().BeEmpty(); + factory.RequestsSent.Should().HaveCount(3); + } + + [Fact] + public async Task DoesNotRetryOn400() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new TestApiRequest(uri, statusCode: 400, responseContent: "bad request")); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().BeEmpty(); + factory.RequestsSent.Should().ContainSingle(); + } + + [Fact] + public async Task HandlesGzipResponse() + { + var applied = new List(); + var factory = new TestRequestFactory(uri => new GzipApiRequest(uri, Body)); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().ContainSingle(); + applied[0].Environment!.Name.Should().Be("production"); + } + + [Fact] + public async Task HandlesNetworkError() + { + var applied = new List(); + var factory = new TestRequestFactory( + uri => new ThrowingApiRequest(uri), + uri => new ThrowingApiRequest(uri), + uri => new ThrowingApiRequest(uri)); + using var source = CreateSource(factory, applied); + + await source.PollAsync(); + + applied.Should().BeEmpty(); + factory.RequestsSent.Should().HaveCount(3); + } + + private static AgentlessConfigurationSource CreateSource( + TestRequestFactory factory, + List applied, + string? environment = null) + => new( + CreateEndpoint(), + factory, + TimeSpan.FromSeconds(30), + configuration => + { + applied.Add(configuration); + return true; + }, + environment, + waitAsync: NoWait); + + private static AgentlessEndpoint CreateEndpoint() + { + AgentlessEndpoint.TryCreate("datadoghq.com", baseUrl: null, out var endpoint, out _).Should().BeTrue(); + return endpoint ?? throw new InvalidOperationException("TryCreate reported success without producing an endpoint."); + } + + private static Task NoWait(TimeSpan delay) => Task.CompletedTask; + + private class ThrowingApiRequest(Uri endpoint) : TestApiRequest(endpoint) + { + public override Task GetAsync() => throw new IOException("The connection was refused"); + } + + private class GzipApiRequest(Uri endpoint, string body) : TestApiRequest(endpoint) + { + public override Task GetAsync() => Task.FromResult(new GzipApiResponse(body)); + } + + private class GzipApiResponse(string body) : IApiResponse + { + public int StatusCode => 200; + + public long ContentLength => -1; + + public string? ContentTypeHeader => "application/json"; + + public string? ContentEncodingHeader => "gzip"; + + public void Dispose() + { + } + + public string? GetHeader(string headerName) => null; + + public Encoding GetCharsetEncoding() => Encoding.UTF8; + + public ContentEncodingType GetContentEncodingType() => ContentEncodingType.GZip; + + public Task GetStreamAsync() + { + var compressed = new MemoryStream(); + using (var gzip = new GZipStream(compressed, CompressionMode.Compress, leaveOpen: true)) + { + var bytes = Encoding.UTF8.GetBytes(body); + gzip.Write(bytes, 0, bytes.Length); + } + + compressed.Position = 0; + return Task.FromResult(compressed); + } + } + + private class DisposingApiRequest(Uri endpoint, string body) : TestApiRequest(endpoint, responseContent: body) + { + public AgentlessConfigurationSource? Source { get; set; } + + public override Task GetAsync() + { + var response = base.GetAsync(); + // Simulate a shutdown arriving after the request completes but before ApplyAsync. + Source?.Dispose(); + return response; + } + } +} diff --git a/tracer/test/Datadog.Trace.Tests/FeatureFlags/UfcConfigurationParserTests.cs b/tracer/test/Datadog.Trace.Tests/FeatureFlags/UfcConfigurationParserTests.cs new file mode 100644 index 000000000000..939151053415 --- /dev/null +++ b/tracer/test/Datadog.Trace.Tests/FeatureFlags/UfcConfigurationParserTests.cs @@ -0,0 +1,132 @@ +// +// 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 System.IO; +using Datadog.Trace.FeatureFlags.Agentless; +using Datadog.Trace.FeatureFlags.Rcm.Model; +using FluentAssertions; +using Xunit; + +namespace Datadog.Trace.Tests.FeatureFlags; + +public class UfcConfigurationParserTests +{ + private const string ValidEnvelope = """ + { "data": { "type": "universal-flag-configuration", + "attributes": { "format": "SERVER", "createdAt": "2025-01-01T00:00:00Z", + "environment": { "name": "production" }, "flags": {} } } } + """; + + private const string Attributes = """ + { "format": "SERVER", "createdAt": "2025-01-01T00:00:00Z", + "environment": { "name": "production" }, "flags": {} } + """; + + [Fact] + public void ParsesValidEnvelope() + { + Parse(ValidEnvelope, out var configuration, out var error) + .Should().BeTrue(); + + error.Should().BeNull(); + configuration.Should().NotBeNull(); + configuration!.Environment!.Name.Should().Be("production"); + configuration.Flags.Should().BeEmpty(); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData("not json")] + [InlineData("{ \"data\": ")] + public void RejectsMalformedJson(string? body) + { + Parse(body, out var configuration, out var error).Should().BeFalse(); + + configuration.Should().BeNull(); + error.Should().Be("Malformed UFC payload"); + } + + [Theory] + // A raw UFC document is rejected too, so every source agrees on one wire format. + [InlineData(Attributes)] + // Wrong resource type + [InlineData("""{ "data": { "type": "wrong-type", "attributes": { "format": "SERVER", "createdAt": "x", "environment": { "name": "prod" }, "flags": {} } } }""")] + // Missing data + [InlineData("""{ "meta": {} }""")] + // data is not an object + [InlineData("""{ "data": "string" }""")] + // data.type is not a string (object) + [InlineData("""{ "data": { "type": { "nested": true }, "attributes": { "format": "SERVER", "createdAt": "x", "environment": { "name": "prod" }, "flags": {} } } }""")] + // data.type is not a string (array) + [InlineData("""{ "data": { "type": [1, 2], "attributes": { "format": "SERVER", "createdAt": "x", "environment": { "name": "prod" }, "flags": {} } } }""")] + public void RejectsInvalidEnvelope(string body) + { + Parse(body, out var configuration, out var error).Should().BeFalse(); + + configuration.Should().BeNull(); + error.Should().Be("Expected a JSON:API Universal Flag Configuration resource"); + } + + [Theory] + // Missing format + [InlineData("""{ "data": { "type": "universal-flag-configuration", "attributes": { "createdAt": "x", "environment": { "name": "prod" }, "flags": {} } } }""")] + // Missing createdAt + [InlineData("""{ "data": { "type": "universal-flag-configuration", "attributes": { "format": "SERVER", "environment": { "name": "prod" }, "flags": {} } } }""")] + // Missing environment + [InlineData("""{ "data": { "type": "universal-flag-configuration", "attributes": { "format": "SERVER", "createdAt": "x", "flags": {} } } }""")] + // Missing flags + [InlineData("""{ "data": { "type": "universal-flag-configuration", "attributes": { "format": "SERVER", "createdAt": "x", "environment": { "name": "prod" } } } }""")] + // flags is not an object + [InlineData("""{ "data": { "type": "universal-flag-configuration", "attributes": { "format": "SERVER", "createdAt": "x", "environment": { "name": "prod" }, "flags": [] } } }""")] + public void RejectsInvalidAttributes(string body) + { + Parse(body, out var configuration, out var error).Should().BeFalse(); + + configuration.Should().BeNull(); + error.Should().Be("Expected a Universal Flag Configuration v1 object"); + } + + [Fact] + public void ParsesFlagsFromEnvelope() + { + var body = """ + { "data": { "type": "universal-flag-configuration", + "attributes": { "format": "SERVER", "createdAt": "2025-01-01T00:00:00Z", + "environment": { "name": "production" }, + "flags": { "test-flag": { "key": "test-flag", "enabled": true, "variationType": "BOOLEAN" } } } } } + """; + + Parse(body, out var configuration, out _).Should().BeTrue(); + + configuration!.Flags.Should().NotBeNull(); + configuration!.Flags!.Should().ContainKey("test-flag"); + configuration!.Flags!["test-flag"].Enabled.Should().BeTrue(); + } + + [Fact] + public void AcceptsANumericEnvironmentName() + { + // The attributes are deserialized straight from the reader, so a scalar of the wrong type is + // coerced rather than rejected. The environment name is an opaque string to us, so a number + // read as its digits is harmless: the request it targets would not have matched anyway. + var body = """{ "data": { "type": "universal-flag-configuration", "attributes": { "format": "SERVER", "createdAt": "x", "environment": { "name": 123 }, "flags": {} } } }"""; + + Parse(body, out var configuration, out var error).Should().BeTrue(); + + error.Should().BeNull(); + configuration!.Environment!.Name.Should().Be("123"); + } + + // The parser reads the response stream directly, so a body under test is handed to it as a reader. + private static bool Parse(string? body, out ServerConfiguration? configuration, out string? error) + { + using var reader = new StringReader(body ?? string.Empty); + return UfcConfigurationParser.TryParse(reader, out configuration, out error); + } +}