Skip to content

Commit e17e173

Browse files
authored
feat: Retry indefinitely after a data source failure instead of stopping permanently in FDv1
fix: Warn and use the documented default for an invalid poll interval or initial reconnect delay
1 parent 1087147 commit e17e173

27 files changed

Lines changed: 2183 additions & 388 deletions

‎contract-tests/async_service.py‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,8 @@ async def handle_status(request: aiohttp.web.Request) -> aiohttp.web.Response:
6969
'migrations',
7070
'persistent-data-store-redis',
7171
'fdv1-fallback',
72+
'retry-conformance-fdv1-streaming',
73+
'retry-conformance-fdv1-polling',
7274
]
7375
}
7476
return aiohttp.web.Response(

‎contract-tests/service.py‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -86,6 +86,8 @@ def status():
8686
'flag-change-listeners',
8787
'flag-value-change-listeners',
8888
'fdv1-fallback',
89+
'retry-conformance-fdv1-streaming',
90+
'retry-conformance-fdv1-polling',
8991
]
9092
}
9193
return json.dumps(body), 200, {'Content-type': 'application/json'}

‎ldclient/async_client.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -389,8 +389,8 @@ async def is_initialized(self) -> bool:
389389
390390
If this returns false, it means that the client has not yet successfully connected to LaunchDarkly.
391391
It might still be in the process of starting up, or it might be attempting to reconnect after an
392-
unsuccessful attempt, or it might have received an unrecoverable error (such as an invalid SDK key)
393-
and given up.
392+
unsuccessful attempt, or it might have received an error that needs to be fixed (such
393+
as an invalid SDK key).
394394
395395
This is a coroutine because determining readiness may query a persistent store.
396396
"""

‎ldclient/async_config.py‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@
1515
from ldclient.config import (
1616
DEFAULT_BASE_URI,
1717
DEFAULT_EVENTS_URI,
18+
DEFAULT_INITIAL_RECONNECT_DELAY,
19+
DEFAULT_POLL_INTERVAL,
1820
DEFAULT_STREAM_URI,
1921
GET_LATEST_FEATURES_PATH,
2022
STREAM_FLAGS_PATH,
@@ -149,11 +151,11 @@ def __init__(
149151
flush_interval: float = 5,
150152
stream_uri: str = DEFAULT_STREAM_URI,
151153
stream: bool = True,
152-
initial_reconnect_delay: float = 1,
154+
initial_reconnect_delay: float = DEFAULT_INITIAL_RECONNECT_DELAY,
153155
defaults: dict = {},
154156
send_events: Optional[bool] = None,
155157
update_processor_class: Optional[Callable[['AsyncConfig', AsyncFeatureStore, AsyncEvent], AsyncUpdateProcessor]] = None,
156-
poll_interval: float = 30,
158+
poll_interval: float = DEFAULT_POLL_INTERVAL,
157159
use_ldd: bool = False,
158160
feature_store: Optional[AsyncFeatureStore] = None,
159161
feature_requester_class=None,
@@ -256,7 +258,7 @@ def __init__(
256258
self.__update_processor_class = update_processor_class
257259
self.__stream = stream
258260
self.__initial_reconnect_delay = initial_reconnect_delay
259-
self.__poll_interval = max(poll_interval, 30.0)
261+
self.__poll_interval = max(poll_interval, DEFAULT_POLL_INTERVAL)
260262
self.__use_ldd = use_ldd
261263
self.__feature_store = AsyncInMemoryFeatureStore() if not feature_store else feature_store
262264
self.__event_processor_class = event_processor_class

‎ldclient/client.py‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -369,10 +369,11 @@ def is_initialized(self) -> bool:
369369
370370
If this returns false, it means the client has not yet obtained any flag data. It might still be
371371
starting up, or attempting to reconnect after an unsuccessful attempt, or it might have received
372-
an unrecoverable error (such as an invalid SDK key) and given up. In this state, feature flag
373-
evaluations will return default values -- unless you are using a persistent store integration and
374-
flag data had already been stored by a successfully connected SDK in the past. You can use
375-
:attr:`data_source_status_provider` to get information on errors, or to wait for a successful retry.
372+
an error that needs to be fixed (such as an invalid SDK key). In this state, feature flag
373+
evaluations will return default values -- unless you are using a persistent store integration
374+
and flag data had already been stored by a successfully connected SDK in the past. You can use
375+
:attr:`data_source_status_provider` to get information on errors, or to wait for a
376+
successful retry.
376377
377378
:return: true if the client is initialized and has flag data available
378379
"""

