From 353e7480021439a918456b1d1dce3d0d03045910 Mon Sep 17 00:00:00 2001
From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com>
Date: Mon, 28 Sep 2026 13:33:20 -0700
Subject: [PATCH] feat: Add a shared file data reloader, poller, and watcher
Adds the Internal/FileLoading namespace with the file reading code that
the flag overrides feature needs: a document parser and merger that keep
each entry's own version and report per-document counts, a reloader that
serializes reloads, debounces change signals, retains the last good
result on failure, retries after a bounded delay, reports an identical
repeated failure once, and skips unchanged content, a stat-based poller,
and a directory watcher that handles files that do not exist yet.
The existing file data source keeps its current behavior and does not use
the shared code. Tests that pin that behavior are added so later changes
to the shared code cannot alter it.
---
.../Internal/FileLoading/FileDataDocument.cs | 184 ++++++
.../Internal/FileLoading/FileDataMerger.cs | 219 +++++++
.../Internal/FileLoading/FileDataPoller.cs | 141 ++++
.../Internal/FileLoading/FileDataReloader.cs | 450 +++++++++++++
.../Internal/FileLoading/FileDataWatcher.cs | 114 ++++
.../FileDataSourceExistingBehaviorTest.cs | 183 ++++++
.../FileLoading/FileDataMergerTest.cs | 159 +++++
.../FileLoading/FileDataParserTest.cs | 142 ++++
.../FileLoading/FileDataPollerTest.cs | 228 +++++++
.../FileLoading/FileDataReloaderTest.cs | 606 ++++++++++++++++++
.../FileLoading/FileDataWatcherTest.cs | 161 +++++
pkgs/sdk/server/test/TestUtils.cs | 24 +
12 files changed, 2611 insertions(+)
create mode 100644 pkgs/sdk/server/src/Internal/FileLoading/FileDataDocument.cs
create mode 100644 pkgs/sdk/server/src/Internal/FileLoading/FileDataMerger.cs
create mode 100644 pkgs/sdk/server/src/Internal/FileLoading/FileDataPoller.cs
create mode 100644 pkgs/sdk/server/src/Internal/FileLoading/FileDataReloader.cs
create mode 100644 pkgs/sdk/server/src/Internal/FileLoading/FileDataWatcher.cs
create mode 100644 pkgs/sdk/server/test/Internal/DataSources/FileDataSourceExistingBehaviorTest.cs
create mode 100644 pkgs/sdk/server/test/Internal/FileLoading/FileDataMergerTest.cs
create mode 100644 pkgs/sdk/server/test/Internal/FileLoading/FileDataParserTest.cs
create mode 100644 pkgs/sdk/server/test/Internal/FileLoading/FileDataPollerTest.cs
create mode 100644 pkgs/sdk/server/test/Internal/FileLoading/FileDataReloaderTest.cs
create mode 100644 pkgs/sdk/server/test/Internal/FileLoading/FileDataWatcherTest.cs
diff --git a/pkgs/sdk/server/src/Internal/FileLoading/FileDataDocument.cs b/pkgs/sdk/server/src/Internal/FileLoading/FileDataDocument.cs
new file mode 100644
index 000000000..ef12149db
--- /dev/null
+++ b/pkgs/sdk/server/src/Internal/FileLoading/FileDataDocument.cs
@@ -0,0 +1,184 @@
+using System;
+using System.Collections.Generic;
+using System.Collections.Immutable;
+using System.Text;
+using System.Text.Json;
+using System.Text.Json.Serialization;
+using LaunchDarkly.Sdk.Server.Internal.Model;
+
+using static LaunchDarkly.Sdk.Internal.JsonConverterHelpers;
+using static LaunchDarkly.Sdk.Json.LdJsonConverters;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ ///
+ /// The parsed form of one data file. A document can hold full flag definitions, simplified
+ /// flag-key-to-value entries, and segment definitions. Each list keeps the order in which the
+ /// file lists its entries.
+ ///
+ internal sealed class FileDataDocument
+ {
+ internal static readonly FileDataDocument Empty = new FileDataDocument(null, null, null);
+
+ internal IReadOnlyList> Flags { get; }
+ internal IReadOnlyList> FlagValues { get; }
+ internal IReadOnlyList> Segments { get; }
+
+ internal FileDataDocument(
+ IReadOnlyList> flags,
+ IReadOnlyList> flagValues,
+ IReadOnlyList> segments
+ )
+ {
+ Flags = flags ?? ImmutableList>.Empty;
+ FlagValues = flagValues ?? ImmutableList>.Empty;
+ Segments = segments ?? ImmutableList>.Empty;
+ }
+ }
+
+ ///
+ /// Parses the file data document format. A document is a JSON object with optional
+ /// flags , flagValues , and segments members. An alternate parser can
+ /// handle other formats, such as YAML.
+ ///
+ internal sealed class FileDataParser
+ {
+ private readonly Func _alternateParser;
+
+ ///
+ /// Constructs a parser.
+ ///
+ /// a function that parses non-JSON content into basic .NET
+ /// collections, or null to accept JSON only
+ internal FileDataParser(Func alternateParser)
+ {
+ _alternateParser = alternateParser;
+ }
+
+ ///
+ /// Parses the content of one file. Throws an exception if the content cannot be parsed.
+ ///
+ /// the file content
+ /// the parsed document
+ internal FileDataDocument Parse(string content)
+ {
+ if (_alternateParser == null)
+ {
+ return ParseJson(content);
+ }
+ if (content.Trim().StartsWith("{"))
+ {
+ try
+ {
+ return ParseJson(content);
+ }
+ catch (Exception)
+ {
+ // The content is not valid JSON. The alternate parser gets a chance to parse it.
+ }
+ }
+ // The alternate parser produces the most basic .NET data structure that can represent
+ // the file content, using types like Dictionary and String. We convert this into a JSON
+ // tree so we can use the JSON deserializer. This is inefficient, but it lets us reuse the
+ // data model deserialization logic.
+ var o = _alternateParser(content);
+ var options = new JsonSerializerOptions();
+ options.Converters.Add(new UntypedDictionaryJsonSerializer());
+ var asJson = JsonSerializer.Serialize(o, options);
+ return ParseJson(asJson);
+ }
+
+ private static FileDataDocument ParseJson(string data)
+ {
+ var r = new Utf8JsonReader(Encoding.UTF8.GetBytes(data));
+ return ParseJson(ref r);
+ }
+
+ private static FileDataDocument ParseJson(ref Utf8JsonReader r)
+ {
+ var flagsBuilder = ImmutableList.CreateBuilder>();
+ var flagValuesBuilder = ImmutableList.CreateBuilder>();
+ var segmentsBuilder = ImmutableList.CreateBuilder>();
+ for (var obj = RequireObject(ref r); obj.Next(ref r);)
+ {
+ switch (obj.Name)
+ {
+ case "flags":
+ for (var subObj = RequireObjectOrNull(ref r); subObj.Next(ref r);)
+ {
+ var key = subObj.Name;
+ var flag = FeatureFlagSerialization.Instance.Read(ref r, null, null);
+ flagsBuilder.Add(new KeyValuePair(key, flag));
+ }
+ break;
+
+ case "flagValues":
+ for (var subObj = RequireObjectOrNull(ref r); subObj.Next(ref r);)
+ {
+ var key = subObj.Name;
+ var value = LdValueConverter.ReadJsonValue(ref r);
+ flagValuesBuilder.Add(new KeyValuePair(key, value));
+ }
+ break;
+
+ case "segments":
+ for (var subObj = RequireObjectOrNull(ref r); subObj.Next(ref r);)
+ {
+ var key = subObj.Name;
+ var segment = SegmentSerialization.Instance.Read(ref r, null, null);
+ segmentsBuilder.Add(new KeyValuePair(key, segment));
+ }
+ break;
+ }
+ }
+ return new FileDataDocument(flagsBuilder.ToImmutable(), flagValuesBuilder.ToImmutable(),
+ segmentsBuilder.ToImmutable());
+ }
+
+ ///
+ /// Constructs a flag that is off and serves the same value for every context. The flag has a
+ /// single variation and that variation as its off variation, so it evaluates with the
+ /// reason. This is the form the override source uses for
+ /// flagValues entries.
+ ///
+ internal static FeatureFlag MakeOffFlagWithValue(string key, LdValue value, int version)
+ {
+ var json = LdValue.BuildObject()
+ .Add("key", key)
+ .Add("version", version)
+ .Add("on", false)
+ .Add("offVariation", 0)
+ .Add("variations", LdValue.ArrayOf(value))
+ .Build()
+ .ToJsonString();
+ return DataModel.Features.Deserialize(json).Item as FeatureFlag;
+ }
+
+ // This custom JSON serializer addresses a problem that can happen when using an external YAML parser.
+ // In JSON, the keys must always be strings, and System.Text.Json will refuse to either serialize or
+ // deserialize anything with non-string keys. But in YAML, the keys can be of any type, so a YAML
+ // parser that is told to deserialize some map-like data without a specific target type may decide to
+ // return the type Dictionary even if the keys really are strings.
+ private class UntypedDictionaryJsonSerializer : JsonConverter>
+ {
+ public override bool CanConvert(Type typeToConvert) =>
+ typeof(IDictionary).IsAssignableFrom(typeToConvert);
+
+ public override IDictionary Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
+ {
+ throw new NotImplementedException();
+ }
+
+ public override void Write(Utf8JsonWriter writer, IDictionary value, JsonSerializerOptions options)
+ {
+ writer.WriteStartObject();
+ foreach (var kv in value)
+ {
+ writer.WritePropertyName(kv.Key.ToString());
+ JsonSerializer.Serialize(writer, kv.Value, options);
+ }
+ writer.WriteEndObject();
+ }
+ }
+ }
+}
diff --git a/pkgs/sdk/server/src/Internal/FileLoading/FileDataMerger.cs b/pkgs/sdk/server/src/Internal/FileLoading/FileDataMerger.cs
new file mode 100644
index 000000000..ad2574476
--- /dev/null
+++ b/pkgs/sdk/server/src/Internal/FileLoading/FileDataMerger.cs
@@ -0,0 +1,219 @@
+using System;
+using System.Collections.Generic;
+using System.Collections.Immutable;
+using LaunchDarkly.Sdk.Server.Internal.Model;
+
+using static LaunchDarkly.Sdk.Server.Subsystems.DataStoreTypes;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ ///
+ /// Determines what happens when the same flag or segment key appears in more than one document.
+ ///
+ internal enum FileDataDuplicateKeysHandling
+ {
+ ///
+ /// A duplicated key makes the merge fail.
+ ///
+ Fail,
+
+ ///
+ /// Only the first occurrence of a duplicated key is used, in the order the documents were given.
+ ///
+ Ignore
+ }
+
+ ///
+ /// An error in loading or combining file data. The message describes the problem.
+ ///
+ internal class FileDataException : Exception
+ {
+ internal FileDataException(string message) : base(message) { }
+
+ internal FileDataException(string message, Exception innerException) : base(message, innerException) { }
+ }
+
+ ///
+ /// Indicates that one of the source files could not be read or parsed. It distinguishes a
+ /// per-file failure from a failure to merge the files' contents.
+ ///
+ internal sealed class FileDataReadException : FileDataException
+ {
+ ///
+ /// The path of the file that failed.
+ ///
+ internal string Path { get; }
+
+ internal FileDataReadException(string path, string description, Exception innerException) :
+ base(description + " [" + path + "]", innerException)
+ {
+ Path = path;
+ }
+ }
+
+ ///
+ /// Counts the entries the merge kept from one document. An entry dropped by the duplicate keys
+ /// handling is not counted.
+ ///
+ internal struct FileDataDocumentSummary
+ {
+ internal int Flags { get; }
+ internal int Segments { get; }
+
+ internal FileDataDocumentSummary(int flags, int segments)
+ {
+ Flags = flags;
+ Segments = segments;
+ }
+ }
+
+ ///
+ /// Describes one configured file after a reload.
+ ///
+ internal struct FileDataFileSummary
+ {
+ internal string Path { get; }
+
+ ///
+ /// False when the file does not exist and missing files are skipped.
+ ///
+ internal bool Present { get; }
+
+ internal int Flags { get; }
+ internal int Segments { get; }
+
+ internal FileDataFileSummary(string path, bool present, int flags, int segments)
+ {
+ Path = path;
+ Present = present;
+ Flags = flags;
+ Segments = segments;
+ }
+ }
+
+ ///
+ /// The merged items from one or more documents. Item values are or
+ /// .
+ ///
+ ///
+ /// All of one document's items precede the next document's, in the order the documents were
+ /// given. Within one document, the file order is kept.
+ ///
+ internal sealed class FileDataMergeResult
+ {
+ internal IReadOnlyList> Flags { get; }
+ internal IReadOnlyList> Segments { get; }
+
+ ///
+ /// For each input document in order, the number of entries the merge kept from it.
+ ///
+ internal IReadOnlyList Documents { get; }
+
+ ///
+ /// Set by the reloader. Describes each configured file in order.
+ ///
+ internal IReadOnlyList Files { get; }
+
+ internal FileDataMergeResult(
+ IReadOnlyList> flags,
+ IReadOnlyList> segments,
+ IReadOnlyList documents,
+ IReadOnlyList files
+ )
+ {
+ Flags = flags;
+ Segments = segments;
+ Documents = documents;
+ Files = files ?? ImmutableList.Empty;
+ }
+
+ internal FileDataMergeResult WithFiles(IReadOnlyList files) =>
+ new FileDataMergeResult(Flags, Segments, Documents, files);
+ }
+
+ ///
+ /// Combines the items of several documents into one set.
+ ///
+ internal static class FileDataMerger
+ {
+ ///
+ /// Combines the items of the given documents, expanding flag-value entries into full flag
+ /// definitions and applying the given duplicate keys handling. An unrecognized handling
+ /// value behaves as .
+ ///
+ /// what to do when a key appears more than once
+ /// the documents, in precedence order
+ /// builds the full flag definition for a flag-value entry
+ /// the merged result
+ /// a key appears more than once and the handling is Fail
+ internal static FileDataMergeResult Merge(
+ FileDataDuplicateKeysHandling duplicateKeysHandling,
+ IReadOnlyList documents,
+ Func expandFlagValue
+ )
+ {
+ var flags = ImmutableList.CreateBuilder>();
+ var segments = ImmutableList.CreateBuilder>();
+ var summaries = ImmutableList.CreateBuilder();
+ var seenFlagKeys = new HashSet();
+ var seenSegmentKeys = new HashSet();
+
+ foreach (var document in documents)
+ {
+ var flagCount = 0;
+ var segmentCount = 0;
+ foreach (var kv in document.Flags)
+ {
+ if (Insert(flags, seenFlagKeys, "flag", kv.Key, kv.Value.Version, kv.Value, duplicateKeysHandling))
+ {
+ flagCount++;
+ }
+ }
+ foreach (var kv in document.FlagValues)
+ {
+ var flag = expandFlagValue(kv.Key, kv.Value);
+ if (Insert(flags, seenFlagKeys, "flag", kv.Key, flag.Version, flag, duplicateKeysHandling))
+ {
+ flagCount++;
+ }
+ }
+ foreach (var kv in document.Segments)
+ {
+ if (Insert(segments, seenSegmentKeys, "segment", kv.Key, kv.Value.Version, kv.Value, duplicateKeysHandling))
+ {
+ segmentCount++;
+ }
+ }
+ summaries.Add(new FileDataDocumentSummary(flagCount, segmentCount));
+ }
+
+ return new FileDataMergeResult(flags.ToImmutable(), segments.ToImmutable(), summaries.ToImmutable(), null);
+ }
+
+ // Adds the entry unless the key was already seen. Returns true if it added the entry.
+ private static bool Insert(
+ ImmutableList>.Builder items,
+ ISet seenKeys,
+ string category,
+ string key,
+ int version,
+ object item,
+ FileDataDuplicateKeysHandling duplicateKeysHandling
+ )
+ {
+ if (seenKeys.Contains(key))
+ {
+ switch (duplicateKeysHandling)
+ {
+ case FileDataDuplicateKeysHandling.Ignore:
+ return false;
+ default:
+ throw new FileDataException(category + " \"" + key + "\" is specified by multiple files");
+ }
+ }
+ items.Add(new KeyValuePair(key, new ItemDescriptor(version, item)));
+ seenKeys.Add(key);
+ return true;
+ }
+ }
+}
diff --git a/pkgs/sdk/server/src/Internal/FileLoading/FileDataPoller.cs b/pkgs/sdk/server/src/Internal/FileLoading/FileDataPoller.cs
new file mode 100644
index 000000000..0dc4a8bcf
--- /dev/null
+++ b/pkgs/sdk/server/src/Internal/FileLoading/FileDataPoller.cs
@@ -0,0 +1,141 @@
+using System;
+using System.Collections.Generic;
+using System.IO;
+using System.Linq;
+using System.Threading;
+using LaunchDarkly.Logging;
+using LaunchDarkly.Sdk.Internal;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ ///
+ /// 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 change 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 , whose debouncing and skip-unchanged
+ /// handling absorb the excess.
+ ///
+ ///
+ internal sealed class FileDataPoller : IDisposable
+ {
+ private struct FileState : IEquatable
+ {
+ internal bool Exists;
+ internal DateTime LastWriteTimeUtc;
+ internal long Length;
+
+ public bool Equals(FileState other) =>
+ Exists == other.Exists && LastWriteTimeUtc == other.LastWriteTimeUtc && Length == other.Length;
+
+ public override bool Equals(object obj) => obj is FileState other && Equals(other);
+
+ public override int GetHashCode() =>
+ Exists.GetHashCode() ^ LastWriteTimeUtc.GetHashCode() ^ Length.GetHashCode();
+ }
+
+ private readonly string[] _paths;
+ private readonly Action _onChange;
+ private readonly Logger _log;
+ private readonly Timer _timer;
+ private FileState[] _last;
+ private int _examining;
+ private volatile bool _disposed;
+
+ ///
+ /// Creates a started poller. It examines the files once before it returns, so only later
+ /// changes invoke the callback. Call Dispose to stop it.
+ ///
+ /// the files to examine
+ /// the time between examinations
+ /// invoked when a change is detected
+ /// receives log output about unexpected errors
+ internal FileDataPoller(IEnumerable paths, TimeSpan interval, Action onChange, Logger log)
+ {
+ _paths = paths.ToArray();
+ _onChange = onChange;
+ _log = log;
+ _last = ObserveAll(_paths);
+ _timer = new Timer(OnTick, null, interval, interval);
+ }
+
+ ///
+ /// Stops the poller. It does not wait for an examination or a callback that is in progress.
+ /// A file system that does not respond must not block shutdown. As a result, the callback
+ /// can run one more time shortly after Dispose returns. Consumers tolerate a late call, as
+ /// they do for a late reload.
+ ///
+ public void Dispose()
+ {
+ _disposed = true;
+ _timer.Dispose();
+ }
+
+ private void OnTick(object state)
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ // Timer callbacks can overlap when an examination or the callback runs longer than the
+ // interval. One examination at a time is enough.
+ if (Interlocked.Exchange(ref _examining, 1) != 0)
+ {
+ return;
+ }
+ try
+ {
+ var current = ObserveAll(_paths);
+ var changed = false;
+ for (var i = 0; i < current.Length; i++)
+ {
+ if (!current[i].Equals(_last[i]))
+ {
+ changed = true;
+ break;
+ }
+ }
+ _last = current;
+ if (changed && !_disposed)
+ {
+ _onChange();
+ }
+ }
+ catch (Exception e)
+ {
+ LogHelpers.LogException(_log, "Unexpected error while examining files for changes", e);
+ }
+ finally
+ {
+ Interlocked.Exchange(ref _examining, 0);
+ }
+ }
+
+ private static FileState[] ObserveAll(string[] paths)
+ {
+ var states = new FileState[paths.Length];
+ for (var i = 0; i < paths.Length; i++)
+ {
+ var info = new FileInfo(paths[i]);
+ if (info.Exists)
+ {
+ states[i] = new FileState
+ {
+ Exists = true,
+ LastWriteTimeUtc = info.LastWriteTimeUtc,
+ Length = info.Length
+ };
+ }
+ }
+ return states;
+ }
+ }
+}
diff --git a/pkgs/sdk/server/src/Internal/FileLoading/FileDataReloader.cs b/pkgs/sdk/server/src/Internal/FileLoading/FileDataReloader.cs
new file mode 100644
index 000000000..3d1315ace
--- /dev/null
+++ b/pkgs/sdk/server/src/Internal/FileLoading/FileDataReloader.cs
@@ -0,0 +1,450 @@
+using System;
+using System.Collections.Generic;
+using System.Collections.Immutable;
+using System.IO;
+using System.Security.Cryptography;
+using System.Text;
+using System.Threading;
+using System.Threading.Tasks;
+using LaunchDarkly.Logging;
+using LaunchDarkly.Sdk.Internal;
+using LaunchDarkly.Sdk.Server.Integrations;
+using LaunchDarkly.Sdk.Server.Internal.Model;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ ///
+ /// Configuration for .
+ ///
+ internal sealed class FileDataReloaderConfig
+ {
+ ///
+ /// The files to load. The order is significant. It determines which file wins under the
+ /// duplicate keys handling.
+ ///
+ internal IReadOnlyList Paths { get; set; }
+
+ ///
+ /// What happens when the same key appears in more than one file.
+ ///
+ internal FileDataDuplicateKeysHandling DuplicateKeysHandling { get; set; }
+
+ ///
+ /// When true, a configured file that does not exist is treated as a file with no content.
+ /// The reload succeeds with the data of the files that exist. When false, a missing file
+ /// fails the reload like any other read error.
+ ///
+ internal bool SkipMissingPaths { get; set; }
+
+ ///
+ /// Receives log output about reloads and failures.
+ ///
+ internal Logger Logger { get; set; }
+
+ ///
+ /// Reads file contents. Defaults to the SDK's standard file reader.
+ ///
+ internal FileDataTypes.IFileReader FileReader { get; set; }
+
+ ///
+ /// Parses non-JSON content, or null to accept JSON only.
+ ///
+ internal Func AlternateParser { get; set; }
+
+ ///
+ /// Builds the full flag definition for a flagValues entry.
+ ///
+ internal Func FlagValueExpander { get; set; }
+
+ ///
+ /// Invoked with each successfully merged result. Calls are serialized. Apply and OnError
+ /// must not call back into Dispose.
+ ///
+ internal Action Apply { get; set; }
+
+ ///
+ /// Invoked when a reload fails, once per distinct failure. With automatic retries, repeats
+ /// of an identical failure do not invoke it again. A success re-arms it. The exception is a
+ /// when a file could not be read or parsed, or a
+ /// for a merge failure. The reloader logs failures itself.
+ ///
+ internal Action OnError { get; set; }
+
+ ///
+ /// How long to wait after a Trigger call for further calls to settle before reloading. This
+ /// coalesces bursts of change notifications into one reload. If zero or negative, each
+ /// Trigger reloads immediately on a worker thread.
+ ///
+ internal TimeSpan DebounceDelay { get; set; }
+
+ ///
+ /// How long to wait after a failed reload before retrying, so that a failure observed while
+ /// a file was being rewritten recovers even if no further change notification arrives. If
+ /// zero or negative, there is no automatic retry.
+ ///
+ internal TimeSpan RetryDelay { get; set; }
+
+ ///
+ /// If true, the Apply call is skipped when the files' raw contents are identical to the
+ /// last successfully applied contents.
+ ///
+ internal bool SkipUnchanged { get; set; }
+ }
+
+ ///
+ /// 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.
+ ///
+ internal sealed class FileDataReloader : IDisposable
+ {
+ ///
+ /// A settle window long enough to coalesce the burst of change notifications produced by a
+ /// single file edit, and short enough to stay responsive.
+ ///
+ internal static readonly TimeSpan DefaultDebounceDelay = TimeSpan.FromMilliseconds(100);
+
+ ///
+ /// 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.
+ ///
+ internal static readonly TimeSpan DefaultRetryDelay = TimeSpan.FromSeconds(1);
+
+ private readonly FileDataReloaderConfig _config;
+ private readonly FileDataParser _parser;
+ private readonly FileDataTypes.IFileReader _fileReader;
+ private readonly Logger _log;
+
+ // Serializes the load work between ReloadNow and the timer callbacks.
+ private readonly object _reloadLock = new object();
+ private byte[] _lastGoodHash; // only touched inside _reloadLock
+ private string _lastErrorMessage; // only touched inside _reloadLock
+
+ // Guards the timers. The timers are created on first use so that a reloader that is never
+ // used holds no scheduled work.
+ private readonly object _timerLock = new object();
+ private Timer _debounceTimer;
+ private Timer _retryTimer;
+ private bool _retryArmed;
+ private int _immediateReloadPending;
+
+ private volatile bool _disposed;
+
+ internal FileDataReloader(FileDataReloaderConfig config)
+ {
+ _config = config;
+ _parser = new FileDataParser(config.AlternateParser);
+ _fileReader = config.FileReader ?? Internal.DataSources.FlagFileReader.Instance;
+ _log = config.Logger ?? Logs.None.Logger("");
+ }
+
+ ///
+ /// Synchronously loads the files and applies the result, or reports the failure. Use it
+ /// for the initial load. A failure here arms the same automatic retry as a failed
+ /// triggered reload.
+ ///
+ internal void ReloadNow()
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ if (!Reload())
+ {
+ // An already armed retry keeps its earlier deadline.
+ ArmRetry(keepExistingDeadline: true);
+ }
+ }
+
+ ///
+ /// 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.
+ ///
+ internal void Trigger()
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ if (_config.DebounceDelay <= TimeSpan.Zero)
+ {
+ if (Interlocked.Exchange(ref _immediateReloadPending, 1) == 0)
+ {
+ Task.Run(() =>
+ {
+ Interlocked.Exchange(ref _immediateReloadPending, 0);
+ ReloadAfterSignal(isRetry: false);
+ });
+ }
+ return;
+ }
+ lock (_timerLock)
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ if (_debounceTimer == null)
+ {
+ _debounceTimer = new Timer(OnDebounceElapsed, null, Timeout.Infinite, Timeout.Infinite);
+ }
+ // Each trigger moves the deadline out again. The reload runs after activity settles.
+ _debounceTimer.Change(_config.DebounceDelay, Timeout.InfiniteTimeSpan);
+ }
+ }
+
+ ///
+ /// 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 OnError shortly after Dispose returns. A
+ /// reload that has not yet reached its callbacks when Dispose is called does not invoke
+ /// them.
+ ///
+ public void Dispose()
+ {
+ _disposed = true;
+ lock (_timerLock)
+ {
+ _debounceTimer?.Dispose();
+ _retryTimer?.Dispose();
+ _debounceTimer = null;
+ _retryTimer = null;
+ _retryArmed = false;
+ }
+ }
+
+ private void OnDebounceElapsed(object state) => RunGuarded(() => ReloadAfterSignal(isRetry: false));
+
+ private void OnRetryElapsed(object state) => RunGuarded(() => ReloadAfterSignal(isRetry: true));
+
+ // Timer callbacks run on thread pool threads. An unhandled exception there ends the process,
+ // so every callback logs instead.
+ private void RunGuarded(Action action)
+ {
+ try
+ {
+ action();
+ }
+ catch (Exception e)
+ {
+ LogHelpers.LogException(_log, "Unexpected error while reloading file data", e);
+ }
+ }
+
+ private void ReloadAfterSignal(bool isRetry)
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ if (isRetry)
+ {
+ _log.Debug("Retrying file data load after earlier failure");
+ }
+ else
+ {
+ _log.Info("Reloading file data after detecting a change");
+ }
+ // A pending retry is superseded by this reload. It either succeeds, or it fails and
+ // arms a fresh retry.
+ DisarmRetry();
+ if (!Reload())
+ {
+ ArmRetry(keepExistingDeadline: false);
+ }
+ }
+
+ private void ArmRetry(bool keepExistingDeadline)
+ {
+ if (_config.RetryDelay <= TimeSpan.Zero)
+ {
+ return;
+ }
+ lock (_timerLock)
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ if (_retryArmed && keepExistingDeadline)
+ {
+ return;
+ }
+ if (_retryTimer == null)
+ {
+ _retryTimer = new Timer(OnRetryElapsed, null, Timeout.Infinite, Timeout.Infinite);
+ }
+ _retryArmed = true;
+ _retryTimer.Change(_config.RetryDelay, Timeout.InfiniteTimeSpan);
+ }
+ }
+
+ private void DisarmRetry()
+ {
+ lock (_timerLock)
+ {
+ _retryArmed = false;
+ _retryTimer?.Change(Timeout.Infinite, Timeout.Infinite);
+ }
+ }
+
+ // Performs one full load of all configured files. Returns true if the load succeeded, which
+ // decides whether a retry gets armed. 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 bool Reload()
+ {
+ lock (_reloadLock)
+ {
+ // A trigger already queued when Dispose was called can still reach here.
+ if (_disposed)
+ {
+ return true;
+ }
+
+ var documents = new List(_config.Paths.Count);
+ var files = new List(_config.Paths.Count);
+ var rawContents = _config.SkipUnchanged ? new StringBuilder() : null;
+ foreach (var path in _config.Paths)
+ {
+ string content;
+ try
+ {
+ content = _fileReader.ReadAllText(path);
+ }
+ catch (FileNotFoundException) when (_config.SkipMissingPaths)
+ {
+ _log.Debug("File {0} does not exist; it contributes no data", path);
+ files.Add(new FileDataFileSummary(path, false, 0, 0));
+ continue;
+ }
+ catch (DirectoryNotFoundException) when (_config.SkipMissingPaths)
+ {
+ _log.Debug("File {0} does not exist; it contributes no data", path);
+ files.Add(new FileDataFileSummary(path, false, 0, 0));
+ continue;
+ }
+ catch (Exception e)
+ {
+ return Fail(new FileDataReadException(path, "unable to read file: " + e.Message, e));
+ }
+ // One read feeds both the hash and the parse, so the skip-unchanged hash can never
+ // disagree with the content that was actually applied.
+ rawContents?.Append(content).Append('\0');
+ FileDataDocument document;
+ try
+ {
+ document = _parser.Parse(content);
+ }
+ catch (Exception e)
+ {
+ return Fail(new FileDataReadException(path, "error parsing file: " + e.Message, e));
+ }
+ documents.Add(document);
+ files.Add(new FileDataFileSummary(path, true, 0, 0));
+ }
+
+ FileDataMergeResult merged;
+ try
+ {
+ merged = FileDataMerger.Merge(_config.DuplicateKeysHandling, documents, _config.FlagValueExpander);
+ }
+ catch (FileDataException e)
+ {
+ return Fail(e);
+ }
+
+ // Documents are the present files in order. Copy their counts onto the file summaries.
+ var filesWithCounts = ImmutableList.CreateBuilder();
+ var next = 0;
+ foreach (var file in files)
+ {
+ if (file.Present)
+ {
+ var summary = merged.Documents[next++];
+ filesWithCounts.Add(new FileDataFileSummary(file.Path, true, summary.Flags, summary.Segments));
+ }
+ else
+ {
+ filesWithCounts.Add(file);
+ }
+ }
+
+ // Dispose may have happened while the files were being read. Deliver nothing in that
+ // case. This check is deliberately not atomic with the delivery below. Dispose never
+ // blocks on a lock shared with callbacks, so a reload that passes this check can rarely
+ // deliver just after Dispose returns.
+ if (_disposed)
+ {
+ return true;
+ }
+
+ // 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 OnError and only
+ // Apply tells it things are good again.
+ var recovering = _lastErrorMessage != null;
+ _lastErrorMessage = null;
+ if (rawContents != null)
+ {
+ var hash = ComputeHash(rawContents.ToString());
+ if (!recovering && HashesEqual(hash, _lastGoodHash))
+ {
+ return true;
+ }
+ _lastGoodHash = hash;
+ }
+ _config.Apply?.Invoke(merged.WithFiles(filesWithCounts.ToImmutable()));
+ return true;
+ }
+ }
+
+ // Called inside _reloadLock.
+ private bool Fail(Exception e)
+ {
+ // Dispose may have happened while the files were being read. Deliver nothing in that
+ // case, and report success so no retry is armed.
+ if (_disposed)
+ {
+ return true;
+ }
+ // With automatic retries, a persistent failure would repeat the same log entry and the
+ // same callback on every attempt. Repeats of an identical failure are demoted to debug
+ // level and do not invoke OnError again.
+ if (e.Message == _lastErrorMessage)
+ {
+ _log.Debug("Unable to load flags: {0}", e.Message);
+ return false;
+ }
+ _lastErrorMessage = e.Message;
+ _log.Error("Unable to load flags: {0}", e.Message);
+ _config.OnError?.Invoke(e);
+ return false;
+ }
+
+ private static byte[] ComputeHash(string contents)
+ {
+ using (var sha = SHA256.Create())
+ {
+ return sha.ComputeHash(Encoding.UTF8.GetBytes(contents));
+ }
+ }
+
+ private static bool HashesEqual(byte[] a, byte[] b)
+ {
+ if (a == null || b == null || a.Length != b.Length)
+ {
+ return false;
+ }
+ for (var i = 0; i < a.Length; i++)
+ {
+ if (a[i] != b[i])
+ {
+ return false;
+ }
+ }
+ return true;
+ }
+ }
+}
diff --git a/pkgs/sdk/server/src/Internal/FileLoading/FileDataWatcher.cs b/pkgs/sdk/server/src/Internal/FileLoading/FileDataWatcher.cs
new file mode 100644
index 000000000..ac78d1044
--- /dev/null
+++ b/pkgs/sdk/server/src/Internal/FileLoading/FileDataWatcher.cs
@@ -0,0 +1,114 @@
+using System;
+using System.Collections.Generic;
+using System.IO;
+using LaunchDarkly.Logging;
+using LaunchDarkly.Sdk.Internal;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ ///
+ /// Detects changes to a set of files through file system change notifications. It watches the
+ /// directory of each file, so a configured file that does not exist yet is picked up when it
+ /// appears, and a deleted file is reported.
+ ///
+ ///
+ /// A single edit often produces several notifications, and a notification can arrive while the
+ /// file is still being written. The watcher reports every notification. Feed it into a
+ /// , whose debouncing and failure retry handle both.
+ ///
+ internal sealed class FileDataWatcher : IDisposable
+ {
+ private readonly HashSet _fullPaths = new HashSet();
+ private readonly List _watchers = new List();
+ private readonly Action _onChange;
+ private readonly Logger _log;
+ private volatile bool _disposed;
+
+ ///
+ /// Creates a started watcher. Throws if a file's directory does not exist or cannot be
+ /// watched. Call Dispose to stop it.
+ ///
+ /// the files to watch
+ /// invoked for each change notification that concerns one of the files
+ /// receives log output about notification errors
+ internal FileDataWatcher(IEnumerable paths, Action onChange, Logger log)
+ {
+ _onChange = onChange;
+ _log = log;
+
+ var directories = new HashSet();
+ foreach (var path in paths)
+ {
+ var fullPath = Path.GetFullPath(path);
+ _fullPaths.Add(fullPath);
+ directories.Add(Path.GetDirectoryName(fullPath));
+ }
+
+ try
+ {
+ foreach (var directory in directories)
+ {
+ var watcher = new FileSystemWatcher(directory)
+ {
+ NotifyFilter = NotifyFilters.LastWrite | NotifyFilters.Size | NotifyFilters.FileName |
+ NotifyFilters.CreationTime,
+ IncludeSubdirectories = false
+ };
+ watcher.Changed += (sender, args) => OnFileEvent(args.FullPath);
+ watcher.Created += (sender, args) => OnFileEvent(args.FullPath);
+ watcher.Deleted += (sender, args) => OnFileEvent(args.FullPath);
+ watcher.Renamed += (sender, args) =>
+ {
+ OnFileEvent(args.OldFullPath);
+ OnFileEvent(args.FullPath);
+ };
+ watcher.Error += OnWatcherError;
+ _watchers.Add(watcher);
+ }
+ foreach (var watcher in _watchers)
+ {
+ watcher.EnableRaisingEvents = true;
+ }
+ }
+ catch (Exception)
+ {
+ Dispose();
+ throw;
+ }
+ }
+
+ public void Dispose()
+ {
+ _disposed = true;
+ foreach (var watcher in _watchers)
+ {
+ watcher.Dispose();
+ }
+ _watchers.Clear();
+ }
+
+ private void OnFileEvent(string fullPath)
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ if (_fullPaths.Contains(Path.GetFullPath(fullPath)))
+ {
+ _onChange();
+ }
+ }
+
+ // An error from the file system, such as a full notification buffer, means that changes may
+ // have been dropped. Treat it as a change so the files are read again.
+ private void OnWatcherError(object sender, ErrorEventArgs args)
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ LogHelpers.LogException(_log, "Error from file system watcher", args.GetException());
+ _onChange();
+ }
+ }
+}
diff --git a/pkgs/sdk/server/test/Internal/DataSources/FileDataSourceExistingBehaviorTest.cs b/pkgs/sdk/server/test/Internal/DataSources/FileDataSourceExistingBehaviorTest.cs
new file mode 100644
index 000000000..b68c28c45
--- /dev/null
+++ b/pkgs/sdk/server/test/Internal/DataSources/FileDataSourceExistingBehaviorTest.cs
@@ -0,0 +1,183 @@
+using System;
+using System.Linq;
+using System.Threading;
+using LaunchDarkly.Logging;
+using LaunchDarkly.Sdk.Server.Integrations;
+using LaunchDarkly.Sdk.Server.Internal.Model;
+using LaunchDarkly.Sdk.Server.Subsystems;
+using Xunit;
+using Xunit.Abstractions;
+
+using static LaunchDarkly.Sdk.Server.Subsystems.DataStoreTypes;
+
+namespace LaunchDarkly.Sdk.Server.Internal.DataSources
+{
+ // These tests pin behavior of the file data source that its other tests do not assert directly,
+ // so that changes to the shared file loading code cannot alter it: every successful load is
+ // applied even when the content is unchanged, the retry delay and its logging, the log messages
+ // and levels for a failed load, missing-file handling, and which file system notifications
+ // trigger a reload.
+ public class FileDataSourceExistingBehaviorTest : BaseTest
+ {
+ private static readonly string ALL_DATA_JSON_FILE = TestUtils.TestFilePath("all-properties.json");
+ private static readonly string FLAG_ONLY_JSON_FILE = TestUtils.TestFilePath("flag-only.json");
+ private static readonly string BAD_FILE = TestUtils.TestFilePath("bad-file.txt");
+
+ private readonly CapturingDataSourceUpdates _updateSink = new CapturingDataSourceUpdates();
+ private readonly FileDataSourceBuilder factory = FileData.DataSource();
+
+ public FileDataSourceExistingBehaviorTest(ITestOutputHelper testOutput) : base(testOutput) { }
+
+ private IDataSource MakeDataSource() =>
+ factory.Build(BasicContext.WithDataSourceUpdates(_updateSink));
+
+ private static int VersionOfFlag1(FullDataSet data) =>
+ data.Data.First(kv => kv.Key == DataModel.Features).Value.Items.First(kv => kv.Key == "flag1").Value.Version;
+
+ private class ScriptedFileReader : FileDataTypes.IFileReader
+ {
+ private int _reads;
+ public int Reads => Volatile.Read(ref _reads);
+
+ public string ReadAllText(string path)
+ {
+ Interlocked.Increment(ref _reads);
+ return @"{""flagValues"": {"; // invalid as JSON and as YAML
+ }
+ }
+
+ private static void WaitUntil(Func condition, string description)
+ {
+ var deadline = DateTime.UtcNow.AddSeconds(15);
+ while (!condition() && DateTime.UtcNow < deadline)
+ {
+ Thread.Sleep(20);
+ }
+ Assert.True(condition(), "timed out waiting for " + description);
+ }
+
+ [Fact]
+ public void IdenticalContentIsReappliedWithANewVersionOnEachChangeNotification()
+ {
+ using (var file = TempFile.Create())
+ {
+ factory.FilePaths(file.Path).AutoUpdate(true);
+ file.SetContentFromPath(FLAG_ONLY_JSON_FILE);
+ using (var fp = MakeDataSource())
+ {
+ fp.Start();
+ Assert.Equal(1, VersionOfFlag1(_updateSink.Inits.ExpectValue()));
+
+ // A rewrite with the same content is a change notification like any other. The
+ // data is applied again, with the next version number.
+ file.SetContentFromPath(FLAG_ONLY_JSON_FILE);
+ AssertHelpers.ExpectPredicate(_updateSink.Inits, data => VersionOfFlag1(data) > 1,
+ "Did not receive a reload of the unchanged file.", TimeSpan.FromSeconds(30));
+ AssertLogMessage(true, LogLevel.Info, "detected file modification, reloading");
+ }
+ }
+ }
+
+ [Fact]
+ public void MissingFileWithoutSkipMissingPathsLogsFailedToLoad()
+ {
+ factory.FilePaths(ALL_DATA_JSON_FILE, "bad-file-path");
+ using (var fp = MakeDataSource())
+ {
+ fp.Start();
+ Assert.False(fp.Initialized);
+ _updateSink.Inits.ExpectNoValue();
+ AssertLogMessageRegex(true, LogLevel.Error, "^Failed to load bad-file-path");
+ }
+ }
+
+ [Fact]
+ public void MalformedFileWithAutoUpdateOffLogsFailedToParse()
+ {
+ factory.FilePaths(BAD_FILE);
+ using (var fp = MakeDataSource())
+ {
+ fp.Start();
+ Assert.False(fp.Initialized);
+ _updateSink.Inits.ExpectNoValue();
+ AssertLogMessageRegex(true, LogLevel.Error, "^Failed to parse .*bad-file.txt");
+ }
+ }
+
+ [Fact]
+ public void DuplicateKeyAcrossFilesFailsTheLoadAndNamesTheKey()
+ {
+ using (var file1 = TempFile.Create())
+ using (var file2 = TempFile.Create())
+ {
+ file1.SetContent(@"{""flagValues"":{""flag1"":""a""}}");
+ file2.SetContent(@"{""flagValues"":{""flag1"":""b""}}");
+ factory.FilePaths(file1.Path, file2.Path);
+ using (var fp = MakeDataSource())
+ {
+ fp.Start();
+ Assert.False(fp.Initialized);
+ _updateSink.Inits.ExpectNoValue();
+ AssertLogMessageRegex(true, LogLevel.Error,
+ "^Failed to load .*: .*in \"features\", key \"flag1\" was already defined");
+ }
+ }
+ }
+
+ [Fact]
+ public void PathInAMissingDirectoryFailsTheLoadEvenWithSkipMissingPaths()
+ {
+ // SkipMissingPaths skips a file that does not exist in an existing directory. A path whose
+ // directory does not exist is a failure.
+ var missingDirectoryPath = "/nonexistent-ld-filedatasource-pin-dir/data.json";
+ factory.FilePaths(ALL_DATA_JSON_FILE, missingDirectoryPath).SkipMissingPaths(true);
+ using (var fp = MakeDataSource())
+ {
+ fp.Start();
+ Assert.False(fp.Initialized);
+ _updateSink.Inits.ExpectNoValue();
+ AssertLogMessageRegex(true, LogLevel.Error, "^Failed to load " + missingDirectoryPath);
+ }
+ }
+
+ [Fact]
+ public void DeletingAWatchedFileDoesNotTriggerAReload()
+ {
+ using (var file = TempFile.Create())
+ {
+ factory.FilePaths(file.Path).AutoUpdate(true);
+ file.SetContentFromPath(FLAG_ONLY_JSON_FILE);
+ using (var fp = MakeDataSource())
+ {
+ fp.Start();
+ _updateSink.Inits.ExpectValue();
+
+ // Only modification, creation, and rename notifications trigger a reload.
+ file.Delete();
+ _updateSink.Inits.ExpectNoValue(TimeSpan.FromMilliseconds(500));
+ AssertLogMessageRegex(false, LogLevel.Error, "Failed to load");
+ AssertLogMessageRegex(false, LogLevel.Warn, "Failed to read");
+ }
+ }
+ }
+
+ [Fact]
+ public void EveryFailedParseAttemptLogsAWarningWithTheRetryDelay()
+ {
+ var reader = new ScriptedFileReader();
+ using (var file = TempFile.Create())
+ {
+ factory.FilePaths(file.Path).AutoUpdate(true).FileReader(reader);
+ using (var fp = MakeDataSource())
+ {
+ fp.Start();
+ WaitUntil(() => reader.Reads >= 3, "3 file reads");
+ var warnings = LogCapture.GetMessages().Where(m =>
+ m.Level == LogLevel.Warn && m.Text.Contains("will retry in 600 ms")).ToList();
+ Assert.True(warnings.Count >= 2, "expected a warning per failed attempt, got " + warnings.Count);
+ Assert.False(fp.Initialized);
+ }
+ }
+ }
+ }
+}
diff --git a/pkgs/sdk/server/test/Internal/FileLoading/FileDataMergerTest.cs b/pkgs/sdk/server/test/Internal/FileLoading/FileDataMergerTest.cs
new file mode 100644
index 000000000..c6a047a25
--- /dev/null
+++ b/pkgs/sdk/server/test/Internal/FileLoading/FileDataMergerTest.cs
@@ -0,0 +1,159 @@
+using System.Collections.Generic;
+using System.Linq;
+using LaunchDarkly.Sdk.Server.Internal.Model;
+using Xunit;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ public class FileDataMergerTest
+ {
+ private static FeatureFlag Expand(string key, LdValue value) =>
+ FileDataParser.MakeOffFlagWithValue(key, value, 0);
+
+ private static FileDataDocument DocWithFlag(FeatureFlag flag) =>
+ new FileDataDocument(new[] { new KeyValuePair(flag.Key, flag) }, null, null);
+
+ private static FileDataDocument DocWithFlagValue(string key, LdValue value) =>
+ new FileDataDocument(null, new[] { new KeyValuePair(key, value) }, null);
+
+ private static FileDataDocument DocWithSegment(Segment segment) =>
+ new FileDataDocument(null, null, new[] { new KeyValuePair(segment.Key, segment) });
+
+ private static FileDataMergeResult Merge(FileDataDuplicateKeysHandling handling, params FileDataDocument[] docs) =>
+ FileDataMerger.Merge(handling, docs, Expand);
+
+ [Fact]
+ public void CombinesDocuments()
+ {
+ var flag1 = new FeatureFlagBuilder("flag1").Version(2).Build();
+ var segment1 = new SegmentBuilder("segment1").Version(4).Build();
+
+ var result = Merge(FileDataDuplicateKeysHandling.Fail,
+ DocWithFlag(flag1), DocWithFlagValue("flag2", LdValue.Of(true)), DocWithSegment(segment1));
+
+ Assert.Equal(2, result.Flags.Count);
+ Assert.Equal("flag1", result.Flags[0].Key);
+ Assert.Equal(2, result.Flags[0].Value.Version);
+ Assert.Same(flag1, result.Flags[0].Value.Item);
+ Assert.Equal("flag2", result.Flags[1].Key);
+ var expanded = Assert.IsType(result.Flags[1].Value.Item);
+ Assert.Equal(new[] { LdValue.Of(true) }, expanded.Variations);
+
+ Assert.Single(result.Segments);
+ Assert.Equal("segment1", result.Segments[0].Key);
+ Assert.Equal(4, result.Segments[0].Value.Version);
+ Assert.Same(segment1, result.Segments[0].Value.Item);
+ }
+
+ [Fact]
+ public void DuplicateFlagFails()
+ {
+ var flagA = new FeatureFlagBuilder("flag1").Version(1).Build();
+ var flagB = new FeatureFlagBuilder("flag1").Version(2).Build();
+ var e = Assert.Throws(() =>
+ Merge(FileDataDuplicateKeysHandling.Fail, DocWithFlag(flagA), DocWithFlag(flagB)));
+ Assert.Equal("flag \"flag1\" is specified by multiple files", e.Message);
+ }
+
+ [Fact]
+ public void UnrecognizedHandlingBehavesAsFail()
+ {
+ var flagA = new FeatureFlagBuilder("flag1").Version(1).Build();
+ var flagB = new FeatureFlagBuilder("flag1").Version(2).Build();
+ Assert.Throws(() =>
+ Merge((FileDataDuplicateKeysHandling)99, DocWithFlag(flagA), DocWithFlag(flagB)));
+ }
+
+ [Fact]
+ public void IgnoreKeepsFirstOccurrence()
+ {
+ var flagA = new FeatureFlagBuilder("flag1").Version(1).Build();
+ var flagB = new FeatureFlagBuilder("flag1").Version(2).Build();
+ var result = Merge(FileDataDuplicateKeysHandling.Ignore, DocWithFlag(flagA), DocWithFlag(flagB));
+ Assert.Single(result.Flags);
+ Assert.Equal(1, result.Flags[0].Value.Version);
+ Assert.Same(flagA, result.Flags[0].Value.Item);
+ }
+
+ [Fact]
+ public void FullFlagAndFlagValueCollide()
+ {
+ var flag = new FeatureFlagBuilder("flag1").Version(1).Build();
+ Assert.Throws(() =>
+ Merge(FileDataDuplicateKeysHandling.Fail, DocWithFlag(flag), DocWithFlagValue("flag1", LdValue.Of(true))));
+ }
+
+ [Fact]
+ public void FlagValueAndFullFlagCollideWithinOneDocument()
+ {
+ var flag = new FeatureFlagBuilder("flag1").Version(1).Build();
+ var doc = new FileDataDocument(
+ new[] { new KeyValuePair("flag1", flag) },
+ new[] { new KeyValuePair("flag1", LdValue.Of(true)) },
+ null);
+ Assert.Throws(() => Merge(FileDataDuplicateKeysHandling.Fail, doc));
+ }
+
+ [Fact]
+ public void DuplicateSegmentFails()
+ {
+ var segment = new SegmentBuilder("segment1").Build();
+ var e = Assert.Throws(() =>
+ Merge(FileDataDuplicateKeysHandling.Fail, DocWithSegment(segment), DocWithSegment(segment)));
+ Assert.Equal("segment \"segment1\" is specified by multiple files", e.Message);
+ }
+
+ [Fact]
+ public void FlagAndSegmentWithTheSameKeyDoNotCollide()
+ {
+ var flag = new FeatureFlagBuilder("same").Build();
+ var segment = new SegmentBuilder("same").Build();
+ var result = Merge(FileDataDuplicateKeysHandling.Fail, DocWithFlag(flag), DocWithSegment(segment));
+ Assert.Single(result.Flags);
+ Assert.Single(result.Segments);
+ }
+
+ [Fact]
+ public void PreservesDocumentOrder()
+ {
+ var expectedKeys = new[] { "flag-a", "flag-b", "flag-c", "flag-d", "flag-e" };
+ var docs = expectedKeys.Select(key => DocWithFlag(new FeatureFlagBuilder(key).Build())).ToArray();
+ var result = Merge(FileDataDuplicateKeysHandling.Fail, docs);
+ Assert.Equal(expectedKeys, result.Flags.Select(kv => kv.Key));
+ }
+
+ [Fact]
+ public void CountsEntriesKeptFromEachDocument()
+ {
+ var first = new FileDataDocument(null, new[]
+ {
+ new KeyValuePair("a", LdValue.Of(true)),
+ new KeyValuePair("shared", LdValue.Of(true))
+ }, null);
+ var second = new FileDataDocument(null, new[]
+ {
+ new KeyValuePair("b", LdValue.Of(true)),
+ new KeyValuePair("shared", LdValue.Of(false))
+ }, new[] { new KeyValuePair("seg", new SegmentBuilder("seg").Build()) });
+
+ var result = Merge(FileDataDuplicateKeysHandling.Ignore, first, second);
+
+ Assert.Equal(2, result.Documents.Count);
+ Assert.Equal(2, result.Documents[0].Flags);
+ Assert.Equal(0, result.Documents[0].Segments);
+ // The duplicate "shared" entry from the second document is dropped and not counted.
+ Assert.Equal(1, result.Documents[1].Flags);
+ Assert.Equal(1, result.Documents[1].Segments);
+ Assert.Empty(result.Files);
+ }
+
+ [Fact]
+ public void NoDocumentsProduceAnEmptyResult()
+ {
+ var result = Merge(FileDataDuplicateKeysHandling.Fail);
+ Assert.Empty(result.Flags);
+ Assert.Empty(result.Segments);
+ Assert.Empty(result.Documents);
+ }
+ }
+}
diff --git a/pkgs/sdk/server/test/Internal/FileLoading/FileDataParserTest.cs b/pkgs/sdk/server/test/Internal/FileLoading/FileDataParserTest.cs
new file mode 100644
index 000000000..e0ef7a9ff
--- /dev/null
+++ b/pkgs/sdk/server/test/Internal/FileLoading/FileDataParserTest.cs
@@ -0,0 +1,142 @@
+using System;
+using System.Linq;
+using System.Text.Json;
+using LaunchDarkly.Sdk.Server.Internal.Model;
+using Xunit;
+using YamlDotNet.Serialization;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ public class FileDataParserTest
+ {
+ private const string FullDocument = @"{
+ ""flags"": {
+ ""flag1"": { ""key"": ""flag1"", ""version"": 7, ""on"": true, ""fallthrough"": { ""variation"": 1 },
+ ""variations"": [ ""a"", ""b"" ] }
+ },
+ ""flagValues"": { ""flag2"": ""value2"", ""flag3"": true },
+ ""segments"": {
+ ""seg1"": { ""key"": ""seg1"", ""version"": 3, ""included"": [ ""user1"" ] }
+ }
+ }";
+
+ private static readonly FileDataParser JsonOnly = new FileDataParser(null);
+
+ [Fact]
+ public void ParsesFlagsFlagValuesAndSegmentsInOrder()
+ {
+ var doc = JsonOnly.Parse(FullDocument);
+
+ Assert.Equal(new[] { "flag1" }, doc.Flags.Select(kv => kv.Key));
+ Assert.Equal(7, doc.Flags[0].Value.Version);
+ Assert.True(doc.Flags[0].Value.On);
+
+ Assert.Equal(new[] { "flag2", "flag3" }, doc.FlagValues.Select(kv => kv.Key));
+ Assert.Equal(LdValue.Of("value2"), doc.FlagValues[0].Value);
+ Assert.Equal(LdValue.Of(true), doc.FlagValues[1].Value);
+
+ Assert.Equal(new[] { "seg1" }, doc.Segments.Select(kv => kv.Key));
+ Assert.Equal(3, doc.Segments[0].Value.Version);
+ Assert.Equal(new[] { "user1" }, doc.Segments[0].Value.Included);
+ }
+
+ [Fact]
+ public void KeepsDocumentVersions()
+ {
+ var doc = JsonOnly.Parse(FullDocument);
+ Assert.Equal(7, doc.Flags[0].Value.Version);
+ Assert.Equal(3, doc.Segments[0].Value.Version);
+ }
+
+ [Fact]
+ public void EmptyObjectIsAnEmptyDocument()
+ {
+ var doc = JsonOnly.Parse("{}");
+ Assert.Empty(doc.Flags);
+ Assert.Empty(doc.FlagValues);
+ Assert.Empty(doc.Segments);
+ }
+
+ [Fact]
+ public void NullSectionsAreEmpty()
+ {
+ var doc = JsonOnly.Parse(@"{""flags"": null, ""flagValues"": null, ""segments"": null}");
+ Assert.Empty(doc.Flags);
+ Assert.Empty(doc.FlagValues);
+ Assert.Empty(doc.Segments);
+ }
+
+ [Fact]
+ public void UnknownTopLevelPropertiesAreIgnored()
+ {
+ var doc = JsonOnly.Parse(@"{""other"": {""x"": 1}, ""flagValues"": {""flag1"": 1}}");
+ Assert.Single(doc.FlagValues);
+ }
+
+ [Theory]
+ [InlineData("")]
+ [InlineData("what is this")]
+ [InlineData(@"{""flagValues""")]
+ [InlineData(@"{""flagValues"": {""flag1"": }}")]
+ [InlineData("[]")]
+ public void MalformedContentThrowsWithoutAlternateParser(string content)
+ {
+ Assert.ThrowsAny(() => JsonOnly.Parse(content));
+ }
+
+ [Fact]
+ public void AlternateParserIsUsedForNonJsonContent()
+ {
+ var yaml = new DeserializerBuilder().WithAttemptingUnquotedStringTypeDeserialization().Build();
+ var parser = new FileDataParser(s => yaml.Deserialize(s));
+
+ var doc = parser.Parse("flagValues:\n flag1: true\n flag2: \"text\"\nsegments:\n seg1:\n key: seg1\n version: 2\n");
+
+ Assert.Equal(new[] { "flag1", "flag2" }, doc.FlagValues.Select(kv => kv.Key));
+ Assert.Equal(LdValue.Of(true), doc.FlagValues[0].Value);
+ Assert.Equal(LdValue.Of("text"), doc.FlagValues[1].Value);
+ Assert.Equal(2, doc.Segments[0].Value.Version);
+ }
+
+ [Fact]
+ public void JsonIsParsedAsJsonWhenAlternateParserIsConfigured()
+ {
+ var parser = new FileDataParser(s => throw new Exception("alternate parser must not be called"));
+ var doc = parser.Parse(FullDocument);
+ Assert.Single(doc.Flags);
+ }
+
+ [Fact]
+ public void AlternateParserGetsContentThatStartsLikeJsonButIsNot()
+ {
+ // YAML flow mappings start with a brace but are not JSON. Both parsers get a chance.
+ var yaml = new DeserializerBuilder().WithAttemptingUnquotedStringTypeDeserialization().Build();
+ var parser = new FileDataParser(s => yaml.Deserialize(s));
+ var doc = parser.Parse("{flagValues: {flag1: yes}}");
+ Assert.Single(doc.FlagValues);
+ }
+
+ [Fact]
+ public void AlternateParserFailureIsThrown()
+ {
+ var parser = new FileDataParser(s => throw new FormatException("bad yaml"));
+ Assert.Throws(() => parser.Parse("not json"));
+ }
+
+ [Fact]
+ public void OffFlagWithValueServesTheValueWithTheOffReason()
+ {
+ var flag = FileDataParser.MakeOffFlagWithValue("flag1", LdValue.Of("x"), 0);
+ Assert.Equal("flag1", flag.Key);
+ Assert.Equal(0, flag.Version);
+ Assert.False(flag.On);
+ Assert.Equal(0, flag.OffVariation);
+ Assert.Equal(new[] { LdValue.Of("x") }, flag.Variations);
+
+ var result = Evaluation.EvaluatorTestUtil.BasicEvaluator.Evaluate(flag, Context.New("any-user"));
+ Assert.Equal(LdValue.Of("x"), result.Result.Value);
+ Assert.Equal(0, result.Result.VariationIndex);
+ Assert.Equal(EvaluationReason.OffReason, result.Result.Reason);
+ }
+ }
+}
diff --git a/pkgs/sdk/server/test/Internal/FileLoading/FileDataPollerTest.cs b/pkgs/sdk/server/test/Internal/FileLoading/FileDataPollerTest.cs
new file mode 100644
index 000000000..91280a16d
--- /dev/null
+++ b/pkgs/sdk/server/test/Internal/FileLoading/FileDataPollerTest.cs
@@ -0,0 +1,228 @@
+using System;
+using System.IO;
+using System.Threading;
+using System.Threading.Tasks;
+using LaunchDarkly.TestHelpers;
+using Xunit;
+using Xunit.Abstractions;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ public class FileDataPollerTest : BaseTest, IDisposable
+ {
+ private static readonly TimeSpan TestTimeout = TimeSpan.FromSeconds(5);
+ private static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(5);
+ private static readonly TimeSpan QuietPeriod = TimeSpan.FromMilliseconds(100);
+
+ private readonly TempDirectory _dir = TempDirectory.Create();
+ private readonly EventSink _changed = new EventSink();
+ private FileDataPoller _poller;
+
+ public FileDataPollerTest(ITestOutputHelper testOutput) : base(testOutput) { }
+
+ public void Dispose()
+ {
+ _poller?.Dispose();
+ _dir.Dispose();
+ }
+
+ private FileDataPoller StartPoller(params string[] paths)
+ {
+ _poller = new FileDataPoller(paths, PollInterval, () => _changed.Enqueue(true), TestLogger);
+ return _poller;
+ }
+
+ private void RequireChange() => _changed.ExpectValue(TestTimeout);
+
+ private void RequireNoChange(TimeSpan duration) => _changed.ExpectNoValue(duration);
+
+ // Rewrites a file and guarantees the observed (modification time, size) state differs from
+ // the previous state, so the poller must detect it regardless of timestamp granularity.
+ private static void WriteWithNewModTime(string path, string content) =>
+ WriteWithModTime(path, content, DateTime.UtcNow.AddSeconds(content.Length));
+
+ // Writes the content and the modification time to a temporary file, then renames it over
+ // the target. The poller observes one change, not one for the content and one for the time.
+ private static void WriteWithModTime(string path, string content, DateTime modTime)
+ {
+ var temp = path + ".tmp";
+ File.WriteAllText(temp, content);
+ File.SetLastWriteTimeUtc(temp, modTime);
+ if (File.Exists(path))
+ {
+ File.Replace(temp, path, null);
+ }
+ else
+ {
+ File.Move(temp, path);
+ }
+ }
+
+ [Fact]
+ public void DetectsModification()
+ {
+ var path = _dir.PathOf("data.json");
+ WriteWithNewModTime(path, "one");
+ StartPoller(path);
+
+ RequireNoChange(QuietPeriod);
+
+ WriteWithNewModTime(path, "two!");
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsSameSizeRewriteWithNewModTime()
+ {
+ var path = _dir.PathOf("data.json");
+ var baseTime = DateTime.UtcNow.AddHours(-1);
+ WriteWithModTime(path, "one", baseTime);
+ StartPoller(path);
+
+ WriteWithModTime(path, "two", baseTime.AddMinutes(1));
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsSizeChangeWithSameModTime()
+ {
+ var path = _dir.PathOf("data.json");
+ var baseTime = DateTime.UtcNow.AddHours(-1);
+ WriteWithModTime(path, "one", baseTime);
+ StartPoller(path);
+
+ WriteWithModTime(path, "four", baseTime);
+ RequireChange();
+ }
+
+ [Fact]
+ public void FiresOncePerChange()
+ {
+ var path = _dir.PathOf("data.json");
+ WriteWithNewModTime(path, "one");
+ StartPoller(path);
+
+ WriteWithNewModTime(path, "two!");
+ RequireChange();
+
+ RequireNoChange(QuietPeriod);
+ }
+
+ [Fact]
+ public void DetectsFileAppearing()
+ {
+ var path = _dir.PathOf("data.json");
+ StartPoller(path);
+
+ RequireNoChange(QuietPeriod);
+
+ WriteWithNewModTime(path, "created");
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsFileDisappearing()
+ {
+ var path = _dir.PathOf("data.json");
+ WriteWithNewModTime(path, "content");
+ StartPoller(path);
+
+ File.Delete(path);
+ RequireChange();
+ }
+
+ [Fact]
+ public void WatchesAllFiles()
+ {
+ var path1 = _dir.PathOf("one.json");
+ var path2 = _dir.PathOf("two.json");
+ WriteWithNewModTime(path1, "one");
+ WriteWithNewModTime(path2, "two");
+ StartPoller(path1, path2);
+
+ WriteWithNewModTime(path2, "two-changed");
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsChangeToFirstOfMultipleFiles()
+ {
+ var path1 = _dir.PathOf("one.json");
+ var path2 = _dir.PathOf("two.json");
+ WriteWithNewModTime(path1, "one");
+ WriteWithNewModTime(path2, "two");
+ StartPoller(path1, path2);
+
+ WriteWithNewModTime(path1, "one-changed");
+ RequireChange();
+ }
+
+ [Fact]
+ public void StopsOnDispose()
+ {
+ var path = _dir.PathOf("data.json");
+ WriteWithNewModTime(path, "one");
+ var poller = StartPoller(path);
+
+ poller.Dispose();
+ poller.Dispose();
+
+ WriteWithNewModTime(path, "two!");
+ RequireNoChange(QuietPeriod);
+ }
+
+ [Fact]
+ public void DisposeReturnsWhileCallbackBlocks()
+ {
+ var path = _dir.PathOf("data.json");
+ WriteWithNewModTime(path, "one");
+ var entered = new ManualResetEventSlim();
+ var release = new ManualResetEventSlim();
+ _poller = new FileDataPoller(new[] { path }, PollInterval, () =>
+ {
+ entered.Set();
+ release.Wait();
+ }, TestLogger);
+
+ WriteWithNewModTime(path, "two!");
+ Assert.True(entered.Wait(TestTimeout), "timed out waiting for the callback to start");
+
+ var disposed = Task.Run(() => _poller.Dispose());
+ Assert.True(disposed.Wait(TestTimeout), "Dispose blocked on a callback in progress");
+
+ release.Set();
+ }
+
+ [Fact]
+ public void CallbackExceptionIsLoggedAndPollingContinues()
+ {
+ var path = _dir.PathOf("data.json");
+ WriteWithNewModTime(path, "one");
+ var calls = 0;
+ _poller = new FileDataPoller(new[] { path }, PollInterval, () =>
+ {
+ if (Interlocked.Increment(ref calls) == 1)
+ {
+ throw new InvalidOperationException("consumer failed");
+ }
+ _changed.Enqueue(true);
+ }, TestLogger);
+
+ WriteWithNewModTime(path, "two!");
+ // The poller logs the exception on its own thread after the callback throws, so the
+ // test waits for the log line rather than for the callback count.
+ var deadline = DateTime.UtcNow + TestTimeout;
+ while (!LogCapture.HasMessageWithRegex(Logging.LogLevel.Error, "Unexpected error while examining files")
+ && DateTime.UtcNow < deadline)
+ {
+ Thread.Sleep(5);
+ }
+ AssertLogMessageRegex(true, Logging.LogLevel.Error, "Unexpected error while examining files");
+ Assert.True(Volatile.Read(ref calls) >= 1, "the callback was not invoked");
+
+ // The exception did not stop the poller: a later change is still reported.
+ WriteWithNewModTime(path, "three");
+ RequireChange();
+ }
+ }
+}
diff --git a/pkgs/sdk/server/test/Internal/FileLoading/FileDataReloaderTest.cs b/pkgs/sdk/server/test/Internal/FileLoading/FileDataReloaderTest.cs
new file mode 100644
index 000000000..a1138124e
--- /dev/null
+++ b/pkgs/sdk/server/test/Internal/FileLoading/FileDataReloaderTest.cs
@@ -0,0 +1,606 @@
+using System;
+using System.IO;
+using System.Linq;
+using System.Threading;
+using System.Threading.Tasks;
+using LaunchDarkly.Logging;
+using LaunchDarkly.Sdk.Server.Integrations;
+using LaunchDarkly.Sdk.Server.Internal.Model;
+using LaunchDarkly.TestHelpers;
+using Xunit;
+using Xunit.Abstractions;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ public class FileDataReloaderTest : BaseTest, IDisposable
+ {
+ private static readonly TimeSpan TestTimeout = TimeSpan.FromSeconds(5);
+
+ private const string Flag1True = @"{""flagValues"": {""flag1"": true}}";
+ private const string Flag1False = @"{""flagValues"": {""flag1"": false}}";
+ private const string Truncated = @"{""flagValues""";
+
+ private readonly TempDirectory _dir = TempDirectory.Create();
+ private readonly string _path;
+ private readonly EventSink _applied = new EventSink();
+ private readonly EventSink _errored = new EventSink();
+ private int _applyCount;
+ private FileDataReloader _reloader;
+
+ public FileDataReloaderTest(ITestOutputHelper testOutput) : base(testOutput)
+ {
+ _path = _dir.PathOf("data.json");
+ }
+
+ public void Dispose()
+ {
+ _reloader?.Dispose();
+ _dir.Dispose();
+ }
+
+ private FileDataReloaderConfig BasicReloaderConfig() =>
+ new FileDataReloaderConfig
+ {
+ Paths = new[] { _path },
+ DuplicateKeysHandling = FileDataDuplicateKeysHandling.Fail,
+ Logger = TestLogger,
+ FlagValueExpander = (key, value) => FileDataParser.MakeOffFlagWithValue(key, value, 0),
+ Apply = RecordApply,
+ OnError = e => _errored.Enqueue(e)
+ };
+
+ private void RecordApply(FileDataMergeResult result)
+ {
+ Interlocked.Increment(ref _applyCount);
+ _applied.Enqueue(result);
+ }
+
+ private int ApplyCount => Volatile.Read(ref _applyCount);
+
+ private FileDataReloader MakeReloader(Action configure = null)
+ {
+ var config = BasicReloaderConfig();
+ configure?.Invoke(config);
+ _reloader = new FileDataReloader(config);
+ return _reloader;
+ }
+
+ private void Write(string content) => File.WriteAllText(_path, content);
+
+ private FileDataMergeResult RequireApplied() => _applied.ExpectValue(TestTimeout);
+
+ private Exception RequireErrored() => _errored.ExpectValue(TestTimeout);
+
+ private void RequireQuiet(TimeSpan duration)
+ {
+ _applied.ExpectNoValue(duration);
+ _errored.ExpectNoValue(TimeSpan.Zero);
+ }
+
+ private static string[] FlagKeys(FileDataMergeResult result) => result.Flags.Select(kv => kv.Key).ToArray();
+
+ private int ErrorLogCount() => LogCapture.GetMessages().Count(m => m.Level == LogLevel.Error);
+
+ [Fact]
+ public void FailsOnMissingPathByDefault()
+ {
+ Write(Flag1True);
+ var missing = _dir.PathOf("missing.json");
+ var reloader = MakeReloader(c => c.Paths = new[] { _path, missing });
+
+ reloader.ReloadNow();
+
+ var e = Assert.IsType(RequireErrored());
+ Assert.Equal(missing, e.Path);
+ Assert.Contains("unable to read file", e.Message);
+ RequireQuiet(TimeSpan.FromMilliseconds(100));
+ }
+
+ [Fact]
+ public void SkipsMissingPathsWhenConfigured()
+ {
+ Write(Flag1True);
+ var second = _dir.PathOf("second.json");
+ var reloader = MakeReloader(c =>
+ {
+ c.Paths = new[] { _path, second };
+ c.SkipMissingPaths = true;
+ c.SkipUnchanged = true;
+ });
+
+ // Step 1: one file exists and one does not. The reload succeeds with the existing file.
+ reloader.ReloadNow();
+ var result = RequireApplied();
+ Assert.Equal(new[] { "flag1" }, FlagKeys(result));
+ Assert.Equal(2, result.Files.Count);
+ Assert.Equal(_path, result.Files[0].Path);
+ Assert.True(result.Files[0].Present);
+ Assert.Equal(1, result.Files[0].Flags);
+ Assert.Equal(0, result.Files[0].Segments);
+ Assert.Equal(second, result.Files[1].Path);
+ Assert.False(result.Files[1].Present);
+ Assert.Equal(0, result.Files[1].Flags);
+
+ // Step 2: the missing file appears. Its data is merged in.
+ File.WriteAllText(second, @"{""flagValues"": {""flag2"": true}}");
+ reloader.ReloadNow();
+ result = RequireApplied();
+ Assert.Equal(new[] { "flag1", "flag2" }, FlagKeys(result));
+ Assert.True(result.Files[1].Present);
+ Assert.Equal(1, result.Files[1].Flags);
+
+ // Step 3: the file is deleted. Its data is gone and the reload still succeeds.
+ File.Delete(second);
+ reloader.ReloadNow();
+ result = RequireApplied();
+ Assert.Equal(new[] { "flag1" }, FlagKeys(result));
+ RequireQuiet(TimeSpan.FromMilliseconds(100));
+ }
+
+ [Fact]
+ public void MissingDirectoryIsAlsoSkippedWhenConfigured()
+ {
+ Write(Flag1True);
+ var inMissingDirectory = Path.Combine(_dir.PathOf("no-such-directory"), "data.json");
+ var reloader = MakeReloader(c =>
+ {
+ c.Paths = new[] { _path, inMissingDirectory };
+ c.SkipMissingPaths = true;
+ });
+
+ reloader.ReloadNow();
+
+ var result = RequireApplied();
+ Assert.Equal(new[] { "flag1" }, FlagKeys(result));
+ Assert.False(result.Files[1].Present);
+ }
+
+ [Fact]
+ public void InitialLoadAppliesTheData()
+ {
+ Write(Flag1True);
+ var reloader = MakeReloader();
+
+ reloader.ReloadNow();
+
+ var result = RequireApplied();
+ Assert.Equal(new[] { "flag1" }, FlagKeys(result));
+ var flag = Assert.IsType(result.Flags[0].Value.Item);
+ Assert.Equal(new[] { LdValue.Of(true) }, flag.Variations);
+ Assert.Single(result.Files);
+ Assert.True(result.Files[0].Present);
+ }
+
+ [Fact]
+ public void ReportsParseFailureAndAppliesNothing()
+ {
+ Write(Truncated);
+ var reloader = MakeReloader();
+
+ reloader.ReloadNow();
+
+ var e = Assert.IsType(RequireErrored());
+ Assert.Equal(_path, e.Path);
+ Assert.Contains("error parsing file", e.Message);
+ AssertLogMessageRegex(true, LogLevel.Error, "Unable to load flags: error parsing file");
+ RequireQuiet(TimeSpan.FromMilliseconds(50));
+ }
+
+ [Fact]
+ public void ReportsMergeFailure()
+ {
+ Write(Flag1True);
+ var second = _dir.PathOf("second.json");
+ File.WriteAllText(second, Flag1False);
+ var reloader = MakeReloader(c => c.Paths = new[] { _path, second });
+
+ reloader.ReloadNow();
+
+ var e = RequireErrored();
+ Assert.IsType(e);
+ Assert.Equal("flag \"flag1\" is specified by multiple files", e.Message);
+ _applied.ExpectNoValue(TimeSpan.FromMilliseconds(50));
+ }
+
+ [Fact]
+ public void ReadFailureIsReportedAsAReadException()
+ {
+ Write(Flag1True);
+ var reader = new ThrowingFileReader();
+ var reloader = MakeReloader(c => c.FileReader = reader);
+
+ reloader.ReloadNow();
+
+ var e = Assert.IsType(RequireErrored());
+ Assert.Equal(_path, e.Path);
+ Assert.Contains("unable to read file", e.Message);
+ Assert.IsType(e.InnerException);
+ }
+
+ [Fact]
+ public void DebounceCoalescesTriggers()
+ {
+ // The settle window is much longer than the whole trigger burst, so that a scheduling
+ // stall during the burst cannot let the debounce fire early and split the reloads.
+ // SkipUnchanged stays off so that every reload is observable as an Apply.
+ Write(Flag1True);
+ var reloader = MakeReloader(c => c.DebounceDelay = TimeSpan.FromMilliseconds(400));
+
+ Write(Flag1False);
+ for (var i = 0; i < 20; i++)
+ {
+ reloader.Trigger();
+ Thread.Sleep(1);
+ }
+
+ RequireApplied();
+ // The burst must coalesce into one reload, with a tolerance of one more: a debounce
+ // tick racing a fresh trigger can produce a single extra serialized reload.
+ Thread.Sleep(600);
+ var extraApplies = ApplyCount - 1;
+ Assert.True(extraApplies <= 1, "trigger burst was not coalesced: " + extraApplies + " extra applies");
+ }
+
+ [Fact]
+ public void DebounceWindowIsExtendedByEachTrigger()
+ {
+ // The debounce is a settle window: each trigger moves the deadline out again. A stream
+ // of notifications spaced closer together than the window must produce no reload while
+ // the stream continues, and exactly one reload after it stops.
+ // The notifications are spaced far inside the window so that a scheduling stall on a
+ // busy machine cannot let the window expire between two of them.
+ Write(Flag1True);
+ var window = TimeSpan.FromMilliseconds(600);
+ var reloader = MakeReloader(c => c.DebounceDelay = window);
+
+ var stop = DateTime.UtcNow + TimeSpan.FromTicks(window.Ticks * 3);
+ while (DateTime.UtcNow < stop)
+ {
+ reloader.Trigger();
+ Thread.Sleep(50);
+ }
+ Assert.True(ApplyCount == 0, "a reload ran while change notifications were still arriving");
+
+ RequireApplied();
+ RequireQuiet(TimeSpan.FromTicks(window.Ticks * 2));
+ Assert.Equal(1, ApplyCount);
+ }
+
+ [Fact]
+ public void TriggerWithoutDebounceReloadsOnAWorkerThread()
+ {
+ Write(Flag1True);
+ var applyThread = -1;
+ var reloader = MakeReloader(c =>
+ {
+ c.DebounceDelay = TimeSpan.Zero;
+ c.Apply = result =>
+ {
+ applyThread = Thread.CurrentThread.ManagedThreadId;
+ _applied.Enqueue(result);
+ };
+ });
+
+ reloader.Trigger();
+
+ RequireApplied();
+ Assert.NotEqual(Thread.CurrentThread.ManagedThreadId, applyThread);
+ }
+
+ [Fact]
+ public void ReportsIdenticalFailureOnlyOnce()
+ {
+ Write(Truncated);
+ var reloader = MakeReloader(c => c.RetryDelay = TimeSpan.FromMilliseconds(10));
+
+ reloader.ReloadNow();
+ RequireErrored();
+
+ // The automatic retries keep failing identically. The error is neither reported to
+ // OnError again nor logged at error level again.
+ RequireQuiet(TimeSpan.FromMilliseconds(200));
+ Assert.Equal(1, ErrorLogCount());
+
+ // A different failure is a new report.
+ Write(@"{""flagValues"": {bad}}");
+ var e = RequireErrored();
+ Assert.Contains("error parsing file", e.Message);
+
+ // Success re-arms reporting: the same failure recurring afterward is reported again.
+ Write(Flag1True);
+ RequireApplied();
+ Write(Truncated);
+ reloader.Trigger();
+ RequireErrored();
+ }
+
+ [Fact]
+ public void RetriesAfterFailureWithoutFurtherTriggers()
+ {
+ Write(Flag1True);
+ var reloader = MakeReloader(c => c.RetryDelay = TimeSpan.FromMilliseconds(20));
+ reloader.ReloadNow();
+ RequireApplied();
+
+ Write(Truncated);
+ reloader.Trigger();
+ RequireErrored();
+
+ // Fix the file without triggering. Only the automatic retry can observe the fix.
+ Write(Flag1False);
+ var result = RequireApplied();
+ var flag = Assert.IsType(result.Flags[0].Value.Item);
+ Assert.Equal(new[] { LdValue.Of(false) }, flag.Variations);
+ }
+
+ [Fact]
+ public void InitialLoadFailureArmsTheRetry()
+ {
+ Write(Truncated);
+ var reloader = MakeReloader(c => c.RetryDelay = TimeSpan.FromMilliseconds(20));
+
+ reloader.ReloadNow();
+ RequireErrored();
+
+ Write(Flag1True);
+ RequireApplied();
+ }
+
+ [Fact]
+ public void StopsRetryingAfterSuccess()
+ {
+ Write(Truncated);
+ var reloader = MakeReloader(c =>
+ {
+ c.RetryDelay = TimeSpan.FromMilliseconds(10);
+ c.SkipUnchanged = true;
+ });
+ reloader.ReloadNow();
+ RequireErrored();
+
+ Write(Flag1True);
+ RequireApplied();
+
+ // After the successful reload there are no further attempts. A changed file with no
+ // trigger must not be picked up.
+ Write(Flag1False);
+ RequireQuiet(TimeSpan.FromMilliseconds(100));
+ }
+
+ [Fact]
+ public void NoRetryWhenRetryDelayIsZero()
+ {
+ Write(Truncated);
+ var reloader = MakeReloader(c => c.RetryDelay = TimeSpan.Zero);
+ reloader.ReloadNow();
+ RequireErrored();
+
+ Write(Flag1True);
+ RequireQuiet(TimeSpan.FromMilliseconds(200));
+ }
+
+ [Fact]
+ public void SkipUnchangedSuppressesIdenticalContent()
+ {
+ Write(Flag1True);
+ var reloader = MakeReloader(c => c.SkipUnchanged = true);
+ reloader.ReloadNow();
+ RequireApplied();
+
+ reloader.Trigger();
+ RequireQuiet(TimeSpan.FromMilliseconds(100));
+
+ Write(Flag1False);
+ reloader.Trigger();
+ RequireApplied();
+ }
+
+ [Fact]
+ public void RecoveryAppliesEvenWhenContentIsUnchanged()
+ {
+ Write(Flag1True);
+ var reloader = MakeReloader(c => c.SkipUnchanged = true);
+ reloader.ReloadNow();
+ RequireApplied();
+
+ // A reload fails. Consumers hear OnError and may move to an interrupted state.
+ File.Delete(_path);
+ reloader.Trigger();
+ RequireErrored();
+
+ // The file comes back with identical content. The success must be applied despite
+ // SkipUnchanged, because only Apply tells the consumer the interruption is over.
+ Write(Flag1True);
+ reloader.Trigger();
+ RequireApplied();
+
+ // Once recovered, identical content skips again.
+ reloader.Trigger();
+ RequireQuiet(TimeSpan.FromMilliseconds(100));
+ }
+
+ [Fact]
+ public void AppliesEveryReloadWhenSkipUnchangedIsOff()
+ {
+ Write(Flag1True);
+ var reloader = MakeReloader();
+ reloader.ReloadNow();
+ RequireApplied();
+ reloader.Trigger();
+ RequireApplied();
+ }
+
+ [Fact]
+ public void MergesMultipleFilesInOrder()
+ {
+ var first = _dir.PathOf("first.json");
+ var second = _dir.PathOf("second.json");
+ File.WriteAllText(first, @"{""flags"": {""flag1"": {""key"": ""flag1"", ""version"": 1}}}");
+ File.WriteAllText(second, @"{""flags"": {""flag1"": {""key"": ""flag1"", ""version"": 2}}}");
+ var reloader = MakeReloader(c =>
+ {
+ c.Paths = new[] { first, second };
+ c.DuplicateKeysHandling = FileDataDuplicateKeysHandling.Ignore;
+ });
+
+ reloader.ReloadNow();
+
+ var result = RequireApplied();
+ Assert.Single(result.Flags);
+ Assert.Equal(1, result.Flags[0].Value.Version);
+ Assert.Equal(2, result.Files.Count);
+ Assert.Equal(1, result.Files[0].Flags);
+ Assert.Equal(0, result.Files[1].Flags);
+ }
+
+ [Fact]
+ public void DoesNothingAfterDispose()
+ {
+ Write(Flag1True);
+ var reloader = MakeReloader(c => c.DebounceDelay = TimeSpan.FromMilliseconds(10));
+ reloader.ReloadNow();
+ RequireApplied();
+
+ reloader.Dispose();
+ reloader.Dispose();
+ reloader.Trigger();
+ reloader.ReloadNow();
+ RequireQuiet(TimeSpan.FromMilliseconds(100));
+ }
+
+ [Fact]
+ public void PendingRetryIsCanceledByDispose()
+ {
+ Write(Truncated);
+ var reloader = MakeReloader(c => c.RetryDelay = TimeSpan.FromMilliseconds(50));
+ reloader.ReloadNow();
+ RequireErrored();
+
+ reloader.Dispose();
+ Write(Flag1True);
+ RequireQuiet(TimeSpan.FromMilliseconds(200));
+ }
+
+ [Fact]
+ public void DisposeDoesNotWaitForAnInFlightReload()
+ {
+ Write(Flag1True);
+ var applyEntered = new ManualResetEventSlim();
+ var applyRelease = new ManualResetEventSlim();
+ var reloader = MakeReloader(c =>
+ {
+ c.DebounceDelay = TimeSpan.Zero;
+ c.Apply = result =>
+ {
+ applyEntered.Set();
+ applyRelease.Wait();
+ _applied.Enqueue(result);
+ };
+ });
+
+ try
+ {
+ reloader.Trigger();
+ Assert.True(applyEntered.Wait(TestTimeout), "timed out waiting for the reload to start");
+
+ // The reload is parked inside Apply. Dispose must return anyway, because a reload wedged
+ // in blocking I/O must not be able to wedge shutdown.
+ var disposed = Task.Run(() => reloader.Dispose());
+ Assert.True(disposed.Wait(TestTimeout), "Dispose blocked on an in-flight reload");
+ }
+ finally
+ {
+ // The parked reload must always be released, or a failing assertion leaves a thread
+ // blocked forever.
+ applyRelease.Set();
+ }
+ RequireApplied();
+ }
+
+ [Fact]
+ public void ReloadsAreSerialized()
+ {
+ Write(Flag1True);
+ var concurrent = 0;
+ var maxConcurrent = 0;
+ var reloader = MakeReloader(c =>
+ {
+ c.DebounceDelay = TimeSpan.Zero;
+ c.Apply = result =>
+ {
+ var now = Interlocked.Increment(ref concurrent);
+ InterlockedMax(ref maxConcurrent, now);
+ Thread.Sleep(20);
+ Interlocked.Decrement(ref concurrent);
+ _applied.Enqueue(result);
+ };
+ });
+
+ var tasks = Enumerable.Range(0, 8).Select(i => Task.Run(() =>
+ {
+ reloader.ReloadNow();
+ reloader.Trigger();
+ })).ToArray();
+ Task.WaitAll(tasks);
+ for (var i = 0; i < 8; i++)
+ {
+ RequireApplied();
+ }
+
+ Assert.Equal(1, maxConcurrent);
+ }
+
+ private static void InterlockedMax(ref int target, int value)
+ {
+ int current;
+ while ((current = Volatile.Read(ref target)) < value)
+ {
+ if (Interlocked.CompareExchange(ref target, value, current) == current)
+ {
+ return;
+ }
+ }
+ }
+
+ [Fact]
+ public void ApplyExceptionOnTimerThreadIsLoggedAndDoesNotStopLaterReloads()
+ {
+ Write(Flag1True);
+ var failNext = true;
+ var reloader = MakeReloader(c =>
+ {
+ c.DebounceDelay = TimeSpan.FromMilliseconds(10);
+ c.Apply = result =>
+ {
+ if (failNext)
+ {
+ failNext = false;
+ throw new InvalidOperationException("consumer failed");
+ }
+ _applied.Enqueue(result);
+ };
+ });
+
+ reloader.Trigger();
+ AssertEventually(() => LogCapture.HasMessageWithRegex(LogLevel.Error, "Unexpected error while reloading file data"));
+
+ reloader.Trigger();
+ RequireApplied();
+ }
+
+ private static void AssertEventually(Func condition)
+ {
+ var deadline = DateTime.UtcNow + TestTimeout;
+ while (!condition() && DateTime.UtcNow < deadline)
+ {
+ Thread.Sleep(10);
+ }
+ Assert.True(condition());
+ }
+
+ private class ThrowingFileReader : FileDataTypes.IFileReader
+ {
+ public string ReadAllText(string path) => throw new IOException("simulated read error");
+ }
+ }
+}
diff --git a/pkgs/sdk/server/test/Internal/FileLoading/FileDataWatcherTest.cs b/pkgs/sdk/server/test/Internal/FileLoading/FileDataWatcherTest.cs
new file mode 100644
index 000000000..650d6c003
--- /dev/null
+++ b/pkgs/sdk/server/test/Internal/FileLoading/FileDataWatcherTest.cs
@@ -0,0 +1,161 @@
+using System;
+using System.IO;
+using LaunchDarkly.TestHelpers;
+using Xunit;
+using Xunit.Abstractions;
+
+namespace LaunchDarkly.Sdk.Server.Internal.FileLoading
+{
+ public class FileDataWatcherTest : BaseTest, IDisposable
+ {
+ private static readonly TimeSpan TestTimeout = TimeSpan.FromSeconds(10);
+ private static readonly TimeSpan QuietPeriod = TimeSpan.FromMilliseconds(300);
+
+ private readonly TempDirectory _dir = TempDirectory.Create();
+ private readonly EventSink _changed = new EventSink();
+ private FileDataWatcher _watcher;
+
+ public FileDataWatcherTest(ITestOutputHelper testOutput) : base(testOutput) { }
+
+ public void Dispose()
+ {
+ _watcher?.Dispose();
+ _dir.Dispose();
+ }
+
+ private FileDataWatcher StartWatcher(params string[] paths)
+ {
+ _watcher = new FileDataWatcher(paths, () => _changed.Enqueue(true), TestLogger);
+ return _watcher;
+ }
+
+ private void RequireChange() => _changed.ExpectValue(TestTimeout);
+
+ [Fact]
+ public void DetectsModification()
+ {
+ var path = _dir.PathOf("data.json");
+ File.WriteAllText(path, "one");
+ StartWatcher(path);
+
+ File.WriteAllText(path, "two");
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsFileAppearing()
+ {
+ var path = _dir.PathOf("data.json");
+ StartWatcher(path);
+
+ File.WriteAllText(path, "created");
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsFileDisappearing()
+ {
+ var path = _dir.PathOf("data.json");
+ File.WriteAllText(path, "one");
+ StartWatcher(path);
+
+ File.Delete(path);
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsFileReplacedByRename()
+ {
+ var path = _dir.PathOf("data.json");
+ var temp = _dir.PathOf("data.json.tmp");
+ File.WriteAllText(path, "one");
+ StartWatcher(path);
+
+ File.WriteAllText(temp, "two");
+ File.Delete(path);
+ File.Move(temp, path);
+ RequireChange();
+ }
+
+ [Fact]
+ public void IgnoresOtherFilesInTheSameDirectory()
+ {
+ var path = _dir.PathOf("data.json");
+ var other = _dir.PathOf("other.json");
+ File.WriteAllText(path, "one");
+ StartWatcher(path);
+
+ File.WriteAllText(other, "unrelated");
+ File.WriteAllText(other, "unrelated again");
+ _changed.ExpectNoValue(QuietPeriod);
+ }
+
+ [Fact]
+ public void WatchesFilesInDifferentDirectories()
+ {
+ var subdir = _dir.PathOf("sub");
+ Directory.CreateDirectory(subdir);
+ var path1 = _dir.PathOf("one.json");
+ var path2 = Path.Combine(subdir, "two.json");
+ File.WriteAllText(path1, "one");
+ File.WriteAllText(path2, "two");
+ StartWatcher(path1, path2);
+
+ File.WriteAllText(path2, "two-changed");
+ RequireChange();
+ }
+
+ [Fact]
+ public void DetectsChangeToFirstOfMultipleFiles()
+ {
+ var path1 = _dir.PathOf("one.json");
+ var path2 = _dir.PathOf("two.json");
+ File.WriteAllText(path1, "one");
+ File.WriteAllText(path2, "two");
+ StartWatcher(path1, path2);
+
+ File.WriteAllText(path1, "one-changed");
+ RequireChange();
+ }
+
+ [Fact]
+ public void RelativePathsAreResolvedAgainstTheCurrentDirectory()
+ {
+ var name = "ld-watcher-test-" + Guid.NewGuid().ToString("N") + ".json";
+ var fullPath = Path.GetFullPath(name);
+ try
+ {
+ StartWatcher(name);
+ File.WriteAllText(fullPath, "created");
+ RequireChange();
+ }
+ finally
+ {
+ _watcher?.Dispose();
+ _watcher = null;
+ File.Delete(fullPath);
+ }
+ }
+
+ [Fact]
+ public void StopsOnDispose()
+ {
+ var path = _dir.PathOf("data.json");
+ File.WriteAllText(path, "one");
+ var watcher = StartWatcher(path);
+
+ watcher.Dispose();
+ watcher.Dispose();
+
+ File.WriteAllText(path, "two");
+ _changed.ExpectNoValue(QuietPeriod);
+ }
+
+ [Fact]
+ public void MissingDirectoryThrows()
+ {
+ var path = Path.Combine(_dir.PathOf("no-such-directory"), "data.json");
+ Assert.ThrowsAny(() => StartWatcher(path));
+ }
+ }
+}
diff --git a/pkgs/sdk/server/test/TestUtils.cs b/pkgs/sdk/server/test/TestUtils.cs
index 932df2adb..457861690 100644
--- a/pkgs/sdk/server/test/TestUtils.cs
+++ b/pkgs/sdk/server/test/TestUtils.cs
@@ -74,6 +74,30 @@ internal static JsonTestValue DataSetAsJson(FullDataSet data)
}
}
+ public class TempDirectory : IDisposable
+ {
+ public string Path { get; }
+
+ public static TempDirectory Create() => new TempDirectory();
+
+ private TempDirectory()
+ {
+ Path = System.IO.Path.Combine(System.IO.Path.GetTempPath(), "ld-test-" + Guid.NewGuid().ToString("N"));
+ Directory.CreateDirectory(Path);
+ }
+
+ public string PathOf(string name) => System.IO.Path.Combine(Path, name);
+
+ public void Dispose()
+ {
+ try
+ {
+ Directory.Delete(Path, true);
+ }
+ catch { }
+ }
+ }
+
public class TempFile : IDisposable
{
public string Path { get; }