Skip to content

Commit 9b8d758

Browse files
committed
fix(openfeature): prevent duplicate feature flag events
1 parent 041bc10 commit 9b8d758

11 files changed

Lines changed: 205 additions & 149 deletions

File tree

‎communication/src/main/java/datadog/communication/BackendApiFactory.java‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -64,13 +64,19 @@ public BackendApi createDirectIntakeApi(Intake intake, boolean responseCompressi
6464
responseCompression);
6565
}
6666

67-
/** Creates an API client that sends data through a compatible local EVP proxy. */
67+
/** Creates an API client that uses the specified retry policy with a compatible local proxy. */
6868
public @Nullable BackendApi createEvpProxyApi(Intake intake) {
6969
return createEvpProxyApi(intake, true);
7070
}
7171

7272
/** Creates an API client that sends data through a compatible local EVP proxy. */
7373
public @Nullable BackendApi createEvpProxyApi(Intake intake, boolean responseCompression) {
74+
return createEvpProxyApi(intake, responseCompression, retryPolicyFactory());
75+
}
76+
77+
/** Creates an API client that sends data through a compatible local EVP proxy. */
78+
public @Nullable BackendApi createEvpProxyApi(
79+
Intake intake, boolean responseCompression, HttpRetryPolicy.Factory retryPolicyFactory) {
7480
DDAgentFeaturesDiscovery featuresDiscovery =
7581
sharedCommunicationObjects.featuresDiscovery(config);
7682
featuresDiscovery.discoverIfOutdated();
@@ -91,7 +97,7 @@ public BackendApi createDirectIntakeApi(Intake intake, boolean responseCompressi
9197
traceId,
9298
evpProxyUrl,
9399
subdomain,
94-
retryPolicyFactory(),
100+
retryPolicyFactory,
95101
sharedCommunicationObjects.agentHttpClient,
96102
responseCompression);
97103
}

‎communication/src/test/java/datadog/communication/BackendApiFactoryTest.java‎

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,11 @@
44
import static org.junit.jupiter.api.Assertions.assertEquals;
55
import static org.junit.jupiter.api.Assertions.assertNotNull;
66
import static org.junit.jupiter.api.Assertions.assertNull;
7+
import static org.junit.jupiter.api.Assertions.assertThrows;
78

89
import datadog.communication.ddagent.DDAgentFeaturesDiscovery;
910
import datadog.communication.ddagent.SharedCommunicationObjects;
11+
import datadog.communication.http.HttpRetryPolicy;
1012
import datadog.metrics.api.Monitoring;
1113
import datadog.trace.api.Config;
1214
import datadog.trace.api.ProtocolVersion;
@@ -61,6 +63,38 @@ void advertisedEvpProxyEndpointSupportsDisabledResponseCompression() throws Exce
6163
}
6264
}
6365

