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 00000000..ef12149d
--- /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 00000000..ad257447
--- /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 00000000..0dc4a8bc
--- /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 00000000..3d1315ac
--- /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 00000000..ac78d104
--- /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 00000000..b68c28c4
--- /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 00000000..c6a047a2
--- /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 00000000..e0ef7a9f
--- /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 00000000..91280a16
--- /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 00000000..a1138124
--- /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 00000000..650d6c00
--- /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 932df2ad..45786169 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; }