From 7e6424285e66604bbaea1b7849349a13f1af701f Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:28:13 -0700 Subject: [PATCH] feat: Add file data loading code with reload, retry, and polling Adds ldclient.impl.integrations.files.filedata, the file reading, parsing, and merging logic that the file-based override source is built on. A document that starts with an opening brace is parsed as JSON and any other document as YAML. Definitions are decoded into the flag and segment models while the file is read, so an invalid definition fails that load. Files are merged in the configured order with a duplicate keys handling of fail or ignore, and the result records how many entries each file supplied. The Reloader owns the reload cycle: it serializes reloads, debounces change signals with a settle window, keeps the last good data by not applying a failed load, retries a failed load after a bounded delay, reports an identical failure once, and skips an application whose file contents did not change. It can treat a configured file that does not exist as a file with no content. The Poller detects changes by comparing modification time and size on an interval, including files that appear or disappear. The Watcher uses the watchdog package, watches the directory of each file so an absent file is picked up when it appears, matches the destination of a move so a file written by rename is detected, retries a directory that does not exist yet, and reacts only to notifications that can change a file's content or presence. The existing file data sources are not changed and keep their current behavior. A new test module pins that behavior: the flagValues expansion and its evaluation reason, the version fallback, the failure messages, the FDv2 status and error kinds, and the polling and watching rules. --- ldclient/impl/integrations/files/filedata.py | 638 ++++++++++++ .../test_file_data_sources_pinned_behavior.py | 430 ++++++++ .../testing/integrations/test_filedata.py | 922 ++++++++++++++++++ 3 files changed, 1990 insertions(+) create mode 100644 ldclient/impl/integrations/files/filedata.py create mode 100644 ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py create mode 100644 ldclient/testing/integrations/test_filedata.py diff --git a/ldclient/impl/integrations/files/filedata.py b/ldclient/impl/integrations/files/filedata.py new file mode 100644 index 00000000..f2388b91 --- /dev/null +++ b/ldclient/impl/integrations/files/filedata.py @@ -0,0 +1,638 @@ +""" +File reading, parsing, and merging logic for components that load flag and segment data from +local files. The file-based override source is built on it. The existing file data sources +keep their own implementation and behavior. + +A data file is a JSON or YAML document with optional ``flags``, ``flagValues``, and +``segments`` members. ``flags`` and ``segments`` hold full definitions keyed by key. +``flagValues`` maps a flag key to a single value. It expands into a full flag definition +that returns that value for every context. +""" + +import hashlib +import json +import os +import threading +import time +from dataclasses import dataclass, field +from enum import Enum +from typing import Any, Callable, Dict, List, Optional, Set, Tuple + +from ldclient.impl.model import FeatureFlag, Segment +from ldclient.impl.repeating_task import RepeatingTask +from ldclient.impl.util import log + +have_yaml = False +try: + import yaml + + have_yaml = True +except ImportError: + pass + +have_watchdog = False +try: + import watchdog + import watchdog.events + import watchdog.observers + + have_watchdog = True +except ImportError: + pass + + +# A settle window long enough to coalesce the burst of change notifications produced by a +# single file edit, and short enough to stay responsive. +DEFAULT_DEBOUNCE_DELAY = 0.1 + +# Bounds how long a failed reload can go uncorrected when no further change notification +# arrives, for example when the failure came from reading a file mid-write. Reading a local +# file is cheap, so this can be short. +DEFAULT_RETRY_DELAY = 1.0 + +# The interval between attempts to watch a directory that could not be watched, for example +# because it does not exist yet. +_WATCH_RETRY_INTERVAL = 1.0 + +# The watchdog event types that can change a file's content or presence. Opening or reading a +# file also produces events, and a reload reads the files, so those must not count as changes. +_CHANGE_EVENT_TYPES = frozenset(["created", "modified", "moved", "deleted", "closed"]) + + +class DuplicateKeysHandling(str, Enum): + """ + Determines what happens when the same flag or segment key appears in more than one file. + """ + + FAIL = "fail" + """A duplicated key causes the load to fail.""" + + IGNORE = "ignore" + """Only the first occurrence of a duplicated key is used, in the order the files were given.""" + + +class FileDataError(Exception): + """Base class for the errors raised while loading file data.""" + + +class FileReadError(FileDataError): + """ + Indicates that one of the source files could not be read or parsed. It distinguishes a + per-file failure from a failure to merge the files' contents. + """ + + def __init__(self, path: str, message: str): + super().__init__("%s [%s]" % (message, path)) + self.path = path + + +class DuplicateKeyError(FileDataError): + """Indicates that the same key appears in more than one file and the handling is FAIL.""" + + +@dataclass +class Document: + """The parsed form of a single data file.""" + + flags: Dict[str, FeatureFlag] = field(default_factory=dict) + flag_values: Dict[str, Any] = field(default_factory=dict) + segments: Dict[str, Segment] = field(default_factory=dict) + + +@dataclass +class DocumentSummary: + """Counts the entries the merge kept from one document.""" + + flags: int = 0 + segments: int = 0 + + +@dataclass +class FileSummary: + """Describes one configured file after a load.""" + + path: str + present: bool = False + """False when the file does not exist and missing files are skipped.""" + flags: int = 0 + segments: int = 0 + + +@dataclass +class MergeResult: + """ + The merged items from one or more documents. The dictionaries preserve document order: + all of one document's items precede the next document's. + """ + + flags: Dict[str, FeatureFlag] = field(default_factory=dict) + segments: Dict[str, Segment] = field(default_factory=dict) + documents: List[DocumentSummary] = field(default_factory=list) + """For each input document in order, the number of entries the merge kept from it.""" + files: List[FileSummary] = field(default_factory=list) + """Set by the loader. Describes each configured file in order.""" + + +def abs_file_paths(paths: List[str]) -> List[str]: + """Converts each of the given paths to an absolute path.""" + return [os.path.abspath(p) for p in paths] + + +def make_flag_with_value(key: str, value: Any) -> FeatureFlag: + """ + Expands a flag-key-to-value entry into a full flag definition that returns the given value + for every context. The flag is off and serves its single variation as the off variation. + """ + return FeatureFlag({"key": key, "version": 1, "on": False, "offVariation": 0, "variations": [value]}) + + +def read_file(path: str) -> Document: + """Reads and parses a single data file, which may be in JSON or YAML format.""" + try: + with open(path, "rb") as f: + raw = f.read() + except OSError as e: + raise FileReadError(path, "unable to read file: %s" % e) + try: + return parse_document(raw) + except Exception as e: + raise FileReadError(path, "error parsing file: %s" % e) + + +def parse_document(raw: bytes) -> Document: + """ + Parses the raw content of a data file. A document whose first non-blank character is ``{`` + is parsed as JSON. Any other document is parsed as YAML, which requires the ``pyyaml`` + package. An empty document is valid and holds no entries. + """ + text = raw.decode("utf-8") + if text.lstrip().startswith("{"): + parsed = json.loads(text) + elif have_yaml: + parsed = yaml.safe_load(text) + else: + raise ValueError("the file is not a JSON object and the pyyaml package is not installed, so it cannot be parsed as YAML") + if parsed is None: + return Document() + if not isinstance(parsed, dict): + raise ValueError("the document must be an object") + + document = Document() + for key, item in _member_items(parsed, "flags").items(): + document.flags[key] = FeatureFlag(_definition(item, key, "flag")) + for key, value in _member_items(parsed, "flagValues").items(): + document.flag_values[key] = value + for key, item in _member_items(parsed, "segments").items(): + document.segments[key] = Segment(_definition(item, key, "segment")) + return document + + +def _member_items(parsed: dict, name: str) -> Dict[str, Any]: + member = parsed.get(name) + if member is None: + return {} + if not isinstance(member, dict): + raise ValueError('"%s" must be an object' % name) + for key in member: + if not isinstance(key, str): + raise ValueError('"%s" has a key that is not a string: %r' % (name, key)) + return member + + +def _definition(item: Any, key: str, kind_name: str) -> dict: + """ + Validates the shape of a flag or segment entry. The entry is stored under its map key. The + file format allows a definition to omit its own key and version, which the model requires, + so those two properties are filled in from the map key and a version of 1. + """ + if not isinstance(item, dict): + raise ValueError('%s "%s" must be an object' % (kind_name, key)) + if "key" not in item: + item["key"] = key + if "version" not in item: + item["version"] = 1 + return item + + +def merge(documents: List[Document], duplicate_keys_handling: DuplicateKeysHandling) -> MergeResult: + """ + Combines the items of the given documents in order, expanding flag-value entries into full + flag definitions and applying the given duplicate keys handling. + """ + result = MergeResult() + seen_flags: Set[str] = set() + seen_segments: Set[str] = set() + + def insert(items: Dict[str, Any], seen: Set[str], kind_name: str, key: str, item: Any) -> bool: + if key in seen: + if duplicate_keys_handling == DuplicateKeysHandling.IGNORE: + return False + raise DuplicateKeyError("%s '%s' is specified by multiple files" % (kind_name, key)) + items[key] = item + seen.add(key) + return True + + for document in documents: + summary = DocumentSummary() + for key, flag in document.flags.items(): + if insert(result.flags, seen_flags, "flag", key, flag): + summary.flags += 1 + for key, value in document.flag_values.items(): + if insert(result.flags, seen_flags, "flag", key, make_flag_with_value(key, value)): + summary.flags += 1 + for key, segment in document.segments.items(): + if insert(result.segments, seen_segments, "segment", key, segment): + summary.segments += 1 + result.documents.append(summary) + return result + + +def load_files(paths: List[str], duplicate_keys_handling: DuplicateKeysHandling, skip_missing_paths: bool = False) -> MergeResult: + """ + Reads, parses, and merges all of the given files in order. Raises :class:`FileReadError` + when a file cannot be read or parsed, :class:`DuplicateKeyError` when a key is duplicated + and the handling is FAIL. When ``skip_missing_paths`` is true, a file that does not exist + contributes no entries instead of failing the load. + """ + merged, _ = _load_files_hashed(paths, duplicate_keys_handling, skip_missing_paths) + return merged + + +def _load_files_hashed(paths: List[str], duplicate_keys_handling: DuplicateKeysHandling, skip_missing_paths: bool) -> Tuple[MergeResult, bytes]: + """ + The load behind :func:`load_files`. It also returns a digest of the raw file contents, so a + caller can tell whether a later load read the same bytes. One read feeds both the digest + and the parse, so the digest can never disagree with the content that was parsed. + """ + documents: List[Document] = [] + files: List[FileSummary] = [] + hasher = hashlib.sha256() + for path in paths: + try: + with open(path, "rb") as f: + raw = f.read() + except FileNotFoundError: + if skip_missing_paths: + log.debug("File %s does not exist; it contributes no data", path) + files.append(FileSummary(path=path)) + continue + raise FileReadError(path, "unable to read file: the file does not exist") + except OSError as e: + raise FileReadError(path, "unable to read file: %s" % e) + hasher.update(raw) + hasher.update(b"\0") + try: + documents.append(parse_document(raw)) + except Exception as e: + raise FileReadError(path, "error parsing file: %s" % e) + files.append(FileSummary(path=path, present=True)) + + merged = merge(documents, duplicate_keys_handling) + # The documents are the present files in order. Copy their counts onto the file summaries. + next_document = 0 + for summary in files: + if summary.present: + summary.flags = merged.documents[next_document].flags + summary.segments = merged.documents[next_document].segments + next_document += 1 + merged.files = files + return merged, hasher.digest() + + +class Reloader: + """ + Owns the reload cycle for a set of data files. It serializes reloads, debounces change + signals, retains the last good result on failure by not calling ``apply``, retries after + failures, and skips applications that would change nothing. + + The worker thread starts on the first :meth:`reload_now` or :meth:`trigger` call, so a + reloader that is constructed but never used does not leak a thread. + """ + + def __init__( + self, + paths: List[str], + duplicate_keys_handling: DuplicateKeysHandling, + apply: Callable[[MergeResult], None], + on_error: Optional[Callable[[Exception], None]] = None, + skip_missing_paths: bool = False, + debounce_delay: float = 0.0, + retry_delay: float = 0.0, + skip_unchanged: bool = False, + ): + """ + :param paths: the files to load, in order. The order determines which file wins under + the duplicate keys handling. + :param duplicate_keys_handling: what to do when the same key appears in more than one file + :param apply: receives each successfully merged result. Calls are serialized. + :param on_error: receives each distinct failure. Repeats of an identical failure do not + call it again until a success re-arms it. Failures are also logged here. + :param skip_missing_paths: when true, a configured file that does not exist contributes + no entries. When false, a missing file fails the load like any other read error. + :param debounce_delay: how long to wait after a trigger for further triggers to settle + before reloading. Zero reloads on every trigger. + :param retry_delay: how long to wait after a failed reload before retrying it + automatically. Zero disables the automatic retry. + :param skip_unchanged: when true, a load whose raw file contents are identical to the + last applied contents does not call ``apply`` + """ + self._paths = list(paths) + self._duplicate_keys_handling = duplicate_keys_handling + self._apply = apply + self._on_error = on_error + self._skip_missing_paths = skip_missing_paths + self._debounce_delay = debounce_delay + self._retry_delay = retry_delay + self._skip_unchanged = skip_unchanged + + # Guards the deadlines, the closed flag, and the worker start. + self._cond = threading.Condition() + self._closed = False + self._started = False + self._debounce_deadline: Optional[float] = None + self._retry_deadline: Optional[float] = None + + # Serializes the load work between reload_now and the worker thread. + self._reload_lock = threading.Lock() + self._last_good_digest: Optional[bytes] = None + self._last_error_message: Optional[str] = None + + def reload_now(self) -> None: + """ + Loads the files synchronously and applies the result or reports the failure. A failure + arms the same automatic retry as a failed triggered reload. + """ + self._ensure_started() + if not self._reload() and self._retry_delay > 0: + with self._cond: + if not self._closed and self._retry_deadline is None: + self._retry_deadline = time.monotonic() + self._retry_delay + self._cond.notify() + + def trigger(self) -> None: + """ + Signals that the files may have changed. A reload happens after the debounce delay. + Signals that arrive while a reload is pending extend the settle window. + """ + self._ensure_started() + with self._cond: + if self._closed: + return + self._debounce_deadline = time.monotonic() + max(self._debounce_delay, 0.0) + self._cond.notify() + + def close(self) -> None: + """ + Stops the reloader. It does not wait for a reload that is already in progress, so such + a reload may still deliver its result shortly after this returns. A reload that has not + yet reached its callbacks does not invoke them. + """ + with self._cond: + self._closed = True + self._cond.notify_all() + + def _ensure_started(self) -> None: + with self._cond: + if self._closed or self._started: + return + self._started = True + thread = threading.Thread(target=self._run, name="ldclient.filedata.reloader", daemon=True) + thread.start() + + def _run(self) -> None: + while True: + is_retry = False + with self._cond: + while True: + if self._closed: + return + now = time.monotonic() + deadlines = [d for d in (self._debounce_deadline, self._retry_deadline) if d is not None] + if len(deadlines) == 0: + self._cond.wait() + continue + next_deadline = min(deadlines) + if next_deadline > now: + self._cond.wait(next_deadline - now) + continue + if self._debounce_deadline is not None and self._debounce_deadline <= now: + # A triggered reload supersedes a pending retry. It either succeeds, or + # it fails and arms a fresh retry below. + self._debounce_deadline = None + self._retry_deadline = None + is_retry = False + else: + self._retry_deadline = None + is_retry = True + break + + if is_retry: + log.debug("Retrying flag data load after earlier failure") + else: + log.info("Reloading flag data after detecting a change") + try: + ok = self._reload() + except Exception as e: + log.exception("Unexpected error while reloading flag data: %s", e) + ok = True + if not ok and self._retry_delay > 0: + with self._cond: + if not self._closed: + self._retry_deadline = time.monotonic() + self._retry_delay + + def _is_closed(self) -> bool: + with self._cond: + return self._closed + + def _reload(self) -> bool: + """ + Performs one full load of all configured files. Returns whether the load succeeded, + which decides whether a retry is armed. A skipped no-op application counts as success. + """ + with self._reload_lock: + if self._is_closed(): + return True + try: + merged, digest = _load_files_hashed(self._paths, self._duplicate_keys_handling, self._skip_missing_paths) + except Exception as e: + return self._fail(e) + + # A close may have happened while the files were being read. Deliver nothing then. + if self._is_closed(): + return True + + # A success right after a failure applies even when the content is unchanged since + # the last success. The consumer heard about the failure and only an application + # tells it that things are good again. + recovering = self._last_error_message is not None + self._last_error_message = None + if self._skip_unchanged and not recovering and digest == self._last_good_digest: + return True + self._last_good_digest = digest + self._apply(merged) + return True + + def _fail(self, err: Exception) -> bool: + if self._is_closed(): + return True + # With automatic retries, a persistent failure would repeat the same log entry and the + # same callback on every attempt. Repeats of an identical failure are logged at debug + # level and do not call on_error again. + message = str(err) + if message == self._last_error_message: + log.debug("Unable to load flag data: %s", err) + return False + self._last_error_message = message + log.error("Unable to load flag data: %s", err) + if self._on_error is not None: + self._on_error(err) + return False + + +FileState = Optional[Tuple[int, int]] + + +def _observe_all(paths: List[str]) -> List[FileState]: + """ + Observes the state of each file: its modification time and size, or None when it does not + exist or cannot be examined. + """ + states: List[FileState] = [] + for path in paths: + try: + info = os.stat(path) + states.append((info.st_mtime_ns, info.st_size)) + except OSError: + states.append(None) + return states + + +class Poller: + """ + Detects changes to a set of files by examining them on a fixed interval. A change to the + modification time or the size of any file invokes the callback. A file that appears or + disappears is also a change. Use it where file system change notifications are not + available or not reliable. + + Detection is generous. The callback can run for a change that does not alter the + effective data. Feed it into a :class:`Reloader`, whose debouncing and skip-unchanged + handling absorb the excess. + """ + + def __init__(self, paths: List[str], interval: float, on_change: Callable[[], None]): + self._paths = list(paths) + self._on_change = on_change + # The files are examined once here, so only later changes invoke the callback. + self._last = _observe_all(self._paths) + self._task = RepeatingTask.at_interval("ldclient.filedata.poll", interval, interval, self._poll) + + def start(self) -> None: + """Starts the polling thread.""" + self._task.start() + + def close(self) -> None: + """ + Stops the poller. It does not wait for an examination or a callback that is in + progress, so the callback can run once more shortly after this returns. + """ + self._task.stop() + + def _poll(self) -> None: + current = _observe_all(self._paths) + changed = current != self._last + self._last = current + if changed: + self._on_change() + + +class Watcher: + """ + Detects changes to a set of files through file system change notifications, using the + ``watchdog`` package. The directory of each file is watched, so a file that does not exist + yet is picked up when it appears. A directory that cannot be watched yet, for example + because it does not exist, is retried on an interval. + + Notifications for the watched paths invoke the callback. The callback can run several times + for one logical edit, so feed it into a :class:`Reloader`. + """ + + def __init__(self, paths: List[str], on_change: Callable[[], None]): + if not have_watchdog: + raise RuntimeError("the watchdog package is required to watch files for changes") + self._on_change = on_change + self._watched_paths: Set[str] = set() + self._lock = threading.Lock() + self._pending_directories: Set[str] = set() + self._retry_task: Optional[RepeatingTask] = None + + directories: List[str] = [] + for path in paths: + absolute = os.path.abspath(path) + real_directory = os.path.realpath(os.path.dirname(absolute)) + self._watched_paths.add(os.path.join(real_directory, os.path.basename(absolute))) + if real_directory not in directories: + directories.append(real_directory) + + watcher = self + + class _Handler(watchdog.events.FileSystemEventHandler): + def on_any_event(self, event): + watcher._handle_event(event) + + self._handler = _Handler() + self._observer = watchdog.observers.Observer() + for directory in directories: + if not self._schedule(directory): + self._pending_directories.add(directory) + try: + self._observer.start() + except Exception as e: + log.error("Unable to start watching files for changes: %s", e) + if len(self._pending_directories) > 0: + self._retry_task = RepeatingTask.at_interval("ldclient.filedata.watch-retry", _WATCH_RETRY_INTERVAL, _WATCH_RETRY_INTERVAL, self._retry_pending) + self._retry_task.start() + + def _schedule(self, directory: str) -> bool: + # The observer accepts a watch on a directory that does not exist and fails later when + # it starts the watch, so the check happens here first. + if not os.path.isdir(directory): + log.warning('Cannot watch directory "%s" for changes yet because it does not exist', directory) + return False + try: + self._observer.schedule(self._handler, directory, recursive=False) + return True + except Exception as e: + log.warning('Cannot watch directory "%s" for changes yet: %s', directory, e) + return False + + def _retry_pending(self) -> None: + with self._lock: + pending = list(self._pending_directories) + for directory in pending: + if self._schedule(directory): + with self._lock: + self._pending_directories.discard(directory) + # Files may have appeared in the directory before the watch was in place. + self._on_change() + with self._lock: + if len(self._pending_directories) == 0 and self._retry_task is not None: + self._retry_task.stop() + + def _handle_event(self, event) -> None: + if getattr(event, "event_type", None) not in _CHANGE_EVENT_TYPES: + return + candidates = [getattr(event, "src_path", None), getattr(event, "dest_path", None)] + for candidate in candidates: + if isinstance(candidate, bytes): + candidate = candidate.decode("utf-8", errors="replace") + if candidate in self._watched_paths: + self._on_change() + return + + def close(self) -> None: + """Stops watching and waits briefly for the observer thread to finish.""" + with self._lock: + if self._retry_task is not None: + self._retry_task.stop() + self._observer.stop() + self._observer.join(timeout=5) diff --git a/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py b/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py new file mode 100644 index 00000000..3950afa5 --- /dev/null +++ b/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py @@ -0,0 +1,430 @@ +""" +Pins the behavior of the two existing file data sources, the FDv1 update processor and the FDv2 +initializer and synchronizer. The file-based override source is built on separate code and +behaves differently in several of these respects. These tests keep the existing sources +observably unchanged: the flagValues expansion and its evaluation reason, the version fallback, +the duplicate key and load failure messages, the FDv2 status and error kinds, and the polling +and watching rules. +""" +import logging +import os +import threading +import time +from typing import Any, Callable, List, Optional + +import pytest + +from ldclient.client import Config, Context, LDClient +from ldclient.datasystem import custom +from ldclient.feature_store import InMemoryFeatureStore +from ldclient.impl.datasource.status import DataSourceUpdateSinkImpl +from ldclient.impl.integrations.files import file_data_sourcev2 +from ldclient.impl.integrations.files.file_data_source import _FileDataSource +from ldclient.impl.integrations.files.file_data_sourcev2 import ( + _FileDataSourceV2, + _PollingAutoUpdaterV2, + _WatchdogAutoUpdaterV2 +) +from ldclient.impl.listeners import Listeners +from ldclient.integrations import Files +from ldclient.interfaces import ( + DataSourceErrorKind, + DataSourceState, + ObjectKind, + Selector +) +from ldclient.testing.mock_components import MockSelectorStore +from ldclient.testing.test_util import SpyListener +from ldclient.versioned_data_kind import FEATURES, SEGMENTS + +have_watchdog = file_data_sourcev2.have_watchdog +watchdog_required = pytest.mark.skipif(not have_watchdog, reason="watchdog is not installed") + +user = Context.create('user') + +DOCUMENT = ''' +{ + "flags": { + "flag1": { + "key": "flag1", + "on": true, + "fallthrough": {"variation": 2}, + "variations": ["fall", "off", "on"] + }, + "flag-versioned": { + "key": "flag-versioned", + "version": 7, + "on": false, + "offVariation": 0, + "variations": ["x"] + } + }, + "flagValues": { + "flag2": "value2" + }, + "segments": { + "seg1": { + "key": "seg1", + "included": ["user1"] + } + } +} +''' + +# The expansion of a flagValues entry: an on flag whose fallthrough serves its single variation. +EXPANDED_FLAG2 = {'key': 'flag2', 'version': 1, 'on': True, 'fallthrough': {'variation': 0}, 'variations': ['value2']} + + +def write_file(path: str, content: str) -> None: + with open(path, 'w') as f: + f.write(content) + + +def make_v1_source(path: str, store: InMemoryFeatureStore, listeners: Optional[Listeners] = None, **kwargs) -> _FileDataSource: + config = Config('SDK_KEY') + if listeners is not None: + config._data_source_update_sink = DataSourceUpdateSinkImpl(store, listeners, Listeners()) + factory = Files.new_data_source(paths=[path], **kwargs) + assert factory is not None + source = factory(config, store, threading.Event()) + assert isinstance(source, _FileDataSource) + return source + + +def make_v2_source(path: str, **kwargs) -> _FileDataSourceV2: + source = Files.new_data_source_v2(paths=[path], **kwargs).build(Config('SDK_KEY')) + assert isinstance(source, _FileDataSourceV2) + return source + + +def v2_changes_by_key(change_set) -> dict: + return {change.key: change for change in change_set.changes} + + +# --------------------------------------------------------------------------- +# flagValues expansion and its evaluation reason +# --------------------------------------------------------------------------- + +def test_v1_expands_flag_values_to_an_on_flag_with_fallthrough(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, DOCUMENT) + store = InMemoryFeatureStore() + source = make_v1_source(path, store) + source.start() + try: + assert store.get(FEATURES, 'flag2').to_json_dict() == EXPANDED_FLAG2 + finally: + source.stop() + + +def test_v1_flag_values_evaluate_with_a_fallthrough_reason(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, DOCUMENT) + config = Config('SDK_KEY', update_processor_class=Files.new_data_source(paths=[path]), send_events=False) + with LDClient(config) as client: + detail = client.variation_detail('flag2', user, 'default') + assert detail.value == 'value2' + assert detail.variation_index == 0 + assert detail.reason == {'kind': 'FALLTHROUGH'} + + +def test_v2_expands_flag_values_to_an_on_flag_with_fallthrough(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, DOCUMENT) + source = make_v2_source(path) + result = source.fetch(MockSelectorStore(Selector.no_selector())) + changes = v2_changes_by_key(result.value.change_set) + assert changes['flag2'].object == EXPANDED_FLAG2 + assert changes['flag2'].version == 1 + + +def test_v2_flag_values_evaluate_with_a_fallthrough_reason(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, DOCUMENT) + datasystem = custom().initializers([Files.new_data_source_v2(paths=[path])]).build() + config = Config('SDK_KEY', datasystem_config=datasystem, send_events=False) + with LDClient(config) as client: + detail = client.variation_detail('flag2', user, 'default') + assert detail.value == 'value2' + assert detail.variation_index == 0 + assert detail.reason == {'kind': 'FALLTHROUGH'} + + +# --------------------------------------------------------------------------- +# Version fallback +# --------------------------------------------------------------------------- + +def test_v1_stamps_version_1_on_entries_without_one_and_keeps_explicit_versions(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, DOCUMENT) + store = InMemoryFeatureStore() + source = make_v1_source(path, store) + source.start() + try: + assert store.get(FEATURES, 'flag1').version == 1 + assert store.get(SEGMENTS, 'seg1').version == 1 + assert store.get(FEATURES, 'flag-versioned').version == 7 + finally: + source.stop() + + +def test_v2_stamps_version_1_on_entries_without_one_and_keeps_explicit_versions(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, DOCUMENT) + result = make_v2_source(path).fetch(MockSelectorStore(Selector.no_selector())) + changes = v2_changes_by_key(result.value.change_set) + assert changes['flag1'].version == 1 + assert changes['flag1'].object['version'] == 1 + assert changes['seg1'].kind == ObjectKind.SEGMENT + assert changes['seg1'].version == 1 + assert changes['seg1'].object['version'] == 1 + assert changes['flag-versioned'].version == 7 + + +# --------------------------------------------------------------------------- +# Failure messages, statuses, and error kinds +# --------------------------------------------------------------------------- + +def test_v1_logs_the_load_failure_with_the_path_and_reports_invalid_data(tmp_path, caplog): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, '{"flagValues":{') + store = InMemoryFeatureStore() + listeners = Listeners() + spy = SpyListener() + listeners.add(spy) + source = make_v1_source(path, store, listeners) + with caplog.at_level(logging.ERROR): + source.start() + try: + assert store.initialized is False + errors = [r.getMessage() for r in caplog.records if r.levelno == logging.ERROR] + assert any(m.startswith('Unable to load flag data from "%s": ' % path) for m in errors), errors + assert len(spy.statuses) == 1 + assert spy.statuses[0].error.kind == DataSourceErrorKind.INVALID_DATA + finally: + source.stop() + + +def test_v1_missing_file_fails_the_load_with_the_same_message(tmp_path, caplog): + path = os.path.join(str(tmp_path), 'missing.json') + store = InMemoryFeatureStore() + source = make_v1_source(path, store) + with caplog.at_level(logging.ERROR): + source.start() + try: + assert source.initialized() is False + errors = [r.getMessage() for r in caplog.records if r.levelno == logging.ERROR] + assert any(m.startswith('Unable to load flag data from "%s": ' % path) for m in errors), errors + finally: + source.stop() + + +def test_v1_duplicate_key_message(tmp_path, caplog): + first = os.path.join(str(tmp_path), 'first.json') + second = os.path.join(str(tmp_path), 'second.json') + write_file(first, '{"flagValues": {"flag1": "a"}}') + write_file(second, '{"flagValues": {"flag1": "b"}}') + store = InMemoryFeatureStore() + config = Config('SDK_KEY') + source = Files.new_data_source(paths=[first, second])(config, store, threading.Event()) + with caplog.at_level(logging.ERROR): + source.start() + try: + assert store.initialized is False + errors = [r.getMessage() for r in caplog.records if r.levelno == logging.ERROR] + assert any('In features, key "flag1" was used more than once' in m for m in errors), errors + finally: + source.stop() + + +def test_v2_fetch_failure_message_names_the_path(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, '{"flagValues":{') + result = make_v2_source(path).fetch(MockSelectorStore(Selector.no_selector())) + assert result.error.startswith('Unable to load flag data from "%s": ' % path) + + +def test_v2_sync_ends_with_off_when_the_initial_load_fails(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, '{"flagValues":{') + source = make_v2_source(path, force_polling=True, poll_interval=0.1) + try: + updates = list(source.sync(MockSelectorStore(Selector.no_selector()))) + assert len(updates) == 1 + assert updates[0].state == DataSourceState.OFF + assert updates[0].change_set is None + assert updates[0].error is not None + assert updates[0].error.kind == DataSourceErrorKind.INVALID_DATA + assert updates[0].error.message.startswith('Unable to load flag data from "%s": ' % path) + finally: + source.stop() + + +def test_v2_sync_reports_invalid_data_when_a_file_becomes_malformed(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, '{"flagValues": {"flag1": true}}') + source = make_v2_source(path, force_polling=True, poll_interval=0.1) + updates: List[Any] = [] + received = threading.Event() + + def collect(): + for update in source.sync(MockSelectorStore(Selector.no_selector())): + updates.append(update) + received.set() + if len(updates) >= 2: + break + + thread = threading.Thread(target=collect, daemon=True) + thread.start() + try: + assert received.wait(5) + assert updates[0].state == DataSourceState.VALID + received.clear() + time.sleep(0.2) + write_file(path, '{"flagValues"') + assert received.wait(5) + assert updates[1].state == DataSourceState.INTERRUPTED + assert updates[1].error.kind == DataSourceErrorKind.INVALID_DATA + assert updates[1].error.message.startswith('Unable to load flag data from "%s": ' % path) + finally: + source.stop() + thread.join(5) + + +# --------------------------------------------------------------------------- +# Polling rules: modification time only, a missing file is not a change +# --------------------------------------------------------------------------- + +class Counter: + def __init__(self): + self.count = 0 + + def __call__(self): + self.count += 1 + + +def set_mtime(path: str, seconds: float) -> None: + os.utime(path, (seconds, seconds)) + + +@pytest.fixture(params=['v1', 'v2']) +def make_poller(request): + pollers = [] + + def factory(paths: List[str], on_change: Callable[[], None]): + # The interval is long, so only the direct _poll calls below examine the files. + poller: Any + if request.param == 'v1': + poller = _FileDataSource.PollingAutoUpdater(paths, on_change, 1000) + else: + poller = _PollingAutoUpdaterV2(paths, on_change, 1000) + pollers.append(poller) + return poller + + yield factory + for poller in pollers: + poller.stop() + + +def test_poller_reloads_when_the_modification_time_changes(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'aaa') + set_mtime(path, 1000000) + counter = Counter() + poller = make_poller([path], counter) + write_file(path, 'bbb') + set_mtime(path, 2000000) + poller._poll() + assert counter.count == 1 + + +def test_poller_ignores_a_size_change_with_the_same_modification_time(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'aaa') + set_mtime(path, 1000000) + counter = Counter() + poller = make_poller([path], counter) + write_file(path, 'aaaa') + set_mtime(path, 1000000) + poller._poll() + assert counter.count == 0 + + +def test_poller_ignores_a_file_that_disappears(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'aaa') + counter = Counter() + poller = make_poller([path], counter) + os.remove(path) + poller._poll() + assert counter.count == 0 + + +def test_poller_reloads_when_a_missing_file_appears(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + counter = Counter() + poller = make_poller([path], counter) + write_file(path, 'aaa') + poller._poll() + assert counter.count == 1 + + +# --------------------------------------------------------------------------- +# Watching rules: any notification on the file's path reloads, a move destination does not +# --------------------------------------------------------------------------- + +def handlers_of(observer) -> list: + return [handler for handlers in observer._handlers.values() for handler in handlers] + + +@pytest.fixture(params=['v1', 'v2']) +def make_watcher(request): + watchers = [] + + def factory(paths: List[str], on_change: Callable[[], None]): + watcher: Any + if request.param == 'v1': + watcher = _FileDataSource.WatchdogAutoUpdater(paths, on_change) + else: + watcher = _WatchdogAutoUpdaterV2(paths, on_change) + watchers.append(watcher) + return watcher + + yield factory + for watcher in watchers: + watcher.stop() + + +@watchdog_required +def test_watcher_reloads_on_any_notification_for_the_file_path(tmp_path, make_watcher): + import watchdog.events + + path = os.path.realpath(os.path.join(str(tmp_path), 'data.json')) + write_file(path, 'aaa') + counter = Counter() + watcher = make_watcher([path], counter) + handlers = handlers_of(watcher._observer) + assert len(handlers) == 1 + handler = handlers[0] + handler.on_any_event(watchdog.events.FileModifiedEvent(path)) + handler.on_any_event(watchdog.events.FileOpenedEvent(path)) + handler.on_any_event(watchdog.events.FileClosedNoWriteEvent(path)) + assert counter.count == 3 + handler.on_any_event(watchdog.events.FileModifiedEvent(os.path.join(os.path.dirname(path), 'other.json'))) + assert counter.count == 3 + + +@watchdog_required +def test_watcher_ignores_a_move_whose_destination_is_the_file_path(tmp_path, make_watcher): + import watchdog.events + + path = os.path.realpath(os.path.join(str(tmp_path), 'data.json')) + temp = os.path.realpath(os.path.join(str(tmp_path), 'data.json.tmp')) + write_file(path, 'aaa') + counter = Counter() + watcher = make_watcher([path], counter) + handler = handlers_of(watcher._observer)[0] + handler.on_any_event(watchdog.events.FileMovedEvent(temp, path)) + assert counter.count == 0 + handler.on_any_event(watchdog.events.FileMovedEvent(path, temp)) + assert counter.count == 1 diff --git a/ldclient/testing/integrations/test_filedata.py b/ldclient/testing/integrations/test_filedata.py new file mode 100644 index 00000000..44e86f3e --- /dev/null +++ b/ldclient/testing/integrations/test_filedata.py @@ -0,0 +1,922 @@ +import os +import threading +import time +from queue import Empty, Queue +from typing import Any, Dict, List + +import pytest + +from ldclient.impl.integrations.files.filedata import ( + Document, + DuplicateKeyError, + DuplicateKeysHandling, + FileReadError, + FileSummary, + MergeResult, + Poller, + Reloader, + Watcher, + abs_file_paths, + have_watchdog, + have_yaml, + load_files, + make_flag_with_value, + merge, + parse_document, + read_file +) +from ldclient.impl.model import FeatureFlag, Segment + +TEST_TIMEOUT = 5.0 + +# Timings for the reloader and poller tests. They are generous multiples of the configured +# delays so the tests stay deterministic on a loaded machine. +SHORT_DELAY = 0.05 +QUIET_PERIOD = 0.3 + + +def write_file(path: str, content: str) -> None: + with open(path, 'w') as f: + f.write(content) + + +def take(queue: Queue, timeout: float = TEST_TIMEOUT) -> Any: + try: + return queue.get(timeout=timeout) + except Empty: + pytest.fail("timed out waiting for a callback") + + +def require_quiet(queue: Queue, duration: float = QUIET_PERIOD) -> None: + try: + item = queue.get(timeout=duration) + except Empty: + return + pytest.fail("received an unexpected callback: %r" % (item,)) + + +# --------------------------------------------------------------------------- +# Parsing +# --------------------------------------------------------------------------- + +def test_parse_json_document(): + document = parse_document(b'{"flags": {"flag1": {"key": "flag1", "version": 3, "on": true}}, "flagValues": {"flag2": "value2"}, "segments": {"seg1": {"key": "seg1", "version": 2, "included": ["user1"]}}}') + assert list(document.flags.keys()) == ['flag1'] + assert isinstance(document.flags['flag1'], FeatureFlag) + assert document.flags['flag1'].version == 3 + assert document.flags['flag1'].on is True + assert document.flag_values == {'flag2': 'value2'} + assert isinstance(document.segments['seg1'], Segment) + assert document.segments['seg1'].included == {'user1'} + + +def test_parse_json_document_with_leading_whitespace_uses_json_parser(): + # Tabs inside strings are invalid JSON but valid YAML, so only the JSON parser rejects this. + with pytest.raises(ValueError): + parse_document(b' \n {"flagValues": {"flag1": "a\tb"}}') + + +def test_parse_yaml_document(): + if not have_yaml: + pytest.skip("pyyaml is not installed") + document = parse_document(b'---\nflags:\n flag1:\n key: flag1\n "on": true\nflagValues:\n flag2: value2\nsegments:\n seg1:\n key: seg1\n') + assert document.flags['flag1'].on is True + assert document.flag_values == {'flag2': 'value2'} + assert list(document.segments.keys()) == ['seg1'] + + +def test_parse_empty_document_has_no_entries(): + document = parse_document(b'') + assert document == Document() + document = parse_document(b'{}') + assert document == Document() + + +def test_parse_fills_in_missing_key_and_version(): + document = parse_document(b'{"flags": {"flag1": {"on": true}}, "segments": {"seg1": {}}}') + assert document.flags['flag1'].key == 'flag1' + assert document.flags['flag1'].version == 1 + assert document.segments['seg1'].key == 'seg1' + assert document.segments['seg1'].version == 1 + + +def test_parse_rejects_documents_that_are_not_objects(): + if have_yaml: + with pytest.raises(ValueError): + parse_document(b'- a\n- b\n') + with pytest.raises(ValueError): + parse_document(b'{"flags": ["not", "an", "object"]}') + with pytest.raises(ValueError): + parse_document(b'{"flags": {"flag1": "not an object"}}') + with pytest.raises(ValueError): + parse_document(b'{"segments": {"seg1": 3}}') + with pytest.raises(ValueError): + parse_document(b'{"flagValues": 3}') + + +def test_parse_rejects_malformed_json(): + with pytest.raises(ValueError): + parse_document(b'{"flagValues"') + + +def test_parse_validates_definition_property_types(): + with pytest.raises(ValueError): + parse_document(b'{"flags": {"flag1": {"key": "flag1", "version": "not a number"}}}') + with pytest.raises(ValueError): + parse_document(b'{"segments": {"seg1": {"key": "seg1", "version": 1, "included": "not a list"}}}') + + +def test_read_file_json_and_errors(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, '{"flagValues": {"flag1": true}}') + assert read_file(path).flag_values == {'flag1': True} + + with pytest.raises(FileReadError) as excinfo: + read_file(os.path.join(str(tmp_path), 'missing.json')) + assert excinfo.value.path == os.path.join(str(tmp_path), 'missing.json') + assert 'unable to read file' in str(excinfo.value) + + write_file(path, '{"flagValues"') + with pytest.raises(FileReadError) as excinfo: + read_file(path) + assert 'error parsing file' in str(excinfo.value) + assert path in str(excinfo.value) + + +def test_abs_file_paths(): + paths = abs_file_paths(['relative/data.json', '/absolute/data.json']) + assert paths[0] == os.path.abspath('relative/data.json') + assert paths[1] == '/absolute/data.json' + + +def test_make_flag_with_value_is_off_and_serves_the_value(): + flag = make_flag_with_value('flag1', 'value1') + assert flag.key == 'flag1' + assert flag.version == 1 + assert flag.on is False + assert flag.off_variation == 0 + assert flag.variations == ['value1'] + assert flag.to_json_dict() == {'key': 'flag1', 'version': 1, 'on': False, 'offVariation': 0, 'variations': ['value1']} + + +# --------------------------------------------------------------------------- +# Merging +# --------------------------------------------------------------------------- + +def doc(json_text: str) -> Document: + return parse_document(json_text.encode('utf-8')) + + +def test_merge_combines_documents(): + result = merge([ + doc('{"flags": {"flag1": {"key": "flag1", "version": 1}}, "segments": {"seg1": {"key": "seg1", "version": 1}}}'), + doc('{"flagValues": {"flag2": "value2"}}'), + ], DuplicateKeysHandling.FAIL) + assert list(result.flags.keys()) == ['flag1', 'flag2'] + assert list(result.segments.keys()) == ['seg1'] + assert result.flags['flag2'].variations == ['value2'] + assert result.flags['flag2'].on is False + + +def test_merge_duplicate_keys_fail(): + documents = [doc('{"flagValues": {"flag1": "a"}}'), doc('{"flags": {"flag1": {"key": "flag1", "version": 1}}}')] + with pytest.raises(DuplicateKeyError) as excinfo: + merge(documents, DuplicateKeysHandling.FAIL) + assert "flag 'flag1' is specified by multiple files" in str(excinfo.value) + + segment_documents = [doc('{"segments": {"seg1": {"key": "seg1"}}}'), doc('{"segments": {"seg1": {"key": "seg1"}}}')] + with pytest.raises(DuplicateKeyError) as excinfo: + merge(segment_documents, DuplicateKeysHandling.FAIL) + assert "segment 'seg1' is specified by multiple files" in str(excinfo.value) + + +def test_merge_duplicate_keys_within_one_document_between_flags_and_flag_values_fail(): + documents = [doc('{"flags": {"flag1": {"key": "flag1", "version": 1}}, "flagValues": {"flag1": "a"}}')] + with pytest.raises(DuplicateKeyError): + merge(documents, DuplicateKeysHandling.FAIL) + + +def test_merge_duplicate_keys_ignore_keeps_first(): + result = merge([ + doc('{"flagValues": {"flag1": "first"}}'), + doc('{"flagValues": {"flag1": "second", "flag2": "other"}}'), + ], DuplicateKeysHandling.IGNORE) + assert result.flags['flag1'].variations == ['first'] + assert list(result.flags.keys()) == ['flag1', 'flag2'] + + +def test_merge_preserves_document_order(): + result = merge([ + doc('{"flagValues": {"b": 1, "a": 2}}'), + doc('{"flagValues": {"c": 3}}'), + doc('{"flags": {"d": {"key": "d", "version": 1}}}'), + ], DuplicateKeysHandling.FAIL) + assert list(result.flags.keys()) == ['b', 'a', 'c', 'd'] + + +def test_merge_counts_entries_kept_from_each_document(): + result = merge([ + doc('{"flagValues": {"flag1": true, "flag2": false}, "segments": {"seg": {"key": "seg"}}}'), + doc('{"flagValues": {"flag2": true, "flag3": true}}'), + doc('{}'), + ], DuplicateKeysHandling.IGNORE) + assert [(d.flags, d.segments) for d in result.documents] == [(2, 1), (1, 0), (0, 0)] + + +# --------------------------------------------------------------------------- +# Loading files +# --------------------------------------------------------------------------- + +def test_load_files_fails_on_missing_path_by_default(tmp_path): + present = os.path.join(str(tmp_path), 'present.json') + missing = os.path.join(str(tmp_path), 'missing.json') + write_file(present, '{"flagValues": {"flag1": true}}') + with pytest.raises(FileReadError) as excinfo: + load_files([present, missing], DuplicateKeysHandling.FAIL) + assert excinfo.value.path == missing + + +def test_load_files_skips_missing_paths_when_configured(tmp_path): + present = os.path.join(str(tmp_path), 'present.json') + missing = os.path.join(str(tmp_path), 'missing.json') + write_file(present, '{"flagValues": {"flag1": true}}') + result = load_files([present, missing], DuplicateKeysHandling.FAIL, skip_missing_paths=True) + assert list(result.flags.keys()) == ['flag1'] + assert result.files == [FileSummary(path=present, present=True, flags=1), FileSummary(path=missing, present=False)] + + +def test_load_files_reports_parse_error_with_path(tmp_path): + path = os.path.join(str(tmp_path), 'bad.json') + write_file(path, '{"flagValues"') + with pytest.raises(FileReadError) as excinfo: + load_files([path], DuplicateKeysHandling.FAIL) + assert excinfo.value.path == path + + +def test_load_files_reports_duplicate_keys(tmp_path): + first = os.path.join(str(tmp_path), 'first.json') + second = os.path.join(str(tmp_path), 'second.json') + write_file(first, '{"flagValues": {"flag1": "first"}}') + write_file(second, '{"flagValues": {"flag1": "second"}}') + with pytest.raises(DuplicateKeyError): + load_files([first, second], DuplicateKeysHandling.FAIL) + result = load_files([first, second], DuplicateKeysHandling.IGNORE) + assert result.flags['flag1'].variations == ['first'] + assert result.files[0].flags == 1 + assert result.files[1].flags == 0 + + +# --------------------------------------------------------------------------- +# Reloader +# --------------------------------------------------------------------------- + +class ReloaderFixture: + def __init__(self, tmp_path, initial_content: str, **kwargs): + self.path = os.path.join(str(tmp_path), 'data.json') + self.applied: Queue = Queue() + self.errored: Queue = Queue() + write_file(self.path, initial_content) + options: Dict[str, Any] = dict( + paths=[self.path], + duplicate_keys_handling=DuplicateKeysHandling.FAIL, + apply=self.applied.put, + on_error=self.errored.put, + ) + options.update(kwargs) + self.reloader = Reloader(**options) + + def write(self, content: str) -> None: + write_file(self.path, content) + + def require_applied(self) -> MergeResult: + return take(self.applied) + + def require_errored(self) -> Exception: + return take(self.errored) + + def require_quiet(self, duration: float = QUIET_PERIOD) -> None: + deadline = time.time() + duration + while time.time() < deadline: + if not self.applied.empty(): + pytest.fail("unexpected apply call") + if not self.errored.empty(): + pytest.fail("unexpected on_error call") + time.sleep(0.01) + + def close(self) -> None: + self.reloader.close() + + +@pytest.fixture +def make_fixture(tmp_path): + fixtures: List[ReloaderFixture] = [] + + def factory(initial_content: str, **kwargs) -> ReloaderFixture: + fixture = ReloaderFixture(tmp_path, initial_content, **kwargs) + fixtures.append(fixture) + return fixture + + yield factory + for fixture in fixtures: + fixture.close() + + +def test_reloader_initial_load(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader.reload_now() + result = f.require_applied() + assert list(result.flags.keys()) == ['flag1'] + assert result.files == [FileSummary(path=f.path, present=True, flags=1)] + + +def test_reloader_fails_on_missing_path_by_default(make_fixture, tmp_path): + missing = os.path.join(str(tmp_path), 'missing.json') + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader = Reloader([f.path, missing], DuplicateKeysHandling.FAIL, apply=f.applied.put, on_error=f.errored.put) + f.reloader.reload_now() + err = f.require_errored() + assert isinstance(err, FileReadError) + assert err.path == missing + f.require_quiet(0.1) + + +def test_reloader_skips_missing_paths_when_configured(make_fixture, tmp_path): + second = os.path.join(str(tmp_path), 'second.json') + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader = Reloader([f.path, second], DuplicateKeysHandling.FAIL, apply=f.applied.put, on_error=f.errored.put, skip_missing_paths=True, skip_unchanged=True) + + # Step 1: one file exists and one does not. The reload succeeds with the existing file. + f.reloader.reload_now() + result = f.require_applied() + assert list(result.flags.keys()) == ['flag1'] + assert result.files == [FileSummary(path=f.path, present=True, flags=1), FileSummary(path=second, present=False)] + + # Step 2: the missing file appears. Its data is merged in. + write_file(second, '{"flagValues": {"flag2": true}}') + f.reloader.reload_now() + result = f.require_applied() + assert list(result.flags.keys()) == ['flag1', 'flag2'] + + # Step 3: the file is deleted. Its data is gone and the reload still succeeds. + os.remove(second) + f.reloader.reload_now() + result = f.require_applied() + assert list(result.flags.keys()) == ['flag1'] + f.require_quiet(0.1) + + +def test_reloader_reports_failure_and_applies_nothing(make_fixture): + f = make_fixture('{"flagValues"') + f.reloader.reload_now() + err = f.require_errored() + assert isinstance(err, FileReadError) + f.require_quiet(0.1) + + +def test_reloader_reports_merge_failure(make_fixture, tmp_path): + second = os.path.join(str(tmp_path), 'second.json') + write_file(second, '{"flagValues": {"flag1": "dup"}}') + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader = Reloader([f.path, second], DuplicateKeysHandling.FAIL, apply=f.applied.put, on_error=f.errored.put) + f.reloader.reload_now() + err = f.require_errored() + assert isinstance(err, DuplicateKeyError) + + +def test_reloader_debounce_coalesces_triggers(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}', debounce_delay=0.2) + for _ in range(20): + f.reloader.trigger() + f.require_applied() + f.require_quiet() + + +def test_reloader_without_debounce_reloads_on_each_trigger(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader.trigger() + f.require_applied() + f.reloader.trigger() + f.require_applied() + + +def test_reloader_debounce_window_is_extended_by_each_trigger(make_fixture): + # The debounce is a settle window: each trigger moves the deadline out again. A stream + # of triggers spaced closer together than the window must produce no reload while the + # stream continues, and exactly one reload after it stops. + window = 0.25 + f = make_fixture('{"flagValues": {"flag1": true}}', debounce_delay=window) + stop = time.time() + 5 * window + while time.time() < stop: + f.reloader.trigger() + time.sleep(window / 5) + assert f.applied.empty(), "a reload ran while triggers were still arriving" + f.require_applied() + f.require_quiet(2 * window) + + +def test_reloader_reports_identical_failure_only_once(make_fixture): + f = make_fixture('{"flagValues"', retry_delay=SHORT_DELAY) + f.reloader.reload_now() + f.require_errored() + # The automatic retries keep failing in the same way. They do not report again. + f.require_quiet() + + # A different failure is reported. + f.write('{"flags": {"flag1": "not an object"}}') + f.require_errored() + f.require_quiet() + + +def test_reloader_unused_spawns_no_thread(make_fixture): + before = threading.active_count() + f = make_fixture('{"flagValues": {"flag1": true}}') + assert threading.active_count() == before + f.reloader.reload_now() + assert threading.active_count() == before + 1 + f.require_applied() + + +def test_reloader_close_does_not_wait_for_in_flight_reload(make_fixture): + entered = threading.Event() + release = threading.Event() + + def blocking_apply(result: MergeResult) -> None: + entered.set() + # The callback stays parked for far longer than close is given, so a close that waits + # for the in-flight reload is detected as a failure rather than hidden by this timeout. + release.wait(TEST_TIMEOUT * 6) + + f = make_fixture('{"flagValues": {"flag1": true}}', apply=blocking_apply) + f.reloader.trigger() + assert entered.wait(TEST_TIMEOUT), "the reload did not start" + + closed = threading.Event() + + def do_close(): + f.reloader.close() + closed.set() + + threading.Thread(target=do_close, daemon=True).start() + assert closed.wait(2.0), "close blocked on an in-flight reload" + release.set() + + +def test_reloader_worker_thread_exits_after_close(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader.reload_now() + f.require_applied() + worker = [t for t in threading.enumerate() if t.name == 'ldclient.filedata.reloader'] + assert len(worker) == 1 + f.reloader.close() + worker[0].join(TEST_TIMEOUT) + assert not worker[0].is_alive() + + +def test_reloader_retries_after_failure_without_further_triggers(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}', retry_delay=SHORT_DELAY) + f.reloader.reload_now() + f.require_applied() + + f.write('{"flagValues"') + f.reloader.trigger() + f.require_errored() + + # Fix the file without triggering. Only the automatic retry can observe the fix. + f.write('{"flagValues": {"flag1": false}}') + f.require_applied() + + +def test_reloader_retries_after_failed_initial_load(make_fixture): + f = make_fixture('{"flagValues"', retry_delay=SHORT_DELAY) + f.reloader.reload_now() + f.require_errored() + f.write('{"flagValues": {"flag1": true}}') + f.require_applied() + + +def test_reloader_stops_retrying_after_success(make_fixture): + f = make_fixture('{"flagValues"', retry_delay=SHORT_DELAY, skip_unchanged=True) + f.reloader.reload_now() + f.require_errored() + + f.write('{"flagValues": {"flag1": true}}') + f.require_applied() + + # After the successful reload there are no further attempts: a changed file is not + # picked up without a trigger. + f.write('{"flagValues": {"flag1": false}}') + f.require_quiet() + + +def test_reloader_does_not_retry_when_retry_delay_is_zero(make_fixture): + f = make_fixture('{"flagValues"') + f.reloader.reload_now() + f.require_errored() + f.write('{"flagValues": {"flag1": true}}') + f.require_quiet() + + +def test_reloader_skip_unchanged(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}', skip_unchanged=True) + f.reloader.reload_now() + f.require_applied() + + f.reloader.trigger() + f.require_quiet() + + f.write('{"flagValues": {"flag1": false}}') + f.reloader.trigger() + f.require_applied() + + +def test_reloader_recovery_applies_even_when_content_unchanged(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}', skip_unchanged=True) + f.reloader.reload_now() + f.require_applied() + + # A reload fails. Consumers hear about it and may move to an interrupted state. + os.remove(f.path) + f.reloader.trigger() + f.require_errored() + + # The file comes back with byte-identical content. The success is applied despite + # skip_unchanged, because only an application tells the consumer the interruption is over. + f.write('{"flagValues": {"flag1": true}}') + f.reloader.trigger() + f.require_applied() + + # Once recovered, identical content skips again. + f.reloader.trigger() + f.require_quiet() + + +def test_reloader_applies_every_reload_when_skip_unchanged_is_off(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader.reload_now() + f.require_applied() + f.reloader.trigger() + f.require_applied() + + +def test_reloader_merges_multiple_files_in_order(make_fixture, tmp_path): + second = os.path.join(str(tmp_path), 'second.json') + write_file(second, '{"flagValues": {"flag1": "second"}}') + f = make_fixture('{"flagValues": {"flag1": "first"}}') + f.reloader = Reloader([f.path, second], DuplicateKeysHandling.IGNORE, apply=f.applied.put, on_error=f.errored.put) + f.reloader.reload_now() + result = f.require_applied() + assert result.flags['flag1'].variations == ['first'] + + +def test_reloader_does_nothing_after_close(make_fixture): + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader.reload_now() + f.require_applied() + f.reloader.close() + f.reloader.close() + f.reloader.trigger() + f.reloader.reload_now() + f.require_quiet(0.1) + + +def test_reloader_serializes_reload_now_against_worker_reloads(make_fixture): + # Concurrent reload_now calls and triggers must not interleave: every application is a + # complete merged result, and the callback never runs on two threads at once. + in_apply = threading.Lock() + overlaps: List[bool] = [] + + def apply(result: MergeResult) -> None: + acquired = in_apply.acquire(blocking=False) + overlaps.append(not acquired) + try: + assert list(result.flags.keys()) == ['flag1'] + time.sleep(0.002) + finally: + if acquired: + in_apply.release() + + f = make_fixture('{"flagValues": {"flag1": true}}', apply=apply) + threads = [threading.Thread(target=lambda: [f.reloader.reload_now() for _ in range(20)]) for _ in range(4)] + for t in threads: + t.start() + for _ in range(20): + f.reloader.trigger() + for t in threads: + t.join(TEST_TIMEOUT) + f.reloader.close() + assert not any(overlaps) + + +def test_reloader_survives_an_apply_callback_that_raises(make_fixture): + calls: Queue = Queue() + + def apply(result: MergeResult) -> None: + calls.put(result) + raise RuntimeError("consumer failure") + + f = make_fixture('{"flagValues": {"flag1": true}}', apply=apply) + f.reloader.trigger() + take(calls) + f.reloader.trigger() + take(calls) + + +# --------------------------------------------------------------------------- +# Poller +# --------------------------------------------------------------------------- + +POLL_INTERVAL = 0.05 + + +class PollerFixture: + def __init__(self, paths: List[str]): + self.changes: Queue = Queue() + self.poller = Poller(paths, POLL_INTERVAL, lambda: self.changes.put(True)) + self.poller.start() + + def require_change(self) -> None: + take(self.changes) + + def require_no_change(self, duration: float = QUIET_PERIOD) -> None: + require_quiet(self.changes, duration) + + +@pytest.fixture +def make_poller(): + pollers: List[PollerFixture] = [] + + def factory(paths: List[str]) -> PollerFixture: + fixture = PollerFixture(paths) + pollers.append(fixture) + return fixture + + yield factory + for fixture in pollers: + fixture.poller.close() + + +def set_mtime(path: str, seconds: float) -> None: + os.utime(path, (seconds, seconds)) + + +def test_poller_detects_modification(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + p = make_poller([path]) + p.require_no_change(0.15) + write_file(path, 'bb') + p.require_change() + + +def test_poller_detects_same_size_rewrite_with_new_mtime(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'aaa') + set_mtime(path, 1000000) + p = make_poller([path]) + write_file(path, 'bbb') + set_mtime(path, 2000000) + p.require_change() + + +def test_poller_detects_size_change_with_same_mtime(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'aaa') + set_mtime(path, 1000000) + p = make_poller([path]) + write_file(path, 'aaaa') + set_mtime(path, 1000000) + p.require_change() + + +def test_poller_fires_once_per_change(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + p = make_poller([path]) + write_file(path, 'bb') + p.require_change() + p.require_no_change() + + +def test_poller_detects_file_appearing(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + p = make_poller([path]) + p.require_no_change(0.15) + write_file(path, 'a') + p.require_change() + + +def test_poller_detects_file_disappearing(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + p = make_poller([path]) + os.remove(path) + p.require_change() + + +def test_poller_watches_all_files(tmp_path, make_poller): + first = os.path.join(str(tmp_path), 'first.json') + second = os.path.join(str(tmp_path), 'second.json') + write_file(first, 'a') + write_file(second, 'a') + p = make_poller([first, second]) + write_file(second, 'bb') + p.require_change() + write_file(first, 'bb') + p.require_change() + + +def test_poller_stops_on_close(tmp_path, make_poller): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + p = make_poller([path]) + p.poller.close() + time.sleep(POLL_INTERVAL * 3) + write_file(path, 'bb') + p.require_no_change() + + +def test_poller_close_returns_while_callback_blocks(tmp_path): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + entered = threading.Event() + release = threading.Event() + + def on_change(): + entered.set() + release.wait(TEST_TIMEOUT) + + poller = Poller([path], POLL_INTERVAL, on_change) + poller.start() + write_file(path, 'bb') + assert entered.wait(TEST_TIMEOUT) + closed = threading.Event() + + def do_close(): + poller.close() + closed.set() + + threading.Thread(target=do_close, daemon=True).start() + assert closed.wait(TEST_TIMEOUT), "close blocked on the callback" + release.set() + + +# --------------------------------------------------------------------------- +# Watcher +# --------------------------------------------------------------------------- + +watchdog_required = pytest.mark.skipif(not have_watchdog, reason="watchdog is not installed") + + +class WatcherFixture: + def __init__(self, paths: List[str]): + self.changes: Queue = Queue() + self.watcher = Watcher(paths, lambda: self.changes.put(True)) + + def require_change(self) -> None: + take(self.changes) + + def require_no_change(self, duration: float = QUIET_PERIOD) -> None: + require_quiet(self.changes, duration) + + def drain(self) -> None: + while True: + try: + self.changes.get(timeout=0.2) + except Empty: + return + + +@pytest.fixture +def make_watcher(): + watchers: List[WatcherFixture] = [] + + def factory(paths: List[str]) -> WatcherFixture: + fixture = WatcherFixture(paths) + watchers.append(fixture) + return fixture + + yield factory + for fixture in watchers: + fixture.watcher.close() + + +@watchdog_required +def test_watcher_detects_modification(tmp_path, make_watcher): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + w = make_watcher([path]) + w.require_no_change(0.15) + write_file(path, 'bb') + w.require_change() + + +@watchdog_required +def test_watcher_ignores_other_files_in_the_directory(tmp_path, make_watcher): + path = os.path.join(str(tmp_path), 'data.json') + other = os.path.join(str(tmp_path), 'other.json') + write_file(path, 'a') + w = make_watcher([path]) + write_file(other, 'bb') + w.require_no_change() + + +@watchdog_required +def test_watcher_detects_absent_file_appearing(tmp_path, make_watcher): + path = os.path.join(str(tmp_path), 'data.json') + w = make_watcher([path]) + w.require_no_change(0.15) + write_file(path, 'a') + w.require_change() + + +@watchdog_required +def test_watcher_detects_file_written_by_rename(tmp_path, make_watcher): + path = os.path.join(str(tmp_path), 'data.json') + temp = os.path.join(str(tmp_path), 'data.json.tmp') + write_file(path, 'a') + w = make_watcher([path]) + write_file(temp, 'bb') + w.drain() + os.replace(temp, path) + w.require_change() + + +@watchdog_required +def test_watcher_matches_the_destination_of_a_move_event(tmp_path, make_watcher): + # A file written by rename arrives as a move event whose destination is the watched path. + import watchdog.events + + path = os.path.join(str(tmp_path), 'data.json') + temp = os.path.join(str(tmp_path), 'data.json.tmp') + real_path = os.path.join(os.path.realpath(str(tmp_path)), 'data.json') + real_temp = os.path.join(os.path.realpath(str(tmp_path)), 'data.json.tmp') + w = make_watcher([path]) + w.watcher._handle_event(watchdog.events.FileMovedEvent(real_temp, real_path)) + w.require_change() + w.watcher._handle_event(watchdog.events.FileMovedEvent(real_path, real_temp)) + w.require_change() + w.watcher._handle_event(watchdog.events.FileMovedEvent(real_temp, real_temp + '.other')) + w.require_no_change() + assert temp not in w.watcher._watched_paths + + +@watchdog_required +def test_watcher_ignores_events_that_do_not_change_the_file(tmp_path, make_watcher): + # Opening and reading a watched file produces notifications too. A reload reads the files, + # so reacting to those would make every reload trigger the next one. + import watchdog.events + + path = os.path.join(str(tmp_path), 'data.json') + real_path = os.path.join(os.path.realpath(str(tmp_path)), 'data.json') + write_file(path, 'a') + w = make_watcher([path]) + w.watcher._handle_event(watchdog.events.FileOpenedEvent(real_path)) + w.watcher._handle_event(watchdog.events.FileClosedNoWriteEvent(real_path)) + w.require_no_change() + w.watcher._handle_event(watchdog.events.FileClosedEvent(real_path)) + w.require_change() + + +@watchdog_required +def test_watcher_does_not_signal_when_the_file_is_only_read(tmp_path, make_watcher): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + w = make_watcher([path]) + for _ in range(3): + with open(path, 'rb') as f: + f.read() + w.require_no_change() + + +@watchdog_required +def test_watcher_detects_file_deletion(tmp_path, make_watcher): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + w = make_watcher([path]) + os.remove(path) + w.require_change() + + +@watchdog_required +def test_watcher_picks_up_directory_that_appears_later(tmp_path, make_watcher): + directory = os.path.join(str(tmp_path), 'later') + path = os.path.join(directory, 'data.json') + w = make_watcher([path]) + w.require_no_change(0.15) + os.mkdir(directory) + write_file(path, 'a') + # The retry that watches the new directory signals a change so the file is read. + w.require_change() + w.drain() + write_file(path, 'bb') + w.require_change() + + +@watchdog_required +def test_watcher_close_stops_notifications(tmp_path, make_watcher): + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + w = make_watcher([path]) + w.watcher.close() + write_file(path, 'bb') + w.require_no_change()