diff --git a/communication/src/main/java/datadog/communication/BackendApiFactory.java b/communication/src/main/java/datadog/communication/BackendApiFactory.java index 2ac0447fc7d..b211a4da6cb 100644 --- a/communication/src/main/java/datadog/communication/BackendApiFactory.java +++ b/communication/src/main/java/datadog/communication/BackendApiFactory.java @@ -8,6 +8,7 @@ import datadog.trace.util.throwable.FatalAgentMisconfigurationError; import javax.annotation.Nullable; import okhttp3.HttpUrl; +import okhttp3.OkHttpClient; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -48,7 +49,13 @@ public BackendApi createDirectIntakeApi(Intake intake) { /** Creates an authenticated API client that sends data directly to a Datadog intake. */ public BackendApi createDirectIntakeApi(Intake intake, boolean responseCompression) { - HttpUrl agentlessUrl = HttpUrl.get(intake.getAgentlessUrl(config)); + return createDirectIntakeApi(intake, responseCompression, true); + } + + /** Creates an authenticated API client that sends data directly to a Datadog intake. */ + public BackendApi createDirectIntakeApi( + Intake intake, boolean responseCompression, boolean followRedirects) { + HttpUrl agentlessUrl = buildDirectIntakeUrl(intake, config); String apiKey = config.getApiKey(); if (apiKey == null || apiKey.isEmpty()) { throw new FatalAgentMisconfigurationError( @@ -60,10 +67,52 @@ public BackendApi createDirectIntakeApi(Intake intake, boolean responseCompressi apiKey, traceId, retryPolicyFactory(), - sharedCommunicationObjects.getIntakeHttpClient(), + directIntakeHttpClient(sharedCommunicationObjects.getIntakeHttpClient(), followRedirects), responseCompression); } + static OkHttpClient directIntakeHttpClient( + final OkHttpClient intakeHttpClient, final boolean followRedirects) { + if (followRedirects) { + return intakeHttpClient; + } + return intakeHttpClient.newBuilder().followRedirects(false).followSslRedirects(false).build(); + } + + private static HttpUrl buildDirectIntakeUrl(Intake intake, Config config) { + if (intake != Intake.EVENT_PLATFORM) { + return HttpUrl.get(intake.getAgentlessUrl(config)); + } + return buildEventPlatformIntakeUrl(config.getSite()); + } + + static HttpUrl buildEventPlatformIntakeUrl(String site) { + if (site == null || site.isEmpty()) { + throw new IllegalArgumentException("Invalid Datadog site"); + } + + String expectedHost = Intake.EVENT_PLATFORM.getUrlPrefix() + "." + site; + HttpUrl url = + new HttpUrl.Builder() + .scheme("https") + .host(expectedHost) + .addPathSegment("api") + .addPathSegment(Intake.EVENT_PLATFORM.getVersion()) + .addPathSegment("") + .build(); + if (!url.isHttps() + || !url.username().isEmpty() + || !url.password().isEmpty() + || !url.host().equalsIgnoreCase(expectedHost) + || url.port() != 443 + || !url.encodedPath().equals("/api/" + Intake.EVENT_PLATFORM.getVersion() + "/") + || url.encodedQuery() != null + || url.encodedFragment() != null) { + throw new IllegalArgumentException("Invalid Datadog site"); + } + return url; + } + /** Creates an API client that uses the specified retry policy with a compatible local proxy. */ public @Nullable BackendApi createEvpProxyApi(Intake intake) { return createEvpProxyApi(intake, true); diff --git a/communication/src/main/java/datadog/communication/ddagent/SharedCommunicationObjects.java b/communication/src/main/java/datadog/communication/ddagent/SharedCommunicationObjects.java index 73e9fbc037c..4b17eea2ba3 100644 --- a/communication/src/main/java/datadog/communication/ddagent/SharedCommunicationObjects.java +++ b/communication/src/main/java/datadog/communication/ddagent/SharedCommunicationObjects.java @@ -11,13 +11,25 @@ import datadog.remoteconfig.DefaultConfigurationPoller; import datadog.trace.api.Config; import datadog.trace.api.civisibility.config.BazelMode; +import datadog.trace.util.AgentProxySelector; import datadog.trace.util.AgentTaskScheduler; import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.net.Proxy; +import java.net.ProxySelector; +import java.net.SocketAddress; +import java.net.URI; import java.security.Security; import java.util.ArrayList; +import java.util.Collections; import java.util.List; +import java.util.Locale; +import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.function.Supplier; +import javax.annotation.Nullable; +import okhttp3.Credentials; import okhttp3.HttpUrl; import okhttp3.OkHttpClient; import org.slf4j.Logger; @@ -44,6 +56,9 @@ public class SharedCommunicationObjects { */ private volatile OkHttpClient intakeHttpClient; + private volatile HttpUrl intakeHttpsProxy; + private volatile Set intakeNoProxyHosts = Collections.emptySet(); + @SuppressFBWarnings("PA_PUBLIC_PRIMITIVE_ATTRIBUTE") public long httpClientTimeout; @@ -78,6 +93,8 @@ public void createRemaining(Config config) { : TimeUnit.SECONDS.toMillis(config.getAgentTimeout()); forceClearTextHttpForIntakeClient = config.isForceClearTextHttpForIntakeClient(); + intakeHttpsProxy = parseHttpsProxy(config.getHttpsProxy()); + intakeNoProxyHosts = config.getNoProxyHosts(); if (agentUrl == null) { agentUrl = parseAgentUrl(config); @@ -269,11 +286,94 @@ public OkHttpClient getIntakeHttpClient() { synchronized (this) { if (this.intakeHttpClient == null) { - this.intakeHttpClient = + OkHttpClient intakeClient = OkHttpUtils.buildHttpClient( forceClearTextHttpForIntakeClient, null, null, httpClientTimeout); + if (intakeHttpsProxy != null) { + final Proxy proxy = + new Proxy( + Proxy.Type.HTTP, + new InetSocketAddress(intakeHttpsProxy.host(), intakeHttpsProxy.port())); + final OkHttpClient.Builder builder = + intakeClient + .newBuilder() + .proxySelector(new IntakeProxySelector(proxy, intakeNoProxyHosts)); + if (!intakeHttpsProxy.username().isEmpty()) { + final String credential = + Credentials.basic(intakeHttpsProxy.username(), intakeHttpsProxy.password()); + builder.proxyAuthenticator( + (route, response) -> + response + .request() + .newBuilder() + .header("Proxy-Authorization", credential) + .build()); + } + intakeClient = builder.build(); + } + this.intakeHttpClient = intakeClient; } return this.intakeHttpClient; } } + + @Nullable + static HttpUrl parseHttpsProxy(@Nullable final String configuredProxy) { + if (configuredProxy == null || configuredProxy.trim().isEmpty()) { + return null; + } + final String candidate = + configuredProxy.contains("://") ? configuredProxy : "http://" + configuredProxy; + final HttpUrl proxy = HttpUrl.parse(candidate); + if (proxy == null || !"http".equalsIgnoreCase(proxy.scheme())) { + log.warn("Ignoring invalid HTTPS proxy configuration"); + return null; + } + return proxy; + } + + static final class IntakeProxySelector extends ProxySelector { + private static final List DIRECT = Collections.singletonList(Proxy.NO_PROXY); + + private final Proxy proxy; + private final Set noProxyHosts; + + IntakeProxySelector(final Proxy proxy, final Set noProxyHosts) { + this.proxy = proxy; + this.noProxyHosts = noProxyHosts; + } + + @Override + public List select(final URI uri) { + final String host = uri.getHost(); + if (host != null && shouldBypassProxy(host)) { + return DIRECT; + } + if ("https".equalsIgnoreCase(uri.getScheme())) { + return Collections.singletonList(proxy); + } + return AgentProxySelector.INSTANCE.select(uri); + } + + @Override + public void connectFailed( + final URI uri, final SocketAddress address, final IOException failure) { + AgentProxySelector.INSTANCE.connectFailed(uri, address, failure); + } + + private boolean shouldBypassProxy(final String host) { + final String normalizedHost = host.toLowerCase(Locale.ROOT); + for (final String configuredHost : noProxyHosts) { + final String normalized = configuredHost.trim().toLowerCase(Locale.ROOT); + if ("*".equals(normalized) + || normalizedHost.equals(normalized) + || (normalized.startsWith(".") + && (normalizedHost.equals(normalized.substring(1)) + || normalizedHost.endsWith(normalized)))) { + return true; + } + } + return false; + } + } } diff --git a/communication/src/test/groovy/datadog/communication/ddagent/SharedCommunicationsObjectsSpecification.groovy b/communication/src/test/groovy/datadog/communication/ddagent/SharedCommunicationsObjectsSpecification.groovy index 6ee88e87115..fe9e3c7cef5 100644 --- a/communication/src/test/groovy/datadog/communication/ddagent/SharedCommunicationsObjectsSpecification.groovy +++ b/communication/src/test/groovy/datadog/communication/ddagent/SharedCommunicationsObjectsSpecification.groovy @@ -90,6 +90,8 @@ class SharedCommunicationsObjectsSpecification extends DDSpecification { 1 * config.isCiVisibilityEnabled() 1 * config.getAgentTimeout() 1 * config.isForceClearTextHttpForIntakeClient() + 1 * config.getHttpsProxy() + 1 * config.getNoProxyHosts() 0 * _ sco.agentUrl.is(url) sco.agentHttpClient.is(okHttpClient) @@ -136,4 +138,46 @@ class SharedCommunicationsObjectsSpecification extends DDSpecification { then: client != null } + + void 'configures standard HTTPS proxy for intake while preserving no-proxy hosts'() { + given: + Config config = Mock() + sco.agentUrl = HttpUrl.get("http://example.com") + sco.agentHttpClient = Mock(OkHttpClient) + sco.monitoring = Monitoring.DISABLED + sco.featuresDiscovery = Mock(DDAgentFeaturesDiscovery) + + when: + sco.createRemaining(config) + def selector = sco.getIntakeHttpClient().proxySelector() + + then: + 1 * config.isCiVisibilityEnabled() >> false + 1 * config.getAgentTimeout() >> 1 + 1 * config.isForceClearTextHttpForIntakeClient() >> false + 1 * config.getHttpsProxy() >> "http://proxy.example:8181" + 1 * config.getNoProxyHosts() >> (["direct.example", ".internal.example"] as Set) + 0 * _ + + and: + def selected = selector.select(new URI("https://event-platform-intake.datadoghq.com")) + selected.size() == 1 + selected[0].type() == Proxy.Type.HTTP + selected[0].address() == new InetSocketAddress("proxy.example", 8181) + selector.select(new URI("https://direct.example")) == [Proxy.NO_PROXY] + selector.select(new URI("https://service.internal.example")) == [Proxy.NO_PROXY] + } + + void 'parses supported HTTPS proxy forms without logging credentials'() { + expect: + SharedCommunicationObjects.parseHttpsProxy(configured)?.toString() == expected + + where: + configured | expected + null | null + "" | null + "proxy.example:8080" | "http://proxy.example:8080/" + "http://user:pass@proxy:3128" | "http://user:pass@proxy:3128/" + "https://unsupported.example" | null + } } diff --git a/communication/src/test/java/datadog/communication/BackendApiFactoryTest.java b/communication/src/test/java/datadog/communication/BackendApiFactoryTest.java index 726c34a7f73..4f6d8d4ab18 100644 --- a/communication/src/test/java/datadog/communication/BackendApiFactoryTest.java +++ b/communication/src/test/java/datadog/communication/BackendApiFactoryTest.java @@ -13,7 +13,9 @@ import datadog.trace.api.Config; import datadog.trace.api.ProtocolVersion; import datadog.trace.api.intake.Intake; +import java.io.IOException; import java.nio.charset.StandardCharsets; +import java.util.Locale; import okhttp3.HttpUrl; import okhttp3.MediaType; import okhttp3.OkHttpClient; @@ -22,11 +24,98 @@ import okhttp3.mockwebserver.MockWebServer; import okhttp3.mockwebserver.RecordedRequest; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.NullAndEmptySource; +import org.junit.jupiter.params.provider.ValueSource; class BackendApiFactoryTest { private static final MediaType JSON = MediaType.parse("application/json"); + @ParameterizedTest + @ValueSource(strings = {"datadoghq.com", "custom.example", "DATADOGHQ.EU"}) + void eventPlatformDirectIntakeUsesExactHttpsHost(String site) { + final HttpUrl url = BackendApiFactory.buildEventPlatformIntakeUrl(site); + + assertEquals("https", url.scheme()); + assertEquals("event-platform-intake." + site.toLowerCase(Locale.ROOT), url.host()); + assertEquals(443, url.port()); + assertEquals("/api/v2/", url.encodedPath()); + assertEquals("", url.username()); + assertEquals("", url.password()); + assertNull(url.encodedQuery()); + assertNull(url.encodedFragment()); + } + + @ParameterizedTest + @NullAndEmptySource + @ValueSource( + strings = { + "datadoghq.com@evil.example", + "datadoghq.com:password@evil.example", + "https://datadoghq.com", + "datadoghq.com:443", + "datadoghq.com:8443", + "datadoghq.com/path", + "datadoghq.com?query=value", + "datadoghq.com#fragment", + "data doghq.com", + " datadoghq.com", + "datadoghq.com ", + "datadoghq.com\\evil.example" + }) + void eventPlatformDirectIntakeRejectsUnsafeSite(String site) { + assertThrows( + IllegalArgumentException.class, () -> BackendApiFactory.buildEventPlatformIntakeUrl(site)); + } + + @ParameterizedTest + @ValueSource(ints = {301, 302, 307, 308}) + void featureFlagDirectIntakeDoesNotFollowRedirects(final int statusCode) throws Exception { + final MockWebServer intake = new MockWebServer(); + final MockWebServer redirectTarget = new MockWebServer(); + final OkHttpClient sharedClient = new OkHttpClient.Builder().build(); + final OkHttpClient directClient = BackendApiFactory.directIntakeHttpClient(sharedClient, false); + redirectTarget.start(); + intake.enqueue( + new MockResponse() + .setResponseCode(statusCode) + .setHeader("Location", redirectTarget.url("/redirected"))); + intake.start(); + try { + final IntakeApi api = + new IntakeApi( + intake.url("/api/v2/"), + "api-key", + "123", + HttpRetryPolicy.Factory.NEVER_RETRY, + directClient, + false); + + assertThrows( + IOException.class, + () -> + api.post( + "flagevaluation", + RequestBody.create(JSON, "{}".getBytes(StandardCharsets.UTF_8)), + stream -> null, + null, + false)); + + final RecordedRequest request = intake.takeRequest(); + assertEquals("api-key", request.getHeader("DD-API-KEY")); + assertEquals(1, intake.getRequestCount()); + assertEquals(0, redirectTarget.getRequestCount()); + } finally { + directClient.dispatcher().executorService().shutdownNow(); + directClient.connectionPool().evictAll(); + sharedClient.dispatcher().executorService().shutdownNow(); + sharedClient.connectionPool().evictAll(); + intake.shutdown(); + redirectTarget.shutdown(); + } + } + @Test void noBackendApiWhenAgentDoesNotAdvertiseEvpProxy() { final FakeFeaturesDiscovery discovery = new FakeFeaturesDiscovery(null); diff --git a/dd-trace-api/src/main/java/datadog/trace/api/config/TracerConfig.java b/dd-trace-api/src/main/java/datadog/trace/api/config/TracerConfig.java index 0ab030189b5..ac68208f8de 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/config/TracerConfig.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/config/TracerConfig.java @@ -29,6 +29,7 @@ public final class TracerConfig { public static final String AGENT_TIMEOUT = "trace.agent.timeout"; public static final String FORCE_CLEAR_TEXT_HTTP_FOR_INTAKE_CLIENT = "force.clear.text.http.for.intake.client"; + public static final String PROXY_HTTPS = "proxy.https"; public static final String PROXY_NO_PROXY = "proxy.no_proxy"; public static final String TRACE_AGENT_PATH = "trace.agent.path"; public static final String TRACE_AGENT_ARGS = "trace.agent.args"; diff --git a/internal-api/src/main/java/datadog/trace/api/Config.java b/internal-api/src/main/java/datadog/trace/api/Config.java index fade2b4c417..a1c04c543a4 100644 --- a/internal-api/src/main/java/datadog/trace/api/Config.java +++ b/internal-api/src/main/java/datadog/trace/api/Config.java @@ -660,6 +660,7 @@ import static datadog.trace.api.config.TracerConfig.PROPAGATION_EXTRACT_LOG_HEADER_NAMES_ENABLED; import static datadog.trace.api.config.TracerConfig.PROPAGATION_STYLE_EXTRACT; import static datadog.trace.api.config.TracerConfig.PROPAGATION_STYLE_INJECT; +import static datadog.trace.api.config.TracerConfig.PROXY_HTTPS; import static datadog.trace.api.config.TracerConfig.PROXY_NO_PROXY; import static datadog.trace.api.config.TracerConfig.REQUEST_HEADER_TAGS; import static datadog.trace.api.config.TracerConfig.REQUEST_HEADER_TAGS_COMMA_ALLOWED; @@ -920,6 +921,7 @@ public static String getHostName() { /** Should be set to {@code true} when running in agentless mode in a JVM without TLS */ private final boolean forceClearTextHttpForIntakeClient; + private final String httpsProxy; private final Set noProxyHosts; private final boolean prioritySamplingEnabled; private final String prioritySamplingForce; @@ -1684,8 +1686,27 @@ private Config(final ConfigProvider configProvider, final InstrumenterConfig ins forceClearTextHttpForIntakeClient = configProvider.getBoolean(FORCE_CLEAR_TEXT_HTTP_FOR_INTAKE_CLIENT, false); - // DD_PROXY_NO_PROXY is specified as a space-separated list of hosts - noProxyHosts = tryMakeImmutableSet(configProvider.getSpacedList(PROXY_NO_PROXY)); + String configuredHttpsProxy = configProvider.getString(PROXY_HTTPS); + if (configuredHttpsProxy == null) { + configuredHttpsProxy = ConfigHelper.env("HTTPS_PROXY"); + if (configuredHttpsProxy == null) { + configuredHttpsProxy = ConfigHelper.env("https_proxy"); + } + } + httpsProxy = configuredHttpsProxy; + + // DD_PROXY_NO_PROXY historically accepted spaces; standard NO_PROXY aliases use commas. + String configuredNoProxyHosts = configProvider.getString(PROXY_NO_PROXY); + if (configuredNoProxyHosts == null) { + configuredNoProxyHosts = ConfigHelper.env("NO_PROXY"); + if (configuredNoProxyHosts == null) { + configuredNoProxyHosts = ConfigHelper.env("no_proxy"); + } + } + noProxyHosts = + configuredNoProxyHosts == null + ? Collections.emptySet() + : parseStringIntoSetOfNonEmptyStrings(configuredNoProxyHosts); prioritySamplingEnabled = configProvider.getBoolean(PRIORITY_SAMPLING, DEFAULT_PRIORITY_SAMPLING_ENABLED); @@ -3626,6 +3647,10 @@ public boolean isForceClearTextHttpForIntakeClient() { return forceClearTextHttpForIntakeClient; } + public String getHttpsProxy() { + return httpsProxy; + } + public Set getNoProxyHosts() { return noProxyHosts; } diff --git a/internal-api/src/test/groovy/datadog/trace/api/ConfigTest.groovy b/internal-api/src/test/groovy/datadog/trace/api/ConfigTest.groovy index ef95d5e902c..a82133e32c0 100644 --- a/internal-api/src/test/groovy/datadog/trace/api/ConfigTest.groovy +++ b/internal-api/src/test/groovy/datadog/trace/api/ConfigTest.groovy @@ -1389,6 +1389,34 @@ class ConfigTest extends DDSpecification { config.serviceName == "what actually wants" } + def "https proxy configuration accepts Datadog and standard environment names"() { + setup: + environmentVariables.set(envName, "http://proxy.example:8080") + + when: + def config = new Config() + + then: + config.httpsProxy == "http://proxy.example:8080" + + where: + envName << ["DD_PROXY_HTTPS", "HTTPS_PROXY", "https_proxy"] + } + + def "no-proxy configuration accepts Datadog and standard environment names"() { + setup: + environmentVariables.set(envName, "event-platform-intake.example, .internal.example local.example") + + when: + def config = new Config() + + then: + config.noProxyHosts == ["event-platform-intake.example", ".internal.example", "local.example"] as Set + + where: + envName << ["DD_PROXY_NO_PROXY", "NO_PROXY", "no_proxy"] + } + def "verify mapping configs on tracer for #mapString"() { setup: System.setProperty(PREFIX + HEADER_TAGS + ".legacy.parsing.enabled", "true") diff --git a/metadata/supported-configurations.json b/metadata/supported-configurations.json index a1459dc3453..b21631ddfcc 100644 --- a/metadata/supported-configurations.json +++ b/metadata/supported-configurations.json @@ -3577,6 +3577,14 @@ "aliases": [] } ], + "DD_PROXY_HTTPS": [ + { + "version": "A", + "type": "string", + "default": null, + "aliases": [] + } + ], "DD_PROXY_NO_PROXY": [ { "version": "A", diff --git a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApi.java b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApi.java index 769ebfd1dd1..533771ab15b 100644 --- a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApi.java +++ b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApi.java @@ -48,15 +48,21 @@ public T post( return selectedApi.post( uri, requestBody, responseParser, requestListener, requestCompression); } catch (final IOException exception) { - if (selectedApi != proxyApi || !isDefinitiveRejection(exception)) { + if (selectedApi != proxyApi) { throw exception; } - final BackendApi directApi = getOrCreateDirectApi(); - if (directApi == null) { - throw exception; + if (isDefinitiveRejection(exception)) { + final BackendApi directApi = getOrCreateDirectApi(); + if (directApi != null) { + return directApi.post( + uri, requestBody, responseParser, requestListener, requestCompression); + } + } else if (isAmbiguousTransportFailure(exception)) { + // The local receiver may have accepted this batch, so only switch future batches. + getOrCreateDirectApi(); } - return directApi.post(uri, requestBody, responseParser, requestListener, requestCompression); + throw exception; } } @@ -98,4 +104,8 @@ private static boolean isDefinitiveRejection(final IOException exception) { } return false; } + + private static boolean isAmbiguousTransportFailure(final IOException exception) { + return !(exception instanceof HttpResponseException); + } } diff --git a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FeatureFlagBackendApiFactory.java b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FeatureFlagBackendApiFactory.java index 0dbf9c74254..451b903ecc7 100644 --- a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FeatureFlagBackendApiFactory.java +++ b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FeatureFlagBackendApiFactory.java @@ -90,7 +90,7 @@ private BackendApi createDirectApi() { } try { return backendApiFactory.createDirectIntakeApi( - Intake.EVENT_PLATFORM, eventType.responseCompressionEnabled()); + Intake.EVENT_PLATFORM, eventType.responseCompressionEnabled(), false); } catch (final IllegalArgumentException exception) { LOGGER.debug( "Cannot configure direct Feature Flagging {} delivery", eventType.logName(), exception); diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApiTest.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApiTest.java index 117c72ba2a1..42ac2ab5129 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApiTest.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/AgentlessFeatureFlagBackendApiTest.java @@ -94,18 +94,23 @@ void doesNotReturnToLocalRouteAfterSwitchingToDirectIntake() throws Exception { @ParameterizedTest @ValueSource(ints = {429, 500}) - void doesNotReplayAmbiguousHttpFailure(final int statusCode) { - assertNoDirectReplay(new HttpResponseException(statusCode, "ambiguous")); + void staysOnLocalRouteAfterNonFallbackHttpFailure(final int statusCode) { + assertNoDirectSwitch(new HttpResponseException(statusCode, "non-fallback")); } @Test - void doesNotReplayTimeout() { - assertNoDirectReplay(new SocketTimeoutException("timed out")); + void switchesFutureBatchToDirectIntakeAfterTimeout() throws Exception { + assertFutureBatchUsesDirectIntake(new SocketTimeoutException("timed out")); } @Test - void doesNotReplayConnectionReset() { - assertNoDirectReplay(new SocketException("connection reset")); + void switchesFutureBatchToDirectIntakeAfterConnectionReset() throws Exception { + assertFutureBatchUsesDirectIntake(new SocketException("connection reset")); + } + + @Test + void switchesFutureBatchToDirectIntakeAfterBrokenPipe() throws Exception { + assertFutureBatchUsesDirectIntake(new SocketException("broken pipe")); } @Test @@ -133,7 +138,7 @@ void doesNotRetryDirectApiCreationWhenFallbackIsUnavailable() { assertEquals(1, directApiCreations.get()); } - private static void assertNoDirectReplay(final IOException failure) { + private static void assertNoDirectSwitch(final IOException failure) { final RecordingBackendApi local = new RecordingBackendApi(failure); final RecordingBackendApi direct = new RecordingBackendApi(); final AtomicInteger directApiCreations = new AtomicInteger(); @@ -149,12 +154,43 @@ private static void assertNoDirectReplay(final IOException failure) { assertThrows( IOException.class, () -> api.post("flagevaluation", requestBody("evaluation"), stream -> null, null, false)); + assertThrows( + IOException.class, + () -> api.post("flagevaluation", requestBody("next"), stream -> null, null, false)); - assertEquals(1, local.calls); + assertEquals(2, local.calls); assertEquals(0, direct.calls); assertEquals(0, directApiCreations.get()); } + private static void assertFutureBatchUsesDirectIntake(final IOException failure) + throws Exception { + final RecordingBackendApi local = new RecordingBackendApi(failure); + final RecordingBackendApi direct = new RecordingBackendApi(); + final AtomicInteger directApiCreations = new AtomicInteger(); + final AgentlessFeatureFlagBackendApi api = + new AgentlessFeatureFlagBackendApi( + local, + () -> { + directApiCreations.incrementAndGet(); + return direct; + }, + "flag evaluation"); + final RequestBody ambiguousBody = requestBody("ambiguous"); + final RequestBody nextBody = requestBody("next"); + + assertThrows( + IOException.class, + () -> api.post("flagevaluation", ambiguousBody, stream -> null, null, false)); + api.post("flagevaluation", nextBody, stream -> null, null, false); + + assertEquals(1, local.calls); + assertEquals(1, direct.calls); + assertEquals(1, directApiCreations.get()); + assertSame(ambiguousBody, local.requestBodies.get(0)); + assertSame(nextBody, direct.requestBodies.get(0)); + } + private static RequestBody requestBody(final String value) { return RequestBody.create(MediaType.parse("application/json"), value); } diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/ExposureWriterTests.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/ExposureWriterTests.java index daaf213ebfb..f70a4ff0fa1 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/ExposureWriterTests.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/ExposureWriterTests.java @@ -168,7 +168,7 @@ void testAgentlessExposureEventWritesDirectlyWithApiKey() throws Exception { new OkHttpClient.Builder().build(), false); when(backendApiFactory.createDirectIntakeApi( - datadog.trace.api.intake.Intake.EVENT_PLATFORM, true)) + eq(datadog.trace.api.intake.Intake.EVENT_PLATFORM), eq(true), eq(false))) .thenReturn(directApi); FeatureFlagBackendApiFactory exposureBackendApiFactory = new FeatureFlagBackendApiFactory(config, backendApiFactory, FeatureFlagEventType.EXPOSURE); @@ -311,7 +311,7 @@ void testAmbiguousExposureBatchIsNotRetriedOrReplayedDirectly() throws Exception when(backendApiFactory.createEvpProxyApi( Intake.EVENT_PLATFORM, true, HttpRetryPolicy.Factory.NEVER_RETRY)) .thenReturn(proxyApi); - when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, true)) + when(backendApiFactory.createDirectIntakeApi(eq(Intake.EVENT_PLATFORM), eq(true), eq(false))) .thenReturn(directApi); when(proxyApi.post(eq("exposures"), any(RequestBody.class), any(), any(), eq(false))) .thenThrow(new SocketTimeoutException("ambiguous timeout")) diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FeatureFlagBackendApiFactoryTest.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FeatureFlagBackendApiFactoryTest.java index 8b117742715..b89cf3ce2b1 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FeatureFlagBackendApiFactoryTest.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FeatureFlagBackendApiFactoryTest.java @@ -32,7 +32,7 @@ void remoteConfigUsesOnlyLocalEvpProxy() { new FeatureFlagBackendApiFactory(config, backendApiFactory, FLAG_EVALUATION).create(); assertSame(proxyApi, selected); - verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false); + verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false); } @Test @@ -45,7 +45,7 @@ void remoteConfigDisablesDeliveryWhenLocalEvpProxyIsUnavailable() { assertNull(selected); verify(backendApiFactory).createEvpProxyApi(Intake.EVENT_PLATFORM, true); - verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, true); + verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, true, false); } @Test @@ -55,7 +55,7 @@ void agentlessPrefersLocalEvpProxyWithDirectFallback() { when(backendApiFactory.createEvpProxyApi( Intake.EVENT_PLATFORM, false, HttpRetryPolicy.Factory.NEVER_RETRY)) .thenReturn(mock(BackendApi.class)); - when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false)) + when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false)) .thenReturn(mock(BackendApi.class)); final BackendApi selected = @@ -64,7 +64,7 @@ void agentlessPrefersLocalEvpProxyWithDirectFallback() { assertInstanceOf(AgentlessFeatureFlagBackendApi.class, selected); verify(backendApiFactory) .createEvpProxyApi(Intake.EVENT_PLATFORM, false, HttpRetryPolicy.Factory.NEVER_RETRY); - verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false); + verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false); } @Test @@ -72,7 +72,7 @@ void agentlessUsesDirectIntakeWhenLocalEvpProxyIsUnavailable() { final Config config = config(CONFIGURATION_SOURCE_AGENTLESS, "api-key"); final BackendApiFactory backendApiFactory = mock(BackendApiFactory.class); final BackendApi directApi = mock(BackendApi.class); - when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false)) + when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false)) .thenReturn(directApi); final BackendApi selected = @@ -92,7 +92,7 @@ void agentlessUsesLocalEvpProxyWhenApiKeyIsUnavailable() { new FeatureFlagBackendApiFactory(config, backendApiFactory, FLAG_EVALUATION).create(); assertSame(proxyApi, selected); - verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false); + verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false); } @Test @@ -115,7 +115,7 @@ void agentlessDisablesDeliveryWhenApiKeyIsEmpty() { new FeatureFlagBackendApiFactory(config, backendApiFactory, EXPOSURE).create(); assertNull(selected); - verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, true); + verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, true, false); } @Test @@ -126,21 +126,21 @@ void agentlessDoesNotValidateDirectUrlWhileLocalRouteIsAvailable() { when(backendApiFactory.createEvpProxyApi( Intake.EVENT_PLATFORM, false, HttpRetryPolicy.Factory.NEVER_RETRY)) .thenReturn(proxyApi); - when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false)) + when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false)) .thenThrow(new IllegalArgumentException("invalid URL")); final BackendApi selected = new FeatureFlagBackendApiFactory(config, backendApiFactory, FLAG_EVALUATION).create(); assertInstanceOf(AgentlessFeatureFlagBackendApi.class, selected); - verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false); + verify(backendApiFactory, never()).createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false); } @Test void agentlessDisablesDeliveryWhenDirectUrlIsInvalidAndLocalRouteIsUnavailable() { final Config config = config(CONFIGURATION_SOURCE_AGENTLESS, "api-key"); final BackendApiFactory backendApiFactory = mock(BackendApiFactory.class); - when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false)) + when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false, false)) .thenThrow(new IllegalArgumentException("invalid URL")); final BackendApi selected = diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java index ba3d03ec6d5..75a767c5c3f 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java @@ -697,7 +697,8 @@ void agentlessWritesFlagEvaluationsDirectlyWhenLocalProxyIsUnavailable() throws HttpRetryPolicy.Factory.NEVER_RETRY, client, false); - when(backendApiFactory.createDirectIntakeApi(Intake.EVENT_PLATFORM, false)) + when(backendApiFactory.createDirectIntakeApi( + eq(Intake.EVENT_PLATFORM), eq(false), eq(false))) .thenReturn(directApi); final FeatureFlagBackendApiFactory featureFlagBackendApiFactory = new FeatureFlagBackendApiFactory(