Skip to content

Commit 6df48d1

Browse files
committed
fix(feature-flags): complete initial poll during activation
1 parent cfaca97 commit 6df48d1

2 files changed

Lines changed: 69 additions & 11 deletions

File tree

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

Lines changed: 27 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,9 @@ final class AgentlessConfigurationSource implements ConfigurationSourceService {
5656
private final RatelimitedLogger ratelimitedLogger;
5757
private final Object lifecycleLock = new Object();
5858
private final AtomicBoolean polling = new AtomicBoolean();
59+
private final AtomicReference<Thread> pollingThread = new AtomicReference<>();
5960
private volatile boolean closed;
61+
private volatile boolean started;
6062
private volatile ScheduledFuture<?> scheduledPoll;
6163
private volatile String etag;
6264

@@ -114,22 +116,38 @@ private AgentlessConfigurationSource(final Config config, final HttpUrl endpoint
114116
@Override
115117
public void init() {
116118
synchronized (lifecycleLock) {
117-
if (closed || scheduledPoll != null) {
119+
if (closed || started) {
118120
return;
119121
}
120-
scheduledPoll =
121-
executor.scheduleWithFixedDelay(
122-
this::pollOnceSafely, 0, pollIntervalMillis, TimeUnit.MILLISECONDS);
122+
started = true;
123+
}
124+
125+
// Complete the first poll cycle on the activation thread. This lets OpenFeature provider
126+
// initialization observe a successful retry before it checks whether configuration is ready.
127+
// No request occurs before application code activates the provider.
128+
pollOnceSafely();
129+
130+
synchronized (lifecycleLock) {
131+
if (!closed) {
132+
scheduledPoll =
133+
executor.scheduleWithFixedDelay(
134+
this::pollOnceSafely,
135+
pollIntervalMillis,
136+
pollIntervalMillis,
137+
TimeUnit.MILLISECONDS);
138+
}
123139
}
124140
}
125141

126142
boolean pollOnce() {
127143
if (closed || !polling.compareAndSet(false, true)) {
128144
return false;
129145
}
146+
pollingThread.set(Thread.currentThread());
130147
try {
131148
return fetchAndApply();
132149
} finally {
150+
pollingThread.compareAndSet(Thread.currentThread(), null);
133151
polling.set(false);
134152
}
135153
}
@@ -142,13 +160,18 @@ public void close() {
142160
return;
143161
}
144162
closed = true;
163+
started = false;
145164
poll = scheduledPoll;
146165
scheduledPoll = null;
147166
}
148167
if (poll != null) {
149168
poll.cancel(true);
150169
}
151170
client.cancel();
171+
final Thread activePollingThread = pollingThread.get();
172+
if (activePollingThread != null) {
173+
activePollingThread.interrupt();
174+
}
152175
executor.shutdownNow();
153176
}
154177

products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/AgentlessConfigurationSourceTest.java

Lines changed: 42 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -847,7 +847,7 @@ void rejectsOverlappingPolls() throws Exception {
847847
}
848848

849849
@Test
850-
void initSchedulesPollAndCloseCancelsFuture() throws Exception {
850+
void initCompletesFirstPollAndCloseCancelsScheduledFuture() throws Exception {
851851
final FakeClient client = new FakeClient(response(200, "etag-a", emptyConfig()));
852852
final AgentlessConfigurationSource service =
853853
new AgentlessConfigurationSource(
@@ -859,12 +859,41 @@ void initSchedulesPollAndCloseCancelsFuture() throws Exception {
859859
FeatureFlaggingGateway.addConfigListener(listener);
860860

861861
service.init();
862-
awaitCalls(client, 1);
862+
assertEquals(1, client.calls.get());
863863
service.close();
864864

865865
verify(listener).accept(any(ServerConfiguration.class));
866866
}
867867

868+
@Test
869+
void initCompletesInitialRetryCycleBeforeReturning() throws Exception {
870+
final List<okhttp3.Request> requests = new ArrayList<>();
871+
final AgentlessConfigurationSource.OkHttpUfcHttpClient client =
872+
scriptedClient(
873+
requests,
874+
delay -> {},
875+
() -> 1.0,
876+
response(500, null, null),
877+
response(200, "etag-a", emptyConfig()));
878+
final AgentlessConfigurationSource service =
879+
new AgentlessConfigurationSource(
880+
HttpUrl.get("http://localhost" + CONFIG_PATH),
881+
config(),
882+
60_000,
883+
client,
884+
Executors.newSingleThreadScheduledExecutor());
885+
FeatureFlaggingGateway.addConfigListener(listener);
886+
887+
try {
888+
service.init();
889+
890+
assertEquals(2, requests.size());
891+
verify(listener).accept(any(ServerConfiguration.class));
892+
} finally {
893+
service.close();
894+
}
895+
}
896+
868897
@Test
869898
void repeatedInitStartsOnlyOnePoller() throws Exception {
870899
final FakeClient client = new FakeClient(response(200, "etag-a", emptyConfig()));
@@ -979,14 +1008,20 @@ void closeInterruptsRetryBackoff() throws Exception {
9791008
final AgentlessConfigurationSource service =
9801009
new AgentlessConfigurationSource(
9811010
HttpUrl.get("http://localhost" + CONFIG_PATH), config(), 30_000, client, executor);
1011+
final ExecutorService runner = Executors.newSingleThreadExecutor();
9821012

983-
service.init();
984-
assertTrue(backoffStarted.await(1, TimeUnit.SECONDS));
1013+
try {
1014+
final Future<?> initialization = runner.submit(service::init);
1015+
assertTrue(backoffStarted.await(1, TimeUnit.SECONDS));
9851016

986-
service.close();
1017+
service.close();
9871018

988-
assertTrue(executor.awaitTermination(1, TimeUnit.SECONDS));
989-
assertEquals(1, requests.size());
1019+
initialization.get(1, TimeUnit.SECONDS);
1020+
assertTrue(executor.awaitTermination(1, TimeUnit.SECONDS));
1021+
assertEquals(1, requests.size());
1022+
} finally {
1023+
runner.shutdownNow();
1024+
}
9901025
}
9911026

9921027
@Test

0 commit comments

Comments
 (0)