From 72d5776e73cde57fae7a8158497153ba033df092 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:28:50 -0700 Subject: [PATCH 01/11] feat: Add file loading code for flag overrides Adds `LaunchDarkly::Impl::FileData`, the file reading, parsing, merging, and reloading code that the file-based override source defined by the OVERRIDE specification needs. That source must tolerate a missing file, a file that is being written, and a burst of change notifications, and it must keep the last good data when a reload fails. The module provides: - Document parsing for JSON and YAML with validation that the document and its flags, flagValues, and segments members are objects. - An ordered merge of several documents with duplicate keys handling of fail or ignore, flagValues expansion into off flags that serve the value, model deserialization, and per-document entry counts. - A reloader that serializes reloads, debounces change signals, retains the last good result when a reload fails, retries after a bounded delay, reports each distinct failure once, and can skip byte-identical reloads. A configured file that does not exist can be treated as a file with no content. - A stat-based poller that compares existence, modification time, and size, so a file that appears or disappears is a change. - A watcher that uses rb-inotify on Linux to watch the directory of each file without descending into subdirectories, and the listen gem elsewhere, so a file that does not exist yet is picked up when it appears. It retries when a directory does not exist. The existing FDv1 and FDv2 file data sources are not changed. They keep their own implementation and behavior. A new spec pins the behaviors of those sources that their other specs did not assert, so that they stay as they are while this code exists beside them. --- lib/ldclient-rb/impl/file_data.rb | 7 + lib/ldclient-rb/impl/file_data/document.rb | 187 +++++++++ lib/ldclient-rb/impl/file_data/merge.rb | 146 +++++++ lib/ldclient-rb/impl/file_data/poller.rb | 81 ++++ lib/ldclient-rb/impl/file_data/reloader.rb | 308 ++++++++++++++ lib/ldclient-rb/impl/file_data/watcher.rb | 231 +++++++++++ spec/impl/file_data/document_spec.rb | 160 ++++++++ spec/impl/file_data/merge_spec.rb | 127 ++++++ spec/impl/file_data/poller_spec.rb | 107 +++++ spec/impl/file_data/reloader_spec.rb | 378 +++++++++++++++++ spec/impl/file_data/watcher_spec.rb | 169 ++++++++ .../file_data_source_compatibility_spec.rb | 387 ++++++++++++++++++ 12 files changed, 2288 insertions(+) create mode 100644 lib/ldclient-rb/impl/file_data.rb create mode 100644 lib/ldclient-rb/impl/file_data/document.rb create mode 100644 lib/ldclient-rb/impl/file_data/merge.rb create mode 100644 lib/ldclient-rb/impl/file_data/poller.rb create mode 100644 lib/ldclient-rb/impl/file_data/reloader.rb create mode 100644 lib/ldclient-rb/impl/file_data/watcher.rb create mode 100644 spec/impl/file_data/document_spec.rb create mode 100644 spec/impl/file_data/merge_spec.rb create mode 100644 spec/impl/file_data/poller_spec.rb create mode 100644 spec/impl/file_data/reloader_spec.rb create mode 100644 spec/impl/file_data/watcher_spec.rb create mode 100644 spec/integrations/file_data_source_compatibility_spec.rb diff --git a/lib/ldclient-rb/impl/file_data.rb b/lib/ldclient-rb/impl/file_data.rb new file mode 100644 index 00000000..eaa8b13b --- /dev/null +++ b/lib/ldclient-rb/impl/file_data.rb @@ -0,0 +1,7 @@ +# frozen_string_literal: true + +require "ldclient-rb/impl/file_data/document" +require "ldclient-rb/impl/file_data/merge" +require "ldclient-rb/impl/file_data/poller" +require "ldclient-rb/impl/file_data/reloader" +require "ldclient-rb/impl/file_data/watcher" diff --git a/lib/ldclient-rb/impl/file_data/document.rb b/lib/ldclient-rb/impl/file_data/document.rb new file mode 100644 index 00000000..005a4752 --- /dev/null +++ b/lib/ldclient-rb/impl/file_data/document.rb @@ -0,0 +1,187 @@ +# frozen_string_literal: true + +require "ldclient-rb/impl/model/serialization" + +require "yaml" + +module LaunchDarkly + module Impl + # + # Shared file reading, parsing, and merging code for the components that load flag and + # segment data from local files. + # + # @private + # + module FileData + # + # Raised when a file cannot be read or parsed. It carries the path so that callers can + # tell a per-file failure from a failure to merge the files' contents. + # + class ReadError < StandardError + # @return [String] + attr_reader :path + + # @return [Boolean] true when the file does not exist + attr_reader :missing + + # + # @param path [String] + # @param message [String] + # @param missing [Boolean] + # + def initialize(path, message, missing: false) + super("#{message} [#{path}]") + @path = path + @missing = missing + end + end + + # + # The parsed form of one data file. A document may contain full flag definitions, flag key + # to value entries, and segment definitions. Every hash has symbol keys. + # + class Document + # @return [Hash{Symbol => Hash}] + attr_reader :flags + + # @return [Hash{Symbol => Object}] + attr_reader :flag_values + + # @return [Hash{Symbol => Hash}] + attr_reader :segments + + # + # @param flags [Hash{Symbol => Hash}] + # @param flag_values [Hash{Symbol => Object}] + # @param segments [Hash{Symbol => Hash}] + # + def initialize(flags: {}, flag_values: {}, segments: {}) + @flags = flags + @flag_values = flag_values + @segments = segments + end + + # + # Parses the content of a data file. The content may be JSON or YAML. JSON is a subset of + # YAML, and the Ruby YAML parser handles it, so one parser serves both formats. + # + # An empty document is a document with no entries. A document that is not a mapping, or + # whose "flags", "flagValues", or "segments" member is not a mapping, is an error. + # + # @param content [String] + # @return [Document] + # @raise [ArgumentError] if the content is not a valid document + # @raise [Psych::SyntaxError] if the content cannot be parsed + # + def self.parse(content) + raw = YAML.safe_load(content) + raw = {} if raw.nil? + raise ArgumentError, "file content must be an object" unless raw.is_a?(Hash) + + data = FileData.symbolize_keys(raw) + Document.new( + flags: section(data, :flags), + flag_values: section(data, :flagValues), + segments: section(data, :segments) + ) + end + + # + # Reads and parses one data file. + # + # @param path [String] + # @return [Document] + # @raise [ReadError] if the file cannot be read or parsed + # + def self.read(path) + content = FileData.read_file(path) + FileData.parse_file(path, content) + end + + private_class_method def self.section(data, name) + value = data[name] + return {} if value.nil? + raise ArgumentError, "\"#{name}\" must be an object" unless value.is_a?(Hash) + + value + end + end + + # + # Reads the raw content of one file. + # + # @param path [String] + # @return [String] + # @raise [ReadError] if the file cannot be read + # + def self.read_file(path) + File.read(path) + rescue Errno::ENOENT => e + raise ReadError.new(path, "unable to read file: #{e.message}", missing: true) + rescue SystemCallError, IOError => e + raise ReadError.new(path, "unable to read file: #{e.message}") + end + + # + # Parses raw content that was read from the given path. + # + # @param path [String] + # @param content [String] + # @return [Document] + # @raise [ReadError] if the content cannot be parsed + # + def self.parse_file(path, content) + Document.parse(content) + rescue StandardError => e + raise ReadError.new(path, "error parsing file: #{e.message}") + end + + # + # Recursively converts hash keys to symbols. The SDK expects all data model objects to + # have symbol keys. + # + # @param value [Object] + # @return [Object] + # + def self.symbolize_keys(value) + case value + when Hash + value.to_h { |k, v| [k.to_s.to_sym, symbolize_keys(v)] } + when Array + value.map { |v| symbolize_keys(v) } + else + value + end + end + + # + # 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 the value as its off variation, so + # an evaluation reports the OFF reason kind. + # + # @param key [String] + # @param value [Object] + # @return [Hash] + # + def self.make_flag_with_value(key, value) + { + key: key, + on: false, + version: 1, + offVariation: 0, + variations: [value], + } + end + + # + # Converts each path to an absolute path. + # + # @param paths [Array, String] + # @return [Array] + # + def self.absolute_paths(paths) + Array(paths).map { |p| File.absolute_path(p.to_s) } + end + end + end +end diff --git a/lib/ldclient-rb/impl/file_data/merge.rb b/lib/ldclient-rb/impl/file_data/merge.rb new file mode 100644 index 00000000..470b8af3 --- /dev/null +++ b/lib/ldclient-rb/impl/file_data/merge.rb @@ -0,0 +1,146 @@ +# frozen_string_literal: true + +require "ldclient-rb/impl/data_store" +require "ldclient-rb/impl/file_data/document" +require "ldclient-rb/impl/model/serialization" + +module LaunchDarkly + module Impl + module FileData + # + # Values for the duplicate keys handling option. They select what happens when the same + # flag or segment key appears in more than one document. + # + module DuplicateKeysHandling + # A duplicated key makes the merge fail. + FAIL = :fail + + # Only the first occurrence of a duplicated key is kept, in the order the documents were given. + IGNORE = :ignore + + ALL = [FAIL, IGNORE].freeze + end + + # + # Raised when documents cannot be combined, for example because a key is duplicated or an + # entry is not an object. + # + class MergeError < StandardError + end + + # + # Counts the entries the merge kept from one document. + # + DocumentSummary = Struct.new(:flags, :segments) + + # + # Describes one configured file after a reload. `present` is false when the file does not + # exist and missing files are skipped. + # + FileSummary = Struct.new(:path, :present, :flags, :segments) + + # + # The merged items from one or more documents. + # + class MergeResult + # @return [Hash{Symbol => LaunchDarkly::Impl::Model::FeatureFlag}] + attr_reader :flags + + # @return [Hash{Symbol => LaunchDarkly::Impl::Model::Segment}] + attr_reader :segments + + # @return [Array] one entry per input document, in order + attr_reader :documents + + # @return [Array] set by the Reloader, one entry per configured file, in order + attr_accessor :files + + def initialize(flags, segments, documents) + @flags = flags + @segments = segments + @documents = documents + @files = [] + end + + # @return [Boolean] + def empty? + @flags.empty? && @segments.empty? + end + end + + # + # Combines the items of the given documents into one set of flags and one set of segments. + # Flag key to value entries expand into full flag definitions. Entries are deserialized into + # the SDK's data model classes, which validate them. The documents are processed in order, + # and the configured duplicate keys handling applies when the same key appears more than once. + # + # Items are keyed by the key under which they appear in the document. An entry that has no + # "key" member receives that key. A missing "version" defaults to 1. + # + # @param documents [Array] + # @param duplicate_keys_handling [Symbol] one of the {DuplicateKeysHandling} values + # @param logger [Logger, nil] receives data model validation messages + # @return [MergeResult] + # @raise [MergeError] if the documents cannot be combined + # + def self.merge(documents, duplicate_keys_handling: DuplicateKeysHandling::FAIL, logger: nil) + flags = {} + segments = {} + summaries = [] + + documents.each do |document| + summary = DocumentSummary.new(0, 0) + + document.flags.each do |key, data| + data = prepare_entry("flag", key, data) + item = Model.deserialize(DataStore::FEATURES, data, logger) + summary.flags += 1 if insert(flags, "flag", key, item, duplicate_keys_handling) + end + + document.flag_values.each do |key, value| + data = make_flag_with_value(key.to_s, value) + item = Model.deserialize(DataStore::FEATURES, data, logger) + summary.flags += 1 if insert(flags, "flag", key, item, duplicate_keys_handling) + end + + document.segments.each do |key, data| + data = prepare_entry("segment", key, data) + item = Model.deserialize(DataStore::SEGMENTS, data, logger) + summary.segments += 1 if insert(segments, "segment", key, item, duplicate_keys_handling) + end + + summaries << summary + end + + MergeResult.new(flags, segments, summaries) + end + + # + # Validates one full flag or segment entry and fills in its key and version. + # + private_class_method def self.prepare_entry(category, key, data) + raise MergeError, "#{category} \"#{key}\" is not an object" unless data.is_a?(Hash) + + data = data.dup + data[:key] = key.to_s if data[:key].nil? + data[:version] = 1 if data[:version].nil? + data + end + + # + # Adds an item unless its key was already seen. Returns true when it added the item. + # + private_class_method def self.insert(items, category, key, item, duplicate_keys_handling) + key = key.to_sym + if items.key?(key) + return false if duplicate_keys_handling == DuplicateKeysHandling::IGNORE + + raise MergeError, "#{category} key \"#{key}\" was used more than once" + end + + items[key] = item + true + end + end + end +end diff --git a/lib/ldclient-rb/impl/file_data/poller.rb b/lib/ldclient-rb/impl/file_data/poller.rb new file mode 100644 index 00000000..b0ce08fc --- /dev/null +++ b/lib/ldclient-rb/impl/file_data/poller.rb @@ -0,0 +1,81 @@ +# frozen_string_literal: true + +require "ldclient-rb/impl/repeating_task" + +module LaunchDarkly + module Impl + module FileData + # + # Detects changes to a set of files by examining them on a fixed interval. Use it where file + # system change notifications are not available or not reliable. 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. A file that cannot be examined counts as absent. + # + # The poller samples the files once per interval and compares only modification time and + # size. A rewrite that keeps both values is not detected. + # + # Detection is generous. The callback can run for a change that does not alter the + # effective data. Feed it into a {Reloader}, whose debouncing and skip-unchanged handling + # absorb the excess. + # + # @private + # + class Poller + # The observed state of one file, or its absence. + FileState = Struct.new(:exists, :mtime, :size) + + ABSENT = FileState.new(false, nil, nil).freeze + private_constant :ABSENT + + # + # Creates and starts a poller. It examines the files once before it returns, so only later + # changes invoke the callback. Call {#stop} to stop it. + # + # @param paths [Array] absolute paths of the files to examine + # @param interval [Numeric] seconds between examinations + # @param on_change [#call] invoked with no arguments when a change is detected + # @param logger [Logger] + # + def initialize(paths, interval, on_change, logger) + @paths = paths + @on_change = on_change + @last = Poller.observe_all(paths) + @task = RepeatingTask.new(interval, interval, method(:examine), logger, "LD/FileDataPoller") + @task.start + end + + # + # Stops the poller and waits for the worker thread to finish. A callback that is already + # running completes first. + # + def stop + @task.stop + end + + # + # Examines every file and returns the observed states, in path order. + # + # @param paths [Array] + # @return [Array] + # + def self.observe_all(paths) + paths.map do |path| + begin + stat = File.stat(path) + FileState.new(true, stat.mtime, stat.size) + rescue SystemCallError + ABSENT + end + end + end + + private def examine + current = Poller.observe_all(@paths) + changed = current != @last + @last = current + @on_change.call if changed + end + end + end + end +end diff --git a/lib/ldclient-rb/impl/file_data/reloader.rb b/lib/ldclient-rb/impl/file_data/reloader.rb new file mode 100644 index 00000000..a080dd9b --- /dev/null +++ b/lib/ldclient-rb/impl/file_data/reloader.rb @@ -0,0 +1,308 @@ +# frozen_string_literal: true + +require "ldclient-rb/impl/file_data/document" +require "ldclient-rb/impl/file_data/merge" +require "ldclient-rb/impl/util" + +require "digest" + +module LaunchDarkly + module Impl + module FileData + # + # 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 no-op applications. + # + # The worker thread starts on the first {#reload_now} or {#trigger} call rather than in the + # constructor. A reloader can be constructed by a component whose lifecycle never uses it, + # and construction alone must not leak a thread. + # + # @private + # + class Reloader + # A settle window long enough to coalesce the burst of change notifications produced by a + # single file edit, and short enough to stay responsive. In seconds. + 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. In seconds. + DEFAULT_RETRY_DELAY = 1.0 + + # + # @param paths [Array] the files to load, already resolved to absolute paths. The + # order is significant: it determines which file wins under the duplicate keys handling. + # @param logger [Logger] + # @param apply [#call] invoked with each successfully merged {MergeResult}. Calls are + # serialized, so implementations do not need their own synchronization against other + # reloads. `apply` and `on_error` must not call back into {#stop}. + # @param on_error [#call, nil] invoked with the error when a reload fails, once per distinct + # failure. With automatic retries, repeats of an identical failure do not re-invoke it. A + # success re-arms it. The error is a {ReadError} when a file could not be read or parsed, + # or a {MergeError} otherwise. The reloader logs failures itself. + # @param duplicate_keys_handling [Symbol] one of the {DuplicateKeysHandling} values + # @param skip_missing_paths [Boolean] when true, a configured file that does not exist is a + # file with no content, and the reload succeeds with the data of the files that exist. + # When false, a missing file fails the reload like any other read error. + # @param debounce_delay [Numeric] seconds to wait after a {#trigger} call for further calls + # to settle before reloading. If zero or negative, each trigger reloads at once. + # @param retry_delay [Numeric] seconds to wait after a failed reload before retrying + # automatically. If zero or negative, there is no automatic retry. + # @param skip_unchanged [Boolean] if true, `apply` is not invoked when the files' raw + # contents are byte-identical to the last successfully applied contents. + # + def initialize(paths:, logger:, apply:, on_error: nil, + duplicate_keys_handling: DuplicateKeysHandling::FAIL, + skip_missing_paths: false, + debounce_delay: DEFAULT_DEBOUNCE_DELAY, + retry_delay: DEFAULT_RETRY_DELAY, + skip_unchanged: false) + @paths = paths + @logger = logger + @apply = apply + @on_error = on_error + @duplicate_keys_handling = duplicate_keys_handling + @skip_missing_paths = skip_missing_paths + @debounce_delay = debounce_delay + @retry_delay = retry_delay + @skip_unchanged = skip_unchanged + + # Guards the scheduling state below and wakes the worker. + @mutex = Mutex.new + @cond = ConditionVariable.new + @worker = nil + @stopped = false + @trigger_pending = false + @retry_requested = false + @debounce_deadline = nil + @retry_deadline = nil + + # Serializes the actual load work between reload_now and the worker. + @reload_mutex = Mutex.new + @last_good_digest = nil + @last_error_message = nil + end + + # + # Synchronously loads the files and applies the result, or reports the failure. Use it for + # the initial load. A failure here schedules the same automatic retry as a failed + # triggered reload. + # + # @return [Boolean] true if the load succeeded + # + def reload_now + ensure_started + ok = reload(retrying: false) + request_retry unless ok + ok + end + + # + # Signals that the files may have changed and a reload should happen after the debounce + # delay. It never blocks. Signals that arrive while a reload is already pending are + # coalesced. + # + def trigger + ensure_started + @mutex.synchronize do + @trigger_pending = true + @cond.signal + end + end + + # + # Stops the reloader. It does not wait for a reload that is already in progress. A reload + # wedged in a blocking file read must not be able to wedge shutdown. Such a reload can + # still deliver its result through `apply` or `on_error` shortly after this method + # returns, and consumers tolerate that. A reload that has not yet reached its callbacks + # when this method is called does not invoke them. + # + def stop + @mutex.synchronize do + @stopped = true + @cond.broadcast + end + end + + private def ensure_started + @mutex.synchronize do + return if @stopped || !@worker.nil? + + @worker = Thread.new { run } + @worker.name = "LD/FileDataReloader" + end + end + + private def request_retry + return unless @retry_delay > 0 + + @mutex.synchronize do + @retry_requested = true + @cond.signal + end + end + + private def run + loop do + action = next_action + break if action == :stop + + if action == :reload + @logger.info { "[LDClient] Reloading flag data after detecting a change" } + else + @logger.debug { "[LDClient] Retrying flag data load after earlier failure" } + end + ok = reload(retrying: action == :retry) + @mutex.synchronize do + # A pending retry is superseded by this reload. The reload either succeeded, or it + # failed and arms a fresh retry here. + @retry_deadline = ok || @retry_delay <= 0 ? nil : monotonic_now + @retry_delay + end + end + rescue => e + Util.log_exception(@logger, "Unexpected error in file data reloader", e) + end + + # + # Waits until a reload is due. Returns :reload for a change-triggered reload, :retry for + # an automatic retry, or :stop. + # + private def next_action + @mutex.synchronize do + loop do + return :stop if @stopped + + now = monotonic_now + if @trigger_pending + @trigger_pending = false + return :reload if @debounce_delay <= 0 + + @debounce_deadline = now + @debounce_delay + end + if @retry_requested + # A synchronous reload_now failed. Arm the retry without reloading again at once. + # An already armed retry keeps its earlier deadline. + @retry_requested = false + @retry_deadline ||= now + @retry_delay + end + if !@debounce_deadline.nil? && now >= @debounce_deadline + @debounce_deadline = nil + return :reload + end + if !@retry_deadline.nil? && now >= @retry_deadline + @retry_deadline = nil + return :retry + end + + deadlines = [@debounce_deadline, @retry_deadline].compact + timeout = deadlines.empty? ? nil : [deadlines.min - now, 0].max + @cond.wait(@mutex, timeout) + end + end + end + + # + # Performs one full load of all configured files and returns whether it succeeded. That + # decides whether a retry is armed, so a skipped no-op application counts as success. The + # whole set is re-read on every reload: entries are combined across files in order, so a + # change to one file can alter which file wins for a key. + # + private def reload(retrying: false) + @reload_mutex.synchronize do + # A trigger already queued when stop was called can still reach here. + return true if stopped? + + documents = [] + files = [] + digest = Digest::SHA256.new + @paths.each do |path| + begin + # One read feeds both the digest and the parse, so the skip-unchanged digest can + # never disagree with the content that was applied. + content = FileData.read_file(path) + rescue ReadError => e + if e.missing && @skip_missing_paths + @logger.debug { "[LDClient] File #{path} does not exist; it contributes no data" } + files << FileSummary.new(path, false, 0, 0) + next + end + return record_failure(e) + end + digest << content << "\0" + begin + documents << FileData.parse_file(path, content) + rescue ReadError => e + return record_failure(e) + end + files << FileSummary.new(path, true, 0, 0) + end + + begin + merged = FileData.merge(documents, + duplicate_keys_handling: @duplicate_keys_handling, + logger: @logger) + rescue => e + return record_failure(e) + end + + # Documents are the present files in order. Copy their counts onto the file summaries. + document_index = 0 + files.each do |file| + next unless file.present + + summary = merged.documents[document_index] + file.flags = summary.flags + file.segments = summary.segments + document_index += 1 + end + merged.files = files + + # stop may have been called while the files were being read. Deliver nothing then. + return true if stopped? + + # A success right after a failure must apply even when the content is unchanged since + # the last success. The consumer heard about the failure through on_error and may + # have moved to an interrupted state. Only apply tells it that things are good again. + recovering = !@last_error_message.nil? + @last_error_message = nil + hexdigest = digest.hexdigest + return true if @skip_unchanged && !recovering && hexdigest == @last_good_digest + + @last_good_digest = hexdigest + @apply.call(merged) + true + end + end + + private def record_failure(error) + # stop may have been called while the files were being read. Deliver nothing then, and + # report success so that no retry is armed. + return true if stopped? + + # 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 re-invoke on_error. + message = error.message + if message == @last_error_message + @logger.debug { "[LDClient] Unable to load flags: #{message}" } + return false + end + + @last_error_message = message + @logger.error { "[LDClient] Unable to load flags: #{message}" } + @on_error&.call(error) + false + end + + private def stopped? + @mutex.synchronize { @stopped } + end + + private def monotonic_now + Process.clock_gettime(Process::CLOCK_MONOTONIC) + end + end + end + end +end diff --git a/lib/ldclient-rb/impl/file_data/watcher.rb b/lib/ldclient-rb/impl/file_data/watcher.rb new file mode 100644 index 00000000..bf079f26 --- /dev/null +++ b/lib/ldclient-rb/impl/file_data/watcher.rb @@ -0,0 +1,231 @@ +# frozen_string_literal: true + +require "ldclient-rb/impl/repeating_task" +require "ldclient-rb/impl/util" + +require "concurrent/atomics" +require "set" + +module LaunchDarkly + module Impl + module FileData + # + # Detects changes to a set of files through file system change notifications. The + # notification mechanism comes from the optional `listen` gem, which the SDK does not depend + # on. Check {Watcher.available?} before constructing a watcher. + # + # On Linux the watcher uses `rb-inotify`, which `listen` depends on, to watch the directory of + # each file without descending into subdirectories. Elsewhere it uses `listen` itself, which + # scans the whole directory tree under each watched directory. + # + # The watcher observes the directory of each file, so a configured file that does not exist + # yet is picked up when it appears. If a directory does not exist, the watcher logs the + # problem and retries once per second until it does. When the watches are in place after a + # retry, the callback runs once, so that a change made while there was no watch is not missed. + # + # @private + # + class Watcher + # Seconds between attempts to set up the watches after a failure. + RETRY_INTERVAL = 1.0 + + INOTIFY_EVENTS = [:create, :modify, :close_write, :attrib, :delete, :moved_to, :moved_from].freeze + private_constant :INOTIFY_EVENTS + + # + # Returns true if a change notification mechanism can be loaded. + # + # @return [Boolean] + # + def self.available? + inotify_available? || listen_available? + end + + # + # Returns true if `rb-inotify` can be loaded. It is a dependency of `listen` on Linux and + # does not load on other platforms. + # + # @return [Boolean] + # + def self.inotify_available? + @inotify_available = load_library("rb-inotify") if @inotify_available.nil? + @inotify_available + end + + # + # Returns true if the `listen` gem can be loaded. + # + # @return [Boolean] + # + def self.listen_available? + @listen_available = load_library("listen") if @listen_available.nil? + @listen_available + end + + private_class_method def self.load_library(name) + require name + true + rescue LoadError, StandardError + false + end + + # + # Creates and starts a watcher. + # + # @param paths [Array] absolute paths of the files to watch + # @param on_change [#call] invoked with no arguments when one of the files changes + # @param logger [Logger] + # + def initialize(paths, on_change, logger) + @paths = paths + @on_change = on_change + @logger = logger + @stopped = Concurrent::AtomicBoolean.new(false) + @lock = Mutex.new + @listener = nil + @retry_task = nil + @last_error_message = nil + + return if try_start + + @retry_task = RepeatingTask.new(RETRY_INTERVAL, RETRY_INTERVAL, method(:retry_start), logger, + "LD/FileDataWatcherRetry") + @retry_task.start + end + + # + # Stops the watcher. No callback runs after this method returns, apart from one that is + # already in progress. + # + def stop + return unless @stopped.make_true + + @retry_task&.stop + listener = @lock.synchronize do + l = @listener + @listener = nil + l + end + listener&.stop + end + + private def retry_start + return if @stopped.value + return unless try_start + + # This runs on the retry task's own thread, which RepeatingTask#stop allows. + @retry_task.stop + @on_change.call unless @stopped.value + end + + # + # Sets up the watches. Returns false, after logging, if that is not possible yet. + # + private def try_start + directories = @paths.map { |p| File.dirname(p) }.uniq + missing = directories.reject { |d| File.directory?(d) } + unless missing.empty? + log_setup_failure("directory does not exist: #{missing.join(', ')}") + return false + end + + listener = Watcher.inotify_available? ? start_inotify : start_listen + + @lock.synchronize do + if @stopped.value + listener.stop + else + @listener = listener + end + end + @last_error_message = nil + true + rescue => e + log_setup_failure(e.message) + false + end + + # + # Watches the real directory of each file, without descending into subdirectories, and + # reports events whose file name is one of the watched names in that directory. + # + private def start_inotify + names_by_directory = {} + @paths.each do |p| + real_directory = File.realpath(File.dirname(p)) + (names_by_directory[real_directory] ||= Set.new) << File.basename(p) + end + + notifier = INotify::Notifier.new + begin + names_by_directory.each do |directory, names| + notifier.watch(directory, *INOTIFY_EVENTS) do |event| + @on_change.call if names.include?(event.name) && !@stopped.value + end + end + rescue + notifier.close + raise + end + InotifyListener.new(notifier, @logger) + end + + # + # Watches the real directory of each file with the `listen` gem, which reports paths under + # the real directory, so the paths to match are built the same way. + # + private def start_listen + directories = @paths.map { |p| File.dirname(p) }.uniq + real_directories = directories.map { |d| File.realpath(d) } + watched = Set.new(@paths.map { |p| File.join(File.realpath(File.dirname(p)), File.basename(p)) }) + + listener = Listen.to(*real_directories) do |modified, added, removed| + changed = (modified + added + removed).any? { |p| watched.include?(p) } + @on_change.call if changed && !@stopped.value + end + listener.start + listener + end + + private def log_setup_failure(message) + if message == @last_error_message + @logger.debug { "[LDClient] Unable to watch data files: #{message}" } + else + @last_error_message = message + @logger.error { "[LDClient] Unable to watch data files: #{message}" } + end + end + + # + # Runs an inotify notifier on its own thread and stops it on request. + # + class InotifyListener + def initialize(notifier, logger) + @notifier = notifier + @thread = Thread.new do + begin + notifier.run + rescue IOError, SystemCallError + # The notifier was closed by stop. + rescue => e + Util.log_exception(logger, "Unexpected error in file data watcher", e) + end + end + @thread.name = "LD/FileDataWatcher" + end + + # + # Stops the notifier and waits briefly for its thread. Closing the notifier ends the + # blocking read that the thread is in. + # + def stop + @notifier.stop + @notifier.close + @thread.join(2) + end + end + private_constant :InotifyListener + end + end + end +end diff --git a/spec/impl/file_data/document_spec.rb b/spec/impl/file_data/document_spec.rb new file mode 100644 index 00000000..e2f565ea --- /dev/null +++ b/spec/impl/file_data/document_spec.rb @@ -0,0 +1,160 @@ +# frozen_string_literal: true + +require "spec_helper" +require "tmpdir" +require "ldclient-rb/impl/file_data" + +module LaunchDarkly + module Impl + module FileData + describe Document do + describe "parse" do + it "parses a JSON document with every section and symbolizes nested keys" do + document = Document.parse(<<~JSON) + { + "flags": { "flag1": { "key": "flag1", "on": true, "rules": [ { "clauses": [ { "op": "in" } ] } ] } }, + "flagValues": { "flag2": "value2" }, + "segments": { "seg1": { "key": "seg1", "included": ["user1"] } } + } + JSON + + expect(document.flags.keys).to eq([:flag1]) + expect(document.flags[:flag1][:rules][0][:clauses][0][:op]).to eq("in") + expect(document.flag_values).to eq({ flag2: "value2" }) + expect(document.segments[:seg1][:included]).to eq(["user1"]) + end + + it "parses a YAML document" do + document = Document.parse(<<~YAML) + --- + flags: + flag1: + key: flag1 + "on": true + flagValues: + flag2: value2 + segments: + seg1: + key: seg1 + YAML + + expect(document.flags[:flag1][:on]).to be true + expect(document.flag_values).to eq({ flag2: "value2" }) + expect(document.segments.keys).to eq([:seg1]) + end + + it "treats an empty document as a document with no entries" do + document = Document.parse("") + + expect(document.flags).to eq({}) + expect(document.flag_values).to eq({}) + expect(document.segments).to eq({}) + end + + it "treats a document with no known sections as a document with no entries" do + document = Document.parse("{}") + + expect(document.flags).to eq({}) + expect(document.flag_values).to eq({}) + expect(document.segments).to eq({}) + end + + it "rejects a document that is not an object" do + expect { Document.parse("\"hello\"") }.to raise_error(ArgumentError, /must be an object/) + expect { Document.parse("[1, 2]") }.to raise_error(ArgumentError, /must be an object/) + end + + it "rejects a section that is not an object" do + expect { Document.parse('{"flags": []}') }.to raise_error(ArgumentError, /"flags" must be an object/) + expect { Document.parse('{"flagValues": 3}') }.to raise_error(ArgumentError, /"flagValues" must be an object/) + expect { Document.parse('{"segments": "x"}') }.to raise_error(ArgumentError, /"segments" must be an object/) + end + + it "symbolizes keys that YAML parsed as numbers" do + document = Document.parse("flagValues:\n 123: true\n") + + expect(document.flag_values).to eq({ "123": true }) + end + + it "raises a syntax error for content that cannot be parsed" do + expect { Document.parse('{"flagValues"') }.to raise_error(Psych::SyntaxError) + end + end + + describe "read" do + around do |example| + Dir.mktmpdir do |dir| + @dir = dir + example.run + end + end + + it "reads and parses a file" do + path = File.join(@dir, "flags.json") + File.write(path, '{"flagValues": {"flag1": 1}}') + + document = Document.read(path) + + expect(document.flag_values).to eq({ flag1: 1 }) + end + + it "reports a missing file as a read error that is marked missing" do + path = File.join(@dir, "no-such-file.json") + + expect { Document.read(path) }.to raise_error(ReadError) do |error| + expect(error.missing).to be true + expect(error.path).to eq(path) + expect(error.message).to include(path) + end + end + + it "reports a path that cannot be read as a read error that is not marked missing" do + expect { Document.read(@dir) }.to raise_error(ReadError) do |error| + expect(error.missing).to be false + expect(error.path).to eq(@dir) + end + end + + it "reports unparseable content as a read error that is not marked missing" do + path = File.join(@dir, "bad.json") + File.write(path, '{"flagValues"') + + expect { Document.read(path) }.to raise_error(ReadError) do |error| + expect(error.missing).to be false + expect(error.message).to include("error parsing file") + expect(error.message).to include(path) + end + end + + it "reports a document that is not an object as a read error" do + path = File.join(@dir, "scalar.yaml") + File.write(path, "just a string\n") + + expect { Document.read(path) }.to raise_error(ReadError, /must be an object/) + end + end + end + + describe "make_flag_with_value" do + it "builds a flag that is off and serves the value as its off variation" do + flag = FileData.make_flag_with_value("flag1", "value1") + + expect(flag).to eq({ + key: "flag1", + on: false, + version: 1, + offVariation: 0, + variations: ["value1"], + }) + end + end + + describe "absolute_paths" do + it "converts relative paths to absolute paths and accepts a single string" do + expect(FileData.absolute_paths("a/b.json")).to eq([File.absolute_path("a/b.json")]) + expect(FileData.absolute_paths(["/x/y.json", "z.json"])).to eq(["/x/y.json", File.absolute_path("z.json")]) + end + end + end + end +end diff --git a/spec/impl/file_data/merge_spec.rb b/spec/impl/file_data/merge_spec.rb new file mode 100644 index 00000000..073af43f --- /dev/null +++ b/spec/impl/file_data/merge_spec.rb @@ -0,0 +1,127 @@ +# frozen_string_literal: true + +require "spec_helper" +require "capturing_logger" +require "ldclient-rb/impl/file_data" + +module LaunchDarkly + module Impl + module FileData + describe "merge" do + def document(flags: {}, flag_values: {}, segments: {}) + Document.new(flags: flags, flag_values: flag_values, segments: segments) + end + + def full_flag(key, version: nil) + data = { key: key, on: true, variations: ["a", "b"], fallthrough: { variation: 1 } } + data[:version] = version unless version.nil? + data + end + + it "combines flags, flag values, and segments from documents in order" do + doc1 = document(flags: { flag1: full_flag("flag1") }, segments: { seg1: { key: "seg1", included: ["u"] } }) + doc2 = document(flag_values: { flag2: "value2" }) + + result = FileData.merge([doc1, doc2]) + + expect(result.flags.keys).to eq([:flag1, :flag2]) + expect(result.flags[:flag1]).to be_a(Model::FeatureFlag) + expect(result.flags[:flag2]).to be_a(Model::FeatureFlag) + expect(result.segments.keys).to eq([:seg1]) + expect(result.segments[:seg1]).to be_a(Model::Segment) + expect(result.documents).to eq([DocumentSummary.new(1, 1), DocumentSummary.new(1, 0)]) + expect(result.empty?).to be false + end + + it "expands a flag value into an off flag that serves the value for every context" do + result = FileData.merge([document(flag_values: { flag2: "value2" })]) + + flag = result.flags[:flag2] + expect(flag.key).to eq("flag2") + expect(flag.on).to be false + expect(flag.version).to eq(1) + expect(flag.variations).to eq(["value2"]) + expect(flag.off_variation).to eq(0) + expect(flag.off_result.reason).to eq(EvaluationReason.off) + end + + it "reports an empty result for no documents" do + expect(FileData.merge([]).empty?).to be true + end + + it "fails when the same flag key appears in two documents" do + docs = [document(flags: { flag1: full_flag("flag1") }), document(flags: { flag1: full_flag("flag1") })] + + expect { FileData.merge(docs) }.to raise_error(MergeError, /flag key "flag1" was used more than once/) + end + + it "fails when the same segment key appears in two documents" do + docs = [document(segments: { seg1: { key: "seg1" } }), document(segments: { seg1: { key: "seg1" } })] + + expect { FileData.merge(docs) }.to raise_error(MergeError, /segment key "seg1" was used more than once/) + end + + it "fails when a key appears in both flags and flagValues of one document" do + docs = [document(flags: { flag1: full_flag("flag1") }, flag_values: { flag1: "x" })] + + expect { FileData.merge(docs) }.to raise_error(MergeError, /flag key "flag1"/) + end + + it "keeps the first document's entry and drops later duplicates with ignore handling" do + doc1 = document(flag_values: { flag1: "first" }) + doc2 = document(flag_values: { flag1: "second", flag2: "other" }) + + result = FileData.merge([doc1, doc2], duplicate_keys_handling: DuplicateKeysHandling::IGNORE) + + expect(result.flags[:flag1].variations).to eq(["first"]) + expect(result.flags[:flag2].variations).to eq(["other"]) + expect(result.documents).to eq([DocumentSummary.new(1, 0), DocumentSummary.new(1, 0)]) + end + + it "defaults a missing version to 1 and keeps a given version" do + doc = document( + flags: { flag1: full_flag("flag1"), flag2: full_flag("flag2", version: 5) }, + segments: { seg1: { key: "seg1" }, seg2: { key: "seg2", version: 9 } } + ) + + result = FileData.merge([doc]) + + expect(result.flags[:flag1].version).to eq(1) + expect(result.flags[:flag2].version).to eq(5) + expect(result.segments[:seg1].version).to eq(1) + expect(result.segments[:seg2].version).to eq(9) + end + + it "fills a missing key member from the key the entry appears under" do + doc = document(flags: { flag1: { on: false, variations: [1] } }, segments: { seg1: {} }) + + result = FileData.merge([doc]) + + expect(result.flags[:flag1].key).to eq("flag1") + expect(result.segments[:seg1].key).to eq("seg1") + end + + it "does not modify the document's own hashes" do + data = { on: false, variations: [1] } + FileData.merge([document(flags: { flag1: data })]) + + expect(data).to eq({ on: false, variations: [1] }) + end + + it "fails when an entry is not an object" do + expect { FileData.merge([document(flags: { flag1: 5 })]) }.to raise_error(MergeError, /flag "flag1" is not an object/) + expect { FileData.merge([document(segments: { seg1: [] })]) }.to raise_error(MergeError, /segment "seg1" is not an object/) + end + + it "passes the logger to model validation" do + logger = CapturingLogger.new + bad = { key: "flag1", on: true, variations: ["a"], fallthrough: { variation: 5 } } + + FileData.merge([document(flags: { flag1: bad })], logger: logger) + + expect(logger.output).to include("Data inconsistency in feature flag \"flag1\"") + end + end + end + end +end diff --git a/spec/impl/file_data/poller_spec.rb b/spec/impl/file_data/poller_spec.rb new file mode 100644 index 00000000..14eb503e --- /dev/null +++ b/spec/impl/file_data/poller_spec.rb @@ -0,0 +1,107 @@ +# frozen_string_literal: true + +require "spec_helper" +require "tmpdir" +require "ldclient-rb/impl/file_data" + +module LaunchDarkly + module Impl + module FileData + describe Poller do + let(:interval) { 0.05 } + + around do |example| + Dir.mktmpdir do |dir| + @dir = dir + example.run + end + end + + def path(name) + File.join(@dir, name) + end + + def wait_for(timeout = 3) + deadline = Time.now + timeout + until yield + return false if Time.now > deadline + sleep 0.01 + end + true + end + + def with_poller(paths) + calls = Concurrent::AtomicFixnum.new(0) + poller = Poller.new(paths, interval, -> { calls.increment }, $null_log) + begin + yield poller, calls + ensure + poller.stop + end + end + + it "invokes the callback when a file's modification time changes" do + File.write(path("a.json"), "{}") + with_poller([path("a.json")]) do |_poller, calls| + File.utime(Time.now + 10, Time.now + 10, path("a.json")) + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "invokes the callback when a file's size changes but its modification time does not" do + File.write(path("a.json"), "{}") + mtime = File.mtime(path("a.json")) + with_poller([path("a.json")]) do |_poller, calls| + File.write(path("a.json"), '{"flagValues": {}}') + File.utime(mtime, mtime, path("a.json")) + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "invokes the callback when a file appears" do + with_poller([path("missing.json")]) do |_poller, calls| + sleep interval * 2 + expect(calls.value).to eq(0) + File.write(path("missing.json"), "{}") + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "invokes the callback when a file disappears" do + File.write(path("a.json"), "{}") + with_poller([path("a.json")]) do |_poller, calls| + File.delete(path("a.json")) + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "watches every configured file" do + File.write(path("a.json"), "{}") + File.write(path("b.json"), "{}") + with_poller([path("a.json"), path("b.json")]) do |_poller, calls| + File.utime(Time.now + 10, Time.now + 10, path("b.json")) + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "does not invoke the callback when nothing changed" do + File.write(path("a.json"), "{}") + with_poller([path("a.json")]) do |_poller, calls| + sleep interval * 6 + expect(calls.value).to eq(0) + end + end + + it "does not invoke the callback after it is stopped" do + File.write(path("a.json"), "{}") + with_poller([path("a.json")]) do |poller, calls| + poller.stop + File.utime(Time.now + 10, Time.now + 10, path("a.json")) + sleep interval * 6 + expect(calls.value).to eq(0) + end + end + end + end + end +end diff --git a/spec/impl/file_data/reloader_spec.rb b/spec/impl/file_data/reloader_spec.rb new file mode 100644 index 00000000..241ad308 --- /dev/null +++ b/spec/impl/file_data/reloader_spec.rb @@ -0,0 +1,378 @@ +# frozen_string_literal: true + +require "spec_helper" +require "capturing_logger" +require "tmpdir" +require "ldclient-rb/impl/file_data" + +module LaunchDarkly + module Impl + module FileData + describe Reloader do + around do |example| + Dir.mktmpdir do |dir| + @dir = dir + example.run + end + end + + def path(name) + File.join(@dir, name) + end + + def write(name, content) + File.write(path(name), content) + path(name) + end + + def values_doc(values) + { flagValues: values }.to_json + end + + def wait_for(timeout = 3) + deadline = Time.now + timeout + until yield + return false if Time.now > deadline + sleep 0.01 + end + true + end + + # Collects apply and on_error calls in a thread-safe way. + class Recorder + attr_reader :applied, :errors + + def initialize + @lock = Mutex.new + @applied = [] + @errors = [] + end + + def apply(merged) + @lock.synchronize { @applied << merged } + end + + def on_error(error) + @lock.synchronize { @errors << error } + end + + def flag_values(index = -1) + @applied[index].flags.transform_values { |flag| flag.variations[0] } + end + end + + def make_reloader(paths, recorder, logger: $null_log, **options) + Reloader.new(paths: paths, logger: logger, apply: recorder.method(:apply), on_error: recorder.method(:on_error), **options) + end + + def with_reloader(paths, logger: $null_log, **options) + recorder = Recorder.new + reloader = make_reloader(paths, recorder, logger: logger, **options) + begin + yield reloader, recorder + ensure + reloader.stop + end + end + + it "does not start a worker thread until it is used" do + with_reloader([write("a.json", "{}")]) do |reloader, _recorder| + expect(reloader.instance_variable_get(:@worker)).to be_nil + + reloader.trigger + + worker = reloader.instance_variable_get(:@worker) + expect(worker).to be_a(Thread) + expect(worker.name).to eq("LD/FileDataReloader") + end + end + + it "applies the merged files and describes each file on a synchronous reload" do + a = write("a.json", values_doc({ flag1: "a" })) + b = write("b.json", { flagValues: { flag2: "b" }, segments: { seg1: { key: "seg1" } } }.to_json) + + with_reloader([a, b]) do |reloader, recorder| + expect(reloader.reload_now).to be true + + expect(recorder.applied.length).to eq(1) + expect(recorder.flag_values).to eq({ flag1: "a", flag2: "b" }) + expect(recorder.applied[0].segments.keys).to eq([:seg1]) + expected_files = [FileSummary.new(a, true, 1, 0), FileSummary.new(b, true, 1, 1)] + expect(recorder.applied[0].files).to eq(expected_files) + expect(recorder.errors).to be_empty + end + end + + it "fails a reload for a missing file by default" do + with_reloader([path("missing.json")], retry_delay: 0) do |reloader, recorder| + expect(reloader.reload_now).to be false + + expect(recorder.applied).to be_empty + expect(recorder.errors.length).to eq(1) + expect(recorder.errors[0]).to be_a(ReadError) + expect(recorder.errors[0].missing).to be true + end + end + + it "treats a missing file as a file with no content when skipping missing paths" do + a = write("a.json", values_doc({ flag1: "a" })) + missing = path("missing.json") + + with_reloader([a, missing], skip_missing_paths: true) do |reloader, recorder| + expect(reloader.reload_now).to be true + + expect(recorder.flag_values).to eq({ flag1: "a" }) + expected_files = [FileSummary.new(a, true, 1, 0), FileSummary.new(missing, false, 0, 0)] + expect(recorder.applied[0].files).to eq(expected_files) + end + end + + it "applies an empty result when every file is missing and missing paths are skipped" do + with_reloader([path("missing.json")], skip_missing_paths: true) do |reloader, recorder| + expect(reloader.reload_now).to be true + + expect(recorder.applied.length).to eq(1) + expect(recorder.applied[0].empty?).to be true + end + end + + it "keeps the last good result when a file cannot be parsed and reports the failure once" do + a = write("a.json", values_doc({ flag1: "a" })) + logger = CapturingLogger.new + + with_reloader([a], logger: logger, retry_delay: 0) do |reloader, recorder| + reloader.reload_now + write("a.json", '{"flagValues"') + + expect(reloader.reload_now).to be false + expect(reloader.reload_now).to be false + + expect(recorder.applied.length).to eq(1) + expect(recorder.errors.length).to eq(1) + expect(recorder.errors[0]).to be_a(ReadError) + expect(logger.output.scan("ERROR").length).to eq(1) + end + end + + it "reports a different failure again" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], retry_delay: 0) do |reloader, recorder| + write("a.json", '{"flagValues"') + reloader.reload_now + write("a.json", '{"flagValues": []}') + reloader.reload_now + + expect(recorder.errors.length).to eq(2) + end + end + + it "fails a reload for duplicate keys across files and applies the first file's entry with ignore handling" do + a = write("a.json", values_doc({ flag1: "first" })) + b = write("b.json", values_doc({ flag1: "second" })) + + with_reloader([a, b], retry_delay: 0) do |reloader, recorder| + expect(reloader.reload_now).to be false + expect(recorder.errors[0]).to be_a(MergeError) + end + + with_reloader([a, b], duplicate_keys_handling: DuplicateKeysHandling::IGNORE) do |reloader, recorder| + expect(reloader.reload_now).to be true + expect(recorder.flag_values).to eq({ flag1: "first" }) + end + end + + it "retries a failed reload after the retry delay and recovers when the file is fixed" do + a = write("a.json", '{"flagValues"') + + with_reloader([a], retry_delay: 0.1) do |reloader, recorder| + expect(reloader.reload_now).to be false + write("a.json", values_doc({ flag1: "fixed" })) + + expect(wait_for { recorder.applied.length == 1 }).to be true + expect(recorder.flag_values).to eq({ flag1: "fixed" }) + end + end + + it "keeps retrying while the failure persists" do + a = write("a.json", '{"flagValues"') + logger = CapturingLogger.new + + with_reloader([a], logger: logger, retry_delay: 0.05) do |reloader, _recorder| + reloader.reload_now + + expect(wait_for { logger.output.scan("Retrying flag data load").length >= 3 }).to be true + end + end + + it "does not retry when the retry delay is zero" do + a = write("a.json", '{"flagValues"') + + with_reloader([a], retry_delay: 0) do |reloader, recorder| + reloader.reload_now + write("a.json", values_doc({ flag1: "fixed" })) + sleep 0.3 + + expect(recorder.applied).to be_empty + end + end + + it "coalesces triggers that arrive within the debounce delay into one reload" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], debounce_delay: 0.2) do |reloader, recorder| + reloader.reload_now + write("a.json", values_doc({ flag1: "b" })) + 5.times { reloader.trigger } + + expect(wait_for { recorder.applied.length == 2 }).to be true + sleep 0.4 + expect(recorder.applied.length).to eq(2) + expect(recorder.flag_values).to eq({ flag1: "b" }) + end + end + + it "restarts the debounce window when a trigger arrives during it" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], debounce_delay: 0.2) do |reloader, recorder| + reloader.reload_now + write("a.json", values_doc({ flag1: "b" })) + reloader.trigger + sleep 0.12 + reloader.trigger + sleep 0.12 + + expect(recorder.applied.length).to eq(1) + expect(wait_for { recorder.applied.length == 2 }).to be true + end + end + + it "reloads at once for each trigger when the debounce delay is zero" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], debounce_delay: 0) do |reloader, recorder| + reloader.reload_now + write("a.json", values_doc({ flag1: "b" })) + reloader.trigger + expect(wait_for(0.5) { recorder.applied.length == 2 }).to be true + write("a.json", values_doc({ flag1: "c" })) + reloader.trigger + expect(wait_for(0.5) { recorder.applied.length == 3 }).to be true + expect(recorder.flag_values).to eq({ flag1: "c" }) + end + end + + it "applies a triggered reload even when the content is unchanged unless told to skip" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], debounce_delay: 0) do |reloader, recorder| + reloader.reload_now + reloader.trigger + + expect(wait_for { recorder.applied.length == 2 }).to be true + end + end + + it "skips a reload whose content is identical to the last applied content" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], debounce_delay: 0, skip_unchanged: true) do |reloader, recorder| + reloader.reload_now + reloader.trigger + sleep 0.3 + expect(recorder.applied.length).to eq(1) + + write("a.json", values_doc({ flag1: "b" })) + reloader.trigger + expect(wait_for { recorder.applied.length == 2 }).to be true + end + end + + it "applies a success that follows a failure even when the content is unchanged" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], retry_delay: 0, skip_unchanged: true) do |reloader, recorder| + reloader.reload_now + write("a.json", '{"flagValues"') + reloader.reload_now + write("a.json", values_doc({ flag1: "a" })) + reloader.reload_now + + expect(recorder.applied.length).to eq(2) + expect(recorder.errors.length).to eq(1) + end + end + + it "runs reloads one at a time even when triggers overlap" do + a = write("a.json", values_doc({ flag1: "a" })) + active = Concurrent::AtomicFixnum.new(0) + max_active = Concurrent::AtomicFixnum.new(0) + applied = Concurrent::AtomicFixnum.new(0) + apply = lambda do |_merged| + current = active.increment + max_active.update { |m| [m, current].max } + sleep 0.1 + active.decrement + applied.increment + end + reloader = Reloader.new(paths: [a], logger: $null_log, apply: apply, debounce_delay: 0) + + begin + threads = Array.new(4) do + Thread.new do + reloader.reload_now + reloader.trigger + end + end + threads.each(&:join) + + expect(wait_for { applied.value >= 5 }).to be true + expect(max_active.value).to eq(1) + ensure + reloader.stop + end + end + + it "does not apply after it is stopped" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], debounce_delay: 0) do |reloader, recorder| + reloader.reload_now + reloader.stop + write("a.json", values_doc({ flag1: "b" })) + reloader.trigger + expect(reloader.reload_now).to be true + sleep 0.2 + + expect(recorder.applied.length).to eq(1) + end + end + + it "does not report a failure after it is stopped" do + a = write("a.json", '{"flagValues"') + + with_reloader([a]) do |reloader, recorder| + reloader.stop + expect(reloader.reload_now).to be true + + expect(recorder.errors).to be_empty + end + end + + it "logs a change-triggered reload at info level" do + a = write("a.json", values_doc({ flag1: "a" })) + logger = CapturingLogger.new + + with_reloader([a], logger: logger, debounce_delay: 0) do |reloader, recorder| + reloader.reload_now + reloader.trigger + expect(wait_for { recorder.applied.length == 2 }).to be true + + expect(logger.output).to include("Reloading flag data after detecting a change") + end + end + end + end + end +end diff --git a/spec/impl/file_data/watcher_spec.rb b/spec/impl/file_data/watcher_spec.rb new file mode 100644 index 00000000..a531abab --- /dev/null +++ b/spec/impl/file_data/watcher_spec.rb @@ -0,0 +1,169 @@ +# frozen_string_literal: true + +require "spec_helper" +require "capturing_logger" +require "tmpdir" +require "ldclient-rb/impl/file_data" + +module LaunchDarkly + module Impl + module FileData + describe Watcher do + before do + skip "the listen gem is not installed" unless Watcher.available? + end + + around do |example| + Dir.mktmpdir do |dir| + @dir = dir + example.run + end + end + + def path(name) + File.join(@dir, name) + end + + def wait_for(timeout = 5) + deadline = Time.now + timeout + until yield + return false if Time.now > deadline + sleep 0.02 + end + true + end + + def with_watcher(paths, logger: $null_log) + calls = Concurrent::AtomicFixnum.new(0) + watcher = Watcher.new(paths, -> { calls.increment }, logger) + begin + yield watcher, calls + ensure + watcher.stop + end + end + + it "reports that the listen gem is available" do + expect(Watcher.available?).to be true + end + + it "invokes the callback when a watched file is modified" do + File.write(path("a.json"), "{}") + with_watcher([path("a.json")]) do |_watcher, calls| + sleep 0.3 + File.write(path("a.json"), '{"flagValues": {}}') + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "invokes the callback when a watched file that did not exist appears" do + with_watcher([path("later.json")]) do |_watcher, calls| + sleep 0.3 + File.write(path("later.json"), "{}") + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "invokes the callback when a watched file is deleted" do + File.write(path("a.json"), "{}") + with_watcher([path("a.json")]) do |_watcher, calls| + sleep 0.3 + File.delete(path("a.json")) + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "ignores other files in the same directory" do + File.write(path("a.json"), "{}") + with_watcher([path("a.json")]) do |_watcher, calls| + sleep 0.3 + File.write(path("other.json"), "{}") + sleep 0.5 + expect(calls.value).to eq(0) + end + end + + it "watches files in more than one directory" do + Dir.mkdir(path("sub")) + File.write(path("a.json"), "{}") + File.write(path("sub/b.json"), "{}") + with_watcher([path("a.json"), path("sub/b.json")]) do |_watcher, calls| + sleep 0.3 + File.write(path("sub/b.json"), '{"flagValues": {}}') + expect(wait_for { calls.value >= 1 }).to be true + end + end + + it "logs when a directory does not exist and starts watching once it appears" do + logger = CapturingLogger.new + missing_dir = path("not-yet") + with_watcher([File.join(missing_dir, "a.json")], logger: logger) do |_watcher, calls| + expect(logger.output).to include("Unable to watch data files") + expect(logger.output).to include("directory does not exist") + sleep 0.2 + expect(calls.value).to eq(0) + + Dir.mkdir(missing_dir) + # The watches are set up on the next retry, and the callback runs once at that point. + expect(wait_for { calls.value >= 1 }).to be true + before = calls.value + + sleep 0.3 + File.write(File.join(missing_dir, "a.json"), "{}") + expect(wait_for { calls.value > before }).to be true + end + end + + it "watches a directory that contains an unreadable subdirectory" do + skip "the current user can read every directory" if Process.uid.zero? + + Dir.mkdir(path("private")) + File.chmod(0o000, path("private")) + begin + File.write(path("a.json"), "{}") + logger = CapturingLogger.new + with_watcher([path("a.json")], logger: logger) do |_watcher, calls| + expect(logger.output).not_to include("Unable to watch data files") + sleep 0.3 + File.write(path("a.json"), '{"flagValues": {}}') + expect(wait_for { calls.value >= 1 }).to be true + end + ensure + File.chmod(0o755, path("private")) + end + end + + it "uses inotify on Linux and runs it on a named thread that stop ends" do + skip "rb-inotify is not available on this platform" unless Watcher.inotify_available? + + File.write(path("a.json"), "{}") + watcher = Watcher.new([path("a.json")], -> {}, $null_log) + threads = Thread.list.select { |t| t.name == "LD/FileDataWatcher" } + expect(threads.length).to eq 1 + + watcher.stop + + expect(threads[0].alive?).to be false + end + + it "does not invoke the callback after it is stopped" do + File.write(path("a.json"), "{}") + with_watcher([path("a.json")]) do |watcher, calls| + sleep 0.3 + watcher.stop + File.write(path("a.json"), '{"flagValues": {}}') + sleep 0.5 + expect(calls.value).to eq(0) + end + end + + it "can be stopped while it is still retrying a missing directory" do + with_watcher([path("not-yet/a.json")]) do |watcher, _calls| + watcher.stop + expect(Thread.list.map(&:name)).not_to include("LD/FileDataWatcherRetry") + end + end + end + end + end +end diff --git a/spec/integrations/file_data_source_compatibility_spec.rb b/spec/integrations/file_data_source_compatibility_spec.rb new file mode 100644 index 00000000..aefa5ba3 --- /dev/null +++ b/spec/integrations/file_data_source_compatibility_spec.rb @@ -0,0 +1,387 @@ +require "spec_helper" +require "capturing_logger" +require "tempfile" +require "ldclient-rb/integrations/file_data" + +# +# These specs pin behaviors of the existing file data sources that their other specs do not +# assert. The file-based override source shares the document format but not the implementation, +# and these sources must keep behaving as they always have. +# +module LaunchDarkly + module Integrations + describe "file data source behavior" do + let(:logger) { CapturingLogger.new } + + before do + @tmp_dir = Dir.mktmpdir + end + + after do + FileUtils.rm_rf(@tmp_dir) + end + + def make_temp_file(content) + file = Tempfile.new('flags', @tmp_dir) + IO.write(file, content) + file + end + + def wait_for(timeout = 5) + deadline = Time.now + timeout + until yield + return false if Time.now > deadline + sleep 0.05 + end + true + end + + def flag_json(key, version: nil) + data = { key: key, on: true, fallthrough: { variation: 0 }, variations: ["a"] } + data[:version] = version unless version.nil? + data + end + + describe FileData, "data_source" do + # Counts every init call so that a test can see each reload. + class CountingFeatureStore < InMemoryFeatureStore + attr_reader :init_count + + def initialize + super + @init_count = 0 + end + + def init(all_data) + @init_count += 1 + super + end + end + + before do + @store = CountingFeatureStore.new + @config = LaunchDarkly::Config.new(logger: logger, feature_store: @store) + executor = SynchronousExecutor.new + @status_broadcaster = LaunchDarkly::Impl::Broadcaster.new(executor, logger) + @flag_change_broadcaster = LaunchDarkly::Impl::Broadcaster.new(executor, logger) + @config.data_source_update_sink = LaunchDarkly::Impl::DataSource::UpdateSink.new(@store, @status_broadcaster, @flag_change_broadcaster) + end + + def with_data_source(options) + ds = FileData.data_source(options).call('', @config) + begin + yield ds + ensure + ds.stop + end + end + + def flag_version(key) + @store.get(Impl::DataStore::FEATURES, key).version + end + + def status + @config.data_source_update_sink.current_status + end + + it "numbers versions per file in load order and keeps counting across reloads" do + file1 = make_temp_file({ flags: { flag1: flag_json("flag1", version: 99) }, flagValues: { value1: true } }.to_json) + file2 = make_temp_file({ segments: { seg1: { key: "seg1" } } }.to_json) + + with_data_source({ paths: [file1.path, file2.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| + ds.start + expect(flag_version("flag1")).to eq 1 + expect(flag_version("value1")).to eq 1 + expect(@store.get(Impl::DataStore::SEGMENTS, "seg1").version).to eq 2 + + sleep 0.2 + IO.write(file2, { segments: { seg1: { key: "seg1" }, seg2: { key: "seg2" } } }.to_json) + expect(wait_for { @store.get(Impl::DataStore::SEGMENTS, "seg2") }).to be true + expect(flag_version("flag1")).to eq 3 + expect(@store.get(Impl::DataStore::SEGMENTS, "seg2").version).to eq 4 + end + end + + it "stores an item under its own key member rather than the key it appears under" do + file = make_temp_file({ flags: { "map-key": flag_json("own-key") } }.to_json) + + with_data_source({ paths: [file.path] }) do |ds| + ds.start + expect(@store.get(Impl::DataStore::FEATURES, "own-key")).not_to be_nil + expect(@store.get(Impl::DataStore::FEATURES, "map-key")).to be_nil + end + end + + it "expands a flag value into a flag that is on and serves the value through its fallthrough" do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + with_data_source({ paths: [file.path] }) do |ds| + ds.start + flag = @store.get(Impl::DataStore::FEATURES, "value1") + expect(flag.on).to be true + expect(flag.fallthrough.variation).to eq 0 + expect(flag.off_variation).to be_nil + expect(flag.variations).to eq ["x"] + end + end + + it "does not retry a failed initial load on its own" do + file = make_temp_file('{"flagValues"') + + with_data_source({ paths: [file.path] }) do |ds| + ds.start + expect(ds.initialized?).to be false + IO.write(file, { flagValues: { value1: "x" } }.to_json) + sleep 1.5 + + expect(ds.initialized?).to be false + expect(@store.init_count).to eq 0 + end + end + + it "logs a failed load with the file path at error level" do + file = make_temp_file('{"flagValues"') + + with_data_source({ paths: [file.path] }) do |ds| + ds.start + expect(logger.output).to match(/ERROR -- : \[LDClient\] Unable to load flag data from "#{Regexp.escape(file.path)}": /) + expect(status.state).to eq Interfaces::DataSource::Status::INITIALIZING + expect(status.last_error.kind).to eq Interfaces::DataSource::ErrorInfo::INVALID_DATA + end + end + + it "treats a missing file as a failed load with the same message" do + missing = File.join(@tmp_dir, "no-such-file.json") + + with_data_source({ paths: [missing] }) do |ds| + ds.start + expect(ds.initialized?).to be false + expect(logger.output).to include("Unable to load flag data from \"#{missing}\"") + end + end + + it "reports a duplicate key with the data kind namespace" do + file1 = make_temp_file({ flags: { flag1: flag_json("flag1") } }.to_json) + file2 = make_temp_file({ flagValues: { flag1: "x" } }.to_json) + + with_data_source({ paths: [file1.path, file2.path] }) do |ds| + ds.start + expect(logger.output).to include('features key "flag1" was used more than once') + expect(@store.init_count).to eq 0 + end + end + + it "compares only the modification time when polling" do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + with_data_source({ paths: [file.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| + ds.start + mtime = File.mtime(file.path) + IO.write(file, { flagValues: { value1: "a much longer value than before" } }.to_json) + File.utime(mtime, mtime, file.path) + sleep 0.5 + + expect(@store.init_count).to eq 1 + end + end + + it "does not react to a deleted file when polling" do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + with_data_source({ paths: [file.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| + ds.start + File.delete(file.path) + sleep 0.5 + + expect(@store.init_count).to eq 1 + expect(status.state).to eq Interfaces::DataSource::Status::VALID + expect(@store.get(Impl::DataStore::FEATURES, "value1")).not_to be_nil + end + end + + it "reloads on every polling interval after the first change" do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + with_data_source({ paths: [file.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| + ds.start + sleep 0.2 + IO.write(file, { flagValues: { value1: "y" } }.to_json) + expect(wait_for { @store.init_count >= 2 }).to be true + + expect(wait_for { @store.init_count >= 4 }).to be true + end + end + + it "uses the listen gem for auto-update when it is available" do + skip "the listen gem is not installed" unless defined?(Listen) + + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + expect(Listen).to receive(:to).and_call_original + + with_data_source({ paths: [file.path], auto_update: true }) do |ds| + ds.start + end + end + + it "does not use the listen gem when polling is forced" do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + expect(Listen).not_to receive(:to) if defined?(Listen) + + with_data_source({ paths: [file.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| + ds.start + expect(Thread.list.map(&:name)).to include("LD/FileDataSource") + end + end + end + + describe Impl::Integrations::FileDataSourceV2 do + def no_selector_store + store = Object.new + store.define_singleton_method(:selector) { Interfaces::DataSystem::Selector.no_selector } + store + end + + def fetch_changes(paths) + source = Impl::Integrations::FileDataSourceV2.new(logger, paths: paths) + begin + result = source.fetch(no_selector_store) + expect(result.success?).to be true + result.value.change_set.changes + ensure + source.stop + end + end + + def without_listen + klass = Impl::Integrations::FileDataSourceV2 + had_listen = klass.class_variable_get(:@@have_listen) + klass.class_variable_set(:@@have_listen, false) + begin + yield + ensure + klass.class_variable_set(:@@have_listen, had_listen) + end + end + + def with_sync(paths, poll_interval: 0.1) + source = Impl::Integrations::FileDataSourceV2.new(logger, paths: paths, poll_interval: poll_interval) + updates = Queue.new + thread = Thread.new { source.sync(no_selector_store) { |update| updates << update } } + begin + initial = updates.pop(timeout: 5) + expect(initial).not_to be_nil + expect(initial.state).to eq Interfaces::DataSource::Status::VALID + yield updates + ensure + source.stop + thread.join(2) + end + end + + it "defaults a missing version to 1 and keeps a given version" do + file = make_temp_file({ flags: { flag1: flag_json("flag1"), flag2: flag_json("flag2", version: 7) }, + flagValues: { value1: "x" }, segments: { seg1: { key: "seg1" } } }.to_json) + + versions = fetch_changes([file.path]).to_h { |change| [change.key, change.version] } + + expect(versions).to eq({ flag1: 1, flag2: 7, value1: 1, seg1: 1 }) + end + + it "expands a flag value into a flag that is on and serves the value through its fallthrough" do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + change = fetch_changes([file.path]).detect { |c| c.key == :value1 } + expect(change.object[:on]).to be true + expect(change.object[:fallthrough]).to eq({ variation: 0 }) + expect(change.object).not_to have_key(:offVariation) + expect(change.object[:variations]).to eq ["x"] + end + + it "stores an item under its own key member rather than the key it appears under" do + file = make_temp_file({ flags: { "map-key": flag_json("own-key") } }.to_json) + + expect(fetch_changes([file.path]).map(&:key)).to eq [:"own-key"] + end + + it "reports a failed load with the file path" do + file = make_temp_file('{"flagValues"') + source = Impl::Integrations::FileDataSourceV2.new(logger, paths: [file.path]) + begin + result = source.fetch(no_selector_store) + + expect(result.success?).to be false + expect(result.error).to start_with("Unable to load flag data from \"#{file.path}\": ") + expect(logger.output).to match(/ERROR -- : \[LDClient\] Unable to load flag data from "#{Regexp.escape(file.path)}": /) + ensure + source.stop + end + end + + it "reports a duplicate key with the section name" do + file1 = make_temp_file({ flags: { flag1: flag_json("flag1") } }.to_json) + file2 = make_temp_file({ flagValues: { flag1: "x" } }.to_json) + source = Impl::Integrations::FileDataSourceV2.new(logger, paths: [file1.path, file2.path]) + begin + result = source.fetch(no_selector_store) + + expect(result.success?).to be false + expect(result.error).to include('In flags, key "flag1" was used more than once') + ensure + source.stop + end + end + + it "uses the listen gem for change detection when it is available" do + skip "the listen gem is not installed" unless defined?(Listen) + + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + expect(Listen).to receive(:to).and_call_original + + with_sync([file.path]) { |_updates| } + end + + it "compares only the modification time when polling" do + without_listen do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + with_sync([file.path]) do |updates| + mtime = File.mtime(file.path) + IO.write(file, { flagValues: { value1: "a much longer value than before" } }.to_json) + File.utime(mtime, mtime, file.path) + + expect(updates.pop(timeout: 0.6)).to be_nil + end + end + end + + it "does not react to a deleted file when polling" do + without_listen do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + with_sync([file.path]) do |updates| + File.delete(file.path) + + expect(updates.pop(timeout: 0.6)).to be_nil + end + end + end + + it "reloads once per change when polling" do + without_listen do + file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + + with_sync([file.path]) do |updates| + sleep 0.2 + IO.write(file, { flagValues: { value1: "y" } }.to_json) + + update = updates.pop(timeout: 5) + expect(update).not_to be_nil + expect(update.state).to eq Interfaces::DataSource::Status::VALID + expect(updates.pop(timeout: 0.6)).to be_nil + end + end + end + end + end + end +end From 821186b32eb30ccdd88af17c49864419c4ab6061 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Thu, 1 Oct 2026 23:06:40 +0000 Subject: [PATCH 02/11] fix: Expand value-only overrides to a flag served by fallthrough --- lib/ldclient-rb/impl/file_data/document.rb | 9 +++++---- spec/impl/file_data/document_spec.rb | 6 +++--- spec/impl/file_data/merge_spec.rb | 8 ++++---- 3 files changed, 12 insertions(+), 11 deletions(-) diff --git a/lib/ldclient-rb/impl/file_data/document.rb b/lib/ldclient-rb/impl/file_data/document.rb index 005a4752..53e820b9 100644 --- a/lib/ldclient-rb/impl/file_data/document.rb +++ b/lib/ldclient-rb/impl/file_data/document.rb @@ -156,8 +156,9 @@ def self.symbolize_keys(value) # # 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 the value as its off variation, so - # an evaluation reports the OFF reason kind. + # value for every context. The flag is on, has the value as its only variation, and + # serves that variation as its fallthrough, so an evaluation reports the FALLTHROUGH + # reason kind. # # @param key [String] # @param value [Object] @@ -166,9 +167,9 @@ def self.symbolize_keys(value) def self.make_flag_with_value(key, value) { key: key, - on: false, + on: true, version: 1, - offVariation: 0, + fallthrough: { variation: 0 }, variations: [value], } end diff --git a/spec/impl/file_data/document_spec.rb b/spec/impl/file_data/document_spec.rb index e2f565ea..7b738322 100644 --- a/spec/impl/file_data/document_spec.rb +++ b/spec/impl/file_data/document_spec.rb @@ -136,14 +136,14 @@ module FileData end describe "make_flag_with_value" do - it "builds a flag that is off and serves the value as its off variation" do + it "builds a flag that is on and serves the value as its fallthrough" do flag = FileData.make_flag_with_value("flag1", "value1") expect(flag).to eq({ key: "flag1", - on: false, + on: true, version: 1, - offVariation: 0, + fallthrough: { variation: 0 }, variations: ["value1"], }) end diff --git a/spec/impl/file_data/merge_spec.rb b/spec/impl/file_data/merge_spec.rb index 073af43f..2039f6c7 100644 --- a/spec/impl/file_data/merge_spec.rb +++ b/spec/impl/file_data/merge_spec.rb @@ -33,16 +33,16 @@ def full_flag(key, version: nil) expect(result.empty?).to be false end - it "expands a flag value into an off flag that serves the value for every context" do + it "expands a flag value into a flag that serves the value by fallthrough for every context" do result = FileData.merge([document(flag_values: { flag2: "value2" })]) flag = result.flags[:flag2] expect(flag.key).to eq("flag2") - expect(flag.on).to be false + expect(flag.on).to be true expect(flag.version).to eq(1) expect(flag.variations).to eq(["value2"]) - expect(flag.off_variation).to eq(0) - expect(flag.off_result.reason).to eq(EvaluationReason.off) + expect(flag.off_variation).to be_nil + expect(flag.fallthrough.variation).to eq(0) end it "reports an empty result for no documents" do From 4546f9bd5bb1970e5f07fa07d90f172ff09eb1c2 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:20:12 +0000 Subject: [PATCH 03/11] fix: Report an exception from apply as a failed reload The result is not remembered as applied, on_error receives the exception, the reload is retried, and the worker thread continues. --- lib/ldclient-rb/impl/file_data/reloader.rb | 18 +++++-- spec/impl/file_data/reloader_spec.rb | 59 ++++++++++++++++++++-- 2 files changed, 68 insertions(+), 9 deletions(-) diff --git a/lib/ldclient-rb/impl/file_data/reloader.rb b/lib/ldclient-rb/impl/file_data/reloader.rb index a080dd9b..a4d5c19a 100644 --- a/lib/ldclient-rb/impl/file_data/reloader.rb +++ b/lib/ldclient-rb/impl/file_data/reloader.rb @@ -36,11 +36,14 @@ class Reloader # @param logger [Logger] # @param apply [#call] invoked with each successfully merged {MergeResult}. Calls are # serialized, so implementations do not need their own synchronization against other - # reloads. `apply` and `on_error` must not call back into {#stop}. + # reloads. An exception from `apply` is a failed reload: it is reported through + # `on_error`, the result is not remembered as applied, and the reload is retried. + # `apply` and `on_error` must not call back into {#stop}. # @param on_error [#call, nil] invoked with the error when a reload fails, once per distinct # failure. With automatic retries, repeats of an identical failure do not re-invoke it. A # success re-arms it. The error is a {ReadError} when a file could not be read or parsed, - # or a {MergeError} otherwise. The reloader logs failures itself. + # a {MergeError} when the documents could not be combined, or the exception that `apply` + # raised. The reloader logs failures itself. # @param duplicate_keys_handling [Symbol] one of the {DuplicateKeysHandling} values # @param skip_missing_paths [Boolean] when true, a configured file that does not exist is a # file with no content, and the reload succeeds with the data of the files that exist. @@ -265,12 +268,19 @@ def stop # the last success. The consumer heard about the failure through on_error and may # have moved to an interrupted state. Only apply tells it that things are good again. recovering = !@last_error_message.nil? - @last_error_message = nil hexdigest = digest.hexdigest return true if @skip_unchanged && !recovering && hexdigest == @last_good_digest + # Nothing is remembered until the consumer has accepted the result. A result that + # apply rejects must not become the baseline that skip-unchanged compares against, + # and must not count as a recovery. + begin + @apply.call(merged) + rescue => e + return record_failure(e) + end + @last_error_message = nil @last_good_digest = hexdigest - @apply.call(merged) true end end diff --git a/spec/impl/file_data/reloader_spec.rb b/spec/impl/file_data/reloader_spec.rb index 241ad308..962d6a06 100644 --- a/spec/impl/file_data/reloader_spec.rb +++ b/spec/impl/file_data/reloader_spec.rb @@ -38,18 +38,26 @@ def wait_for(timeout = 3) true end - # Collects apply and on_error calls in a thread-safe way. + # Collects apply and on_error calls in a thread-safe way. It raises from the first `reject` + # apply calls, as a consumer that cannot use a result does. class Recorder attr_reader :applied, :errors - def initialize + def initialize(reject: 0) @lock = Mutex.new @applied = [] @errors = [] + @reject = reject end def apply(merged) - @lock.synchronize { @applied << merged } + @lock.synchronize do + if @reject > 0 + @reject -= 1 + raise ArgumentError, "the result was rejected" + end + @applied << merged + end end def on_error(error) @@ -65,8 +73,8 @@ def make_reloader(paths, recorder, logger: $null_log, **options) Reloader.new(paths: paths, logger: logger, apply: recorder.method(:apply), on_error: recorder.method(:on_error), **options) end - def with_reloader(paths, logger: $null_log, **options) - recorder = Recorder.new + def with_reloader(paths, logger: $null_log, reject: 0, **options) + recorder = Recorder.new(reject: reject) reloader = make_reloader(paths, recorder, logger: logger, **options) begin yield reloader, recorder @@ -304,6 +312,47 @@ def with_reloader(paths, logger: $null_log, **options) end end + it "reports an exception from apply as a failed reload and applies the same content on the retry" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], reject: 1, retry_delay: 0.1) do |reloader, recorder| + expect(reloader.reload_now).to be false + expect(recorder.applied).to be_empty + expect(recorder.errors.length).to eq(1) + expect(recorder.errors[0]).to be_a(ArgumentError) + + # The retry applies the content that was rejected, so a consumer that recovers sees it. + expect(wait_for { recorder.applied.length == 1 }).to be true + expect(recorder.flag_values).to eq({ flag1: "a" }) + end + end + + it "does not skip content as unchanged after apply rejected it" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], reject: 1, retry_delay: 0, skip_unchanged: true) do |reloader, recorder| + expect(reloader.reload_now).to be false + + # The same content again. It was never applied, so it is not skipped as unchanged. + expect(reloader.reload_now).to be true + expect(recorder.flag_values).to eq({ flag1: "a" }) + end + end + + it "reloads on a later trigger after apply raised during a triggered reload" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], reject: 1, debounce_delay: 0, retry_delay: 0) do |reloader, recorder| + reloader.trigger + expect(wait_for { recorder.errors.length == 1 }).to be true + + write("a.json", values_doc({ flag1: "b" })) + reloader.trigger + expect(wait_for { recorder.applied.length == 1 }).to be true + expect(recorder.flag_values).to eq({ flag1: "b" }) + end + end + it "runs reloads one at a time even when triggers overlap" do a = write("a.json", values_doc({ flag1: "a" })) active = Concurrent::AtomicFixnum.new(0) From 4f0af3921c25e093169c43466a9d517ae0151458 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:20:12 +0000 Subject: [PATCH 04/11] fix: Watch a data file directory again after it is deleted The inotify watcher tears its watches down when a watched directory is deleted or moved, logs it, retries setup once per second, and signals one change when the watches are in place again. --- lib/ldclient-rb/impl/file_data/watcher.rb | 81 ++++++++++++++++++----- spec/impl/file_data/watcher_spec.rb | 24 +++++++ 2 files changed, 88 insertions(+), 17 deletions(-) diff --git a/lib/ldclient-rb/impl/file_data/watcher.rb b/lib/ldclient-rb/impl/file_data/watcher.rb index bf079f26..625950f0 100644 --- a/lib/ldclient-rb/impl/file_data/watcher.rb +++ b/lib/ldclient-rb/impl/file_data/watcher.rb @@ -19,9 +19,10 @@ module FileData # scans the whole directory tree under each watched directory. # # The watcher observes the directory of each file, so a configured file that does not exist - # yet is picked up when it appears. If a directory does not exist, the watcher logs the - # problem and retries once per second until it does. When the watches are in place after a - # retry, the callback runs once, so that a change made while there was no watch is not missed. + # yet is picked up when it appears. If a directory does not exist, or is deleted while it is + # watched, the watcher logs the problem and retries once per second until it exists. When the + # watches are in place after a retry, the callback runs once, so that a change made while + # there was no watch is not missed. # # @private # @@ -29,9 +30,14 @@ class Watcher # Seconds between attempts to set up the watches after a failure. RETRY_INTERVAL = 1.0 - INOTIFY_EVENTS = [:create, :modify, :close_write, :attrib, :delete, :moved_to, :moved_from].freeze + INOTIFY_EVENTS = [:create, :modify, :close_write, :attrib, :delete, :moved_to, :moved_from, + :delete_self, :move_self].freeze private_constant :INOTIFY_EVENTS + # The inotify event flags that mean the watched directory itself is gone. + INOTIFY_DIRECTORY_LOST = [:delete_self, :move_self].freeze + private_constant :INOTIFY_DIRECTORY_LOST + # # Returns true if a change notification mechanism can be loaded. # @@ -86,11 +92,7 @@ def initialize(paths, on_change, logger) @retry_task = nil @last_error_message = nil - return if try_start - - @retry_task = RepeatingTask.new(RETRY_INTERVAL, RETRY_INTERVAL, method(:retry_start), logger, - "LD/FileDataWatcherRetry") - @retry_task.start + schedule_retry unless try_start end # @@ -100,24 +102,62 @@ def initialize(paths, on_change, logger) def stop return unless @stopped.make_true - @retry_task&.stop - listener = @lock.synchronize do - l = @listener + listener, retry_task = @lock.synchronize do + pair = [@listener, @retry_task] @listener = nil - l + @retry_task = nil + pair end + retry_task&.stop listener&.stop end + # + # Starts the task that attempts to set up the watches once per second, unless it already + # runs or the watcher is stopped. + # + private def schedule_retry + @lock.synchronize do + return if @stopped.value || !@retry_task.nil? + + @retry_task = RepeatingTask.new(RETRY_INTERVAL, RETRY_INTERVAL, method(:retry_start), @logger, + "LD/FileDataWatcherRetry") + @retry_task.start + end + end + private def retry_start return if @stopped.value return unless try_start + retry_task = @lock.synchronize do + task = @retry_task + @retry_task = nil + task + end # This runs on the retry task's own thread, which RepeatingTask#stop allows. - @retry_task.stop + retry_task&.stop @on_change.call unless @stopped.value end + # + # Handles the loss of a watched directory. The notification mechanism reports it on its + # own thread, and the watches end with the directory, so they are torn down and set up + # again through the same retry as at start, once the directory exists. + # + private def directory_lost(directory) + listener = @lock.synchronize do + return if @stopped.value || @listener.nil? + + l = @listener + @listener = nil + l + end + @logger.warn { "[LDClient] Directory #{directory} no longer exists; its data files are watched again when it exists" } + listener.stop + schedule_retry + end + # # Sets up the watches. Returns false, after logging, if that is not possible yet. # @@ -160,7 +200,13 @@ def stop begin names_by_directory.each do |directory, names| notifier.watch(directory, *INOTIFY_EVENTS) do |event| - @on_change.call if names.include?(event.name) && !@stopped.value + next if @stopped.value + + if (event.flags & INOTIFY_DIRECTORY_LOST).empty? + @on_change.call if names.include?(event.name) + else + directory_lost(directory) + end end end rescue @@ -216,12 +262,13 @@ def initialize(notifier, logger) # # Stops the notifier and waits briefly for its thread. Closing the notifier ends the - # blocking read that the thread is in. + # blocking read that the thread is in. A callback can call this on the notifier's own + # thread, which then ends when the callback returns. # def stop @notifier.stop @notifier.close - @thread.join(2) + @thread.join(2) unless Thread.current == @thread end end private_constant :InotifyListener diff --git a/spec/impl/file_data/watcher_spec.rb b/spec/impl/file_data/watcher_spec.rb index a531abab..b3b0e64c 100644 --- a/spec/impl/file_data/watcher_spec.rb +++ b/spec/impl/file_data/watcher_spec.rb @@ -2,6 +2,7 @@ require "spec_helper" require "capturing_logger" +require "fileutils" require "tmpdir" require "ldclient-rb/impl/file_data" @@ -114,6 +115,29 @@ def with_watcher(paths, logger: $null_log) end end + it "logs when a watched directory is deleted and watches it again once it exists again" do + logger = CapturingLogger.new + Dir.mkdir(path("sub")) + File.write(path("sub/a.json"), "{}") + with_watcher([path("sub/a.json")], logger: logger) do |_watcher, calls| + sleep 0.3 + FileUtils.rm_rf(path("sub")) + expect(wait_for { logger.output.include?("WARN") }).to be true + expect(logger.output).to match(/WARN.*#{Regexp.escape(path('sub'))}/) + before = calls.value + + Dir.mkdir(path("sub")) + File.write(path("sub/a.json"), "{}") + # The watches are set up on the next retry, and the callback runs once at that point. + expect(wait_for { calls.value > before }).to be true + before = calls.value + + sleep 0.3 + File.write(path("sub/a.json"), '{"flagValues": {}}') + expect(wait_for { calls.value > before }).to be true + end + end + it "watches a directory that contains an unreadable subdirectory" do skip "the current user can read every directory" if Process.uid.zero? From 9b7d0d05122b5c301457ce2a89cdba71e9df4c8b Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:20:12 +0000 Subject: [PATCH 05/11] fix: Report an entry that cannot be deserialized as a merge error --- lib/ldclient-rb/impl/file_data/merge.rb | 19 +++++++++++++++---- spec/impl/file_data/merge_spec.rb | 5 +++++ 2 files changed, 20 insertions(+), 4 deletions(-) diff --git a/lib/ldclient-rb/impl/file_data/merge.rb b/lib/ldclient-rb/impl/file_data/merge.rb index 470b8af3..270f84a3 100644 --- a/lib/ldclient-rb/impl/file_data/merge.rb +++ b/lib/ldclient-rb/impl/file_data/merge.rb @@ -22,8 +22,8 @@ module DuplicateKeysHandling end # - # Raised when documents cannot be combined, for example because a key is duplicated or an - # entry is not an object. + # Raised when documents cannot be combined, for example because a key is duplicated, an + # entry is not an object, or an entry cannot be deserialized into the data model. # class MergeError < StandardError end @@ -93,7 +93,7 @@ def self.merge(documents, duplicate_keys_handling: DuplicateKeysHandling::FAIL, document.flags.each do |key, data| data = prepare_entry("flag", key, data) - item = Model.deserialize(DataStore::FEATURES, data, logger) + item = deserialize(DataStore::FEATURES, "flag", key, data, logger) summary.flags += 1 if insert(flags, "flag", key, item, duplicate_keys_handling) end @@ -105,7 +105,7 @@ def self.merge(documents, duplicate_keys_handling: DuplicateKeysHandling::FAIL, document.segments.each do |key, data| data = prepare_entry("segment", key, data) - item = Model.deserialize(DataStore::SEGMENTS, data, logger) + item = deserialize(DataStore::SEGMENTS, "segment", key, data, logger) summary.segments += 1 if insert(segments, "segment", key, item, duplicate_keys_handling) end @@ -127,6 +127,17 @@ def self.merge(documents, duplicate_keys_handling: DuplicateKeysHandling::FAIL, data end + # + # Deserializes one entry into its data model class. The model classes raise for a value + # that does not fit the schema, and that failure is reported as a merge error that names + # the entry. + # + private_class_method def self.deserialize(kind, category, key, data, logger) + Model.deserialize(kind, data, logger) + rescue => e + raise MergeError, "#{category} \"#{key}\" is not valid: #{e.message}" + end + # # Adds an item unless its key was already seen. Returns true when it added the item. # diff --git a/spec/impl/file_data/merge_spec.rb b/spec/impl/file_data/merge_spec.rb index 2039f6c7..9fd23cb5 100644 --- a/spec/impl/file_data/merge_spec.rb +++ b/spec/impl/file_data/merge_spec.rb @@ -113,6 +113,11 @@ def full_flag(key, version: nil) expect { FileData.merge([document(segments: { seg1: [] })]) }.to raise_error(MergeError, /segment "seg1" is not an object/) end + it "fails when an entry cannot be deserialized into the data model" do + expect { FileData.merge([document(flags: { flag1: { rules: 5 } })]) }.to raise_error(MergeError, /flag "flag1"/) + expect { FileData.merge([document(segments: { seg1: { rules: 5 } })]) }.to raise_error(MergeError, /segment "seg1"/) + end + it "passes the logger to model validation" do logger = CapturingLogger.new bad = { key: "flag1", on: true, variations: ["a"], fallthrough: { variation: 5 } } From 753eed3976826ff93ae30ef4415f80a13743d16c Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:20:12 +0000 Subject: [PATCH 06/11] test: Make the file data specs portable across platforms and runtimes Compare absolute paths through File.absolute_path, close temp files before deleting them, and backdate modification times instead of sleeping before a write that the pollers must see. --- spec/impl/file_data/document_spec.rb | 4 +++- .../file_data_source_compatibility_spec.rb | 16 +++++++++++++--- 2 files changed, 16 insertions(+), 4 deletions(-) diff --git a/spec/impl/file_data/document_spec.rb b/spec/impl/file_data/document_spec.rb index 7b738322..b5a20703 100644 --- a/spec/impl/file_data/document_spec.rb +++ b/spec/impl/file_data/document_spec.rb @@ -151,8 +151,10 @@ module FileData describe "absolute_paths" do it "converts relative paths to absolute paths and accepts a single string" do + absolute = File.absolute_path("/x/y.json") + expect(FileData.absolute_paths("a/b.json")).to eq([File.absolute_path("a/b.json")]) - expect(FileData.absolute_paths(["/x/y.json", "z.json"])).to eq(["/x/y.json", File.absolute_path("z.json")]) + expect(FileData.absolute_paths([absolute, "z.json"])).to eq([absolute, File.absolute_path("z.json")]) end end end diff --git a/spec/integrations/file_data_source_compatibility_spec.rb b/spec/integrations/file_data_source_compatibility_spec.rb index aefa5ba3..473f9dea 100644 --- a/spec/integrations/file_data_source_compatibility_spec.rb +++ b/spec/integrations/file_data_source_compatibility_spec.rb @@ -27,6 +27,14 @@ def make_temp_file(content) file end + # Moves a file's modification time into the past. The pollers compare modification times, + # so a write that follows at once registers as a change even where the file system records + # the time in whole seconds. + def backdate(file) + past = Time.now - 10 + File.utime(past, past, file.path) + end + def wait_for(timeout = 5) deadline = Time.now + timeout until yield @@ -87,6 +95,7 @@ def status it "numbers versions per file in load order and keeps counting across reloads" do file1 = make_temp_file({ flags: { flag1: flag_json("flag1", version: 99) }, flagValues: { value1: true } }.to_json) file2 = make_temp_file({ segments: { seg1: { key: "seg1" } } }.to_json) + backdate(file2) with_data_source({ paths: [file1.path, file2.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| ds.start @@ -94,7 +103,6 @@ def status expect(flag_version("value1")).to eq 1 expect(@store.get(Impl::DataStore::SEGMENTS, "seg1").version).to eq 2 - sleep 0.2 IO.write(file2, { segments: { seg1: { key: "seg1" }, seg2: { key: "seg2" } } }.to_json) expect(wait_for { @store.get(Impl::DataStore::SEGMENTS, "seg2") }).to be true expect(flag_version("flag1")).to eq 3 @@ -190,6 +198,7 @@ def status with_data_source({ paths: [file.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| ds.start + file.close File.delete(file.path) sleep 0.5 @@ -201,10 +210,10 @@ def status it "reloads on every polling interval after the first change" do file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + backdate(file) with_data_source({ paths: [file.path], auto_update: true, force_polling: true, poll_interval: 0.1 }) do |ds| ds.start - sleep 0.2 IO.write(file, { flagValues: { value1: "y" } }.to_json) expect(wait_for { @store.init_count >= 2 }).to be true @@ -359,6 +368,7 @@ def with_sync(paths, poll_interval: 0.1) file = make_temp_file({ flagValues: { value1: "x" } }.to_json) with_sync([file.path]) do |updates| + file.close File.delete(file.path) expect(updates.pop(timeout: 0.6)).to be_nil @@ -369,9 +379,9 @@ def with_sync(paths, poll_interval: 0.1) it "reloads once per change when polling" do without_listen do file = make_temp_file({ flagValues: { value1: "x" } }.to_json) + backdate(file) with_sync([file.path]) do |updates| - sleep 0.2 IO.write(file, { flagValues: { value1: "y" } }.to_json) update = updates.pop(timeout: 5) From 6c688e21c40dd173c5753e5be987e3f4940849b2 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:41:57 +0000 Subject: [PATCH 07/11] test: Run the deleted directory example only with inotify --- spec/impl/file_data/watcher_spec.rb | 2 ++ 1 file changed, 2 insertions(+) diff --git a/spec/impl/file_data/watcher_spec.rb b/spec/impl/file_data/watcher_spec.rb index b3b0e64c..5ffb0c4c 100644 --- a/spec/impl/file_data/watcher_spec.rb +++ b/spec/impl/file_data/watcher_spec.rb @@ -116,6 +116,8 @@ def with_watcher(paths, logger: $null_log) end it "logs when a watched directory is deleted and watches it again once it exists again" do + skip "rb-inotify is not available on this platform" unless Watcher.inotify_available? + logger = CapturingLogger.new Dir.mkdir(path("sub")) File.write(path("sub/a.json"), "{}") From 6f3e67eedf56688f49603574d6aa1ffde4720276 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:42:59 +0000 Subject: [PATCH 08/11] fix: Start the watcher retry again when a directory is lost while the retry finishes --- lib/ldclient-rb/impl/file_data/watcher.rb | 4 ++++ spec/impl/file_data/watcher_spec.rb | 27 +++++++++++++++++++++++ 2 files changed, 31 insertions(+) diff --git a/lib/ldclient-rb/impl/file_data/watcher.rb b/lib/ldclient-rb/impl/file_data/watcher.rb index 625950f0..4a26197d 100644 --- a/lib/ldclient-rb/impl/file_data/watcher.rb +++ b/lib/ldclient-rb/impl/file_data/watcher.rb @@ -137,6 +137,10 @@ def stop end # This runs on the retry task's own thread, which RepeatingTask#stop allows. retry_task&.stop + # The new watches can report the loss of their directory before the retry task is + # released above. That report finds the task still present and leaves the watches to + # it, so when they are gone the retry is started again here. + schedule_retry if @lock.synchronize { @listener.nil? } @on_change.call unless @stopped.value end diff --git a/spec/impl/file_data/watcher_spec.rb b/spec/impl/file_data/watcher_spec.rb index 5ffb0c4c..a0204148 100644 --- a/spec/impl/file_data/watcher_spec.rb +++ b/spec/impl/file_data/watcher_spec.rb @@ -140,6 +140,33 @@ def with_watcher(paths, logger: $null_log) end end + it "keeps retrying when a directory is lost again as soon as its watches are set up" do + skip "rb-inotify is not available on this platform" unless Watcher.inotify_available? + + with_watcher([path("sub/a.json")]) do |watcher, calls| + # Take the directory away as soon as the watches are set up, and wait until the loss + # has been handled, so that it is reported while the retry that set up the watches + # is still finishing. + lost_once = false + allow(watcher).to receive(:try_start).and_wrap_original do |original| + started = original.call + if started && !lost_once + lost_once = true + FileUtils.rm_rf(path("sub")) + wait_for { Thread.list.none? { |t| t.name == "LD/FileDataWatcher" } } + end + started + end + + Dir.mkdir(path("sub")) + expect(wait_for { calls.value >= 1 }).to be true + + # The retry that follows the second loss sets the watches up once the directory exists. + Dir.mkdir(path("sub")) + expect(wait_for { calls.value >= 2 }).to be true + end + end + it "watches a directory that contains an unreadable subdirectory" do skip "the current user can read every directory" if Process.uid.zero? From 2e31db3a2d191e29cde6774fcb853b8677deec2c Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:42:59 +0000 Subject: [PATCH 09/11] test: Cover the deduplication of a repeated rejection by apply --- spec/impl/file_data/reloader_spec.rb | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/spec/impl/file_data/reloader_spec.rb b/spec/impl/file_data/reloader_spec.rb index 962d6a06..08a734a9 100644 --- a/spec/impl/file_data/reloader_spec.rb +++ b/spec/impl/file_data/reloader_spec.rb @@ -353,6 +353,20 @@ def with_reloader(paths, logger: $null_log, reject: 0, **options) end end + it "reports a repeated rejection by apply once" do + a = write("a.json", values_doc({ flag1: "a" })) + + with_reloader([a], reject: 2, retry_delay: 0) do |reloader, recorder| + expect(reloader.reload_now).to be false + expect(reloader.reload_now).to be false + + # The second rejection repeats the first, so on_error is not invoked again. + expect(recorder.errors.length).to eq(1) + expect(reloader.reload_now).to be true + expect(recorder.applied.length).to eq(1) + end + end + it "runs reloads one at a time even when triggers overlap" do a = write("a.json", values_doc({ flag1: "a" })) active = Concurrent::AtomicFixnum.new(0) From ba5b65550405bd518b006199aa456b3c560a4f87 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 01:50:56 +0000 Subject: [PATCH 10/11] test: Wait for the legacy V2 source's change detection without Queue#pop timeouts Queue#pop accepts a timeout only from Ruby 3.2 on. The sync examples also wait for the poller thread or the listen call before they act, because the source starts change detection after it yields the initial data. --- .../file_data_source_compatibility_spec.rb | 42 +++++++++++++++---- 1 file changed, 34 insertions(+), 8 deletions(-) diff --git a/spec/integrations/file_data_source_compatibility_spec.rb b/spec/integrations/file_data_source_compatibility_spec.rb index 473f9dea..5ba5a65b 100644 --- a/spec/integrations/file_data_source_compatibility_spec.rb +++ b/spec/integrations/file_data_source_compatibility_spec.rb @@ -44,6 +44,12 @@ def wait_for(timeout = 5) true end + # Returns the next item from the queue, or nil when none arrives within the timeout. + # Queue#pop accepts a timeout only from Ruby 3.2 on, and the gem supports earlier runtimes. + def pop_with_timeout(queue, timeout) + wait_for(timeout) { !queue.empty? } ? queue.pop : nil + end + def flag_json(key, version: nil) data = { key: key, on: true, fallthrough: { variation: 0 }, variations: ["a"] } data[:version] = version unless version.nil? @@ -261,9 +267,13 @@ def fetch_changes(paths) end end + def listen_enabled? + Impl::Integrations::FileDataSourceV2.class_variable_get(:@@have_listen) + end + def without_listen klass = Impl::Integrations::FileDataSourceV2 - had_listen = klass.class_variable_get(:@@have_listen) + had_listen = listen_enabled? klass.class_variable_set(:@@have_listen, false) begin yield @@ -272,14 +282,23 @@ def without_listen end end + def poller_threads + Thread.list.select { |t| t.name == "LD/FileDataSourceV2" } + end + def with_sync(paths, poll_interval: 0.1) source = Impl::Integrations::FileDataSourceV2.new(logger, paths: paths, poll_interval: poll_interval) updates = Queue.new + pollers_before = poller_threads thread = Thread.new { source.sync(no_selector_store) { |update| updates << update } } begin - initial = updates.pop(timeout: 5) + initial = pop_with_timeout(updates, 5) expect(initial).not_to be_nil expect(initial.state).to eq Interfaces::DataSource::Status::VALID + # The source starts change detection after it yields the initial data. A polling + # source has recorded the files' modification times once its poller thread exists, + # so a change made after that is one the poller can see. + expect(wait_for { (poller_threads - pollers_before).any? }).to be true unless listen_enabled? yield updates ensure source.stop @@ -344,9 +363,16 @@ def with_sync(paths, poll_interval: 0.1) skip "the listen gem is not installed" unless defined?(Listen) file = make_temp_file({ flagValues: { value1: "x" } }.to_json) - expect(Listen).to receive(:to).and_call_original + listened = Queue.new + allow(Listen).to receive(:to).and_wrap_original do |original, *args, **options, &block| + listened << true + original.call(*args, **options, &block) + end - with_sync([file.path]) { |_updates| } + with_sync([file.path]) do |_updates| + # The source starts the listener after it yields the initial data. + expect(pop_with_timeout(listened, 5)).to be true + end end it "compares only the modification time when polling" do @@ -358,7 +384,7 @@ def with_sync(paths, poll_interval: 0.1) IO.write(file, { flagValues: { value1: "a much longer value than before" } }.to_json) File.utime(mtime, mtime, file.path) - expect(updates.pop(timeout: 0.6)).to be_nil + expect(pop_with_timeout(updates, 0.6)).to be_nil end end end @@ -371,7 +397,7 @@ def with_sync(paths, poll_interval: 0.1) file.close File.delete(file.path) - expect(updates.pop(timeout: 0.6)).to be_nil + expect(pop_with_timeout(updates, 0.6)).to be_nil end end end @@ -384,10 +410,10 @@ def with_sync(paths, poll_interval: 0.1) with_sync([file.path]) do |updates| IO.write(file, { flagValues: { value1: "y" } }.to_json) - update = updates.pop(timeout: 5) + update = pop_with_timeout(updates, 5) expect(update).not_to be_nil expect(update.state).to eq Interfaces::DataSource::Status::VALID - expect(updates.pop(timeout: 0.6)).to be_nil + expect(pop_with_timeout(updates, 0.6)).to be_nil end end end From 73afbaad70bff6018c4e5be71e563b59cce3f27f Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Thu, 8 Oct 2026 16:51:37 +0000 Subject: [PATCH 11/11] fix: Watch each data file directory independently The watcher set up its watches all or nothing. If any directory of the configured paths did not exist, it logged the problem, watched nothing, and retried once per second until every directory existed. A configured file in a directory that never exists is a valid steady state for the override source, which treats a missing file as contributing nothing, so one permanently missing directory meant that changes to files in the directories that do exist were never detected. Each directory is now watched on its own. At start, the directories that exist are watched at once and the others are remembered as missing and retried once per second. A retry adds the watches for the directories that exist by then, and the callback runs once when it added at least one, whether or not other directories are still missing. The directories that are still missing, and any watch that could not be set up, are logged together, so that an unchanged situation repeats at debug level as before. When a watched directory is lost, only its watch is removed and it goes back to the missing set; the other directories keep their watches. With rb-inotify, one notifier and one thread take the watches as they are added, and a lost directory's watch is closed. With the listen gem, there is one listener per directory, and stop ends each of them. The spec that lost a directory as soon as its watches were set up waited for the inotify thread to end as its sign that the loss had been handled. The thread outlives one directory now, so it waits for the warning in the log instead. New specs cover a directory missing at start beside one that exists, a lost directory watched again while another is still missing, and stop while a directory is missing. --- lib/ldclient-rb/impl/file_data/watcher.rb | 225 ++++++++++++++-------- spec/impl/file_data/watcher_spec.rb | 89 ++++++++- 2 files changed, 233 insertions(+), 81 deletions(-) diff --git a/lib/ldclient-rb/impl/file_data/watcher.rb b/lib/ldclient-rb/impl/file_data/watcher.rb index 4a26197d..73c0add6 100644 --- a/lib/ldclient-rb/impl/file_data/watcher.rb +++ b/lib/ldclient-rb/impl/file_data/watcher.rb @@ -19,10 +19,10 @@ module FileData # scans the whole directory tree under each watched directory. # # The watcher observes the directory of each file, so a configured file that does not exist - # yet is picked up when it appears. If a directory does not exist, or is deleted while it is - # watched, the watcher logs the problem and retries once per second until it exists. When the - # watches are in place after a retry, the callback runs once, so that a change made while - # there was no watch is not missed. + # yet is picked up when it appears. Each directory is watched on its own: one that does not + # exist, or that is deleted while it is watched, is logged and retried once per second until + # it exists, while the other directories stay watched. When a retry sets up a watch, the + # callback runs once, so that a change made while there was no watch is not missed. # # @private # @@ -38,6 +38,11 @@ class Watcher INOTIFY_DIRECTORY_LOST = [:delete_self, :move_self].freeze private_constant :INOTIFY_DIRECTORY_LOST + # The watch on one real directory: the configured directories that resolve to it, the + # names of the configured files in it, and a callable that removes the watch. + WatchedDirectory = Struct.new(:directories, :names, :close) + private_constant :WatchedDirectory + # # Returns true if a change notification mechanism can be loaded. # @@ -88,11 +93,18 @@ def initialize(paths, on_change, logger) @logger = logger @stopped = Concurrent::AtomicBoolean.new(false) @lock = Mutex.new - @listener = nil @retry_task = nil @last_error_message = nil + # The directories of the configured paths, in configuration order. A directory is in + # @missing until it is watched, and then in @watched under its real path. The lock + # guards both, @retry_task, and @inotify. + @directories = paths.map { |p| File.dirname(p) }.uniq + @missing = Set.new(@directories) + @watched = {} + @inotify = nil - schedule_retry unless try_start + try_start + schedule_retry unless @lock.synchronize { @missing.empty? } end # @@ -102,19 +114,21 @@ def initialize(paths, on_change, logger) def stop return unless @stopped.make_true - listener, retry_task = @lock.synchronize do - pair = [@listener, @retry_task] - @listener = nil + closers, retry_task, inotify = @lock.synchronize do + state = [@watched.values.map(&:close), @retry_task, @inotify] + @watched.clear @retry_task = nil - pair + @inotify = nil + state end retry_task&.stop - listener&.stop + closers.each(&:call) + inotify&.stop end # - # Starts the task that attempts to set up the watches once per second, unless it already - # runs or the watcher is stopped. + # Starts the task that attempts to watch the missing directories once per second, unless + # it already runs or the watcher is stopped. # private def schedule_retry @lock.synchronize do @@ -128,113 +142,157 @@ def stop private def retry_start return if @stopped.value - return unless try_start + added = try_start + # The task ends when nothing is missing. A directory lost meanwhile is back in @missing + # before its loss report calls schedule_retry, so either the task is kept here or that + # call starts a new one. retry_task = @lock.synchronize do + next nil unless @missing.empty? + task = @retry_task @retry_task = nil task end # This runs on the retry task's own thread, which RepeatingTask#stop allows. retry_task&.stop - # The new watches can report the loss of their directory before the retry task is - # released above. That report finds the task still present and leaves the watches to - # it, so when they are gone the retry is started again here. - schedule_retry if @lock.synchronize { @listener.nil? } - @on_change.call unless @stopped.value + @on_change.call if added && !@stopped.value end # # Handles the loss of a watched directory. The notification mechanism reports it on its - # own thread, and the watches end with the directory, so they are torn down and set up - # again through the same retry as at start, once the directory exists. + # own thread. The watch ends with the directory, so it is removed and set up again through + # the same retry as at start, once the directory exists. The other directories keep their + # watches. # - private def directory_lost(directory) - listener = @lock.synchronize do - return if @stopped.value || @listener.nil? + private def directory_lost(real_directory) + entry = @lock.synchronize do + return if @stopped.value + + e = @watched.delete(real_directory) + return if e.nil? - l = @listener - @listener = nil - l + @missing.merge(e.directories) + e end - @logger.warn { "[LDClient] Directory #{directory} no longer exists; its data files are watched again when it exists" } - listener.stop + @logger.warn { "[LDClient] Directory #{real_directory} no longer exists; its data files are watched again when it exists" } + entry.close.call schedule_retry end # - # Sets up the watches. Returns false, after logging, if that is not possible yet. + # Watches each missing directory that exists now. Returns true if at least one watch was + # added. The directories that are still missing, and any watch that could not be set up, + # are logged together, so that an unchanged situation repeats at debug level. # private def try_start - directories = @paths.map { |p| File.dirname(p) }.uniq - missing = directories.reject { |d| File.directory?(d) } - unless missing.empty? - log_setup_failure("directory does not exist: #{missing.join(', ')}") - return false + added = false + missing = [] + problems = [] + @lock.synchronize { @directories.select { |d| @missing.include?(d) } }.each do |directory| + unless File.directory?(directory) + missing << directory + next + end + begin + added = true if watch_directory(directory) + rescue => e + problems << e.message + end end + problems.unshift("directory does not exist: #{missing.join(', ')}") unless missing.empty? + if problems.empty? + @last_error_message = nil + else + log_setup_failure(problems.join("; ")) + end + added + end - listener = Watcher.inotify_available? ? start_inotify : start_listen - + # + # Watches the real directory of a configured directory, and records the names of the + # configured files in it. Two configured directories can resolve to the same real + # directory, which then has one watch for the names in both. Returns false if the watcher + # is stopped. + # + private def watch_directory(directory) + real_directory = File.realpath(directory) + names = @paths.select { |p| File.dirname(p) == directory }.map { |p| File.basename(p) } @lock.synchronize do - if @stopped.value - listener.stop - else - @listener = listener + return false if @stopped.value + + entry = @watched[real_directory] + if entry.nil? + entry = WatchedDirectory.new(Set.new, Set.new, add_watch(real_directory)) + @watched[real_directory] = entry end + entry.directories << directory + entry.names.merge(names) + @missing.delete(directory) end - @last_error_message = nil true - rescue => e - log_setup_failure(e.message) - false end # - # Watches the real directory of each file, without descending into subdirectories, and - # reports events whose file name is one of the watched names in that directory. + # Adds the notification mechanism's watch on a real directory and returns a callable that + # removes it. Called with the lock held, so that stop sees either no watch or a recorded + # one. # - private def start_inotify - names_by_directory = {} - @paths.each do |p| - real_directory = File.realpath(File.dirname(p)) - (names_by_directory[real_directory] ||= Set.new) << File.basename(p) - end + private def add_watch(real_directory) + Watcher.inotify_available? ? add_inotify_watch(real_directory) : start_listen(real_directory) + end - notifier = INotify::Notifier.new - begin - names_by_directory.each do |directory, names| - notifier.watch(directory, *INOTIFY_EVENTS) do |event| - next if @stopped.value - - if (event.flags & INOTIFY_DIRECTORY_LOST).empty? - @on_change.call if names.include?(event.name) - else - directory_lost(directory) - end - end + # + # Adds a watch on the directory itself, without descending into subdirectories, to the one + # inotify notifier, which is created with the first watch. + # + private def add_inotify_watch(real_directory) + @inotify ||= InotifyListener.new(INotify::Notifier.new, @logger) + watch = @inotify.watch(real_directory) { |event| inotify_event(real_directory, event) } + lambda do + begin + watch.close + rescue SystemCallError + # The kernel already removed the watch along with the directory. end - rescue - notifier.close - raise end - InotifyListener.new(notifier, @logger) + end + + private def inotify_event(real_directory, event) + return if @stopped.value + + if (event.flags & INOTIFY_DIRECTORY_LOST).empty? + @on_change.call if watched_name?(real_directory, event.name) + else + directory_lost(real_directory) + end end # - # Watches the real directory of each file with the `listen` gem, which reports paths under - # the real directory, so the paths to match are built the same way. + # Watches a real directory with the `listen` gem, which reports paths under the real + # directory, so a reported path is matched by its directory and name. # - private def start_listen - directories = @paths.map { |p| File.dirname(p) }.uniq - real_directories = directories.map { |d| File.realpath(d) } - watched = Set.new(@paths.map { |p| File.join(File.realpath(File.dirname(p)), File.basename(p)) }) + private def start_listen(real_directory) + listener = Listen.to(real_directory) do |modified, added, removed| + next if @stopped.value - listener = Listen.to(*real_directories) do |modified, added, removed| - changed = (modified + added + removed).any? { |p| watched.include?(p) } + changed = (modified + added + removed).any? do |p| + File.dirname(p) == real_directory && watched_name?(real_directory, File.basename(p)) + end @on_change.call if changed && !@stopped.value end listener.start - listener + -> { listener.stop } + end + + # + # Returns true if the name is one of the configured files in the watched real directory. + # + private def watched_name?(real_directory, name) + @lock.synchronize do + entry = @watched[real_directory] + !entry.nil? && entry.names.include?(name) + end end private def log_setup_failure(message) @@ -247,7 +305,8 @@ def stop end # - # Runs an inotify notifier on its own thread and stops it on request. + # Runs an inotify notifier on its own thread, takes watches for it over time, and stops it + # on request. # class InotifyListener def initialize(notifier, logger) @@ -264,6 +323,14 @@ def initialize(notifier, logger) @thread.name = "LD/FileDataWatcher" end + # + # Adds a watch on a directory and returns the notifier's watch object, whose `close` + # removes the watch again. + # + def watch(directory, &callback) + @notifier.watch(directory, *INOTIFY_EVENTS, &callback) + end + # # Stops the notifier and waits briefly for its thread. Closing the notifier ends the # blocking read that the thread is in. A callback can call this on the notifier's own diff --git a/spec/impl/file_data/watcher_spec.rb b/spec/impl/file_data/watcher_spec.rb index a0204148..d5357456 100644 --- a/spec/impl/file_data/watcher_spec.rb +++ b/spec/impl/file_data/watcher_spec.rb @@ -44,6 +44,19 @@ def with_watcher(paths, logger: $null_log) end end + # Returns the call count once it has stopped changing for a moment. One edit can produce a + # burst of notifications, and a count taken in the middle of the burst would make the rest + # of it look like a later signal. + def settled(calls) + count = calls.value + loop do + sleep 0.2 + break if calls.value == count + count = calls.value + end + count + end + it "reports that the listen gem is available" do expect(Watcher.available?).to be true end @@ -143,7 +156,8 @@ def with_watcher(paths, logger: $null_log) it "keeps retrying when a directory is lost again as soon as its watches are set up" do skip "rb-inotify is not available on this platform" unless Watcher.inotify_available? - with_watcher([path("sub/a.json")]) do |watcher, calls| + logger = CapturingLogger.new + with_watcher([path("sub/a.json")], logger: logger) do |watcher, calls| # Take the directory away as soon as the watches are set up, and wait until the loss # has been handled, so that it is reported while the retry that set up the watches # is still finishing. @@ -153,7 +167,9 @@ def with_watcher(paths, logger: $null_log) if started && !lost_once lost_once = true FileUtils.rm_rf(path("sub")) - wait_for { Thread.list.none? { |t| t.name == "LD/FileDataWatcher" } } + # The loss has been handled once it is logged. The inotify thread cannot mark it, + # as it used to: each directory has its own watch now, and the thread outlives one. + wait_for { logger.output.include?("no longer exists") } end started end @@ -216,6 +232,75 @@ def with_watcher(paths, logger: $null_log) expect(Thread.list.map(&:name)).not_to include("LD/FileDataWatcherRetry") end end + + it "watches the directories that exist while another one is missing" do + logger = CapturingLogger.new + Dir.mkdir(path("sub")) + File.write(path("sub/a.json"), "{}") + missing_dir = path("not-yet") + with_watcher([path("sub/a.json"), File.join(missing_dir, "b.json")], logger: logger) do |_watcher, calls| + expect(logger.output).to include("directory does not exist: #{missing_dir}") + sleep 0.3 + # The directory that exists is watched at once, not only when every directory exists. + File.write(path("sub/a.json"), '{"flagValues": {}}') + expect(wait_for { calls.value >= 1 }).to be true + before = settled(calls) + + Dir.mkdir(missing_dir) + # The watch is set up on the next retry, and the callback runs once at that point. + expect(wait_for { calls.value > before }).to be true + # Nothing is missing any more, so the retry ends and does not signal again. + expect(wait_for { Thread.list.none? { |t| t.name == "LD/FileDataWatcherRetry" } }).to be true + before = settled(calls) + sleep 0.3 + expect(calls.value).to eq(before) + + File.write(File.join(missing_dir, "b.json"), "{}") + expect(wait_for { calls.value > before }).to be true + end + end + + it "watches a lost directory again while another directory is still missing" do + skip "rb-inotify is not available on this platform" unless Watcher.inotify_available? + + logger = CapturingLogger.new + Dir.mkdir(path("sub")) + File.write(path("sub/a.json"), "{}") + with_watcher([path("sub/a.json"), path("not-yet/b.json")], logger: logger) do |_watcher, calls| + FileUtils.rm_rf(path("sub")) + expect(wait_for { logger.output.match?(/WARN.*#{Regexp.escape(path('sub'))}/) }).to be true + before = settled(calls) + + Dir.mkdir(path("sub")) + File.write(path("sub/a.json"), "{}") + # The directory is watched again on the next retry, which signals once at that point, + # although the other directory is still missing. + expect(wait_for { calls.value > before }).to be true + before = settled(calls) + + File.write(path("sub/a.json"), '{"flagValues": {}}') + expect(wait_for { calls.value > before }).to be true + expect(Dir.exist?(path("not-yet"))).to be false + expect(logger.output).to include("directory does not exist: #{path('not-yet')}") + end + end + + it "stops while a directory is missing and runs no callback afterwards" do + Dir.mkdir(path("sub")) + File.write(path("sub/a.json"), "{}") + with_watcher([path("sub/a.json"), path("not-yet/b.json")]) do |watcher, calls| + sleep 0.3 + watcher.stop + expect(Thread.list.map(&:name)).not_to include("LD/FileDataWatcherRetry") + expect(Thread.list.map(&:name)).not_to include("LD/FileDataWatcher") if Watcher.inotify_available? + + Dir.mkdir(path("not-yet")) + File.write(path("not-yet/b.json"), "{}") + File.write(path("sub/a.json"), '{"flagValues": {}}') + sleep 0.5 + expect(calls.value).to eq(0) + end + end end end end