66+
@Test
67+
void explicitNoRetryProxyPolicyDoesNotReplayAmbiguousFailure() throws Exception {
68+
final MockWebServer agent = new MockWebServer();
69+
agent.enqueue(new MockResponse().setResponseCode(500).setBody("ambiguous"));
70+
agent.enqueue(new MockResponse().setResponseCode(200).setBody("{}"));
71+
agent.start();
72+
try {
73+
final FakeFeaturesDiscovery discovery = new FakeFeaturesDiscovery(V4_EVP_PROXY_ENDPOINT);
74+
final BackendApiFactory factory =
75+
new BackendApiFactory(
76+
Config.get(), sharedCommunicationObjects(discovery, agent.url("/")));
77+
final BackendApi api =
78+
factory.createEvpProxyApi(
79+
Intake.EVENT_PLATFORM, false, HttpRetryPolicy.Factory.NEVER_RETRY);
80+
81+
assertNotNull(api);
82+
assertThrows(
83+
HttpResponseException.class,
84+
() ->
85+
api.post(
86+
"flagevaluation",
87+
RequestBody.create(JSON, "{}".getBytes(StandardCharsets.UTF_8)),
88+
stream -> null,
89+
null,
90+
false));
91+
92+
assertEquals(1, agent.getRequestCount());
93+
} finally {
94+
agent.shutdown();
95+
}
96+
}
97+
6498
private static SharedCommunicationObjects sharedCommunicationObjects(
6599
final DDAgentFeaturesDiscovery discovery, final HttpUrl agentUrl) {
66100
final TestSharedCommunicationObjects sco = new TestSharedCommunicationObjects(discovery);

‎products/feature-flagging/feature-flagging-lib/src/jmh/java/com/datadog/featureflag/FlagEvaluationEnqueueContentionBenchmark.java‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import datadog.trace.api.Config;
88
import datadog.trace.api.featureflag.FeatureFlaggingGateway;
99
import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent;
10+
import datadog.trace.api.intake.Intake;
1011
import de.thetaphi.forbiddenapis.SuppressForbidden;
1112
import java.util.HashMap;
1213
import java.util.Map;
@@ -90,7 +91,13 @@ public void setUp() {
9091
final Config config = Config.get();
9192
final BackendApiFactory factory = new BackendApiFactory(config, null);
9293
// Capacity well above what the batch-draining consumer should ever let build up.
93-
writer = new FlagEvaluationWriterImpl(1 << 20, Long.MAX_VALUE, NANOSECONDS, factory, config);
94+
writer =
95+
new FlagEvaluationWriterImpl(
96+
1 << 20,
97+
Long.MAX_VALUE,
98+
NANOSECONDS,
99+
() -> factory.createBackendApi(Intake.EVENT_PLATFORM, false),
100+
config);
94101
}
95102

96103
/**

‎products/feature-flagging/feature-flagging-lib/src/jmh/java/com/datadog/featureflag/FlagEvaluationHotPathBenchmark.java‎

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
import datadog.communication.BackendApiFactory;
77
import datadog.trace.api.Config;
88
import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent;
9+
import datadog.trace.api.intake.Intake;
910
import java.util.HashMap;
1011
import java.util.Map;
1112
import org.openjdk.jmh.annotations.Benchmark;
@@ -79,10 +80,18 @@ public void setUp() {
7980
final BackendApiFactory factory = new BackendApiFactory(config, null);
8081
final Map<String, String> ddContext = new HashMap<>();
8182
ddContext.put("service", "bench-service");
82-
handler = FlagEvaluationWriterImpl.createHandlerForTest(factory, ddContext);
83+
handler =
84+
FlagEvaluationWriterImpl.createHandlerForTest(
85+
() -> factory.createBackendApi(Intake.EVENT_PLATFORM, false), ddContext);
8386

8487
// Capacity large enough that the benchmark never overflows within a measurement window.
85-
writer = new FlagEvaluationWriterImpl(1 << 20, Long.MAX_VALUE, NANOSECONDS, factory, config);
88+
writer =
89+
new FlagEvaluationWriterImpl(
90+
1 << 20,
91+
Long.MAX_VALUE,
92+
NANOSECONDS,
93+
() -> factory.createBackendApi(Intake.EVENT_PLATFORM, false),
94+
config);
8695
}
8796

8897
/**

‎products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/ExposureWriterImpl.java‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -187,9 +187,12 @@ protected void flushIfNecessary() {
187187
}
188188
try {
189189
evpPublisher.post(EXPOSURES_ROUTE, payload);
190-
this.buffer.clear();
191190
} catch (Exception e) {
192191
LOGGER.debug("Could not submit exposures", e);
192+
} finally {
193+
// Best-effort delivery must not retry an ambiguously accepted batch. A later definitive
194+
// proxy rejection could otherwise replay the same exposures through direct intake.
195+
this.buffer.clear();
193196
}
194197
}
195198
}

‎products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FeatureFlagBackendApiFactory.java‎

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import datadog.communication.BackendApi;
66
import datadog.communication.BackendApiFactory;
77
import datadog.communication.ddagent.SharedCommunicationObjects;
8+
import datadog.communication.http.HttpRetryPolicy;
89
import datadog.trace.api.Config;
910
import datadog.trace.api.intake.Intake;
1011
import javax.annotation.Nullable;
@@ -38,9 +39,17 @@ final class FeatureFlagBackendApiFactory {
3839

3940
@Nullable
4041
BackendApi create() {
42+
final boolean directFallbackAvailable =
43+
CONFIGURATION_SOURCE_AGENTLESS.equals(config.getFeatureFlaggingConfigurationSource())
44+
&& hasDirectCredentials();
4145
final BackendApi proxyApi =
42-
backendApiFactory.createEvpProxyApi(
43-
Intake.EVENT_PLATFORM, eventType.responseCompressionEnabled());
46+
directFallbackAvailable
47+
? backendApiFactory.createEvpProxyApi(
48+
Intake.EVENT_PLATFORM,
49+
eventType.responseCompressionEnabled(),
50+
HttpRetryPolicy.Factory.NEVER_RETRY)
51+
: backendApiFactory.createEvpProxyApi(
52+
Intake.EVENT_PLATFORM, eventType.responseCompressionEnabled());
4453
if (!CONFIGURATION_SOURCE_AGENTLESS.equals(config.getFeatureFlaggingConfigurationSource())) {
4554
if (proxyApi == null) {
4655
LOGGER.warn(
@@ -51,7 +60,7 @@ BackendApi create() {
5160
}
5261

5362
if (proxyApi != null) {
54-
if (hasDirectCredentials()) {
63+
if (directFallbackAvailable) {
5564
return new AgentlessFeatureFlagBackendApi(
5665
proxyApi, this::createDirectApi, eventType.logName());
5766
}

‎products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java‎

Lines changed: 17 additions & 117 deletions
Original file line numberDiff line numberDiff line change
@@ -7,14 +7,12 @@
77
import datadog.common.queue.MessagePassingBlockingQueue;
88
import datadog.common.queue.Queues;
99
import datadog.communication.BackendApi;
10-
import datadog.communication.BackendApiFactory;
1110
import datadog.communication.EvpProxy;
1211
import datadog.communication.ddagent.SharedCommunicationObjects;
1312
import datadog.trace.api.Config;
1413
import datadog.trace.api.featureflag.FeatureFlaggingGateway;
1514
import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent;
1615
import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationWriter;
17-
import datadog.trace.api.intake.Intake;
1816
import datadog.trace.api.telemetry.CoreMetricCollector;
1917
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
2018
import java.util.ArrayList;
@@ -33,8 +31,8 @@
3331
* EVP flagevaluation writer for Java.
3432
*
3533
* <p>Uses the same EVP publisher path as ExposureWriterImpl, with two-tier aggregation replacing
36-
* the single-exposure buffer. Routes to the Agent-advertised EVP proxy endpoint for
37-
* /api/v2/flagevaluation.
34+
* the single-exposure buffer. Uses a local EVP proxy when available. Agentless mode can use direct
35+
* intake when no compatible local route is available.
3836
*
3937
* <p>Two-tier aggregation contract: Full key: (flagKey, variant, allocationKey, runtimeDefault,
4038
* errorMessage, targetingKey, canonical-context-key). Degraded key: (flagKey, variant,
@@ -102,51 +100,15 @@ private static void countMetric(final String metricName, final long value, final
102100
new ConcurrentHashMap<>();
103101

104102
public FlagEvaluationWriterImpl(final SharedCommunicationObjects sco, final Config config) {
105-
this(DEFAULT_CAPACITY, FLUSH_INTERVAL_SECONDS, SECONDS, sco, config);
106-
}
107-
108-
FlagEvaluationWriterImpl(
109-
final int capacity,
110-
final long flushInterval,
111-
final TimeUnit timeUnit,
112-
final SharedCommunicationObjects sco,
113-
final Config config) {
114-
this(
115-
capacity,
116-
flushInterval,
117-
timeUnit,
118-
new FeatureFlagBackendApiFactory(config, sco, FeatureFlagEventType.FLAG_EVALUATION),
119-
config);
120-
}
121-
122-
/** Package-private constructor allowing a BackendApiFactory to be injected for tests. */
123-
FlagEvaluationWriterImpl(
124-
final int capacity,
125-
final long flushInterval,
126-
final TimeUnit timeUnit,
127-
final BackendApiFactory backendApiFactory,
128-
final Config config) {
129103
this(
130-
capacity,
131-
flushInterval,
132-
timeUnit,
133-
() ->
134-
backendApiFactory.createBackendApi(
135-
Intake.EVENT_PLATFORM,
136-
FeatureFlagEventType.FLAG_EVALUATION.responseCompressionEnabled()),
104+
DEFAULT_CAPACITY,
105+
FLUSH_INTERVAL_SECONDS,
106+
SECONDS,
107+
new FeatureFlagBackendApiFactory(config, sco, FeatureFlagEventType.FLAG_EVALUATION)::create,
137108
config);
138109
}
139110

140111
FlagEvaluationWriterImpl(
141-
final int capacity,
142-
final long flushInterval,
143-
final TimeUnit timeUnit,
144-
final FeatureFlagBackendApiFactory backendApiFactory,
145-
final Config config) {
146-
this(capacity, flushInterval, timeUnit, backendApiFactory::create, config);
147-
}
148-
149-
private FlagEvaluationWriterImpl(
150112
final int capacity,
151113
final long flushInterval,
152114
final TimeUnit timeUnit,
@@ -162,7 +124,8 @@ private FlagEvaluationWriterImpl(
162124
FeatureFlagEvpContext.from(config),
163125
droppedQueueOverflow,
164126
contextTruncatedCounts,
165-
this::close);
127+
this::close,
128+
FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES);
166129
this.serializerThread = newAgentThread(FEATURE_FLAG_EVALUATION_PROCESSOR, serializer);
167130
}
168131

