diff --git a/ldclient/impl/datasource/async_streaming.py b/ldclient/impl/datasource/async_streaming.py index 7c52af45..fb3c5bc5 100644 --- a/ldclient/impl/datasource/async_streaming.py +++ b/ldclient/impl/datasource/async_streaming.py @@ -27,7 +27,7 @@ classify_http_status, for_streaming ) -from ldclient.impl.util import http_error_description, log +from ldclient.impl.util import http_error_description, log, monotonic_seconds from ldclient.interfaces import ( AsyncUpdateProcessor, DataSourceErrorInfo, @@ -61,7 +61,7 @@ def __init__(self, config, store, ready, diagnostic_accumulator, sse_factory: Op self._sse_factory = sse_factory self._owned_session = None self._sse: Any = None - self._connection_attempt_start_time: Optional[float] = None + self._connection_attempt_started_monotonic: Optional[float] = None self._runner = AsyncTaskRunner() self._started = False self._retry = retry_state or for_streaming(config.initial_reconnect_delay) @@ -86,7 +86,7 @@ async def _run(self): self._running = True try: self._sse = self._sse_factory.create(self._uri, self._config.initial_reconnect_delay, sdk_managed_retry=True) - self._connection_attempt_start_time = time.time() + self._connection_attempt_started_monotonic = monotonic_seconds() async for action in self._sse.all: if isinstance(action, Start): # interrupt() is a no-op when the connection has already gone, so @@ -94,9 +94,9 @@ async def _run(self): self._interrupted_by_sdk = False # On reconnect after an error the timer was cleared; reset it here. - # For the initial connect the pre-loop timestamp is already set. - if self._connection_attempt_start_time is None: - self._connection_attempt_start_time = time.time() + # For the initial connect the pre-loop stamp is already set. + if self._connection_attempt_started_monotonic is None: + self._connection_attempt_started_monotonic = monotonic_seconds() elif isinstance(action, Event): message_ok = False message_handled = False @@ -121,7 +121,7 @@ async def _run(self): if message_ok: self._record_stream_init(False) - self._connection_attempt_start_time = None + self._connection_attempt_started_monotonic = None if self._data_source_update_sink is not None: self._data_source_update_sink.update_status(DataSourceState.VALID, None) @@ -160,10 +160,10 @@ async def _close_owned_session(self): self._owned_session = None def _record_stream_init(self, failed: bool): - if self._diagnostic_accumulator and self._connection_attempt_start_time: + if self._diagnostic_accumulator and self._connection_attempt_started_monotonic is not None: current_time = int(time.time() * 1000) - elapsed = current_time - int(self._connection_attempt_start_time * 1000) - self._diagnostic_accumulator.record_stream_init(current_time, elapsed if elapsed >= 0 else 0, failed) + elapsed = int((monotonic_seconds() - self._connection_attempt_started_monotonic) * 1000) + self._diagnostic_accumulator.record_stream_init(current_time, elapsed, failed) async def stop(self): log.info("Stopping AsyncStreamingUpdateProcessor") @@ -276,9 +276,9 @@ async def _handle_error(self, error: Exception) -> bool: if delay > 0: await asyncio.sleep(delay) - # Read after the wait, so a clock change during it cannot skew the + # Set after the wait, so the backoff delay is not counted in the # stream-init latency we report. - self._connection_attempt_start_time = time.time() + self._connection_attempt_started_monotonic = monotonic_seconds() return self._running # magic methods for "with" statement (used in testing) diff --git a/ldclient/impl/datasource/streaming.py b/ldclient/impl/datasource/streaming.py index 2e48e918..9e425e3f 100644 --- a/ldclient/impl/datasource/streaming.py +++ b/ldclient/impl/datasource/streaming.py @@ -29,7 +29,7 @@ classify_http_status, for_streaming ) -from ldclient.impl.util import http_error_description, log +from ldclient.impl.util import http_error_description, log, monotonic_seconds from ldclient.interfaces import ( DataSourceErrorInfo, DataSourceErrorKind, @@ -62,7 +62,7 @@ def __init__(self, config, store, ready, diagnostic_accumulator, retry_state: Op self._running = False self._ready = ready self._diagnostic_accumulator = diagnostic_accumulator - self._connection_attempt_start_time: Optional[float] = None + self._connection_attempt_started_monotonic: Optional[float] = None self._retry = retry_state or for_streaming(config.initial_reconnect_delay) self._sse: Optional[SSEClient] = None self._stop_event = ThreadEvent() @@ -79,7 +79,7 @@ def run(self): self._sse.close() return - self._connection_attempt_start_time = time.time() + self._connection_attempt_started_monotonic = monotonic_seconds() try: for action in self._sse.all: if isinstance(action, Start): @@ -111,7 +111,7 @@ def run(self): if message_ok: self._record_stream_init(False) - self._connection_attempt_start_time = None + self._connection_attempt_started_monotonic = None if self._data_source_update_sink is not None: self._data_source_update_sink.update_status(DataSourceState.VALID, None) @@ -137,10 +137,10 @@ def run(self): self._sse.close() def _record_stream_init(self, failed: bool): - if self._diagnostic_accumulator and self._connection_attempt_start_time: + if self._diagnostic_accumulator and self._connection_attempt_started_monotonic is not None: current_time = int(time.time() * 1000) - elapsed = current_time - int(self._connection_attempt_start_time * 1000) - self._diagnostic_accumulator.record_stream_init(current_time, elapsed if elapsed >= 0 else 0, failed) + elapsed = int((monotonic_seconds() - self._connection_attempt_started_monotonic) * 1000) + self._diagnostic_accumulator.record_stream_init(current_time, elapsed, failed) def _create_sse_client(self) -> SSEClient: # We don't want the stream to use the same read timeout as the rest of the SDK. @@ -261,9 +261,9 @@ def _handle_error(self, error: Exception) -> bool: interrupted = self._stop_event.wait(min(delay, TIMEOUT_MAX)) - # Read after the wait, so a clock change during it cannot skew the + # Set after the wait, so the backoff delay is not counted in the # stream-init latency we report. - self._connection_attempt_start_time = time.time() + self._connection_attempt_started_monotonic = monotonic_seconds() return not interrupted # magic methods for "with" statement (used in testing) diff --git a/ldclient/impl/datasystem/async_fdv2.py b/ldclient/impl/datasystem/async_fdv2.py index 612fe4cb..d6b6c3e3 100644 --- a/ldclient/impl/datasystem/async_fdv2.py +++ b/ldclient/impl/datasystem/async_fdv2.py @@ -33,6 +33,7 @@ DataSourceStatusProviderImpl, DataStoreStatusProviderImpl, _FDv2Base, + _StateAgeTracker, fallback_condition, recovery_condition ) @@ -545,6 +546,7 @@ async def _consume_synchronizer_results( :return: the ConditionDirective describing how to proceed """ action_queue: AsyncQueue = AsyncQueue() + state_age = _StateAgeTracker() timer = AsyncRepeatingTask.at_interval( label="AsyncFDv2-sync-cond-timer", interval=10, @@ -578,9 +580,10 @@ async def reader(): if update == "check": # Check condition periodically current_status = self._data_source_status_provider.status - if check_recovery and recovery_condition(current_status): + seconds_in_state = state_age.seconds_in_state(current_status) + if check_recovery and recovery_condition(current_status, seconds_in_state): return ConditionDirective.RECOVER - if fallback_condition(current_status): + if fallback_condition(current_status, seconds_in_state): return ConditionDirective.FALLBACK continue diff --git a/ldclient/impl/datasystem/fdv2.py b/ldclient/impl/datasystem/fdv2.py index 75f8d28c..1fbcbd33 100644 --- a/ldclient/impl/datasystem/fdv2.py +++ b/ldclient/impl/datasystem/fdv2.py @@ -11,6 +11,7 @@ DataSourceStatusProviderImpl, DataStoreStatusProviderImpl, _FDv2Base, + _StateAgeTracker, fallback_condition, recovery_condition ) @@ -536,6 +537,7 @@ def _consume_synchronizer_results( :return: the ConditionDirective describing how to proceed """ action_queue: Queue = Queue() + state_age = _StateAgeTracker() timer = RepeatingTask.at_interval( label="FDv2-sync-cond-timer", interval=10, @@ -570,9 +572,10 @@ def reader(self: 'FDv2'): if update == "check": # Check condition periodically current_status = self._data_source_status_provider.status - if check_recovery and recovery_condition(current_status): + seconds_in_state = state_age.seconds_in_state(current_status) + if check_recovery and recovery_condition(current_status, seconds_in_state): return ConditionDirective.RECOVER - if fallback_condition(current_status): + if fallback_condition(current_status, seconds_in_state): return ConditionDirective.FALLBACK continue diff --git a/ldclient/impl/datasystem/fdv2_common.py b/ldclient/impl/datasystem/fdv2_common.py index b88837d8..f692919f 100644 --- a/ldclient/impl/datasystem/fdv2_common.py +++ b/ldclient/impl/datasystem/fdv2_common.py @@ -9,13 +9,13 @@ import time from copy import copy from enum import Enum -from typing import Callable, Optional +from typing import Callable, Optional, Tuple from ldclient.impl.datasystem import DataAvailability, DiagnosticAccumulator from ldclient.impl.datasystem.store import _StoreBase from ldclient.impl.listeners import Listeners from ldclient.impl.rwlock import ReadWriteLock -from ldclient.impl.util import log +from ldclient.impl.util import log, monotonic_seconds from ldclient.interfaces import ( DataSourceErrorInfo, DataSourceState, @@ -135,37 +135,69 @@ class ConditionDirective(str, Enum): """ -def fallback_condition(status: DataSourceStatus) -> bool: +class _StateAgeTracker: + """ + Measures how long the data source has been in its current state, on the + monotonic timeline, by observing statuses as they are sampled. + + The identity key is ``(state, since)``: a real transition changes ``since`` + even when the state repeats (an interrupted-valid-interrupted flap between + two samples), while a same-state error update preserves it, so a stream of + error reports cannot keep resetting the measured age. ``since`` is used + only for equality, never arithmetic, so wall-clock steps cannot distort + the measurement. + + Ages are measured from first observation, so they under-count the true + time in state by up to one sampling period. + """ + + def __init__(self) -> None: + self.__key: Optional[Tuple[DataSourceState, float]] = None + self.__observed_at = 0.0 + + def seconds_in_state(self, status: DataSourceStatus) -> float: + key = (status.state, status.since) + if key != self.__key: + self.__key = key + self.__observed_at = monotonic_seconds() + return monotonic_seconds() - self.__observed_at + + +def fallback_condition(status: DataSourceStatus, seconds_in_state: float) -> bool: """ Determine if we should fallback to the next synchronizer in the list. This applies at any position in the synchronizers list. :param status: Current data source status + :param seconds_in_state: monotonic seconds spent in ``status.state``, as + measured by a :class:`_StateAgeTracker` :return: True if fallback condition is met """ interrupted_at_runtime = ( status.state == DataSourceState.INTERRUPTED - and time.time() - status.since > 60 # 1 minute + and seconds_in_state > 60 # 1 minute ) cannot_initialize = ( status.state == DataSourceState.INITIALIZING - and time.time() - status.since > 10 # 10 seconds + and seconds_in_state > 10 # 10 seconds ) return interrupted_at_runtime or cannot_initialize -def recovery_condition(status: DataSourceStatus) -> bool: +def recovery_condition(status: DataSourceStatus, seconds_in_state: float) -> bool: """ Determine if we should try to recover to the first (preferred) synchronizer. This only applies when not already at the first synchronizer (index > 0). :param status: Current data source status + :param seconds_in_state: monotonic seconds spent in ``status.state``, as + measured by a :class:`_StateAgeTracker` :return: True if recovery condition is met """ healthy_for_too_long = ( status.state == DataSourceState.VALID - and time.time() - status.since > 300 # 5 minutes + and seconds_in_state > 300 # 5 minutes ) return healthy_for_too_long diff --git a/ldclient/impl/util.py b/ldclient/impl/util.py index b7d1713c..deec78c8 100644 --- a/ldclient/impl/util.py +++ b/ldclient/impl/util.py @@ -14,6 +14,14 @@ def current_time_millis() -> int: return int(time.time() * 1000) +def monotonic_seconds() -> float: + """ + Seconds on the monotonic clock. Use this for measuring intervals and + durations; use the wall clock for functionality tied to time outside the process. + """ + return time.monotonic() + + def timedelta_millis(delta: timedelta) -> float: return delta / timedelta(milliseconds=1) diff --git a/ldclient/testing/impl/datasource/test_async_streaming.py b/ldclient/testing/impl/datasource/test_async_streaming.py index 0c8865c3..302a2a3c 100644 --- a/ldclient/testing/impl/datasource/test_async_streaming.py +++ b/ldclient/testing/impl/datasource/test_async_streaming.py @@ -908,3 +908,28 @@ async def close(): assert DataSourceState.OFF in order assert order.index(DataSourceState.OFF) < order.index('closed') + + +def test_stream_init_duration_unaffected_by_wall_clock_step(monkeypatch): + processor = object.__new__(AsyncStreamingUpdateProcessor) + recorded = [] + + class _Recorder: + def record_stream_init(self, timestamp, duration, failed): + recorded.append((timestamp, duration, failed)) + + processor._diagnostic_accumulator = _Recorder() + processor._connection_attempt_started_monotonic = time.monotonic() + + real_time = time.time() + # A forward wall step between connect start and finish used to inflate the + # reported duration by the step size (a backward step was clamped to zero, + # silently losing the measurement). + monkeypatch.setattr(time, "time", lambda: real_time + 3600) + processor._record_stream_init(False) + + assert len(recorded) == 1 + _, duration, failed = recorded[0] + assert failed is False + assert duration >= 0 + assert duration < 1000 diff --git a/ldclient/testing/impl/datasource/test_streaming.py b/ldclient/testing/impl/datasource/test_streaming.py index 9150a36d..0fa3b803 100644 --- a/ldclient/testing/impl/datasource/test_streaming.py +++ b/ldclient/testing/impl/datasource/test_streaming.py @@ -1011,3 +1011,28 @@ def await_item(store, kind, key, expected_item): if current_item == expected_item: return assert False, 'expected %s = %s but value was still %s after %d seconds' % (key, json.dumps(expected_item), json.dumps(current_item), update_wait) + + +def test_stream_init_duration_unaffected_by_wall_clock_step(monkeypatch): + processor = object.__new__(StreamingUpdateProcessor) + recorded = [] + + class _Recorder: + def record_stream_init(self, timestamp, duration, failed): + recorded.append((timestamp, duration, failed)) + + processor._diagnostic_accumulator = _Recorder() + processor._connection_attempt_started_monotonic = time.monotonic() + + real_time = time.time() + # A forward wall step between connect start and finish used to inflate the + # reported duration by the step size (a backward step was clamped to zero, + # silently losing the measurement). + monkeypatch.setattr(time, "time", lambda: real_time + 3600) + processor._record_stream_init(False) + + assert len(recorded) == 1 + _, duration, failed = recorded[0] + assert failed is False + assert duration >= 0 + assert duration < 1000 diff --git a/ldclient/testing/impl/datasystem/test_fdv2_conditions.py b/ldclient/testing/impl/datasystem/test_fdv2_conditions.py new file mode 100644 index 00000000..8cdfdcd4 --- /dev/null +++ b/ldclient/testing/impl/datasystem/test_fdv2_conditions.py @@ -0,0 +1,122 @@ +""" +Tests for the FDv2 fallback/recovery conditions and the monotonic +time-in-state measurement in :class:`_StateAgeTracker`. +""" + +import time + +import pytest + +from ldclient.impl.datasystem import fdv2_common +from ldclient.impl.datasystem.fdv2_common import ( + _StateAgeTracker, + fallback_condition, + recovery_condition +) +from ldclient.interfaces import DataSourceState, DataSourceStatus + + +def _status(state: DataSourceState, since: float = 0.0) -> DataSourceStatus: + return DataSourceStatus(state, since or time.time(), None) + + +class _FakeMonotonic: + def __init__(self, now: float = 1000.0): + self.now = now + + def __call__(self) -> float: + return self.now + + +@pytest.fixture +def mono(monkeypatch): + clock = _FakeMonotonic() + monkeypatch.setattr(fdv2_common, "monotonic_seconds", clock) + return clock + + +@pytest.mark.parametrize( + "state,seconds,expected", + [ + (DataSourceState.INTERRUPTED, 59, False), + (DataSourceState.INTERRUPTED, 60, False), # strictly greater than + (DataSourceState.INTERRUPTED, 61, True), + (DataSourceState.INITIALIZING, 9, False), + (DataSourceState.INITIALIZING, 10, False), # strictly greater than + (DataSourceState.INITIALIZING, 11, True), + (DataSourceState.VALID, 10_000, False), + (DataSourceState.OFF, 10_000, False), + ], +) +def test_fallback_condition(state, seconds, expected): + assert fallback_condition(_status(state), seconds) is expected + + +@pytest.mark.parametrize( + "state,seconds,expected", + [ + (DataSourceState.VALID, 299, False), + (DataSourceState.VALID, 300, False), # strictly greater than + (DataSourceState.VALID, 301, True), + (DataSourceState.INTERRUPTED, 10_000, False), + (DataSourceState.INITIALIZING, 10_000, False), + (DataSourceState.OFF, 10_000, False), + ], +) +def test_recovery_condition(state, seconds, expected): + assert recovery_condition(_status(state), seconds) is expected + + +def test_tracker_measures_age_from_first_observation(mono): + tracker = _StateAgeTracker() + status = _status(DataSourceState.INTERRUPTED, since=500.0) + + assert tracker.seconds_in_state(status) == 0.0 + + mono.now = 1030.0 + assert tracker.seconds_in_state(status) == 30.0 + + +def test_tracker_keeps_age_across_same_state_error_updates(mono): + tracker = _StateAgeTracker() + # Same (state, since), new object: how a same-state error update looks. + first = _status(DataSourceState.INTERRUPTED, since=500.0) + second = _status(DataSourceState.INTERRUPTED, since=500.0) + + tracker.seconds_in_state(first) + mono.now = 1030.0 + assert tracker.seconds_in_state(second) == 30.0 + + +def test_tracker_resets_on_state_change(mono): + tracker = _StateAgeTracker() + tracker.seconds_in_state(_status(DataSourceState.INITIALIZING, since=500.0)) + + mono.now = 1030.0 + assert tracker.seconds_in_state(_status(DataSourceState.VALID, since=530.0)) == 0.0 + + mono.now = 1045.0 + assert tracker.seconds_in_state(_status(DataSourceState.VALID, since=530.0)) == 15.0 + + +def test_tracker_resets_on_a_flap_through_another_state(mono): + tracker = _StateAgeTracker() + # INTERRUPTED -> VALID -> INTERRUPTED between two samples: the state + # matches the last sample but ``since`` differs, so the age must reset. + tracker.seconds_in_state(_status(DataSourceState.INTERRUPTED, since=500.0)) + + mono.now = 1070.0 + assert tracker.seconds_in_state(_status(DataSourceState.INTERRUPTED, since=560.0)) == 0.0 + + +def test_wall_clock_step_cannot_trigger_a_spurious_fallback(monkeypatch): + tracker = _StateAgeTracker() + status = DataSourceStatus(DataSourceState.INTERRUPTED, time.time(), None) + + real_time = time.time() + # A forward wall step used to make the interrupted state look an hour old, + # which would have triggered an immediate synchronizer fallback. + monkeypatch.setattr(time, "time", lambda: real_time + 3600) + + seconds = tracker.seconds_in_state(status) + assert fallback_condition(status, seconds) is False