Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 12 additions & 12 deletions ldclient/impl/datasource/async_streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand All @@ -86,17 +86,17 @@ 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
# clear a stale flag here rather than swallow the next real close.
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
Expand All @@ -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)
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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)
Expand Down
18 changes: 9 additions & 9 deletions ldclient/impl/datasource/streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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()
Expand All @@ -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):
Expand Down Expand Up @@ -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)
Expand All @@ -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.
Expand Down Expand Up @@ -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)
Expand Down
7 changes: 5 additions & 2 deletions ldclient/impl/datasystem/async_fdv2.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
DataSourceStatusProviderImpl,
DataStoreStatusProviderImpl,
_FDv2Base,
_StateAgeTracker,
fallback_condition,
recovery_condition
)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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

Expand Down
7 changes: 5 additions & 2 deletions ldclient/impl/datasystem/fdv2.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
DataSourceStatusProviderImpl,
DataStoreStatusProviderImpl,
_FDv2Base,
_StateAgeTracker,
fallback_condition,
recovery_condition
)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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

Expand Down
46 changes: 39 additions & 7 deletions ldclient/impl/datasystem/fdv2_common.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down
8 changes: 8 additions & 0 deletions ldclient/impl/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

could we use time.monotonic_ns() instead? (refer to https://docs.python.org/3/library/time.html#time.monotonic)



def timedelta_millis(delta: timedelta) -> float:
return delta / timedelta(milliseconds=1)

Expand Down
25 changes: 25 additions & 0 deletions ldclient/testing/impl/datasource/test_async_streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
25 changes: 25 additions & 0 deletions ldclient/testing/impl/datasource/test_streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading
Loading