‎ldclient/config.py‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,11 @@
3636
DEFAULT_EVENTS_URI = 'https://events.launchdarkly.com'
3737
DEFAULT_STREAM_URI = 'https://stream.launchdarkly.com'
3838

39+
# Defaults, in seconds, for the two configurable data source intervals. The
40+
# poll interval is also its own minimum.
41+
DEFAULT_INITIAL_RECONNECT_DELAY = 1
42+
DEFAULT_POLL_INTERVAL = 30
43+
3944

4045
class BigSegmentsConfig:
4146
"""Configuration options related to Big Segments.
@@ -295,11 +300,11 @@ def __init__(
295300
flush_interval: float = 5,
296301
stream_uri: str = DEFAULT_STREAM_URI,
297302
stream: bool = True,
298-
initial_reconnect_delay: float = 1,
303+
initial_reconnect_delay: float = DEFAULT_INITIAL_RECONNECT_DELAY,
299304
defaults: dict = {},
300305
send_events: Optional[bool] = None,
301306
update_processor_class: Optional[Callable[['Config', FeatureStore, Event], UpdateProcessor]] = None,
302-
poll_interval: float = 30,
307+
poll_interval: float = DEFAULT_POLL_INTERVAL,
303308
use_ldd: bool = False,
304309
feature_store: Optional[FeatureStore] = None,
305310
feature_requester_class=None,
@@ -402,7 +407,7 @@ def __init__(
402407
self.__update_processor_class = update_processor_class
403408
self.__stream = stream
404409
self.__initial_reconnect_delay = initial_reconnect_delay
405-
self.__poll_interval = max(poll_interval, 30.0)
410+
self.__poll_interval = max(poll_interval, DEFAULT_POLL_INTERVAL)
406411
self.__use_ldd = use_ldd
407412
self.__feature_store = InMemoryFeatureStore() if not feature_store else feature_store
408413
self.__event_processor_class = event_processor_class

‎ldclient/impl/aio/transport.py‎

Lines changed: 29 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -118,11 +118,16 @@ def __init__(self, config, session: Optional[aiohttp.ClientSession] = None, prox
118118
self._http_options = http_options if http_options is not None else config.http
119119
self._proxy = proxy if proxy is not None else (self._http_options.http_proxy or None)
120120

121-
def create(self, url: str, initial_retry_delay: float, query_params=None) -> AsyncSSEClient:
122-
"""Builds an SSE client for the given stream URL. Headers, timeouts,
123-
proxy settings, and the retry/backoff policy come from the SDK config.
124-
``query_params`` is an optional zero-argument callable evaluated on
125-
each (re)connect to produce additional query string parameters."""
121+
def create(self, url: str, initial_retry_delay: float, query_params=None, sdk_managed_retry: bool = False) -> AsyncSSEClient:
122+
"""Builds an SSE client for the given stream URL. Headers, timeouts and
123+
proxy settings come from the SDK config. ``query_params`` is an
124+
optional zero-argument callable evaluated on each (re)connect to
125+
produce additional query string parameters.
126+
127+
``sdk_managed_retry`` moves the delay between connection attempts to
128+
the caller. The SSE client then never waits, and
129+
``initial_retry_delay`` is ignored. When it is false, the SSE client
130+
backs off on its own."""
126131
base_headers = _base_headers(self._config, ASYNC_USER_AGENT)
127132
aiohttp_request_options: dict = {
128133
"timeout": aiohttp.ClientTimeout(
@@ -134,6 +139,24 @@ def create(self, url: str, initial_retry_delay: float, query_params=None) -> Asy
134139
proxy = self._proxy or _get_proxy_url(url)
135140
if proxy:
136141
aiohttp_request_options["proxy"] = proxy
142+
if sdk_managed_retry:
143+
# The SSE client's retry is disabled; the SDK owns the delay. The base
144+
# strategy must be passed: omitting it selects the library's backoff.
145+
retry_options: dict = {
146+
"initial_retry_delay": 0,
147+
"retry_delay_strategy": RetryDelayStrategy(),
148+
"retry_delay_reset_threshold": 0,
149+
}
150+
else:
151+
retry_options = {
152+
"initial_retry_delay": initial_retry_delay,
153+
"retry_delay_strategy": RetryDelayStrategy.default(
154+
max_delay=MAX_RETRY_DELAY,
155+
backoff_multiplier=2,
156+
jitter_multiplier=JITTER_RATIO,
157+
),
158+
"retry_delay_reset_threshold": BACKOFF_RESET_INTERVAL,
159+
}
137160
return AsyncSSEClient(
138161
connect=AsyncConnectStrategy.http(
139162
url=url,
@@ -143,12 +166,6 @@ def create(self, url: str, initial_retry_delay: float, query_params=None) -> Asy
143166
query_params=query_params,
144167
),
145168
error_strategy=ErrorStrategy.always_continue(), # we'll make error-handling decisions when we see a Fault
146-
initial_retry_delay=initial_retry_delay,
147-
retry_delay_strategy=RetryDelayStrategy.default(
148-
max_delay=MAX_RETRY_DELAY,
149-
backoff_multiplier=2,
150-
jitter_multiplier=JITTER_RATIO,
151-
),
152-
retry_delay_reset_threshold=BACKOFF_RESET_INTERVAL,
153169
logger=log,
170+
**retry_options,
154171
)

‎ldclient/impl/datasource/async_polling.py‎

Lines changed: 54 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -9,11 +9,16 @@
99

1010
from ldclient.async_config import AsyncConfig
1111
from ldclient.impl.aio.concurrency import AsyncEvent, AsyncRepeatingTask
12-
from ldclient.impl.datasource.datasource_common import sink_or_store
12+
from ldclient.impl.datasource.datasource_common import async_sink_or_store
13+
from ldclient.impl.retry import (
14+
FailureKind,
15+
RetryState,
16+
classify_http_status,
17+
for_polling
18+
)
1319
from ldclient.impl.util import (
1420
UnsuccessfulResponseException,
15-
http_error_message,
16-
is_http_error_recoverable,
21+
http_error_description,
1722
log
1823
)
1924
from ldclient.interfaces import (
@@ -27,13 +32,22 @@
2732

2833

2934
class AsyncPollingUpdateProcessor(AsyncUpdateProcessor):
30-
def __init__(self, config: AsyncConfig, requester: AsyncFeatureRequester, store: AsyncFeatureStore, ready: AsyncEvent):
35+
"""Polls LaunchDarkly for flag data on its own background task.
36+
37+
The loop reads its wait from the retry state, which ``_fetch_and_store``
38+
updates, so a failure can push the next poll further out than the poll
39+
interval. See :mod:`ldclient.impl.retry`.
40+
"""
41+
42+
def __init__(self, config: AsyncConfig, requester: AsyncFeatureRequester, store: AsyncFeatureStore, ready: AsyncEvent, retry_state: Optional[RetryState] = None):
3143
self._config = config
3244
self._data_source_update_sink = config.data_source_update_sink
3345
self._requester = requester
3446
self._store = store
3547
self._ready = ready
36-
self._task = AsyncRepeatingTask.at_interval("ldclient.datasource.polling", config.poll_interval, 0, self._fetch_and_store)
48+
self._retry = retry_state or for_polling(config.poll_interval)
49+
# No initial delay: the first poll is immediate.
50+
self._task = AsyncRepeatingTask("ldclient.datasource.polling", self._retry, 0, self._fetch_and_store)
3751

3852
def start(self):
3953
log.info("Starting AsyncPollingUpdateProcessor with request interval: " + str(self._config.poll_interval))
@@ -43,48 +57,53 @@ def initialized(self):
4357
return self._ready.is_set() and self._store.initialized
4458

4559
async def stop(self):
46-
self.__stop_with_error_info(None)
47-
# Wait for the current poll to finish before closing the transport, so we do
48-
# not close it while a request is still using it. The close is in a finally
49-
# so an owned transport is still released if stop() is cancelled mid-wait.
50-
try:
51-
await self._task.wait_stopped()
52-
finally:
53-
await self._requester.close()
54-
55-
def __stop_with_error_info(self, error: Optional[DataSourceErrorInfo]):
5660
log.info("Stopping AsyncPollingUpdateProcessor")
5761
self._task.stop()
5862

59-
if self._data_source_update_sink is None:
60-
return
63+
if self._data_source_update_sink is not None:
64+
self._data_source_update_sink.update_status(DataSourceState.OFF, None)
6165

62-
self._data_source_update_sink.update_status(DataSourceState.OFF, error)
66+
# Do not close the transport while an in-flight request still uses it.
67+
try:
68+
await self._task.wait_stopped()
69+
finally:
70+
await self._requester.close()
6371

64-
async def _fetch_and_store(self):
72+
async def _fetch_and_store(self) -> None:
73+
"""Makes one poll request and records the outcome on the retry state."""
6574
try:
6675
all_data = await self._requester.get_all_data()
67-
await sink_or_store(self._data_source_update_sink, self._store).init(all_data)
76+
await async_sink_or_store(self._data_source_update_sink, self._store).init(all_data)
77+
78+
if self._data_source_update_sink is not None:
79+
self._data_source_update_sink.update_status(DataSourceState.VALID, None)
80+
81+
# Report the status before signaling readiness, so a caller that
82+
# wakes on readiness cannot still read INITIALIZING.
6883
if not self._ready.is_set() and self._store.initialized:
6984
log.info("AsyncPollingUpdateProcessor initialized ok")
7085
self._ready.set()
7186

72-
if self._data_source_update_sink is not None:
73-
self._data_source_update_sink.update_status(DataSourceState.VALID, None)
87+
self._retry.record_success()
88+
return
7489
except UnsuccessfulResponseException as e:
90+
kind = classify_http_status(e.status)
7591
error_info = DataSourceErrorInfo(DataSourceErrorKind.ERROR_RESPONSE, e.status, time.time(), str(e))
92+
description = "Received %s for polling request" % http_error_description(e.status)
93+
level = log.error if kind is FailureKind.UNEXPECTED else log.warning
94+
stacktrace = None
95+
except Exception as e:
96+
kind = FailureKind.NORMAL
97+
error_info = DataSourceErrorInfo(DataSourceErrorKind.UNKNOWN, 0, time.time(), str(e))
98+
description = "Error encountered when updating flags: %s" % e
99+
level = log.error
100+
# The exception is passed explicitly: by the time the message is
101+
# logged, the handler has exited and exc_info() is empty.
102+
stacktrace = e
76103

77-
http_error_message_result = http_error_message(e.status, "polling request")
78-
if not is_http_error_recoverable(e.status):
79-
log.error(http_error_message_result)
80-
self._ready.set() # if client is initializing, make it stop waiting; has no effect if already inited
81-
self.__stop_with_error_info(error_info)
82-
else:
83-
log.warning(http_error_message_result)
104+
self._retry.record_failure(kind)
105+
delay = self._retry.next_delay
106+
level("%s - will retry in %.1fs" % (description, delay), exc_info=stacktrace)
84107

85-
if self._data_source_update_sink is not None:
86-
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, error_info)
87-
except Exception as e:
88-
log.exception('Error: Exception encountered when updating flags. %s' % e)
89-
if self._data_source_update_sink is not None:
90-
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, DataSourceErrorInfo(DataSourceErrorKind.UNKNOWN, 0, time.time(), str(e)))
108+
if self._data_source_update_sink is not None:
109+
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, error_info)

‎ldclient/impl/datasource/async_status.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,11 @@ def update_status(self, new_state: DataSourceState, new_error: Optional[DataSour
7272

7373
old_status = self.__status
7474

75+
# OFF is terminal. A poll or stream connection that was still in
76+
# flight when the data source stopped must not report after it.
77+
if old_status.state == DataSourceState.OFF:
78+
return
79+
7580
if new_state == DataSourceState.INTERRUPTED and old_status.state == DataSourceState.INITIALIZING:
7681
new_state = DataSourceState.INITIALIZING
7782

0 commit comments

Comments
 (0)