@@ -342,70 +305,6 @@ static class FlagEvaluationSerializingHandler implements Runnable {
342305
private final AtomicBoolean shutdownRequested = new AtomicBoolean(false);
343306
private final CountDownLatch finalFlushDone = new CountDownLatch(1);
344307

345-
FlagEvaluationSerializingHandler(
346-
final BackendApiFactory backendApiFactory,
347-
final MessagePassingBlockingQueue<FlagEvalEvent> queue,
348-
final long flushInterval,
349-
final TimeUnit timeUnit,
350-
final Map<String, String> context,
351-
final AtomicLong droppedQueueOverflow,
352-
final ConcurrentHashMap<String, AtomicLong> contextTruncatedCounts,
353-
final Runnable errorCallback) {
354-
this(
355-
() -> backendApiFactory.createBackendApi(Intake.EVENT_PLATFORM, false),
356-
queue,
357-
flushInterval,
358-
timeUnit,
359-
context,
360-
droppedQueueOverflow,
361-
contextTruncatedCounts,
362-
errorCallback,
363-
FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES);
364-
}
365-
366-
FlagEvaluationSerializingHandler(
367-
final BackendApiFactory backendApiFactory,
368-
final MessagePassingBlockingQueue<FlagEvalEvent> queue,
369-
final long flushInterval,
370-
final TimeUnit timeUnit,
371-
final Map<String, String> context,
372-
final AtomicLong droppedQueueOverflow,
373-
final ConcurrentHashMap<String, AtomicLong> contextTruncatedCounts,
374-
final Runnable errorCallback,
375-
final int payloadSizeLimitBytes) {
376-
this(
377-
() -> backendApiFactory.createBackendApi(Intake.EVENT_PLATFORM, false),
378-
queue,
379-
flushInterval,
380-
timeUnit,
381-
context,
382-
droppedQueueOverflow,
383-
contextTruncatedCounts,
384-
errorCallback,
385-
payloadSizeLimitBytes);
386-
}
387-
388-
FlagEvaluationSerializingHandler(
389-
final Supplier<BackendApi> backendApiSupplier,
390-
final MessagePassingBlockingQueue<FlagEvalEvent> queue,
391-
final long flushInterval,
392-
final TimeUnit timeUnit,
393-
final Map<String, String> context,
394-
final AtomicLong droppedQueueOverflow,
395-
final ConcurrentHashMap<String, AtomicLong> contextTruncatedCounts,
396-
final Runnable errorCallback) {
397-
this(
398-
backendApiSupplier,
399-
queue,
400-
flushInterval,
401-
timeUnit,
402-
context,
403-
droppedQueueOverflow,
404-
contextTruncatedCounts,
405-
errorCallback,
406-
FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES);
407-
}
408-
409308
FlagEvaluationSerializingHandler(
410309
final Supplier<BackendApi> backendApiSupplier,
411310
final MessagePassingBlockingQueue<FlagEvalEvent> queue,
@@ -622,16 +521,17 @@ private boolean shouldFlush() {
622521
*/
623522
static class SerializingHandlerForTest extends FlagEvaluationSerializingHandler {
624523

625-
SerializingHandlerForTest(final BackendApiFactory factory, final Map<String, String> context) {
626-
this(factory, context, FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES);
524+
SerializingHandlerForTest(
525+
final Supplier<BackendApi> backendApiSupplier, final Map<String, String> context) {
526+
this(backendApiSupplier, context, FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES);
627527
}
628528

629529
SerializingHandlerForTest(
630-
final BackendApiFactory factory,
530+
final Supplier<BackendApi> backendApiSupplier,
631531
final Map<String, String> context,
632532
final int payloadSizeLimitBytes) {
633533
super(
634-
factory,
534+
backendApiSupplier,
635535
Queues.mpscBlockingConsumerArrayQueue(DEFAULT_CAPACITY),
636536
Long.MAX_VALUE, // effectively never auto-flush
637537
TimeUnit.NANOSECONDS,
@@ -695,14 +595,14 @@ int fullTierSizeForTest() {
695595

696596
/** Factory method for test use - creates a SerializingHandlerForTest. */
697597
static SerializingHandlerForTest createHandlerForTest(
698-
final BackendApiFactory factory, final Map<String, String> context) {
699-
return new SerializingHandlerForTest(factory, context);
598+
final Supplier<BackendApi> backendApiSupplier, final Map<String, String> context) {
599+
return new SerializingHandlerForTest(backendApiSupplier, context);
700600
}
701601

702602
static SerializingHandlerForTest createHandlerForTest(
703-
final BackendApiFactory factory,
603+
final Supplier<BackendApi> backendApiSupplier,
704604
final Map<String, String> context,
705605
final int payloadSizeLimitBytes) {
706-
return new SerializingHandlerForTest(factory, context, payloadSizeLimitBytes);
606+
return new SerializingHandlerForTest(backendApiSupplier, context, payloadSizeLimitBytes);
707607
}
708608
}

0 commit comments

Comments
 (0)