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
Expand Up @@ -6,6 +6,7 @@
#nullable enable

using System.Collections.Generic;
using Datadog.Trace.Configuration;

namespace Datadog.Trace.Agent.DiscoveryService;

Expand All @@ -30,7 +31,8 @@ public AgentConfiguration(
List<string>? peerTags = null,
int obfuscationVersion = 0,
AgentTraceFilterConfig? traceFilterConfig = null,
List<string>? featureFlags = null)
List<string>? featureFlags = null,
bool eventPlatformProxySupportsEvpOriginHeaders = false)
{
ConfigurationEndpoint = configurationEndpoint;
DebuggerEndpoint = debuggerEndpoint;
Expand All @@ -51,6 +53,7 @@ public AgentConfiguration(
ObfuscationVersion = obfuscationVersion;
TraceFilterConfig = traceFilterConfig ?? AgentTraceFilterConfig.Empty;
FeatureFlags = featureFlags;
EventPlatformProxySupportsEvpOriginHeaders = eventPlatformProxySupportsEvpOriginHeaders;
}

public string? ConfigurationEndpoint { get; }
Expand Down Expand Up @@ -83,6 +86,8 @@ public AgentConfiguration(

public string? EventPlatformProxyEndpoint { get; }

public bool EventPlatformProxySupportsEvpOriginHeaders { get; }

public string? TelemetryProxyEndpoint { get; }

public string? TracerFlareEndpoint { get; }
Expand All @@ -102,4 +107,8 @@ public AgentConfiguration(
public AgentTraceFilterConfig TraceFilterConfig { get; }

public List<string>? 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; }
}
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -62,15 +64,15 @@ 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 =>
{
if (changes.UpdatedExporter is { } exporter)
{
var newFactory = CreateApiRequestFactory(exporter, containerMetadata.ContainerId, tcpTimeout);
Interlocked.Exchange(ref _apiRequestFactory!, new(newFactory));
UpdateRequestFactory(newFactory, exporter);
}
});
}
Expand All @@ -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;
Expand Down Expand Up @@ -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;
Comment on lines +179 to +183

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Revoke capabilities held by existing subscribers

When the Agent endpoint changes after discovery has completed, clearing _configuration only prevents future subscribers from receiving the old snapshot; already-registered subscribers are not notified and continue using capabilities learned from the previous Agent. For example, RemoteConfigurationManager.SetRcmEnabled retains its prior value indefinitely if discovery against the replacement Agent fails, so clearing the cache does not actually enforce the endpoint boundary for normal production subscribers. Publish an invalidated snapshot or otherwise revoke the old capabilities when replacing the factory.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Checked against base 0610b13b and head b16a7dcc. Existing subscribers retaining capability state is pre-existing; this PR fences cached replay and publication from obsolete endpoint generations. The RC example is narrower: exporter replacement creates a new RemoteConfigurationApi with no endpoint, and the probe confirmed zero requests before fresh discovery. A real SpanEventsManager does retain its old capability, so universal subscriber invalidation remains a separate policy/lifecycle follow-up. Leaving this open for that discussion rather than claiming it is fixed.

Comment thread
leoromanovsky marked this conversation as resolved.
}
}

/// <inheritdoc cref="IDiscoveryService.SubscribeToChanges"/>
public void SubscribeToChanges(Action<AgentConfiguration> callback)
Expand Down Expand Up @@ -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<Action<AgentConfiguration>> 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
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -373,6 +404,10 @@ private async Task ProcessDiscoveryResponse(IApiResponse response)
}

var discoveredEndpoints = (jObject["endpoints"] as JArray)?.Values<string>().ToArray();
var evpProxyAllowedHeaders = (jObject["evp_proxy_allowed_headers"] as JArray)?.Values<string>().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;
Expand Down Expand Up @@ -442,8 +477,6 @@ private async Task ProcessDiscoveryResponse(IApiResponse response)
}
}

var existingConfiguration = _configuration;

var newConfig = new AgentConfiguration(
configurationEndpoint: configurationEndpoint,
debuggerEndpoint: debuggerEndpoint,
Expand All @@ -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()
Expand All @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -46,7 +47,8 @@ public void TriggerChange(
containerTagsHash: containerTagsHash,
clientDropP0: clientDropP0,
spanMetaStructs: spanMetaStructs,
spanEvents: spanEvents));
spanEvents: spanEvents,
eventPlatformProxySupportsEvpOriginHeaders: eventPlatformProxySupportsEvpOriginHeaders));

public void TriggerChange(AgentConfiguration config)
{
Expand Down
101 changes: 101 additions & 0 deletions tracer/test/Datadog.Trace.Tests/Agent/DiscoveryServiceTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<AgentConfiguration>();
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<AgentConfiguration>();
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<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
var release = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
var notifications = new List<AgentConfiguration>();
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()
{
Expand Down Expand Up @@ -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()
{
Expand Down Expand Up @@ -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<bool> started, TaskCompletionSource<bool> release, string body) : TestApiRequest(endpoint)
{
public override async Task<IApiResponse> 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()
Expand Down
Loading