From 2074b74a61b2d52ac81ecf93e002af9abda31dd6 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Fri, 25 Sep 2026 15:29:38 -0700 Subject: [PATCH] feat: Add the file-based override source Adds ldclient.integrations.overrides.FileOverrideSourceBuilder, the file-based override source described by the OVERRIDE specification. The source reads one or more JSON or YAML files in the file data source document format, with optional flags, flagValues, and segments members, and supplies each successful load to the SDK's override sink as a full snapshot. Files are combined in the configured order. The duplicate keys handling is fail by default, which rejects the reload and keeps the previously loaded overrides, or ignore, which keeps the first file's entry. A configured file that does not exist contributes no overrides, so a file can be created later and deleting a file removes its overrides. A file that exists but cannot be read or parsed fails that reload, the last good overrides stay in effect, the failure is logged, and the load is retried after a bounded delay and on the next detected change. Change detection is one of two modes: polling, the default, examines the files once per second by default with a one second minimum, and watching reacts to file system notifications through the watchdog package. Watching without the watchdog package and a builder with no paths are construction errors. The initial load completes during client construction. Every applied change is logged at Info level with the overrides in effect and what each configured file supplied. --- docs/api-integrations.rst | 10 + .../impl/integrations/overrides/__init__.py | 0 .../overrides/file_override_source.py | 148 ++++++ ldclient/integrations/overrides.py | 152 ++++++ .../integrations/test_file_override_source.py | 441 ++++++++++++++++++ 5 files changed, 751 insertions(+) create mode 100644 ldclient/impl/integrations/overrides/__init__.py create mode 100644 ldclient/impl/integrations/overrides/file_override_source.py create mode 100644 ldclient/integrations/overrides.py create mode 100644 ldclient/testing/integrations/test_file_override_source.py diff --git a/docs/api-integrations.rst b/docs/api-integrations.rst index 8d8146ff..594d7ed8 100644 --- a/docs/api-integrations.rst +++ b/docs/api-integrations.rst @@ -8,3 +8,13 @@ ldclient.integrations module :members: :special-members: __init__ :show-inheritance: + +ldclient.integrations.overrides module +-------------------------------------- + +The entry point for this feature is :class:`ldclient.integrations.overrides.FileOverrideSourceBuilder`. + +.. automodule:: ldclient.integrations.overrides + :members: + :special-members: __init__ + :show-inheritance: diff --git a/ldclient/impl/integrations/overrides/__init__.py b/ldclient/impl/integrations/overrides/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/ldclient/impl/integrations/overrides/file_override_source.py b/ldclient/impl/integrations/overrides/file_override_source.py new file mode 100644 index 00000000..f9367701 --- /dev/null +++ b/ldclient/impl/integrations/overrides/file_override_source.py @@ -0,0 +1,148 @@ +""" +The file-based override source. It reads flag and segment overrides from one or more local +files and reloads them when the files change. +""" + +import threading +from enum import Enum +from typing import List, Optional, Union + +from ldclient.impl.integrations.files.filedata import ( + DEFAULT_DEBOUNCE_DELAY, + DEFAULT_RETRY_DELAY, + DuplicateKeysHandling, + FileSummary, + MergeResult, + Poller, + Reloader, + Watcher +) +from ldclient.impl.util import log +from ldclient.interfaces import OverrideSink, OverrideSource + + +class ChangeDetection(str, Enum): + """ + Selects how the file-based override source learns that a file changed. The two modes are + alternatives. Flag overrides are currently experimental and subject to change. + """ + + POLLING = "polling" + """ + The source examines the files on a fixed interval and reloads when the modification time + or the size of a file changes. Polling works on every file system, including network + mounts and directories whose contents are swapped through symbolic links, as Kubernetes + does for mounted ConfigMaps. It is the default. + """ + + WATCHING = "watching" + """ + The source reloads in response to file system change notifications, using the ``watchdog`` + package. It reacts faster than polling. It depends on notifications, which some file + systems do not deliver reliably. + """ + + +class _FileOverrideSource(OverrideSource): + """ + Reads overrides from the configured files and supplies each successful load to the sink as + a full snapshot. A configured file that does not exist contributes no overrides. A file + that cannot be read or parsed fails that load, the last good snapshot stays in effect, and + the load is retried. + """ + + def __init__(self, paths: List[str], duplicate_keys_handling: DuplicateKeysHandling, change_detection: ChangeDetection, poll_interval: float): + self._paths = paths + self._duplicate_keys_handling = duplicate_keys_handling + self._change_detection = change_detection + self._poll_interval = poll_interval + self._reloader: Optional[Reloader] = None + self._change_detector: Optional[Union[Poller, Watcher]] = None + self._lock = threading.Lock() + self._closed = False + + def start(self, sink: OverrideSink) -> None: + def apply(merged: MergeResult) -> None: + sink.set_overrides(merged.flags, merged.segments) + _log_overrides_in_effect(merged) + + reloader = Reloader( + self._paths, + self._duplicate_keys_handling, + apply=apply, + skip_missing_paths=True, + debounce_delay=DEFAULT_DEBOUNCE_DELAY, + retry_delay=DEFAULT_RETRY_DELAY, + skip_unchanged=True, + ) + with self._lock: + if self._closed: + return + self._reloader = reloader + + # The initial load runs synchronously, so overrides present in the files are in effect + # by the time the client constructor returns. A file that does not exist yet + # contributes no overrides. A file that cannot be read or parsed is not fatal: the + # client runs with no overrides, the failure is logged, and the retry recovers once + # the file is readable. + reloader.reload_now() + + with self._lock: + if self._closed: + return + if self._change_detection == ChangeDetection.WATCHING: + self._change_detector = Watcher(self._paths, reloader.trigger) + else: + poller = Poller(self._paths, self._poll_interval, reloader.trigger) + poller.start() + self._change_detector = poller + + def close(self) -> None: + with self._lock: + if self._closed: + return + self._closed = True + change_detector = self._change_detector + reloader = self._reloader + self._change_detector = None + self._reloader = None + if change_detector is not None: + change_detector.close() + if reloader is not None: + reloader.close() + + +def _log_overrides_in_effect(merged: MergeResult) -> None: + """ + Reports the overrides now in effect and what each configured file supplied. The reloader + applies a snapshot only when the content changed, so this logs each change once. + """ + details = "; ".join(_file_summary_text(summary) for summary in merged.files) + if len(merged.flags) == 0 and len(merged.segments) == 0: + log.info("Flag overrides: none in effect (%s)", details) + return + log.info("Flag overrides in effect: %s (%s)", _counts_text(len(merged.flags), len(merged.segments)), details) + + +def _file_summary_text(summary: FileSummary) -> str: + if not summary.present: + return "%s: absent" % summary.path + if summary.flags == 0 and summary.segments == 0: + return "%s: no entries" % summary.path + return "%s: %s" % (summary.path, _counts_text(summary.flags, summary.segments)) + + +def _counts_text(flags: int, segments: int) -> str: + """Formats flag and segment counts, for example "2 flags, 1 segment".""" + parts = [] + if flags > 0: + parts.append(_pluralize(flags, "flag")) + if segments > 0: + parts.append(_pluralize(segments, "segment")) + return ", ".join(parts) + + +def _pluralize(count: int, noun: str) -> str: + if count == 1: + return "1 %s" % noun + return "%d %ss" % (count, noun) diff --git a/ldclient/integrations/overrides.py b/ldclient/integrations/overrides.py new file mode 100644 index 00000000..6933253b --- /dev/null +++ b/ldclient/integrations/overrides.py @@ -0,0 +1,152 @@ +""" +Override sources for the SDK's flag override capability. Flag overrides are currently +experimental and subject to change. + +Overrides are flag and segment definitions that take precedence over data received from +LaunchDarkly at evaluation time, on a per-key basis. They exist for resilience during an +incident. An operator can force one or more flags to a known state on a running application, +whether or not the application can reach LaunchDarkly. The override stays in effect until the +operator removes it. Flags not present in the override data are completely unaffected. + +This module currently provides one source: :class:`FileOverrideSourceBuilder`, which reads +overrides from local files and reloads them as the files change. Configure it with the data +system builder: +:: + + from ldclient import Config, datasystem + from ldclient.integrations.overrides import FileOverrideSourceBuilder + + source = FileOverrideSourceBuilder(['/etc/launchdarkly/overrides.json']) + config = Config(sdk_key, datasystem_config=datasystem.default().overrides(source).build()) + +An evaluation that an override affects is marked. The marking is direct or transitive. It +applies when the evaluated flag, a prerequisite at any depth, or a segment read during the +evaluation came from the override layer. The evaluation reason's ``overrideAffected`` property +reports the marking. Marked evaluations appear in analytics summary events only, under separate +counters, so LaunchDarkly can distinguish them. They produce no individual evaluation events. +""" + +from typing import List, Optional, Union + +from ldclient.config import DataSourceBuilderConfig, OverrideSourceBuilder +from ldclient.impl.integrations.files.filedata import ( + DuplicateKeysHandling, + abs_file_paths, + have_watchdog +) +from ldclient.impl.integrations.overrides.file_override_source import ( + ChangeDetection, + _FileOverrideSource +) +from ldclient.impl.util import log +from ldclient.interfaces import OverrideSource + + +class FileOverrideSourceBuilder(OverrideSourceBuilder): + """ + A builder for a file-based override source. Flag overrides are currently experimental and + subject to change. + + The source reads flag and segment overrides from one or more local files and reloads them + as the files change. The files use the same document format as the file data source: + each file is a JSON or YAML document with optional ``flags``, ``flagValues``, and + ``segments`` members. ``flagValues`` entries are expanded into full flag definitions that + return the given value for every context. YAML requires the ``pyyaml`` package. + + When multiple files are configured, their entries are combined in the configured order. + The order determines which file wins under the duplicate keys handling. A reload replaces + the entire override set, so removing an entry from the files removes the override. A + configured file that does not exist contributes no overrides: deleting a file removes its + overrides, and deleting every file removes them all. A file that exists but cannot be read + or parsed makes that whole reload fail. The previously loaded overrides stay in effect, the + source logs the failure, retries after a short delay, and recovers on its own once the file + is readable again. + + Whenever the set of overrides in effect changes, including at startup, the source logs the + overrides in effect and what each configured file supplied, at Info level. + + By default the source polls the files for changes once per second. See + :meth:`change_detection` and :meth:`poll_interval`. + """ + + DEFAULT_POLL_INTERVAL = 1.0 + """ + The interval, in seconds, at which the source examines the files for changes in polling + mode when no interval was specified. Because the source reads local files rather than + contacting a service, a short interval keeps an override responsive during an incident at + negligible cost. + """ + + MINIMUM_POLL_INTERVAL = 1.0 + """ + The shortest allowed polling interval, in seconds. A configured interval below this is + raised to it. The minimum exists only to prevent a pathological tight loop over the file + system. + """ + + def __init__(self, paths: Union[str, List[str]]): + """ + :param paths: the files to load overrides from, as a single path or a list of paths. + The order is significant: it determines which file wins under the duplicate keys + handling when the same key appears in more than one file. Relative paths are + resolved against the current working directory when the source is built. + """ + self.__paths: List[str] = [paths] if isinstance(paths, str) else list(paths) + self.__duplicate_keys_handling = DuplicateKeysHandling.FAIL + self.__change_detection = ChangeDetection.POLLING + self.__poll_interval = self.DEFAULT_POLL_INTERVAL + + def duplicate_keys_handling(self, handling: Union[DuplicateKeysHandling, str]) -> 'FileOverrideSourceBuilder': + """ + Specifies how to handle the same key appearing in more than one file. The default is + :attr:`DuplicateKeysHandling.FAIL`, which treats the reload as failed and keeps the + previously loaded overrides. :attr:`DuplicateKeysHandling.IGNORE` keeps the entry from + the first configured file that defines the key and discards the others. + + :param handling: the handling, as the enum or its string value (``"fail"`` or ``"ignore"``) + """ + self.__duplicate_keys_handling = DuplicateKeysHandling(handling) + return self + + def change_detection(self, mode: Union[ChangeDetection, str]) -> 'FileOverrideSourceBuilder': + """ + Selects how the source detects file changes. The default is + :attr:`ChangeDetection.POLLING`. The two modes are alternatives, so setting one replaces + the other. :attr:`ChangeDetection.WATCHING` requires the ``watchdog`` package. + + :param mode: the mode, as the enum or its string value (``"polling"`` or ``"watching"``) + """ + self.__change_detection = ChangeDetection(mode) + return self + + def poll_interval(self, seconds: float) -> 'FileOverrideSourceBuilder': + """ + Sets the interval between examinations of the files in polling mode. Watching mode + ignores it. The default is :attr:`DEFAULT_POLL_INTERVAL`. An interval below + :attr:`MINIMUM_POLL_INTERVAL` is raised to the minimum. + + :param seconds: the interval in seconds + """ + self.__poll_interval = seconds + return self + + def build(self, config: DataSourceBuilderConfig) -> OverrideSource: # pylint: disable=unused-argument + """ + Builds the override source. This is called internally by the SDK. It raises + ``ValueError`` when no file paths were specified or when watching mode was selected + without the ``watchdog`` package. + """ + if len(self.__paths) == 0: + raise ValueError("no file paths were specified for the file-based override source") + if self.__change_detection == ChangeDetection.WATCHING and not have_watchdog: + raise ValueError("the file-based override source cannot watch files for changes because the watchdog package is not installed; install it or use polling") + + poll_interval = self.__poll_interval + if self.__change_detection == ChangeDetection.POLLING and poll_interval < self.MINIMUM_POLL_INTERVAL: + log.warning("Poll interval %s is below the minimum for the file-based override source; using %s", poll_interval, self.MINIMUM_POLL_INTERVAL) + poll_interval = self.MINIMUM_POLL_INTERVAL + + return _FileOverrideSource(abs_file_paths(self.__paths), self.__duplicate_keys_handling, self.__change_detection, poll_interval) + + +__all__ = ['ChangeDetection', 'DuplicateKeysHandling', 'FileOverrideSourceBuilder'] diff --git a/ldclient/testing/integrations/test_file_override_source.py b/ldclient/testing/integrations/test_file_override_source.py new file mode 100644 index 00000000..ab940a85 --- /dev/null +++ b/ldclient/testing/integrations/test_file_override_source.py @@ -0,0 +1,441 @@ +""" +Tests for the file-based override source: its builder, initial load, multi-file merging, +absent files, change detection in both modes, failure retention, and its Info log. +""" +import logging +import os +import time +from queue import Empty, Queue +from typing import Any, Dict, List, Optional, Tuple + +import pytest + +from ldclient.client import Config, Context, LDClient +from ldclient.config import Config as SDKConfig +from ldclient.datasystem import custom +from ldclient.impl.integrations.files import filedata +from ldclient.impl.integrations.overrides.file_override_source import ( + _FileOverrideSource +) +from ldclient.integrations.overrides import ( + ChangeDetection, + DuplicateKeysHandling, + FileOverrideSourceBuilder +) +from ldclient.testing.mock_components import HangingSynchronizer +from ldclient.testing.stub_util import MockEventProcessor + +TEST_TIMEOUT = 10.0 +QUIET_PERIOD = 0.3 + +watchdog_required = pytest.mark.skipif(not filedata.have_watchdog, reason="watchdog is not installed") +yaml_required = pytest.mark.skipif(not filedata.have_yaml, reason="pyyaml is not installed") + +Snapshot = Tuple[Dict[str, Any], Dict[str, Any]] + + +class CapturingSink: + """Records every override snapshot it receives.""" + + def __init__(self): + self.snapshots: Queue = Queue() + + def set_overrides(self, flags, segments) -> None: + self.snapshots.put((dict(flags), dict(segments))) + + def require_snapshot(self, timeout: float = TEST_TIMEOUT) -> Snapshot: + try: + return self.snapshots.get(timeout=timeout) + except Empty: + pytest.fail("timed out waiting for an override snapshot") + + def require_no_snapshot(self, duration: float = QUIET_PERIOD) -> None: + try: + snapshot = self.snapshots.get(timeout=duration) + except Empty: + return + pytest.fail("received an unexpected override snapshot: %r" % (snapshot,)) + + +def write_file(path: str, content: str) -> None: + with open(path, 'w') as f: + f.write(content) + + +def flag_values(snapshot: Snapshot) -> Dict[str, Any]: + """The single variation of each flag in a snapshot, keyed by flag key.""" + return {key: flag.variations[0] for key, flag in snapshot[0].items()} + + +@pytest.fixture +def sources(): + started: List[_FileOverrideSource] = [] + yield started + for source in started: + source.close() + + +def build_source(sources, paths, configure=None) -> Tuple[_FileOverrideSource, CapturingSink]: + builder = FileOverrideSourceBuilder(paths) + if configure is not None: + configure(builder) + source = builder.build(SDKConfig('SDK_KEY')) + assert isinstance(source, _FileOverrideSource) + sink = CapturingSink() + source.start(sink) + sources.append(source) + return source, sink + + +def info_lines(caplog) -> List[str]: + return [r.getMessage() for r in caplog.records if r.levelno == logging.INFO] + + +def require_info_line(caplog, expected: str) -> None: + """Waits for an Info line equal to the expected text. The source logs after it hands the snapshot to the sink.""" + deadline = time.time() + TEST_TIMEOUT + while True: + if expected in info_lines(caplog): + return + if time.time() > deadline: + pytest.fail("timed out waiting for the Info line %r; Info output: %r" % (expected, info_lines(caplog))) + time.sleep(0.01) + + +# --------------------------------------------------------------------------- +# Builder +# --------------------------------------------------------------------------- + +def test_builder_requires_paths(): + with pytest.raises(ValueError): + FileOverrideSourceBuilder([]).build(SDKConfig('SDK_KEY')) + + +def test_builder_accepts_a_single_path_string(tmp_path): + path = os.path.join(str(tmp_path), 'overrides.json') + source = FileOverrideSourceBuilder(path).build(SDKConfig('SDK_KEY')) + assert isinstance(source, _FileOverrideSource) + assert source._paths == [path] + + +def test_builder_resolves_relative_paths(): + source = FileOverrideSourceBuilder(['relative/overrides.json']).build(SDKConfig('SDK_KEY')) + assert isinstance(source, _FileOverrideSource) + assert source._paths == [os.path.abspath('relative/overrides.json')] + + +def test_builder_polls_by_default(tmp_path): + path = os.path.join(str(tmp_path), 'overrides.json') + source = FileOverrideSourceBuilder([path]).build(SDKConfig('SDK_KEY')) + assert isinstance(source, _FileOverrideSource) + assert source._change_detection == ChangeDetection.POLLING + assert source._poll_interval == FileOverrideSourceBuilder.DEFAULT_POLL_INTERVAL + assert source._duplicate_keys_handling == DuplicateKeysHandling.FAIL + + +def test_builder_rejects_unknown_change_detection(): + with pytest.raises(ValueError) as excinfo: + FileOverrideSourceBuilder(['a']).change_detection('notify') + assert 'notify' in str(excinfo.value) + + +def test_builder_rejects_unknown_duplicate_keys_handling(): + with pytest.raises(ValueError): + FileOverrideSourceBuilder(['a']).duplicate_keys_handling('merge') + + +def test_builder_accepts_string_values(): + builder = FileOverrideSourceBuilder(['a']).change_detection('polling').duplicate_keys_handling('ignore') + source = builder.build(SDKConfig('SDK_KEY')) + assert isinstance(source, _FileOverrideSource) + assert source._change_detection == ChangeDetection.POLLING + assert source._duplicate_keys_handling == DuplicateKeysHandling.IGNORE + + +def test_builder_clamps_poll_interval_to_the_minimum(caplog): + with caplog.at_level(logging.WARNING): + source = FileOverrideSourceBuilder(['a']).poll_interval(0.001).build(SDKConfig('SDK_KEY')) + assert isinstance(source, _FileOverrideSource) + assert source._poll_interval == FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL + assert any('below the minimum' in r.getMessage() for r in caplog.records) + + source = FileOverrideSourceBuilder(['a']).poll_interval(2.5).build(SDKConfig('SDK_KEY')) + assert isinstance(source, _FileOverrideSource) + assert source._poll_interval == 2.5 + + +def test_builder_rejects_watching_without_watchdog(monkeypatch): + import ldclient.integrations.overrides as module + monkeypatch.setattr(module, 'have_watchdog', False) + with pytest.raises(ValueError) as excinfo: + FileOverrideSourceBuilder(['a']).change_detection(ChangeDetection.WATCHING).build(SDKConfig('SDK_KEY')) + assert 'watchdog' in str(excinfo.value) + + +# --------------------------------------------------------------------------- +# Loading +# --------------------------------------------------------------------------- + +def test_source_loads_initial_data_synchronously(tmp_path, sources): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{"flagValues": {"flag1": true, "flag2": "x"}, "segments": {"seg": {"key": "seg"}}}') + source, sink = build_source(sources, [path]) + # The snapshot was delivered before start returned. + snapshot = sink.snapshots.get_nowait() + assert flag_values(snapshot) == {'flag1': True, 'flag2': 'x'} + assert list(snapshot[1].keys()) == ['seg'] + assert snapshot[0]['flag1'].is_override is False, "the source supplies unmarked definitions; the SDK marks them" + + +@yaml_required +def test_source_loads_yaml(tmp_path, sources): + path = os.path.join(str(tmp_path), 'overrides.yaml') + write_file(path, 'flagValues:\n yaml-flag: "override-value"\n') + _, sink = build_source(sources, [path]) + assert flag_values(sink.require_snapshot()) == {'yaml-flag': 'override-value'} + + +def test_source_merges_files_in_configured_order(tmp_path, sources): + first = os.path.join(str(tmp_path), 'first.json') + second = os.path.join(str(tmp_path), 'second.json') + write_file(first, '{"flagValues": {"shared": "first", "only-first": 1}}') + write_file(second, '{"flagValues": {"shared": "second", "only-second": 2}}') + _, sink = build_source(sources, [first, second], lambda b: b.duplicate_keys_handling(DuplicateKeysHandling.IGNORE)) + assert flag_values(sink.require_snapshot()) == {'shared': 'first', 'only-first': 1, 'only-second': 2} + + +def test_source_duplicate_keys_fail_by_default(tmp_path, sources, caplog): + first = os.path.join(str(tmp_path), 'first.json') + second = os.path.join(str(tmp_path), 'second.json') + write_file(first, '{"flagValues": {"shared": "first"}}') + write_file(second, '{"flagValues": {"shared": "second"}}') + with caplog.at_level(logging.ERROR): + _, sink = build_source(sources, [first, second]) + # The load failed, so no snapshot was supplied and the failure was logged. + sink.require_no_snapshot() + assert any("is specified by multiple files" in r.getMessage() for r in caplog.records) + + +def test_source_starts_with_missing_file(tmp_path, sources): + path = os.path.join(str(tmp_path), 'not-yet.json') + _, sink = build_source(sources, [path], lambda b: b.poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL)) + # A missing file contributes no overrides. The initial snapshot is empty. + assert sink.require_snapshot() == ({}, {}) + # Once the file appears, the change detection picks it up. + write_file(path, '{"flagValues": {"flag1": true}}') + assert flag_values(sink.require_snapshot()) == {'flag1': True} + + +def test_source_missing_file_contributes_no_entries(tmp_path, sources): + first = os.path.join(str(tmp_path), 'first.json') + second = os.path.join(str(tmp_path), 'second.json') + write_file(first, '{"flagValues": {"from-first": true}}') + + # Step 1: one configured file exists and one does not. The existing file applies. + _, sink = build_source(sources, [first, second], lambda b: b.poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL)) + assert flag_values(sink.require_snapshot()) == {'from-first': True} + + # Step 2: the second file appears. Both apply. + write_file(second, '{"flagValues": {"from-second": true}}') + assert flag_values(sink.require_snapshot()) == {'from-first': True, 'from-second': True} + + # Step 3: the second file is deleted. Its overrides are removed. + os.remove(second) + assert flag_values(sink.require_snapshot()) == {'from-first': True} + + # Step 4: the last file is deleted. The layer is cleared. + os.remove(first) + assert sink.require_snapshot() == ({}, {}) + + +def test_source_logs_overrides_in_effect_on_each_change(tmp_path, sources, caplog): + first = os.path.join(str(tmp_path), 'first.json') + second = os.path.join(str(tmp_path), 'second.json') + write_file(first, '{"flagValues": {"flag1": true, "flag2": false}, "segments": {"seg": {"key": "seg"}}}') + + with caplog.at_level(logging.INFO): + # Step 1: at startup, one file supplies entries and the other is absent. + _, sink = build_source(sources, [first, second], lambda b: b.poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL)) + sink.require_snapshot() + require_info_line(caplog, "Flag overrides in effect: 2 flags, 1 segment (%s: 2 flags, 1 segment; %s: absent)" % (first, second)) + + # Step 2: the absent file appears with one entry. + write_file(second, '{"flagValues": {"flag3": true}}') + sink.require_snapshot() + require_info_line(caplog, "Flag overrides in effect: 3 flags, 1 segment (%s: 2 flags, 1 segment; %s: 1 flag)" % (first, second)) + + # Step 3: both files are deleted. Nothing is in effect. + os.remove(first) + os.remove(second) + sink.require_snapshot() + require_info_line(caplog, "Flag overrides: none in effect (%s: absent; %s: absent)" % (first, second)) + + +def test_source_logs_none_in_effect_at_startup_without_files(tmp_path, sources, caplog): + path = os.path.join(str(tmp_path), 'overrides.json') + with caplog.at_level(logging.INFO): + _, sink = build_source(sources, [path]) + sink.require_snapshot() + require_info_line(caplog, "Flag overrides: none in effect (%s: absent)" % path) + + +def test_source_logs_a_present_file_with_no_entries(tmp_path, sources, caplog): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{}') + with caplog.at_level(logging.INFO): + _, sink = build_source(sources, [path]) + sink.require_snapshot() + require_info_line(caplog, "Flag overrides: none in effect (%s: no entries)" % path) + + +# --------------------------------------------------------------------------- +# Change detection +# --------------------------------------------------------------------------- + +@watchdog_required +def test_watching_mode_is_quiet_when_file_is_absent(tmp_path, sources): + path = os.path.join(str(tmp_path), 'not-yet.json') + _, sink = build_source(sources, [path], lambda b: b.change_detection(ChangeDetection.WATCHING)) + assert sink.require_snapshot() == ({}, {}) + sink.require_no_snapshot() + write_file(path, '{"flagValues": {"flag1": true}}') + assert flag_values(sink.require_snapshot()) == {'flag1': True} + + +@watchdog_required +def test_watching_mode_reloads_on_change(tmp_path, sources): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{"flagValues": {"flag1": true}}') + source, sink = build_source(sources, [path], lambda b: b.change_detection(ChangeDetection.WATCHING)) + assert isinstance(source._change_detector, filedata.Watcher) + assert flag_values(sink.require_snapshot()) == {'flag1': True} + write_file(path, '{"flagValues": {"flag1": false}}') + assert flag_values(sink.require_snapshot()) == {'flag1': False} + + +def test_polling_mode_reloads_on_change(tmp_path, sources): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{"flagValues": {"flag1": true}}') + source, sink = build_source(sources, [path], lambda b: b.change_detection(ChangeDetection.POLLING).poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL)) + assert isinstance(source._change_detector, filedata.Poller) + assert flag_values(sink.require_snapshot()) == {'flag1': True} + write_file(path, '{"flagValues": {"flag1": false}}') + assert flag_values(sink.require_snapshot()) == {'flag1': False} + + +@pytest.mark.parametrize("mode", [ChangeDetection.POLLING, pytest.param(ChangeDetection.WATCHING, marks=watchdog_required)]) +def test_source_retains_last_good_data_across_malformed_edit(tmp_path, sources, caplog, mode): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{"flagValues": {"flag1": true}}') + with caplog.at_level(logging.ERROR): + _, sink = build_source(sources, [path], lambda b: b.change_detection(mode).poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL)) + sink.require_snapshot() + + # A malformed edit produces no snapshot: the previously applied overrides stay in + # effect because the sink is never called. The failure is logged. + write_file(path, '{"flagValues"') + sink.require_no_snapshot(1.5) + deadline = time.time() + TEST_TIMEOUT + while not any('Unable to load flag data' in r.getMessage() for r in caplog.records): + assert time.time() < deadline, "the load failure was not logged" + time.sleep(0.01) + + # Fixing the file recovers, through the change notification or the failure retry. + write_file(path, '{"flagValues": {"flag1": false}}') + assert flag_values(sink.require_snapshot()) == {'flag1': False} + + +def test_source_retries_a_failed_load_without_a_change_signal(tmp_path, sources, caplog): + # The fixed content has the same size as the malformed content and the file keeps its + # modification time, so the poller sees no change. Only the automatic retry can recover. + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{"flagValues": {"flag1": true}}') + with caplog.at_level(logging.ERROR): + _, sink = build_source(sources, [path], lambda b: b.poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL)) + sink.require_snapshot() + + malformed = '{"flagValues": {"flag1": false}' + fixed = '{"flagValues": {"flag1":false}}' + assert len(malformed) == len(fixed) + write_file(path, malformed) + deadline = time.time() + TEST_TIMEOUT + while not any('Unable to load flag data' in r.getMessage() for r in caplog.records): + assert time.time() < deadline, "the load failure was not logged" + time.sleep(0.01) + observed = os.stat(path) + write_file(path, fixed) + os.utime(path, ns=(observed.st_atime_ns, observed.st_mtime_ns)) + assert os.stat(path).st_size == observed.st_size + assert flag_values(sink.require_snapshot()) == {'flag1': False} + + +def test_source_does_not_supply_snapshots_after_close(tmp_path, sources): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{"flagValues": {"flag1": true}}') + source, sink = build_source(sources, [path], lambda b: b.poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL)) + sink.require_snapshot() + source.close() + write_file(path, '{"flagValues": {"flag1": false}}') + sink.require_no_snapshot(2.5) + + +def test_source_close_is_idempotent(tmp_path, sources): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{}') + source, _ = build_source(sources, [path]) + source.close() + source.close() + + +def test_source_start_after_close_does_nothing(tmp_path): + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{"flagValues": {"flag1": true}}') + source = FileOverrideSourceBuilder([path]).build(SDKConfig('SDK_KEY')) + source.close() + sink = CapturingSink() + source.start(sink) + sink.require_no_snapshot(0.2) + + +# --------------------------------------------------------------------------- +# Through the client +# --------------------------------------------------------------------------- + +def test_file_overrides_end_to_end(tmp_path): + # An operator writes, edits, and empties an override file. Evaluations follow without any + # client restart, even though the client never obtains data from LaunchDarkly. + path = os.path.join(str(tmp_path), 'overrides.json') + write_file(path, '{}') + user = Context.create('user-key') + source = FileOverrideSourceBuilder([path]).poll_interval(FileOverrideSourceBuilder.MINIMUM_POLL_INTERVAL) + datasystem = custom().synchronizers(HangingSynchronizer().builder).overrides(source).build() + config = Config(sdk_key='SDK_KEY', datasystem_config=datasystem, event_processor_class=MockEventProcessor) + + def eventually(predicate, message: str) -> None: + deadline = time.time() + TEST_TIMEOUT + while not predicate(): + assert time.time() < deadline, message + time.sleep(0.05) + + with LDClient(config, start_wait=0) as client: + # Not initialized and no override present: the default is served. + detail = client.variation_detail('overridden-flag', user, False) + assert detail.value is False + assert detail.reason == {'kind': 'ERROR', 'errorKind': 'CLIENT_NOT_READY'} + + # An operator adds an override. The running client picks it up. + write_file(path, '{"flagValues": {"overridden-flag": true}}') + eventually(lambda: client.variation('overridden-flag', user, False) is True, "the override was not picked up") + + # The override changes value. + write_file(path, '{"flagValues": {"overridden-flag": false}}') + + def changed() -> bool: + detail = client.variation_detail('overridden-flag', user, True) + return detail.value is False and detail.reason.get('overrideAffected') is True + + eventually(changed, "the changed override was not picked up") + + # The override is removed. The not-initialized short-circuit returns. + write_file(path, '{}') + eventually(lambda: client.variation_detail('overridden-flag', user, False).reason.get('errorKind') == 'CLIENT_NOT_READY', "the removed override was not picked up")