From b16a7dcc32da62fcb7e049d823a068d00b7fea7c Mon Sep 17 00:00:00 2001 From: Leo Romanovsky Date: Thu, 1 Oct 2026 21:55:05 +0000 Subject: [PATCH] Bind Agent discovery capabilities to the configured endpoint Invalidate cached capabilities on endpoint replacement and reject stale responses before publishing configuration or hashes. Report EVP identity-header forwarding support without changing product routing. Environment: Datadog workspace --- .../DiscoveryService/AgentConfiguration.cs | 11 +- .../DiscoveryService/DiscoveryService.cs | 83 +++++++++----- .../Agent/DiscoveryServiceMock.cs | 6 +- .../Agent/DiscoveryServiceTests.cs | 101 ++++++++++++++++++ 4 files changed, 172 insertions(+), 29 deletions(-) diff --git a/tracer/src/Datadog.Trace/Agent/DiscoveryService/AgentConfiguration.cs b/tracer/src/Datadog.Trace/Agent/DiscoveryService/AgentConfiguration.cs index 0ba79525cf22..052b0fa43dcb 100644 --- a/tracer/src/Datadog.Trace/Agent/DiscoveryService/AgentConfiguration.cs +++ b/tracer/src/Datadog.Trace/Agent/DiscoveryService/AgentConfiguration.cs @@ -6,6 +6,7 @@ #nullable enable using System.Collections.Generic; +using Datadog.Trace.Configuration; namespace Datadog.Trace.Agent.DiscoveryService; @@ -30,7 +31,8 @@ public AgentConfiguration( List? peerTags = null, int obfuscationVersion = 0, AgentTraceFilterConfig? traceFilterConfig = null, - List? featureFlags = null) + List? featureFlags = null, + bool eventPlatformProxySupportsEvpOriginHeaders = false) { ConfigurationEndpoint = configurationEndpoint; DebuggerEndpoint = debuggerEndpoint; @@ -51,6 +53,7 @@ public AgentConfiguration( ObfuscationVersion = obfuscationVersion; TraceFilterConfig = traceFilterConfig ?? AgentTraceFilterConfig.Empty; FeatureFlags = featureFlags; + EventPlatformProxySupportsEvpOriginHeaders = eventPlatformProxySupportsEvpOriginHeaders; } public string? ConfigurationEndpoint { get; } @@ -83,6 +86,8 @@ public AgentConfiguration( public string? EventPlatformProxyEndpoint { get; } + public bool EventPlatformProxySupportsEvpOriginHeaders { get; } + public string? TelemetryProxyEndpoint { get; } public string? TracerFlareEndpoint { get; } @@ -102,4 +107,8 @@ public AgentConfiguration( public AgentTraceFilterConfig TraceFilterConfig { get; } public List? FeatureFlags { get; } + + // Bind capabilities to the immutable exporter snapshot used for discovery, including UDS + // and named pipes. Internal so record diagnostics do not print endpoint configuration. + internal ExporterSettings? DiscoverySettings { get; init; } } diff --git a/tracer/src/Datadog.Trace/Agent/DiscoveryService/DiscoveryService.cs b/tracer/src/Datadog.Trace/Agent/DiscoveryService/DiscoveryService.cs index a4f5b96c7d14..51c54fa24823 100644 --- a/tracer/src/Datadog.Trace/Agent/DiscoveryService/DiscoveryService.cs +++ b/tracer/src/Datadog.Trace/Agent/DiscoveryService/DiscoveryService.cs @@ -35,6 +35,8 @@ internal sealed class DiscoveryService : IDiscoveryService private const string SupportedDataStreamsEndpoint = "v0.1/pipeline_stats"; private const string SupportedEventPlatformProxyEndpointV2 = "evp_proxy/v2"; private const string SupportedEventPlatformProxyEndpointV4 = "evp_proxy/v4"; + private const string EvpOriginHeader = "DD-EVP-ORIGIN"; + private const string EvpOriginVersionHeader = "DD-EVP-ORIGIN-VERSION"; private const string SupportedTelemetryProxyEndpoint = "telemetry/proxy"; private const string SupportedTracerFlareEndpoint = "tracer_flare/v1"; @@ -62,7 +64,7 @@ public DiscoveryService( int initialRetryDelayMs, int maxRetryDelayMs, int recheckIntervalMs) - : this(CreateApiRequestFactory(settings.InitialExporterSettings, containerMetadata.ContainerId, tcpTimeout), serviceRemappingHash, initialRetryDelayMs, maxRetryDelayMs, recheckIntervalMs) + : this(CreateApiRequestFactory(settings.InitialExporterSettings, containerMetadata.ContainerId, tcpTimeout), serviceRemappingHash, initialRetryDelayMs, maxRetryDelayMs, recheckIntervalMs, exporterSettings: settings.InitialExporterSettings) { // Create as a "managed" service that can update the request factory _settingSubscription = settings.SubscribeToChanges(changes => @@ -70,7 +72,7 @@ public DiscoveryService( if (changes.UpdatedExporter is { } exporter) { var newFactory = CreateApiRequestFactory(exporter, containerMetadata.ContainerId, tcpTimeout); - Interlocked.Exchange(ref _apiRequestFactory!, new(newFactory)); + UpdateRequestFactory(newFactory, exporter); } }); } @@ -82,9 +84,10 @@ internal DiscoveryService( int initialRetryDelayMs, int maxRetryDelayMs, int recheckIntervalMs, - bool autoStartLoop = true) + bool autoStartLoop = true, + ExporterSettings? exporterSettings = null) { - _apiRequestFactory = new(apiRequestFactory); + _apiRequestFactory = new(apiRequestFactory, exporterSettings); _serviceRemappingHash = serviceRemappingHash; _initialRetryDelayMs = initialRetryDelayMs; _maxRetryDelayMs = maxRetryDelayMs; @@ -165,7 +168,21 @@ public static DiscoveryService CreateUnmanaged( serviceRemappingHash, initialRetryDelayMs, maxRetryDelayMs, - recheckIntervalMs); + recheckIntervalMs, + exporterSettings: exporterSettings); + + [TestingAndPrivateOnly] + internal void UpdateRequestFactory(IApiRequestFactory factory, ExporterSettings exporterSettings) + { + lock (_lock) + { + _apiRequestFactory = new(factory, exporterSettings); + // A different Agent must be polled even if it returns the same body/hash. Never + // replay the previous Agent's cached capabilities to a new subscriber. + _configuration = null; + _configurationHash = null; + } + } /// public void SubscribeToChanges(Action callback) @@ -211,11 +228,29 @@ public void SetCurrentConfigStateHash(string configStateHash) Interlocked.Exchange(ref _agentConfigStateHashUnixTime, DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()); } - private void NotifySubscribers(AgentConfiguration newConfig) + private void NotifySubscribers(AgentConfiguration newConfig, ApiFactoryHolder requestFactory, string configurationHash, string? containerTagsHash) { List> subscribers; lock (_lock) { + if (!ReferenceEquals(requestFactory, _apiRequestFactory)) + { + // A response that started before an endpoint change cannot validate the new Agent. + return; + } + + if (containerTagsHash is not null) + { + _serviceRemappingHash.UpdateContainerTagsHash(containerTagsHash); + } + + _configurationHash = configurationHash; + if (newConfig.Equals(_configuration)) + { + return; + } + + Log.Debug("Discovery configuration updated, notifying subscribers: {Configuration}", newConfig); subscribers = _agentChangeCallbacks.ToList(); // Setting the configuration immediately after grabbing // the subscribers ensures subscribers receive the @@ -262,7 +297,7 @@ private void NotifySubscribers(AgentConfiguration newConfig) using var response = await api.GetAsync().ConfigureAwait(false); if (response.StatusCode is >= 200 and < 300) { - await ProcessDiscoveryResponse(response).ConfigureAwait(false); + await ProcessDiscoveryResponse(response, requestFactory).ConfigureAwait(false); return null; } @@ -313,14 +348,10 @@ internal bool RequireRefresh(string? currentHash, DateTimeOffset utcNow) return Volatile.Read(ref _agentConfigStateHashUnixTime) + _recheckIntervalMs < utcNow.ToUnixTimeMilliseconds(); } - private async Task ProcessDiscoveryResponse(IApiResponse response) + private async Task ProcessDiscoveryResponse(IApiResponse response, ApiFactoryHolder requestFactory) { // Extract and store container tags hash from response headers var containerTagsHash = response.GetHeader(AgentHttpHeaderNames.ContainerTagsHash); - if (containerTagsHash != null) - { - _serviceRemappingHash.UpdateContainerTagsHash(containerTagsHash); - } // Grab the original stream var stream = await response.GetStreamAsync().ConfigureAwait(false); @@ -373,6 +404,10 @@ private async Task ProcessDiscoveryResponse(IApiResponse response) } var discoveredEndpoints = (jObject["endpoints"] as JArray)?.Values().ToArray(); + var evpProxyAllowedHeaders = (jObject["evp_proxy_allowed_headers"] as JArray)?.Values().ToArray(); + var eventPlatformProxySupportsEvpOriginHeaders = + evpProxyAllowedHeaders?.Any(header => string.Equals(header?.Trim(), EvpOriginHeader, StringComparison.OrdinalIgnoreCase)) == true + && evpProxyAllowedHeaders.Any(header => string.Equals(header?.Trim(), EvpOriginVersionHeader, StringComparison.OrdinalIgnoreCase)); string? configurationEndpoint = null; string? debuggerEndpoint = null; string? debuggerV2Endpoint = null; @@ -442,8 +477,6 @@ private async Task ProcessDiscoveryResponse(IApiResponse response) } } - var existingConfiguration = _configuration; - var newConfig = new AgentConfiguration( configurationEndpoint: configurationEndpoint, debuggerEndpoint: debuggerEndpoint, @@ -456,24 +489,20 @@ private async Task ProcessDiscoveryResponse(IApiResponse response) eventPlatformProxyEndpoint: eventPlatformProxyEndpoint, telemetryProxyEndpoint: telemetryProxyEndpoint, tracerFlareEndpoint: tracerFlareEndpoint, - containerTagsHash: _serviceRemappingHash.ContainerTagsHash, // either the value just received, or the one we stored before (prevents overriding with null) + containerTagsHash: containerTagsHash ?? _serviceRemappingHash.ContainerTagsHash, clientDropP0: clientDropP0, spanMetaStructs: spanMetaStructs, spanEvents: spanEvents, peerTags: peerTags!, obfuscationVersion: obfuscationVersion, traceFilterConfig: traceFilterConfig, - featureFlags: featureFlags!); - - // Save the hash, whether the details we care about changed or not - _configurationHash = HexString.ToHexString(sha256.Hash); - - // AgentConfiguration is a record, so this compares by value - if (existingConfiguration is null || !newConfig.Equals(existingConfiguration)) + featureFlags: featureFlags!, + eventPlatformProxySupportsEvpOriginHeaders: eventPlatformProxySupportsEvpOriginHeaders) { - Log.Debug("Discovery configuration updated, notifying subscribers: {Configuration}", newConfig); - NotifySubscribers(newConfig); - } + DiscoverySettings = requestFactory.ExporterSettings, + }; + + NotifySubscribers(newConfig, requestFactory, HexString.ToHexString(sha256.Hash), containerTagsHash); } public Task DisposeAsync() @@ -497,8 +526,10 @@ private static IApiRequestFactory CreateApiRequestFactory(ExporterSettings expor httpHeaderHelper: containerId is null ? MinimalAgentHeaderHelper.Instance : new MinimalWithContainerIdAgentHeaderHelper(containerId)); } - private sealed class ApiFactoryHolder(IApiRequestFactory apiFactory) + private sealed class ApiFactoryHolder(IApiRequestFactory apiFactory, ExporterSettings? exporterSettings) { + public ExporterSettings? ExporterSettings { get; } = exporterSettings; + public IApiRequestFactory ApiFactory { get; } = apiFactory; public Uri Uri { get; } = apiFactory.GetEndpoint("info"); diff --git a/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceMock.cs b/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceMock.cs index 696ebf84bc15..b19b7e2dc8b2 100644 --- a/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceMock.cs +++ b/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceMock.cs @@ -29,7 +29,8 @@ public void TriggerChange( string containerTagsHash = "containerTagsHash", bool clientDropP0 = true, bool spanMetaStructs = true, - bool spanEvents = true) + bool spanEvents = true, + bool eventPlatformProxySupportsEvpOriginHeaders = true) => TriggerChange( new AgentConfiguration( configurationEndpoint: configurationEndpoint, @@ -46,7 +47,8 @@ public void TriggerChange( containerTagsHash: containerTagsHash, clientDropP0: clientDropP0, spanMetaStructs: spanMetaStructs, - spanEvents: spanEvents)); + spanEvents: spanEvents, + eventPlatformProxySupportsEvpOriginHeaders: eventPlatformProxySupportsEvpOriginHeaders)); public void TriggerChange(AgentConfiguration config) { diff --git a/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceTests.cs b/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceTests.cs index 139104a114ac..c69aee0a3aef 100644 --- a/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceTests.cs +++ b/tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceTests.cs @@ -10,6 +10,7 @@ using System.Threading.Tasks; using Datadog.Trace.Agent; using Datadog.Trace.Agent.DiscoveryService; +using Datadog.Trace.Configuration; using Datadog.Trace.PlatformHelpers; using Datadog.Trace.TestHelpers; using Datadog.Trace.TestHelpers.TransportHelpers; @@ -29,6 +30,73 @@ public class DiscoveryServiceTests private static readonly ServiceRemappingHash DisabledServiceRemappingHash = new(null); + [Fact] + public async Task EndpointReplacementClearsCachedCapabilitiesAndPublishesEqualConfiguration() + { + var settingsA = new ExporterSettings(); + var settingsB = new ExporterSettings(); + var notifications = new List(); + var factoryA = new TestRequestFactory(uri => new TestApiRequest(uri, responseContent: GetConfig())); + var factoryB = new TestRequestFactory(uri => new TestApiRequest(uri, responseContent: GetConfig())); + await using var discovery = new DiscoveryService(factoryA, DisabledServiceRemappingHash, 1, 1, RecheckIntervalMs, autoStartLoop: false, exporterSettings: settingsA); + discovery.SubscribeToChanges(notifications.Add); + await discovery.RunOneIterationAsync(null); + discovery.SetCurrentConfigStateHash(discovery.ConfigStateHash); + discovery.RequireRefresh(discovery.ConfigStateHash, DateTimeOffset.UtcNow).Should().BeFalse(); + + discovery.UpdateRequestFactory(factoryB, settingsB); + + discovery.ConfigStateHash.Should().BeNull(); + var newSubscriber = new List(); + discovery.SubscribeToChanges(newSubscriber.Add); + newSubscriber.Should().BeEmpty("cached capabilities belonged to the previous endpoint"); + await discovery.RunOneIterationAsync(null); + + factoryB.RequestsSent.Should().ContainSingle(); + notifications.Should().HaveCount(2, "even an identical /info body belongs to a new endpoint"); + notifications[0].DiscoverySettings.Should().BeSameAs(settingsA); + notifications[1].DiscoverySettings.Should().BeSameAs(settingsB); + newSubscriber.Should().ContainSingle().Which.DiscoverySettings.Should().BeSameAs(settingsB); + notifications[1].ToString().Should().NotContain(nameof(AgentConfiguration.DiscoverySettings)); + } + + [Fact] + public async Task LateResponseFromPreviousEndpointCannotOverwriteNewConfigurationOrHashes() + { + var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var notifications = new List(); + var serviceHash = new ServiceRemappingHash("service:test"); + var factoryA = new TestRequestFactory(uri => new DelayedDiscoveryRequest(uri, started, release, GetConfig(version: "old-agent"))); + var factoryB = new TestRequestFactory(uri => new TestApiRequest(uri, responseContent: GetConfig(version: "new-agent"), responseHeaders: new() { { AgentHttpHeaderNames.ContainerTagsHash, "new-tags" } })); + await using var discovery = new DiscoveryService(factoryA, serviceHash, 1, 1, RecheckIntervalMs, autoStartLoop: false, exporterSettings: new ExporterSettings()); + discovery.SubscribeToChanges(notifications.Add); + var oldResponse = discovery.RunOneIterationAsync(null); + await started.Task; + try + { + var settingsB = new ExporterSettings(); + discovery.UpdateRequestFactory(factoryB, settingsB); + await discovery.RunOneIterationAsync(null); + var currentHash = discovery.ConfigStateHash; + + release.SetResult(true); + await oldResponse; + + notifications.Should().ContainSingle().Which.AgentVersion.Should().Be("new-agent"); + discovery.ConfigStateHash.Should().Be(currentHash); + serviceHash.ContainerTagsHash.Should().Be("new-tags"); + AgentConfiguration cached = null; + discovery.SubscribeToChanges(config => cached = config); + cached.DiscoverySettings.Should().BeSameAs(settingsB); + } + finally + { + release.TrySetResult(true); + await oldResponse; + } + } + [Fact] public async Task HandlesFlakyConfiguration() { @@ -76,9 +144,32 @@ public async Task ReturnsDeserializedConfig() config.StatsEndpoint.Should().NotBeNullOrEmpty(); config.DataStreamsMonitoringEndpoint.Should().NotBeNullOrEmpty(); config.EventPlatformProxyEndpoint.Should().Be(evpProxyEndpoint); + config.EventPlatformProxySupportsEvpOriginHeaders.Should().BeFalse(); await ds.DisposeAsync(); } + [Theory] + [InlineData("null", false)] + [InlineData("[]", false)] + [InlineData("[\"DD-EVP-ORIGIN\"]", false)] + [InlineData("[\"DD-EVP-ORIGIN-VERSION\"]", false)] + [InlineData("[\" dd-evp-origin-version \",\"dd-evp-origin\"]", true)] + public async Task ReportsWhetherEvpProxyCanForwardLogicalProducerIdentity(string allowedHeaders, bool expected) + { + AgentConfiguration config = null; + var response = $"{{\"endpoints\":[\"/evp_proxy/v4/\"],\"evp_proxy_allowed_headers\":{allowedHeaders}}}"; + var factory = new TestRequestFactory(x => new TestApiRequest(x, responseContent: response)); + + await using var ds = new DiscoveryService(factory, DisabledServiceRemappingHash, InitialRetryDelayMs, MaxRetryDelayMs, RecheckIntervalMs, autoStartLoop: false); + ds.SubscribeToChanges(x => config = x); + + await ds.RunOneIterationAsync(previousRetryDuration: null); + + config.Should().NotBeNull(); + config.EventPlatformProxyEndpoint.Should().Be("evp_proxy/v4"); + config.EventPlatformProxySupportsEvpOriginHeaders.Should().Be(expected); + } + [Fact] public async Task CalculatesConfigStateHash() { @@ -389,6 +480,16 @@ public async Task ExtractsContainerTagsHashFromResponseHeader() private string GetConfig(bool dropP0 = true, string version = null) => JsonConvert.SerializeObject(new MockTracerAgent.AgentConfiguration() { ClientDropP0s = dropP0, AgentVersion = version }); + internal sealed class DelayedDiscoveryRequest(Uri endpoint, TaskCompletionSource started, TaskCompletionSource release, string body) : TestApiRequest(endpoint) + { + public override async Task GetAsync() + { + started.TrySetResult(true); + await release.Task; + return new TestApiResponse(200, body, "application/json", headers: new() { { AgentHttpHeaderNames.ContainerTagsHash, "old-tags" } }); + } + } + internal class ThrowingRequest : TestApiRequest { public ThrowingRequest()