diff --git a/ldclient/impl/integrations/files/filedata.py b/ldclient/impl/integrations/files/filedata.py new file mode 100644 index 00000000..a4f71d6a --- /dev/null +++ b/ldclient/impl/integrations/files/filedata.py @@ -0,0 +1,745 @@ +""" +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. + Flag overrides are currently experimental and subject to change. + """ + + 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 on, has the value as its only variation, and serves that + variation as its fallthrough. + """ + return FeatureFlag({"key": key, "version": 1, "on": True, "fallthrough": {"variation": 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, which may be JSON or YAML. The parser is chosen the + way the file data sources choose it: the YAML parser reads both formats when the ``pyyaml`` + package is installed, and the standard JSON parser is used when it is not. An empty + document is valid and holds no entries. + """ + text = raw.decode("utf-8") + if have_yaml: + parsed = yaml.safe_load(text) # pyyaml correctly parses JSON too + else: + parsed = json.loads(text) + 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. An + exception from it is reported like a load failure and retried, and the result it + rejected is not remembered as the last good one. + :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. An exception + from it is logged and does not prevent the retry. + :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 + if self._skip_unchanged and not recovering and digest == self._last_good_digest: + return True + + # Nothing is remembered until the consumer has accepted the result. A result the + # consumer rejects must not become the baseline that skip-unchanged compares + # against, and must not count as a recovery. + try: + self._apply(merged) + except Exception as e: + return self._fail(e) + self._last_error_message = None + self._last_good_digest = digest + 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: + try: + self._on_error(err) + except Exception as e: + # A consumer that raises while it reports a failure must not stop the retry. + log.error("Error while reporting a flag data load failure: %s", e) + 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() + + +def _event_path(value: Any) -> Optional[str]: + """The path carried by a watchdog event as text, or None when the event has none.""" + if isinstance(value, bytes): + return value.decode("utf-8", errors="replace") + if isinstance(value, str): + return value + return None + + +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, for example because it does not exist yet or because + the notification mechanism cannot be started, is attempted again on an interval. A watched + directory that is deleted has its watch dropped and set up again once the directory exists. + When a retry sets up a watch, the callback runs once, so that changes made while the watch + was not in place are picked up. + + 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._directories: Set[str] = set() + 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))) + self._directories.add(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() + + # Guards the state below. It is never held while the observer is called, because the + # observer holds its own lock while it dispatches an event to the handler. + self._lock = threading.Lock() + self._closed = False + self._observer_running = False + # The watch on each directory that currently has one. + self._watches: Dict[str, Any] = {} + # The directories that have no watch yet. The retry schedule runs while this is not empty. + self._pending_directories: Set[str] = set(self._directories) + self._retry_task: Optional[RepeatingTask] = None + # Serializes the watch setup against close, so that a setup in progress cannot start + # anything after close has stopped the observer. + self._setup_lock = threading.Lock() + + self._set_up_watches(is_retry=False) + + def _set_up_watches(self, is_retry: bool) -> None: + """ + Starts the observer when it is not running, then sets up a watch on each pending + directory. Whatever cannot be set up stays pending and is attempted again on the retry + schedule. When a retry sets up a watch, the callback runs once, because the files may + have changed while the watch was not in place. + """ + with self._setup_lock: + with self._lock: + if self._closed: + return + running = self._observer_running + if not running: + try: + self._observer.start() + except Exception as e: + log.error("Unable to start watching files for changes: %s", e) + self._ensure_retry_scheduled() + return + with self._lock: + self._observer_running = True + + with self._lock: + pending = sorted(self._pending_directories) + armed = False + for directory in pending: + watch = self._schedule(directory) + if watch is None: + continue + with self._lock: + self._pending_directories.discard(directory) + self._watches[directory] = watch + armed = True + + with self._lock: + if len(self._pending_directories) > 0: + self._ensure_retry_scheduled_locked() + if is_retry and armed: + self._on_change() + + def _schedule(self, directory: str) -> Optional[Any]: + """Sets up a watch on a directory. Returns the watch, or None when the directory cannot be watched yet.""" + # 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 None + try: + return self._observer.schedule(self._handler, directory, recursive=False) + except Exception as e: + log.warning('Cannot watch directory "%s" for changes yet: %s', directory, e) + return None + + def _ensure_retry_scheduled(self) -> None: + with self._lock: + self._ensure_retry_scheduled_locked() + + def _ensure_retry_scheduled_locked(self) -> None: + if self._closed or self._retry_task is not None: + return + self._retry_task = RepeatingTask.at_interval("ldclient.filedata.watch-retry", _WATCH_RETRY_INTERVAL, _WATCH_RETRY_INTERVAL, self._retry_pending) + self._retry_task.start() + + def _retry_pending(self) -> None: + self._set_up_watches(is_retry=True) + with self._lock: + # The schedule ends once everything is watched. It is set up again when a directory + # becomes pending later. The check and the stop happen under the lock, so a + # directory that becomes pending at the same time either keeps this schedule or + # starts a new one. + if len(self._pending_directories) == 0 and self._observer_running and self._retry_task is not None: + self._retry_task.stop() + self._retry_task = None + + def _directory_lost(self, directory: str) -> None: + """ + Drops the watch on a directory that no longer exists and arms the retry that sets it up + again once the directory exists. + """ + with self._lock: + if self._closed: + return + stale = self._watches.pop(directory, None) + self._pending_directories.add(directory) + self._ensure_retry_scheduled_locked() + log.warning('Directory "%s" no longer exists. It is watched again when it exists.', directory) + if stale is not None: + # The observer keeps a watch whose directory is gone, and a later watch on the same + # path would reuse it, so it is removed. + try: + self._observer.unschedule(stale) + except Exception as e: + log.debug('Unable to remove the watch on directory "%s": %s', directory, e) + + def _handle_event(self, event) -> None: + event_type = getattr(event, "event_type", None) + if event_type not in _CHANGE_EVENT_TYPES: + return + src_path = _event_path(getattr(event, "src_path", None)) + if event_type in ("deleted", "moved") and getattr(event, "is_directory", False) and src_path in self._directories: + self._directory_lost(src_path) + return + for candidate in (src_path, _event_path(getattr(event, "dest_path", None))): + 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._setup_lock: + with self._lock: + self._closed = True + retry_task = self._retry_task + self._retry_task = None + running = self._observer_running + if retry_task is not None: + retry_task.stop() + self._observer.stop() + # A thread that never started cannot be joined. + if running: + 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..0037261a --- /dev/null +++ b/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py @@ -0,0 +1,434 @@ +""" +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] + events = [watchdog.events.FileModifiedEvent(path), watchdog.events.FileOpenedEvent(path)] + # Older watchdog versions report no event for a read-only close. + closed_no_write = getattr(watchdog.events, 'FileClosedNoWriteEvent', None) + if closed_no_write is not None: + events.append(closed_no_write(path)) + for event in events: + handler.on_any_event(event) + assert counter.count == len(events) + handler.on_any_event(watchdog.events.FileModifiedEvent(os.path.join(os.path.dirname(path), 'other.json'))) + assert counter.count == len(events) + + +@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..34e8aceb --- /dev/null +++ b/ldclient/testing/integrations/test_filedata.py @@ -0,0 +1,1036 @@ +import logging +import os +import threading +import time +from queue import Empty, Queue +from typing import Any, Dict, List, Set +from unittest import mock + +import pytest + +try: + import yaml +except ImportError: # pragma: no cover + yaml = None # type: ignore + +from ldclient.impl.integrations.files import filedata +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_chooses_the_parser_the_way_the_file_data_source_does(): + # The parser depends only on whether pyyaml is installed, never on the content. A tab inside + # a string is invalid JSON but valid YAML, so with pyyaml this document parses, and a parser + # chosen from the leading brace would reject it. + if not have_yaml: + pytest.skip("pyyaml is not installed") + document = parse_document(b' \n {"flagValues": {"flag1": "a\tb"}}') + assert document.flag_values == {'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 pyyaml the YAML parser reports the error, as in the file data sources. + expected = (ValueError, yaml.YAMLError) if have_yaml else (ValueError,) + with pytest.raises(expected): + 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] == os.path.abspath('/absolute/data.json') + + +def test_make_flag_with_value_is_on_and_serves_the_value_by_fallthrough(): + flag = make_flag_with_value('flag1', 'value1') + assert flag.key == 'flag1' + assert flag.version == 1 + assert flag.on is True + assert flag.off_variation is None + assert flag.fallthrough.variation == 0 + assert flag.variations == ['value1'] + assert flag.to_json_dict() == {'key': 'flag1', 'version': 1, 'on': True, 'fallthrough': {'variation': 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 True + assert result.flags['flag2'].fallthrough.variation == 0 + + +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 reloader_threads() -> Set[threading.Thread]: + return {t for t in threading.enumerate() if t.name == 'ldclient.filedata.reloader'} + + +def test_reloader_unused_spawns_no_thread(make_fixture): + before = reloader_threads() + f = make_fixture('{"flagValues": {"flag1": true}}') + assert reloader_threads() == before + f.reloader.reload_now() + assert len(reloader_threads() - 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): + before = reloader_threads() + f = make_fixture('{"flagValues": {"flag1": true}}') + f.reloader.reload_now() + f.require_applied() + started = reloader_threads() - before + assert len(started) == 1 + f.reloader.close() + worker = started.pop() + worker.join(TEST_TIMEOUT) + assert not worker.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] = [] + applied_keys: List[List[str]] = [] + + def apply(result: MergeResult) -> None: + acquired = in_apply.acquire(blocking=False) + overlaps.append(not acquired) + try: + applied_keys.append(list(result.flags.keys())) + 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) + assert len(applied_keys) > 0 + assert all(keys == ['flag1'] for keys in applied_keys) + + +def test_reloader_reports_an_apply_failure_and_retries_it(make_fixture): + # The consumer rejects every result until it is released. Each call records whether the + # result was accepted. + accept = threading.Event() + calls: Queue = Queue() + + def apply(result: MergeResult) -> None: + accepted = accept.is_set() + calls.put(accepted) + if not accepted: + raise RuntimeError("consumer failure") + + f = make_fixture('{"flagValues": {"flag1": true}}', apply=apply, retry_delay=SHORT_DELAY, skip_unchanged=True) + + # The synchronous load returns normally. The rejection is reported like a load failure. + f.reloader.reload_now() + assert take(calls) is False + assert isinstance(f.require_errored(), RuntimeError) + + # The automatic retry offers the unchanged content again until the consumer accepts it. + accept.set() + while take(calls) is False: + pass + + # Once the consumer has accepted the result, the retries stop. + with pytest.raises(Empty): + calls.get(timeout=QUIET_PERIOD) + + +def test_reloader_retries_when_the_error_callback_raises(make_fixture): + def on_error(err: Exception) -> None: + raise RuntimeError("reporter failure") + + f = make_fixture('{"flagValues"', on_error=on_error, retry_delay=SHORT_DELAY) + + # The failing report does not escape the synchronous load. + f.reloader.reload_now() + + # The automatic retry still runs and picks up the corrected file. + f.write('{"flagValues": {"flag1": true}}') + f.require_applied() + + +# --------------------------------------------------------------------------- +# 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)) + # Older watchdog versions report no event for a read-only close. + closed_no_write = getattr(watchdog.events, 'FileClosedNoWriteEvent', None) + if closed_no_write is not None: + w.watcher._handle_event(closed_no_write(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_watches_a_directory_again_after_it_is_deleted_and_recreated(tmp_path, make_watcher, monkeypatch): + monkeypatch.setattr(filedata, '_WATCH_RETRY_INTERVAL', SHORT_DELAY) + directory = os.path.join(str(tmp_path), 'later') + path = os.path.join(directory, 'data.json') + os.mkdir(directory) + write_file(path, 'a') + w = make_watcher([path]) + + # The file and then its directory are removed. + os.remove(path) + os.rmdir(directory) + w.drain() + + # The directory comes back with the file. The retry that watches it again signals a change. + os.mkdir(directory) + write_file(path, 'bb') + w.require_change() + w.drain() + + # Notifications from the recreated directory are delivered. + write_file(path, 'ccc') + w.require_change() + + +@watchdog_required +def test_watcher_retries_when_the_observer_cannot_start(tmp_path, make_watcher, monkeypatch, caplog): + import watchdog.observers + + monkeypatch.setattr(filedata, '_WATCH_RETRY_INTERVAL', SHORT_DELAY) + path = os.path.join(str(tmp_path), 'data.json') + write_file(path, 'a') + + # The observer cannot start while the watcher is constructed. + with caplog.at_level(logging.ERROR): + with mock.patch.object(watchdog.observers.Observer, 'start', side_effect=RuntimeError("cannot start thread")): + w = make_watcher([path]) + assert len([r for r in caplog.records if r.levelno == logging.ERROR]) > 0 + + # Once the observer can start, the retry sets up the watch and signals a change, and later + # edits are detected. + w.require_change() + w.drain() + write_file(path, 'bb') + w.require_change() + + +@watchdog_required +def test_watcher_close_succeeds_when_the_observer_never_started(tmp_path): + import watchdog.observers + + path = os.path.join(str(tmp_path), 'data.json') + with mock.patch.object(watchdog.observers.Observer, 'start', side_effect=RuntimeError("cannot start thread")): + watcher = Watcher([path], lambda: None) + watcher.close() + + +@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()