Skip to content

Commit e1f29df

Browse files
committed
fix: Validate data source intervals in the retry state and bound the wait
1 parent 1b50827 commit e1f29df

15 files changed

Lines changed: 271 additions & 135 deletions

‎ldclient/async_config.py‎

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,6 @@
3030
from ldclient.impl.util import (
3131
log,
3232
validate_application_info,
33-
validate_positive_finite,
3433
validate_sdk_key_format
3534
)
3635
from ldclient.interfaces import (
@@ -258,13 +257,8 @@ def __init__(
258257
self.__stream_uri = stream_uri.rstrip('/')
259258
self.__update_processor_class = update_processor_class
260259
self.__stream = stream
261-
self.__initial_reconnect_delay = validate_positive_finite(
262-
initial_reconnect_delay, DEFAULT_INITIAL_RECONNECT_DELAY, 'initial_reconnect_delay', log
263-
)
264-
self.__poll_interval = max(
265-
validate_positive_finite(poll_interval, DEFAULT_POLL_INTERVAL, 'poll_interval', log),
266-
DEFAULT_POLL_INTERVAL,
267-
)
260+
self.__initial_reconnect_delay = initial_reconnect_delay
261+
self.__poll_interval = max(poll_interval, DEFAULT_POLL_INTERVAL)
268262
self.__use_ldd = use_ldd
269263
self.__feature_store = AsyncInMemoryFeatureStore() if not feature_store else feature_store
270264
self.__event_processor_class = event_processor_class

‎ldclient/config.py‎

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,6 @@
1515
from ldclient.impl.util import (
1616
log,
1717
validate_application_info,
18-
validate_positive_finite,
1918
validate_sdk_key_format
2019
)
2120
from ldclient.interfaces import (
@@ -407,13 +406,8 @@ def __init__(
407406
self.__stream_uri = stream_uri.rstrip('/')
408407
self.__update_processor_class = update_processor_class
409408
self.__stream = stream
410-
self.__initial_reconnect_delay = validate_positive_finite(
411-
initial_reconnect_delay, DEFAULT_INITIAL_RECONNECT_DELAY, 'initial_reconnect_delay', log
412-
)
413-
self.__poll_interval = max(
414-
validate_positive_finite(poll_interval, DEFAULT_POLL_INTERVAL, 'poll_interval', log),
415-
DEFAULT_POLL_INTERVAL,
416-
)
409+
self.__initial_reconnect_delay = initial_reconnect_delay
410+
self.__poll_interval = max(poll_interval, DEFAULT_POLL_INTERVAL)
417411
self.__use_ldd = use_ldd
418412
self.__feature_store = InMemoryFeatureStore() if not feature_store else feature_store
419413
self.__event_processor_class = event_processor_class

‎ldclient/impl/aio/transport.py‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -140,8 +140,7 @@ def create(self, url: str, initial_retry_delay: float, query_params=None, sdk_ma
140140
if proxy:
141141
aiohttp_request_options["proxy"] = proxy
142142
if sdk_managed_retry:
143-
# A zero base delay plus the no-op base strategy holds
144-
# next_retry_delay at zero, so the SSE client never sleeps.
143+
# The SSE client's retry is disabled; the SDK owns the delay.
145144
retry_options: dict = {
146145
"initial_retry_delay": 0,
147146
"retry_delay_strategy": RetryDelayStrategy(),

‎ldclient/impl/datasource/async_streaming.py‎

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -126,11 +126,8 @@ async def _run(self):
126126
log.info("AsyncStreamingUpdateProcessor initialized ok.")
127127
self._ready.set()
128128
elif isinstance(action, Fault):
129-
# A Fault with no error means the connection closed cleanly.
130-
# If we asked for that close, we have already recorded the
131-
# failure behind it and must not record it twice. Otherwise
132-
# the server closed a connection it normally leaves open,
133-
# which is a connection failure the SDK backs off from.
129+
# A Fault with no error is a clean close. An interrupt the
130+
# SDK asked for is not a failure.
134131
if action.error is None:
135132
if self._interrupted_by_sdk:
136133
self._interrupted_by_sdk = False

‎ldclient/impl/datasource/streaming.py‎

Lines changed: 5 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import json
22
import time
3+
from threading import TIMEOUT_MAX
34
from threading import Event as ThreadEvent
45
from threading import Thread
56
from typing import Callable, Optional
@@ -107,11 +108,8 @@ def run(self):
107108
log.info("StreamingUpdateProcessor initialized ok.")
108109
self._ready.set()
109110
elif isinstance(action, Fault):
110-
# A Fault with no error means the connection closed cleanly. If
111-
# we asked for that close, we have already recorded the failure
112-
# behind it and must not record it twice. Otherwise the server
113-
# closed a connection it normally leaves open, which is a
114-
# connection failure the SDK backs off from.
111+
# A Fault with no error is a clean close. An interrupt the SDK
112+
# asked for is not a failure.
115113
if action.error is None:
116114
if self._interrupted_by_sdk:
117115
self._interrupted_by_sdk = False
@@ -139,13 +137,7 @@ def _create_sse_client(self) -> SSEClient:
139137
url=self._uri, headers=http_factory.base_headers, pool=stream_http_factory.create_pool_manager(1, self._uri), urllib3_request_options={"timeout": stream_http_factory.timeout}
140138
),
141139
error_strategy=ErrorStrategy.always_continue(), # we'll make error-handling decisions when we see a Fault
142-
# The SDK owns the retry delay, so the SSE client must never wait.
143-
# A zero base delay plus the no-op base strategy holds
144-
# next_retry_delay at zero, which is what these three arguments
145-
# are for. The SSE client hands us the Fault before it would
146-
# sleep, so we classify the failure and wait ourselves in
147-
# _handle_error. Our wait is interruptible, which matters because
148-
# the extended regime can ask for an hour.
140+
# The SSE client's retry is disabled; the SDK owns the delay.
149141
initial_retry_delay=0,
150142
retry_delay_strategy=RetryDelayStrategy(),
151143
retry_delay_reset_threshold=0,
@@ -253,7 +245,7 @@ def _handle_error(self, error: Exception) -> bool:
253245
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, error_info)
254246

255247
self._connection_attempt_start_time = time.time() + delay
256-
return not self._stop_event.wait(delay)
248+
return not self._stop_event.wait(min(delay, TIMEOUT_MAX))
257249

258250
# magic methods for "with" statement (used in testing)
259251
def __enter__(self):

‎ldclient/impl/repeating_task.py‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
from threading import Event, Thread
1+
from threading import TIMEOUT_MAX, Event, Thread
22
from typing import Any, Callable
33

44
from ldclient.impl.delay import DelaySource, FixedDelay
@@ -66,7 +66,7 @@ def stop(self):
6666

6767
def _run(self):
6868
if self.__initial_delay > 0:
69-
if self.__stop.wait(self.__initial_delay):
69+
if self.__stop.wait(min(self.__initial_delay, TIMEOUT_MAX)):
7070
return
7171
stopped = self.__stop.is_set()
7272
while not stopped:
@@ -77,4 +77,4 @@ def _run(self):
7777
# The wait starts when the callback returns, so a slow callback
7878
# never shortens it.
7979
delay = self.__delays.next_delay
80-
stopped = self.__stop.wait(delay) if delay > 0 else self.__stop.is_set()
80+
stopped = self.__stop.wait(min(delay, TIMEOUT_MAX)) if delay > 0 else self.__stop.is_set()

‎ldclient/impl/retry.py‎

Lines changed: 29 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818

1919
# currently excluded from documentation - see docs/README.md
2020

21+
import math
2122
import random
2223
import time
2324
from enum import Enum
@@ -27,7 +28,7 @@
2728
DEFAULT_INITIAL_RECONNECT_DELAY,
2829
DEFAULT_POLL_INTERVAL
2930
)
30-
from ldclient.impl.util import log, validate_positive_finite
31+
from ldclient.impl.util import log
3132

3233
# The delay ceiling of the normal regime for streaming, in seconds.
3334
NORMAL_STREAMING_CEILING_DELAY = 30
@@ -53,6 +54,23 @@
5354
_MAX_BACKOFF_EXPONENT = 30
5455

5556

57+
def _usable_delay(value: float, default: float, name: str, ceiling: float = math.inf) -> float:
58+
"""
59+
Returns the delay to use, clamped to the ceiling. A value that is
60+
not a positive, finite number of seconds is replaced by the default.
61+
62+
:param value: the configured number of seconds
63+
:param default: the value to use when ``value`` is not usable
64+
:param name: the option name, for the warning message
65+
:param ceiling: the longest delay allowed
66+
"""
67+
68+
if value > 0 and math.isfinite(value):
69+
return min(value, ceiling)
70+
log.warning("%s must be a positive, finite number of seconds; using the default of %ss" % (name, default))
71+
return default
72+
73+
5674
class FailureKind(Enum):
5775
"""How a failure is classified, which decides how long the next wait is."""
5876

@@ -251,19 +269,12 @@ def for_streaming(initial_reconnect_delay: float) -> RetryState:
251269
"""
252270
Builds the retry state for a streaming data source.
253271
254-
Streaming's operating cadence is zero, so there is no delay during
255-
healthy operation. Stream failures use either the normal or extended
256-
initial delay to determine their backoff wait. A stream returns to
257-
healthy operation after establishing a successful connection with no
258-
failures during the ``STREAMING_RESET_INTERVAL``.
259-
260-
``Config`` validates the configured delay, so this guard only catches a
261-
state built without it.
262-
263-
The extended regime never starts below the configured delay.
272+
Streaming's cadence is zero, so a healthy stream never waits. An invalid
273+
delay value is replaced by the documented default; one longer than a
274+
ceiling raises that bound rather than being cut down to it.
264275
"""
265-
initial_reconnect_delay = validate_positive_finite(
266-
initial_reconnect_delay, DEFAULT_INITIAL_RECONNECT_DELAY, 'initial_reconnect_delay', log
276+
initial_reconnect_delay = _usable_delay(
277+
initial_reconnect_delay, DEFAULT_INITIAL_RECONNECT_DELAY, 'initial_reconnect_delay'
267278
)
268279
return RetryState(
269280
normal_initial_delay=initial_reconnect_delay,
@@ -279,16 +290,12 @@ def for_polling(poll_interval: float) -> RetryState:
279290
"""
280291
Builds the retry state for a polling data source.
281292
282-
The poll interval is polling's operating cadence, so no wait is ever
283-
shorter than it. In the normal regime the delay bounds are the poll
284-
interval itself, which means a normal failure simply polls again on
285-
schedule. Polling is healthy on any successful poll, and resets after two
286-
in a row.
287-
288-
``Config`` validates and clamps the poll interval, so this guard only
289-
catches a state built without it.
293+
The poll interval is polling's cadence and its normal ceiling, so a normal
294+
failure waits the interval rather than backing off past it. An invalid
295+
interval is replaced by the documented default. No wait is ever shorter
296+
than the interval, so the cadence wins over the extended ceiling.
290297
"""
291-
poll_interval = validate_positive_finite(poll_interval, DEFAULT_POLL_INTERVAL, 'poll_interval', log)
298+
poll_interval = _usable_delay(poll_interval, DEFAULT_POLL_INTERVAL, 'poll_interval')
292299
return RetryState(
293300
normal_initial_delay=poll_interval,
294301
normal_ceiling_delay=poll_interval,

‎ldclient/impl/util.py‎

Lines changed: 4 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,4 @@
11
import logging
2-
import math
32
import re
43
import sys
54
import time
@@ -61,25 +60,6 @@ def validate_application_value(value: Any, name: str, logger: logging.Logger) ->
6160
return value
6261

6362

64-
def validate_positive_finite(value: float, default: float, name: str, logger: logging.Logger) -> float:
65-
"""
66-
Validates that a number of seconds is positive and finite.
67-
68-
A non-finite value makes later arithmetic produce NaN, and a non-positive
69-
one makes a wait no wait at all.
70-
71-
:param value: the number of seconds to validate
72-
:param default: the value to use when ``value`` is not usable
73-
:param name: the option name, for the warning message
74-
:param logger: the logger to use for logging warnings
75-
:return: ``value``, or ``default`` if ``value`` is not positive and finite
76-
"""
77-
if value > 0 and math.isfinite(value):
78-
return value
79-
logger.warning("%s must be a positive, finite number of seconds; using the default of %ss" % (name, default))
80-
return default
81-
82-
8363
def validate_sdk_key_format(sdk_key: str, logger: logging.Logger) -> str:
8464
"""
8565
Validates that an SDK key does not contain invalid characters and is not too long for our systems.
@@ -155,14 +135,11 @@ def throw_if_unsuccessful_response(resp):
155135

156136
def is_http_error_recoverable(status):
157137
"""
158-
Reports whether a component that treats some statuses as fatal should
159-
keep going.
160-
161138
Deprecated. Use :func:`ldclient.impl.retry.classify_http_status` instead.
162139
"""
163140
if status >= 400 and status < 500:
164-
return status in _RETRYABLE_STATUSES # all other 4xx besides these are treated as fatal
165-
return True
141+
return status in _RETRYABLE_STATUSES # all other 4xx besides these are unrecoverable
142+
return True # all other errors are recoverable
166143

167144

168145
def http_error_description(status):
@@ -171,11 +148,8 @@ def http_error_description(status):
171148

172149
def http_error_message(status, context, retryable_message="will retry"):
173150
"""
174-
Builds the log message for an HTTP failure in a component that stops on
175-
some statuses.
176-
177-
Deprecated. The FDv1 data sources build their own message instead, so
178-
that it can report the real retry delay.
151+
Deprecated. The FDv1 data sources build their own message, so that it can
152+
report the real retry delay.
179153
"""
180154
return "Received %s for %s - %s" % (http_error_description(status), context, retryable_message if is_http_error_recoverable(status) else "giving up permanently")
181155

‎ldclient/interfaces.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1013,7 +1013,8 @@ class DataSourceState(Enum):
10131013
"""
10141014
Indicates that the data source has been permanently shut down.
10151015
1016-
This means the SDK client was explicitly shut down, or that its configuration could not be parsed.
1016+
This could be because the SDK client was explicitly shut down, because its configuration could not
1017+
be parsed, or because the data source encountered a condition it will not retry.
10171018
"""
10181019

10191020

‎ldclient/testing/impl/datasource/test_async_polling.py‎

Lines changed: 30 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import asyncio
66
import logging
77
import ssl
8+
import time
89
from unittest.mock import AsyncMock, MagicMock, patch
910

1011
import aiohttp
@@ -58,9 +59,12 @@ def make_config(**kwargs):
5859
)
5960

6061

62+
ONE_HOUR = 60 * 60
63+
64+
6165
def fast_retry_state(delay=0.001):
62-
"""A retry state with tiny delays, so a test does not have to wait out the
63-
real extended-regime delay of five minutes."""
66+
"""A retry state whose every delay is ``delay``: small enough to skip the
67+
real extended-regime wait, or large enough to prove a stop interrupts one."""
6468
return RetryState(
6569
normal_initial_delay=delay,
6670
normal_ceiling_delay=delay,
@@ -469,6 +473,30 @@ async def close():
469473

470474
assert order == ['poll_done', 'transport_closed']
471475

476+
@pytest.mark.asyncio
477+
async def test_an_extended_regime_wait_is_cut_short_by_stop(self):
478+
"""Shutdown must not sit through an hour-long backoff. The in-flight-poll
479+
case is test_stop_cancels_polling_task_cleanly; this one stops while the
480+
task is waiting between polls."""
481+
retry = fast_retry_state(ONE_HOUR)
482+
processor = make_processor(retry_state=retry)
483+
processor._requester.get_all_data = AsyncMock(
484+
side_effect=UnsuccessfulResponseException(401)
485+
)
486+
487+
processor.start()
488+
# Confirm the wait under test really is long before measuring the stop.
489+
deadline = time.time() + 2
490+
while retry.next_delay <= 60 and time.time() < deadline:
491+
await asyncio.sleep(0.01)
492+
assert retry.next_delay > 60, "the wait under test should be minutes long"
493+
494+
started = time.time()
495+
await processor.stop()
496+
elapsed = time.time() - started
497+
498+
assert elapsed < 2, "stop() took %.2fs" % elapsed
499+
472500
@pytest.mark.asyncio
473501
async def test_stop_cancels_polling_task_cleanly(self):
474502
store = MockAsyncFeatureStore()

0 commit comments

Comments
 (0)