From ee3234c6d885e36b60c5b872687b114a65891f5a 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 1/4] 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 | 637 ++++++++++++ .../test_file_data_sources_pinned_behavior.py | 430 ++++++++ .../testing/integrations/test_filedata.py | 933 ++++++++++++++++++ 3 files changed, 2000 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..a6536164 --- /dev/null +++ b/ldclient/impl/integrations/files/filedata.py @@ -0,0 +1,637 @@ +""" +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, 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. + :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..4cf20945 --- /dev/null +++ b/ldclient/testing/integrations/test_filedata.py @@ -0,0 +1,933 @@ +import os +import threading +import time +from queue import Empty, Queue +from typing import Any, Dict, List + +import pytest + +try: + import yaml +except ImportError: # pragma: no cover + yaml = None # type: ignore + +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] == '/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() From 27bde6f8bf6e93090d56727d57d161c503a9d98e Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Thu, 1 Oct 2026 23:09:04 +0000 Subject: [PATCH 2/4] fix: Expand value-only overrides to a flag served by fallthrough --- ldclient/impl/integrations/files/filedata.py | 5 +++-- ldclient/testing/integrations/test_filedata.py | 12 +++++++----- 2 files changed, 10 insertions(+), 7 deletions(-) diff --git a/ldclient/impl/integrations/files/filedata.py b/ldclient/impl/integrations/files/filedata.py index a6536164..9f7cb19e 100644 --- a/ldclient/impl/integrations/files/filedata.py +++ b/ldclient/impl/integrations/files/filedata.py @@ -141,9 +141,10 @@ def abs_file_paths(paths: List[str]) -> List[str]: 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. + 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": False, "offVariation": 0, "variations": [value]}) + return FeatureFlag({"key": key, "version": 1, "on": True, "fallthrough": {"variation": 0}, "variations": [value]}) def read_file(path: str) -> Document: diff --git a/ldclient/testing/integrations/test_filedata.py b/ldclient/testing/integrations/test_filedata.py index 4cf20945..4ccc967e 100644 --- a/ldclient/testing/integrations/test_filedata.py +++ b/ldclient/testing/integrations/test_filedata.py @@ -160,14 +160,15 @@ def test_abs_file_paths(): assert paths[1] == '/absolute/data.json' -def test_make_flag_with_value_is_off_and_serves_the_value(): +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 False - assert flag.off_variation == 0 + 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': False, 'offVariation': 0, 'variations': ['value1']} + assert flag.to_json_dict() == {'key': 'flag1', 'version': 1, 'on': True, 'fallthrough': {'variation': 0}, 'variations': ['value1']} # --------------------------------------------------------------------------- @@ -186,7 +187,8 @@ def test_merge_combines_documents(): 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 + assert result.flags['flag2'].on is True + assert result.flags['flag2'].fallthrough.variation == 0 def test_merge_duplicate_keys_fail(): From f748973b70e5e1bc10743ffd2dd8740f0dd3c8a8 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:32:24 +0000 Subject: [PATCH 3/4] fix: Report a rejected file data result as a failure and recover a lost watch The reloader treats an exception from apply like a load failure: it is reported, retried, and the result is not remembered as the last good one. An exception from the error callback is logged and does not stop the retry. The watcher retries when the observer cannot start, watches a directory again after it is deleted and recreated, and closes without error when the observer never started. --- ldclient/impl/integrations/files/filedata.py | 203 +++++++++++++----- .../testing/integrations/test_filedata.py | 139 ++++++++++-- 2 files changed, 275 insertions(+), 67 deletions(-) diff --git a/ldclient/impl/integrations/files/filedata.py b/ldclient/impl/integrations/files/filedata.py index 9f7cb19e..a4f71d6a 100644 --- a/ldclient/impl/integrations/files/filedata.py +++ b/ldclient/impl/integrations/files/filedata.py @@ -62,6 +62,7 @@ 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" @@ -324,9 +325,12 @@ def __init__( :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 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. + 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 @@ -465,11 +469,18 @@ def _reload(self) -> bool: # 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 + + # 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 - self._apply(merged) return True def _fail(self, err: Exception) -> bool: @@ -485,7 +496,11 @@ def _fail(self, err: Exception) -> bool: 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) + 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 @@ -545,12 +560,26 @@ def _poll(self) -> None: 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 yet, for example - because it does not exist, is retried on an interval. + 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`. @@ -561,17 +590,12 @@ def __init__(self, paths: List[str], on_change: Callable[[], None]): 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] = [] + 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))) - if real_directory not in directories: - directories.append(real_directory) + self._directories.add(real_directory) watcher = self @@ -581,58 +605,141 @@ def on_any_event(self, 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: + # 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 False + return None try: - self._observer.schedule(self._handler, directory, recursive=False) - return True + 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 False + return None - def _retry_pending(self) -> None: + def _ensure_retry_scheduled(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() + 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: - if len(self._pending_directories) == 0 and self._retry_task is not None: + # 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: - if getattr(event, "event_type", None) not in _CHANGE_EVENT_TYPES: + event_type = getattr(event, "event_type", None) + if event_type 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") + 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._lock: - if self._retry_task is not None: - self._retry_task.stop() - self._observer.stop() - self._observer.join(timeout=5) + 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_filedata.py b/ldclient/testing/integrations/test_filedata.py index 4ccc967e..34e8aceb 100644 --- a/ldclient/testing/integrations/test_filedata.py +++ b/ldclient/testing/integrations/test_filedata.py @@ -1,8 +1,10 @@ +import logging import os import threading import time from queue import Empty, Queue -from typing import Any, Dict, List +from typing import Any, Dict, List, Set +from unittest import mock import pytest @@ -11,6 +13,7 @@ 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, @@ -157,7 +160,7 @@ def test_read_file_json_and_errors(tmp_path): 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' + assert paths[1] == os.path.abspath('/absolute/data.json') def test_make_flag_with_value_is_on_and_serves_the_value_by_fallthrough(): @@ -440,12 +443,16 @@ def test_reloader_reports_identical_failure_only_once(make_fixture): 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 = threading.active_count() + before = reloader_threads() f = make_fixture('{"flagValues": {"flag1": true}}') - assert threading.active_count() == before + assert reloader_threads() == before f.reloader.reload_now() - assert threading.active_count() == before + 1 + assert len(reloader_threads() - before) == 1 f.require_applied() @@ -475,14 +482,16 @@ def do_close(): 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() - worker = [t for t in threading.enumerate() if t.name == 'ldclient.filedata.reloader'] - assert len(worker) == 1 + started = reloader_threads() - before + assert len(started) == 1 f.reloader.close() - worker[0].join(TEST_TIMEOUT) - assert not worker[0].is_alive() + worker = started.pop() + worker.join(TEST_TIMEOUT) + assert not worker.is_alive() def test_reloader_retries_after_failure_without_further_triggers(make_fixture): @@ -597,12 +606,13 @@ def test_reloader_serializes_reload_now_against_worker_reloads(make_fixture): # 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: - assert list(result.flags.keys()) == ['flag1'] + applied_keys.append(list(result.flags.keys())) time.sleep(0.002) finally: if acquired: @@ -618,20 +628,51 @@ def apply(result: MergeResult) -> None: 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_survives_an_apply_callback_that_raises(make_fixture): +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: - calls.put(result) - raise RuntimeError("consumer failure") + accepted = accept.is_set() + calls.put(accepted) + if not accepted: + raise RuntimeError("consumer failure") - f = make_fixture('{"flagValues": {"flag1": true}}', apply=apply) - f.reloader.trigger() - take(calls) - f.reloader.trigger() - take(calls) + 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() # --------------------------------------------------------------------------- @@ -884,7 +925,10 @@ def test_watcher_ignores_events_that_do_not_change_the_file(tmp_path, make_watch 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)) + # 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() @@ -925,6 +969,63 @@ def test_watcher_picks_up_directory_that_appears_later(tmp_path, make_watcher): 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') From b309a878d05e81a015ca97c915c7a08e4b0d2a4a Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:32:24 +0000 Subject: [PATCH 4/4] test: Build the read-only close event only when the watchdog version has it --- .../test_file_data_sources_pinned_behavior.py | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py b/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py index 3950afa5..0037261a 100644 --- a/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py +++ b/ldclient/testing/integrations/test_file_data_sources_pinned_behavior.py @@ -406,12 +406,16 @@ def test_watcher_reloads_on_any_notification_for_the_file_path(tmp_path, make_wa 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 + 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 == 3 + assert counter.count == len(events) @watchdog_required