diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataPoller.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataPoller.java
new file mode 100644
index 00000000..2cc5e488
--- /dev/null
+++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataPoller.java
@@ -0,0 +1,134 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import java.io.Closeable;
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.attribute.BasicFileAttributes;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Objects;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+/**
+ * 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, alone or together with them.
+ * 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 whose attributes cannot be read counts as
+ * absent.
+ *
+ * The poller samples the files once per interval and compares only the modification time and the
+ * 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 {@link FileDataReloader}, whose debouncing and skip-unchanged handling
+ * absorb the excess.
+ */
+final class FileDataPoller implements Closeable {
+ private final List paths;
+ private final Runnable onChange;
+ private final ScheduledThreadPoolExecutor executor;
+ private final AtomicBoolean closed = new AtomicBoolean(false);
+ private List last;
+
+ /**
+ * Creates a started poller. It examines the files once before it returns, so only later changes
+ * invoke the callback. Call {@link #close()} to stop it.
+ *
+ * @param paths the files to examine
+ * @param interval the time between examinations
+ * @param onChange called when any file changed since the previous examination
+ */
+ FileDataPoller(List paths, Duration interval, Runnable onChange) {
+ this.paths = new ArrayList<>(paths);
+ this.onChange = onChange;
+ this.last = observeAll(this.paths);
+ ThreadFactory threadFactory = runnable -> {
+ Thread t = new Thread(runnable, "LaunchDarkly-FileDataPoller");
+ t.setDaemon(true);
+ return t;
+ };
+ this.executor = new ScheduledThreadPoolExecutor(1, threadFactory);
+ this.executor.setExecuteExistingDelayedTasksAfterShutdownPolicy(false);
+ long millis = Math.max(interval.toMillis(), 1);
+ this.executor.scheduleWithFixedDelay(this::examine, millis, millis, TimeUnit.MILLISECONDS);
+ }
+
+ /**
+ * 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 close returns. Consumers tolerate a late call, as they do for a
+ * late reload.
+ */
+ @Override
+ public void close() {
+ if (closed.getAndSet(true)) {
+ return;
+ }
+ executor.shutdownNow();
+ }
+
+ private void examine() {
+ if (closed.get()) {
+ return;
+ }
+ List current = observeAll(paths);
+ boolean changed = !current.equals(last);
+ last = current;
+ if (changed && !closed.get()) {
+ onChange.run();
+ }
+ }
+
+ static List observeAll(List paths) {
+ List states = new ArrayList<>(paths.size());
+ for (Path path : paths) {
+ states.add(observe(path));
+ }
+ return states;
+ }
+
+ private static FileState observe(Path path) {
+ try {
+ BasicFileAttributes attributes = Files.readAttributes(path, BasicFileAttributes.class);
+ return new FileState(true, attributes.lastModifiedTime().toMillis(), attributes.size());
+ } catch (IOException | RuntimeException e) {
+ return FileState.ABSENT;
+ }
+ }
+
+ /**
+ * The observed state of one file, or its absence.
+ */
+ static final class FileState {
+ static final FileState ABSENT = new FileState(false, 0, 0);
+
+ final boolean exists;
+ final long modifiedTimeMillis;
+ final long size;
+
+ FileState(boolean exists, long modifiedTimeMillis, long size) {
+ this.exists = exists;
+ this.modifiedTimeMillis = modifiedTimeMillis;
+ this.size = size;
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ if (other instanceof FileState) {
+ FileState o = (FileState) other;
+ return exists == o.exists && modifiedTimeMillis == o.modifiedTimeMillis && size == o.size;
+ }
+ return false;
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(exists, modifiedTimeMillis, size);
+ }
+ }
+}
diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataReloader.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataReloader.java
new file mode 100644
index 00000000..259ee185
--- /dev/null
+++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataReloader.java
@@ -0,0 +1,324 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.logging.LDLogger;
+import com.launchdarkly.logging.LogValues;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FileDataException;
+import com.launchdarkly.sdk.server.integrations.OverrideFileLoader.LoadResult;
+
+import java.io.Closeable;
+import java.time.Duration;
+import java.util.Arrays;
+import java.util.concurrent.ScheduledFuture;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+/**
+ * Owns the reload cycle for the files of an override source. It serializes reloads, debounces
+ * change signals, retains the last good result on failure by not calling the apply callback,
+ * retries after failures, and skips applications whose content did not change.
+ *
+ * The reload work itself is delegated to a {@link Loader}. The loader reads every configured file
+ * and returns the merged result, or throws {@link FileDataException} when a file cannot be read or
+ * parsed or when the files cannot be merged.
+ *
+ * The file data source does not use this class. It keeps its own reload behavior.
+ */
+final class FileDataReloader implements Closeable {
+ /**
+ * A settle window long enough to coalesce the burst of change notifications produced by a
+ * single file edit, and short enough to stay responsive.
+ */
+ static final Duration DEFAULT_DEBOUNCE_DELAY = Duration.ofMillis(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.
+ */
+ static final Duration DEFAULT_RETRY_DELAY = Duration.ofSeconds(1);
+
+ /**
+ * Performs one full load of all configured files.
+ */
+ interface Loader {
+ LoadResult load() throws FileDataException;
+ }
+
+ /**
+ * Receives the outcome of each reload.
+ *
+ * Both methods are called with the reload lock held, so calls never overlap with each other
+ * or with another reload. Neither method may call back into {@link FileDataReloader#close()}.
+ */
+ interface Handler {
+ /**
+ * Called with each successfully merged result.
+ *
+ * @param result the merged data
+ */
+ void apply(LoadResult result);
+
+ /**
+ * Called when a reload fails, once per distinct failure. With automatic retries, repeats of an
+ * identical failure do not call this again. A success re-arms it. The reloader logs failures
+ * itself, so implementations only need to update their own state.
+ *
+ * @param e the failure
+ */
+ void onError(FileDataException e);
+ }
+
+ private final Loader loader;
+ private final Handler handler;
+ private final LDLogger logger;
+ private final long debounceDelayMillis;
+ private final long retryDelayMillis;
+ private final boolean skipUnchanged;
+
+ // Serializes the load work between reloadNow and the worker thread.
+ private final Object reloadLock = new Object();
+ private byte[] lastGoodHash = null;
+ private String lastErrorMessage = null;
+
+ // Guards the worker and the two timers.
+ private final Object timerLock = new Object();
+ private ScheduledThreadPoolExecutor worker = null;
+ private ScheduledFuture> debounceFuture = null;
+ private ScheduledFuture> retryFuture = null;
+
+ private final AtomicBoolean closed = new AtomicBoolean(false);
+
+ /**
+ * Creates a reloader. The caller performs the initial load with {@link #reloadNow()}, routes
+ * change signals to {@link #trigger()}, and calls {@link #close()} when finished.
+ *
+ * The worker thread starts on the first {@link #trigger()} call or the first failed load that
+ * arms a retry, not here. Construction alone does not start a thread.
+ *
+ * @param loader performs each full load
+ * @param handler receives results and failures
+ * @param logger receives log output about reloads and failures
+ * @param debounceDelay how long to wait after a trigger for further triggers to settle before
+ * reloading. Zero or negative reloads on every trigger without waiting.
+ * @param retryDelay how long to wait after a failed reload before retrying without a trigger.
+ * Zero or negative disables the automatic retry.
+ * @param skipUnchanged if true, a successful load whose raw file contents are identical to the
+ * last applied contents does not call the apply callback
+ */
+ FileDataReloader(
+ Loader loader,
+ Handler handler,
+ LDLogger logger,
+ Duration debounceDelay,
+ Duration retryDelay,
+ boolean skipUnchanged
+ ) {
+ this.loader = loader;
+ this.handler = handler;
+ this.logger = logger;
+ this.debounceDelayMillis = debounceDelay == null ? 0 : debounceDelay.toMillis();
+ this.retryDelayMillis = retryDelay == null ? 0 : retryDelay.toMillis();
+ this.skipUnchanged = skipUnchanged;
+ }
+
+ /**
+ * 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.
+ */
+ void reloadNow() {
+ if (!reload() && retryDelayMillis > 0) {
+ armRetry();
+ }
+ }
+
+ /**
+ * 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 extend the settle
+ * window, so a burst of signals produces one reload after the burst ends.
+ */
+ void trigger() {
+ synchronized (timerLock) {
+ if (closed.get()) {
+ return;
+ }
+ ScheduledThreadPoolExecutor executor = ensureWorker();
+ if (debounceFuture != null) {
+ debounceFuture.cancel(false);
+ }
+ debounceFuture = executor.schedule(() -> runReload(false), Math.max(debounceDelayMillis, 0),
+ TimeUnit.MILLISECONDS);
+ }
+ }
+
+ /**
+ * 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 may still
+ * deliver its result shortly after close returns, and consumers tolerate that. A reload that
+ * has not yet reached its callbacks when close is called does not invoke them.
+ */
+ @Override
+ public void close() {
+ if (closed.getAndSet(true)) {
+ return;
+ }
+ synchronized (timerLock) {
+ if (debounceFuture != null) {
+ debounceFuture.cancel(false);
+ debounceFuture = null;
+ }
+ if (retryFuture != null) {
+ retryFuture.cancel(false);
+ retryFuture = null;
+ }
+ if (worker != null) {
+ worker.shutdownNow();
+ }
+ }
+ }
+
+ /**
+ * Reports whether the worker thread has been started. Visible for tests.
+ *
+ * @return true if a worker exists
+ */
+ boolean hasWorker() {
+ synchronized (timerLock) {
+ return worker != null;
+ }
+ }
+
+ /**
+ * Formats a failure for logging. The exception's own description requires a cause, and a
+ * duplicate key failure has none, so this builds the text from the parts that are present.
+ *
+ * @param e the failure
+ * @return the message, followed by the cause in brackets when there is one
+ */
+ static String describe(FileDataException e) {
+ StringBuilder s = new StringBuilder();
+ if (e.getMessage() != null) {
+ s.append(e.getMessage());
+ }
+ if (e.getCause() != null) {
+ if (s.length() > 0) {
+ s.append(" ");
+ }
+ s.append("[").append(e.getCause().toString()).append("]");
+ }
+ return s.toString();
+ }
+
+ private ScheduledThreadPoolExecutor ensureWorker() {
+ if (worker == null) {
+ ThreadFactory threadFactory = runnable -> {
+ Thread t = new Thread(runnable, "LaunchDarkly-FileDataReloader");
+ t.setDaemon(true);
+ return t;
+ };
+ ScheduledThreadPoolExecutor executor = new ScheduledThreadPoolExecutor(1, threadFactory);
+ executor.setRemoveOnCancelPolicy(true);
+ executor.setExecuteExistingDelayedTasksAfterShutdownPolicy(false);
+ worker = executor;
+ }
+ return worker;
+ }
+
+ private void armRetry() {
+ synchronized (timerLock) {
+ if (closed.get()) {
+ return;
+ }
+ // An already armed retry keeps its earlier deadline.
+ if (retryFuture != null && !retryFuture.isDone()) {
+ return;
+ }
+ retryFuture = ensureWorker().schedule(() -> runReload(true), retryDelayMillis, TimeUnit.MILLISECONDS);
+ }
+ }
+
+ // Runs on the worker thread for a debounced or a retried reload.
+ private void runReload(boolean isRetry) {
+ try {
+ if (isRetry) {
+ logger.debug("Retrying flag data load after earlier failure");
+ } else {
+ logger.info("Reloading flag data after detecting a change");
+ }
+ synchronized (timerLock) {
+ // This reload supersedes a pending retry. It either succeeds, or it fails and arms a fresh
+ // retry below.
+ if (retryFuture != null) {
+ retryFuture.cancel(false);
+ retryFuture = null;
+ }
+ }
+ if (!reload() && retryDelayMillis > 0) {
+ armRetry();
+ }
+ } catch (RuntimeException e) {
+ // The executor would otherwise swallow the exception and the worker would go quiet.
+ logger.error("Unexpected error while reloading flag data: {}", LogValues.exceptionSummary(e));
+ logger.debug(LogValues.exceptionTrace(e));
+ }
+ }
+
+ // Performs one full load of all configured files. It reports whether 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 because entries are combined across files in order, so a
+ // change to one file can alter which file wins for a key.
+ private boolean reload() {
+ synchronized (reloadLock) {
+ // A trigger already queued when close was called can still reach here.
+ if (closed.get()) {
+ return true;
+ }
+
+ LoadResult result;
+ try {
+ result = loader.load();
+ } catch (FileDataException e) {
+ return fail(e);
+ }
+
+ // Close may have happened while the files were being read. Deliver nothing in that case.
+ if (closed.get()) {
+ 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 and may have moved to an interrupted
+ // state. Only apply tells it that things are good again.
+ boolean recovering = lastErrorMessage != null;
+ lastErrorMessage = null;
+ byte[] hash = result.getContentHash();
+ if (skipUnchanged && !recovering && hash != null && Arrays.equals(hash, lastGoodHash)) {
+ return true;
+ }
+ lastGoodHash = hash;
+ handler.apply(result);
+ return true;
+ }
+ }
+
+ private boolean fail(FileDataException e) {
+ // Close may have happened while the files were being read. Deliver nothing and report success
+ // so that no retry is armed.
+ if (closed.get()) {
+ return true;
+ }
+ // With automatic retries, a persistent failure would repeat the same log entry and the same
+ // callback on every attempt. Repeats of an identical failure are logged at debug level and do
+ // not call onError again. A consumer therefore sees one report per distinct failure.
+ String description = describe(e);
+ if (description.equals(lastErrorMessage)) {
+ logger.debug("Unable to load flags: {}", description);
+ return false;
+ }
+ lastErrorMessage = description;
+ logger.error("Unable to load flags: {}", description);
+ handler.onError(e);
+ return false;
+ }
+}
diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataSourceParsing.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataSourceParsing.java
index 81ff8cdc..bae5309d 100644
--- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataSourceParsing.java
+++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataSourceParsing.java
@@ -193,6 +193,31 @@ private FlagFactory() {}
static ItemDescriptor flagFromJson(LDValue jsonTree, int version) {
return FEATURES.deserialize(replaceVersion(jsonTree, version).toJsonString());
}
+
+ /**
+ * Constructs a flag from raw JSON and keeps the version that the document specifies. The
+ * override source uses this. The file data source keeps assigning its load version.
+ */
+ static ItemDescriptor flagFromJson(LDValue jsonTree) {
+ return FEATURES.deserialize(jsonTree.toJsonString());
+ }
+
+ /**
+ * Constructs a flag that is off and serves the given value as its single variation for every
+ * context. The document supplies no version, so the flag has version zero. This is the shape
+ * that the override source uses for value-only entries. The file data source keeps
+ * {@link #flagWithValue(String, LDValue, int)}.
+ */
+ static ItemDescriptor offFlagWithValue(String key, LDValue jsonValue) {
+ LDValue o = LDValue.buildObject()
+ .put("key", key)
+ .put("version", 0)
+ .put("on", false)
+ .put("offVariation", 0)
+ .put("variations", LDValue.buildArray().add(jsonValue).build())
+ .build();
+ return FEATURES.deserialize(o.toJsonString());
+ }
/**
* Constructs a flag that always returns the same value. This is done by giving it a single
@@ -214,6 +239,14 @@ static ItemDescriptor flagWithValue(String key, LDValue jsonValue, int version)
static ItemDescriptor segmentFromJson(LDValue jsonTree, int version) {
return SEGMENTS.deserialize(replaceVersion(jsonTree, version).toJsonString());
}
+
+ /**
+ * Constructs a segment from raw JSON and keeps the version that the document specifies. The
+ * override source uses this.
+ */
+ static ItemDescriptor segmentFromJson(LDValue jsonTree) {
+ return SEGMENTS.deserialize(jsonTree.toJsonString());
+ }
private static LDValue replaceVersion(LDValue objectValue, int version) {
ObjectBuilder b = LDValue.buildObject();
diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataWatcher.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataWatcher.java
new file mode 100644
index 00000000..172f1a0f
--- /dev/null
+++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataWatcher.java
@@ -0,0 +1,155 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.logging.LDLogger;
+import com.launchdarkly.logging.LogValues;
+
+import java.io.Closeable;
+import java.io.IOException;
+import java.nio.file.ClosedWatchServiceException;
+import java.nio.file.FileSystems;
+import java.nio.file.Path;
+import java.nio.file.WatchEvent;
+import java.nio.file.WatchKey;
+import java.nio.file.WatchService;
+import java.nio.file.Watchable;
+import java.util.HashSet;
+import java.util.Set;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import static java.nio.file.StandardWatchEventKinds.ENTRY_CREATE;
+import static java.nio.file.StandardWatchEventKinds.ENTRY_DELETE;
+import static java.nio.file.StandardWatchEventKinds.ENTRY_MODIFY;
+import static java.nio.file.StandardWatchEventKinds.OVERFLOW;
+
+/**
+ * Watches a set of files for changes on a worker thread and invokes a callback when one of them
+ * is created, modified, or deleted.
+ *
+ * The Java file system API reports changes at the directory level, so the watcher registers the
+ * parent directory of every file. A file that does not exist yet is reported when it appears, as
+ * long as its directory exists when the watcher is created. A directory that does not exist is an
+ * error from {@link #create(Iterable, LDLogger)}.
+ *
+ * The callback is invoked for every relevant notification, including the burst of notifications
+ * that a single edit can produce, and once more right after the watcher starts so that a change
+ * made between an initial load and the start of watching is not missed. Feed it into a
+ * {@link FileDataReloader}, which coalesces the burst and skips no-op reloads.
+ */
+final class FileDataWatcher implements Closeable, Runnable {
+ private final WatchService watchService;
+ private final Set watchedFilePaths;
+ private final LDLogger logger;
+ private final Thread thread;
+ private final AtomicBoolean stopped = new AtomicBoolean(false);
+ private volatile Runnable onChange;
+
+ /**
+ * Creates a watcher for the given files. Nothing is watched until {@link #start(Runnable)}.
+ *
+ * @param filePaths the files to watch, absolute or relative to the working directory
+ * @param logger the logger
+ * @return the watcher
+ * @throws IOException if the watch service cannot be created or a parent directory cannot be
+ * registered, for example because it does not exist
+ */
+ static FileDataWatcher create(Iterable filePaths, LDLogger logger) throws IOException {
+ Set directoryPaths = new HashSet<>();
+ Set absoluteFilePaths = new HashSet<>();
+ for (Path p : filePaths) {
+ Path absolutePath = p.toAbsolutePath().normalize();
+ absoluteFilePaths.add(absolutePath);
+ directoryPaths.add(absolutePath.getParent());
+ }
+ WatchService ws = FileSystems.getDefault().newWatchService();
+ try {
+ for (Path d : directoryPaths) {
+ d.register(ws, ENTRY_CREATE, ENTRY_MODIFY, ENTRY_DELETE);
+ }
+ } catch (IOException | RuntimeException e) {
+ ws.close();
+ throw e;
+ }
+ return new FileDataWatcher(ws, absoluteFilePaths, logger);
+ }
+
+ private FileDataWatcher(WatchService watchService, Set watchedFilePaths, LDLogger logger) {
+ this.watchService = watchService;
+ this.watchedFilePaths = watchedFilePaths;
+ this.logger = logger;
+ this.thread = new Thread(this, "LaunchDarkly-FileDataWatcher");
+ this.thread.setDaemon(true);
+ }
+
+ /**
+ * Starts the worker thread and then signals one change, because a file may have changed between
+ * the caller's initial load and the registration of the watches.
+ *
+ * @param onChange called when a watched file changed
+ */
+ void start(Runnable onChange) {
+ this.onChange = onChange;
+ thread.start();
+ signal();
+ }
+
+ @Override
+ public void run() {
+ while (!stopped.get()) {
+ WatchKey key;
+ try {
+ key = watchService.take(); // blocks until a change is available or the thread is interrupted
+ } catch (InterruptedException e) {
+ continue; // if stopped, the loop condition ends the thread
+ } catch (ClosedWatchServiceException e) {
+ return;
+ }
+ boolean watchedFileWasChanged = false;
+ for (WatchEvent> event : key.pollEvents()) {
+ if (event.kind() == OVERFLOW) {
+ // Notifications were dropped, so a watched file may have changed.
+ watchedFileWasChanged = true;
+ break;
+ }
+ Watchable w = key.watchable();
+ Object context = event.context();
+ if (w instanceof Path && context instanceof Path) {
+ Path absolutePath = ((Path) w).resolve((Path) context);
+ if (watchedFilePaths.contains(absolutePath)) {
+ watchedFileWasChanged = true;
+ break;
+ }
+ }
+ }
+ key.reset(); // without this, the watch on this key stops working
+ if (watchedFileWasChanged && !stopped.get()) {
+ signal();
+ }
+ }
+ }
+
+ private void signal() {
+ try {
+ onChange.run();
+ } catch (RuntimeException e) {
+ // COVERAGE: there is no way to simulate this condition in a unit test
+ logger.warn("Unexpected exception when reloading file data: {}", LogValues.exceptionSummary(e));
+ }
+ }
+
+ /**
+ * Stops watching and releases the watch service. It does not wait for the worker thread.
+ */
+ @Override
+ public void close() {
+ if (stopped.getAndSet(true)) {
+ return;
+ }
+ thread.interrupt();
+ try {
+ watchService.close();
+ } catch (IOException e) {
+ // COVERAGE: there is no way to simulate this condition in a unit test
+ logger.debug("Error closing file watch service: {}", LogValues.exceptionSummary(e));
+ }
+ }
+}
diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoader.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoader.java
new file mode 100644
index 00000000..77dca424
--- /dev/null
+++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoader.java
@@ -0,0 +1,247 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.google.common.collect.ImmutableList;
+import com.launchdarkly.sdk.LDValue;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FileDataException;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FlagFactory;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FlagFileParser;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FlagFileRep;
+import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.DataKind;
+import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.ItemDescriptor;
+import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.KeyedItems;
+
+import java.io.ByteArrayInputStream;
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.NoSuchFileException;
+import java.nio.file.Path;
+import java.security.MessageDigest;
+import java.security.NoSuchAlgorithmException;
+import java.util.AbstractMap;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static com.launchdarkly.sdk.server.DataModel.FEATURES;
+import static com.launchdarkly.sdk.server.DataModel.SEGMENTS;
+
+/**
+ * Reads the files of the file-based override source and merges them into one snapshot.
+ *
+ * This loader is separate from the file data source's loader because the override source has
+ * different rules: a configured file that does not exist contributes no entries, entries keep the
+ * versions that the documents specify, a value-only entry becomes a flag that is off and serves the
+ * value, and every failure to read or parse a file is reported as a file data error. The file data
+ * source keeps its own behavior.
+ */
+final class OverrideFileLoader {
+ /**
+ * Describes one configured file after a load.
+ */
+ static final class FileSummary {
+ private final Path path;
+ private final boolean present;
+ private final int flags;
+ private final int segments;
+
+ FileSummary(Path path, boolean present, int flags, int segments) {
+ this.path = path;
+ this.present = present;
+ this.flags = flags;
+ this.segments = segments;
+ }
+
+ Path getPath() {
+ return path;
+ }
+
+ /**
+ * False when the file does not exist.
+ */
+ boolean isPresent() {
+ return present;
+ }
+
+ /**
+ * The number of flag entries the merge kept from this file. An entry dropped by the duplicate keys
+ * handling is not counted.
+ */
+ int getFlags() {
+ return flags;
+ }
+
+ /**
+ * The number of segment entries the merge kept from this file.
+ */
+ int getSegments() {
+ return segments;
+ }
+ }
+
+ /**
+ * The merged result of one full load.
+ */
+ static final class LoadResult {
+ private final Iterable>> data;
+ private final List files;
+ private final int flagCount;
+ private final int segmentCount;
+ private final byte[] contentHash;
+
+ LoadResult(
+ Iterable>> data,
+ List files,
+ int flagCount,
+ int segmentCount,
+ byte[] contentHash
+ ) {
+ this.data = data;
+ this.files = Collections.unmodifiableList(new ArrayList<>(files));
+ this.flagCount = flagCount;
+ this.segmentCount = segmentCount;
+ this.contentHash = contentHash;
+ }
+
+ Iterable>> getData() {
+ return data;
+ }
+
+ /**
+ * One summary per configured file, in configuration order.
+ */
+ List getFiles() {
+ return files;
+ }
+
+ int getFlagCount() {
+ return flagCount;
+ }
+
+ int getSegmentCount() {
+ return segmentCount;
+ }
+
+ /**
+ * A digest of the raw content of every file that was read. Two loads with an equal digest produced
+ * the same data. The digest may be null when the runtime cannot compute one.
+ */
+ byte[] getContentHash() {
+ return contentHash;
+ }
+ }
+
+ private final List paths;
+ private final FileData.DuplicateKeysHandling duplicateKeysHandling;
+
+ /**
+ * Creates a loader.
+ *
+ * @param paths the files to read, in precedence order
+ * @param duplicateKeysHandling what to do when the same key appears in more than one file
+ */
+ OverrideFileLoader(List paths, FileData.DuplicateKeysHandling duplicateKeysHandling) {
+ this.paths = new ArrayList<>(paths);
+ this.duplicateKeysHandling = duplicateKeysHandling;
+ }
+
+ /**
+ * Reads every configured file in order and merges the entries.
+ *
+ * @return the merged result with a summary of each file
+ * @throws FileDataException if a file that exists cannot be read or parsed, or the files cannot be merged
+ */
+ LoadResult load() throws FileDataException {
+ MessageDigest digest = newDigest();
+ Map> data = new HashMap<>();
+ List files = new ArrayList<>(paths.size());
+ for (Path path : paths) {
+ byte[] raw;
+ try {
+ raw = Files.readAllBytes(path);
+ } catch (NoSuchFileException e) {
+ // Absence is a valid state, not a failure.
+ files.add(new FileSummary(path, false, 0, 0));
+ continue;
+ } catch (IOException e) {
+ throw new FileDataException("file " + path + ": unable to read file", e);
+ }
+ // One read feeds both the digest and the parse, so the skip-unchanged digest can never disagree
+ // with the content that was applied.
+ if (digest != null) {
+ digest.update(raw);
+ digest.update((byte) 0);
+ }
+ int flagsAdded = 0;
+ int segmentsAdded = 0;
+ try {
+ FlagFileParser parser = FlagFileParser.selectForContent(raw);
+ FlagFileRep fileContents = parser.parse(new ByteArrayInputStream(raw));
+ if (fileContents.flags != null) {
+ for (Map.Entry e : fileContents.flags.entrySet()) {
+ if (add(data, FEATURES, e.getKey(), FlagFactory.flagFromJson(e.getValue()))) {
+ flagsAdded++;
+ }
+ }
+ }
+ if (fileContents.flagValues != null) {
+ for (Map.Entry e : fileContents.flagValues.entrySet()) {
+ if (add(data, FEATURES, e.getKey(), FlagFactory.offFlagWithValue(e.getKey(), e.getValue()))) {
+ flagsAdded++;
+ }
+ }
+ }
+ if (fileContents.segments != null) {
+ for (Map.Entry e : fileContents.segments.entrySet()) {
+ if (add(data, SEGMENTS, e.getKey(), FlagFactory.segmentFromJson(e.getValue()))) {
+ segmentsAdded++;
+ }
+ }
+ }
+ } catch (FileDataException e) {
+ throw new FileDataException("file " + path + ": " + e.getMessage(), e.getCause());
+ } catch (IOException e) {
+ throw new FileDataException("file " + path + ": cannot read document", e);
+ } catch (RuntimeException e) {
+ // A document can be valid JSON or YAML and still hold a flag or segment that the data model
+ // cannot accept. For the override source that is a parse failure of the file, so the last good
+ // overrides stay in effect.
+ throw new FileDataException("file " + path + ": cannot parse flag or segment data", e);
+ }
+ files.add(new FileSummary(path, true, flagsAdded, segmentsAdded));
+ }
+
+ ImmutableList.Builder>> all = ImmutableList.builder();
+ for (Map.Entry> e : data.entrySet()) {
+ all.add(new AbstractMap.SimpleEntry<>(e.getKey(), new KeyedItems<>(ImmutableList.copyOf(e.getValue().entrySet()))));
+ }
+ Map flags = data.get(FEATURES);
+ Map segments = data.get(SEGMENTS);
+ return new LoadResult(all.build(), files, flags == null ? 0 : flags.size(),
+ segments == null ? 0 : segments.size(), digest == null ? null : digest.digest());
+ }
+
+ // Adds an entry unless the key was already added. Reports whether it added the entry.
+ private boolean add(Map> data, DataKind kind, String key, ItemDescriptor item)
+ throws FileDataException {
+ Map items = data.computeIfAbsent(kind, k -> new HashMap<>());
+ if (items.containsKey(key)) {
+ if (duplicateKeysHandling == FileData.DuplicateKeysHandling.IGNORE) {
+ return false;
+ }
+ throw new FileDataException("in " + kind.getName() + ", key \"" + key + "\" was already defined", null);
+ }
+ items.put(key, item);
+ return true;
+ }
+
+ private static MessageDigest newDigest() {
+ try {
+ return MessageDigest.getInstance("SHA-256");
+ } catch (NoSuchAlgorithmException e) {
+ // COVERAGE: every Java runtime provides SHA-256
+ return null;
+ }
+ }
+}
diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataLoadingBehaviorTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataLoadingBehaviorTest.java
new file mode 100644
index 00000000..87b1ffd4
--- /dev/null
+++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataLoadingBehaviorTest.java
@@ -0,0 +1,80 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.logging.LDLogger;
+import com.launchdarkly.sdk.server.datasources.Synchronizer;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceBase.DataBuilder;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceBase.DataLoader;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FileDataException;
+import com.launchdarkly.sdk.server.subsystems.SerializationException;
+import com.launchdarkly.testhelpers.TempDir;
+import com.launchdarkly.testhelpers.TempFile;
+
+import org.junit.Test;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.instanceOf;
+import static org.junit.Assert.fail;
+
+/**
+ * Pins the error behavior of the file data source's loader. A document that parses but holds a flag
+ * that the data model rejects is not a file data error: the deserialization exception escapes as it
+ * always has. The exception description requires a cause. The override source has its own loader
+ * with different rules.
+ */
+@SuppressWarnings("javadoc")
+public class FileDataLoadingBehaviorTest {
+ private static final String MODEL_REJECTED_DOCUMENT = "{\"flags\":{\"f1\":{\"key\":\"f1\",\"rules\":5}}}";
+
+ @Test
+ public void modelRejectedDocumentEscapesTheLoaderAsARuntimeException() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ try (TempFile file = dir.tempFile(".json")) {
+ file.setContents(MODEL_REJECTED_DOCUMENT);
+ DataLoader loader = new DataLoader(FileData.dataSource().filePaths(file.getPath()).sources);
+ try {
+ loader.load(new DataBuilder(FileData.DuplicateKeysHandling.FAIL));
+ fail("expected exception");
+ } catch (FileDataException e) {
+ fail("the loader must not convert a model rejection into a file data error");
+ } catch (RuntimeException e) {
+ assertThat(e, instanceOf(SerializationException.class));
+ }
+ }
+ }
+ }
+
+ @Test
+ public void modelRejectedDocumentEscapesTheSynchronizerAsARuntimeException() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ try (TempFile file = dir.tempFile(".json")) {
+ file.setContents(MODEL_REJECTED_DOCUMENT);
+ try (Synchronizer synchronizer = FileData.synchronizer().filePaths(file.getPath())
+ .build(TestDataSourceBuildInputs.create(LDLogger.none()))) {
+ try {
+ synchronizer.next();
+ fail("expected exception");
+ } catch (RuntimeException e) {
+ assertThat(e, instanceOf(SerializationException.class));
+ }
+ }
+ }
+ }
+ }
+
+ @Test
+ public void exceptionDescriptionRequiresACause() {
+ FileDataException withCause = new FileDataException("message", new RuntimeException("boom"), null);
+ assertThat(withCause.getDescription(), equalTo("message [java.lang.RuntimeException: boom]"));
+
+ // A duplicate key failure carries no cause. The file data source builds its own description for
+ // that case rather than calling this method.
+ FileDataException withoutCause = new FileDataException("message", null, null);
+ try {
+ withoutCause.getDescription();
+ fail("expected exception");
+ } catch (NullPointerException e) {
+ // current behavior
+ }
+ }
+}
diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataPollerTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataPollerTest.java
new file mode 100644
index 00000000..b243d968
--- /dev/null
+++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataPollerTest.java
@@ -0,0 +1,182 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.testhelpers.TempDir;
+
+import org.junit.Test;
+
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.attribute.FileTime;
+import java.time.Duration;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.lessThan;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+@SuppressWarnings("javadoc")
+public class FileDataPollerTest {
+ private static final Duration INTERVAL = Duration.ofMillis(20);
+ private static final long WAIT_MILLIS = 5000;
+ private static final long QUIET_MILLIS = 250;
+
+ private static void awaitCount(AtomicInteger counter, int expected) throws InterruptedException {
+ long deadline = System.currentTimeMillis() + WAIT_MILLIS;
+ while (counter.get() < expected) {
+ if (System.currentTimeMillis() > deadline) {
+ throw new AssertionError("timed out waiting for " + expected + " change signals, got " + counter.get());
+ }
+ Thread.sleep(5);
+ }
+ }
+
+ private static Path createFile(TempDir dir, String name, String contents, long modifiedMillis) throws Exception {
+ Path path = dir.getPath().resolve(name);
+ Files.write(path, contents.getBytes());
+ Files.setLastModifiedTime(path, FileTime.fromMillis(modifiedMillis));
+ return path;
+ }
+
+ private static void rewrite(Path path, String contents, long modifiedMillis) throws Exception {
+ Files.write(path, contents.getBytes());
+ Files.setLastModifiedTime(path, FileTime.fromMillis(modifiedMillis));
+ }
+
+ @Test
+ public void detectsModifiedTimeChange() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = createFile(dir, "a.json", "{}", 1000_000);
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet)) {
+ Thread.sleep(QUIET_MILLIS);
+ assertEquals(0, changes.get());
+
+ rewrite(file, "{}", 2000_000); // same size, new modified time
+ awaitCount(changes, 1);
+ }
+ }
+ }
+
+ @Test
+ public void detectsSizeChangeWithSameModifiedTime() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = createFile(dir, "a.json", "{}", 1000_000);
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet)) {
+ rewrite(file, "{\"flagValues\":{}}", 1000_000);
+ awaitCount(changes, 1);
+ }
+ }
+ }
+
+ @Test
+ public void firesOncePerChange() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = createFile(dir, "a.json", "{}", 1000_000);
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet)) {
+ rewrite(file, "{}", 2000_000);
+ awaitCount(changes, 1);
+ Thread.sleep(QUIET_MILLIS);
+ assertEquals(1, changes.get());
+ }
+ }
+ }
+
+ @Test
+ public void detectsFileAppearing() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("later.json");
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet)) {
+ Thread.sleep(QUIET_MILLIS);
+ assertEquals(0, changes.get());
+
+ Files.write(file, "{}".getBytes());
+ awaitCount(changes, 1);
+ }
+ }
+ }
+
+ @Test
+ public void detectsFileDisappearing() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = createFile(dir, "a.json", "{}", 1000_000);
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet)) {
+ Files.delete(file);
+ awaitCount(changes, 1);
+ }
+ }
+ }
+
+ @Test
+ public void watchesEveryConfiguredFile() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path first = createFile(dir, "a.json", "{}", 1000_000);
+ Path second = createFile(dir, "b.json", "{}", 1000_000);
+ List paths = Arrays.asList(first, second);
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataPoller poller = new FileDataPoller(paths, INTERVAL, changes::incrementAndGet)) {
+ rewrite(second, "{}", 2000_000);
+ awaitCount(changes, 1);
+ rewrite(first, "{}", 2000_000);
+ awaitCount(changes, 2);
+ }
+ }
+ }
+
+ @Test
+ public void stopsOnClose() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = createFile(dir, "a.json", "{}", 1000_000);
+ AtomicInteger changes = new AtomicInteger();
+ FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet);
+ poller.close();
+ poller.close(); // idempotent
+
+ rewrite(file, "{}", 2000_000);
+ Thread.sleep(QUIET_MILLIS);
+ assertEquals(0, changes.get());
+ }
+ }
+
+ @Test
+ public void closeReturnsWhileCallbackBlocks() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = createFile(dir, "a.json", "{}", 1000_000);
+ CountDownLatch entered = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, () -> {
+ entered.countDown();
+ try {
+ release.await(WAIT_MILLIS, TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ });
+ rewrite(file, "{}", 2000_000);
+ assertTrue(entered.await(WAIT_MILLIS, TimeUnit.MILLISECONDS));
+
+ long start = System.currentTimeMillis();
+ poller.close();
+ assertThat(System.currentTimeMillis() - start, lessThan(1000L));
+ release.countDown();
+ }
+ }
+
+ @Test
+ public void unreadableAttributesCountAsAbsent() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path missing = dir.getPath().resolve("nope").resolve("missing.json");
+ List states = FileDataPoller.observeAll(Collections.singletonList(missing));
+ assertEquals(FileDataPoller.FileState.ABSENT, states.get(0));
+ }
+ }
+}
diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataReloaderTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataReloaderTest.java
new file mode 100644
index 00000000..4c2c1704
--- /dev/null
+++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataReloaderTest.java
@@ -0,0 +1,512 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.logging.LDLogger;
+import com.launchdarkly.logging.LogCapture;
+import com.launchdarkly.logging.Logs;
+import com.launchdarkly.sdk.server.integrations.OverrideFileLoader.LoadResult;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FileDataException;
+
+import org.junit.After;
+import org.junit.Test;
+
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Supplier;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.greaterThanOrEqualTo;
+import static org.hamcrest.Matchers.hasItem;
+import static org.hamcrest.Matchers.lessThan;
+import static org.hamcrest.Matchers.startsWith;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+@SuppressWarnings("javadoc")
+public class FileDataReloaderTest {
+ private static final Duration SHORT = Duration.ofMillis(50);
+ private static final long WAIT_MILLIS = 5000;
+
+ private final LogCapture logCapture = Logs.capture();
+ private final LDLogger logger = LDLogger.withAdapter(logCapture, "");
+ private final List reloaders = new ArrayList<>();
+
+ @After
+ public void closeReloaders() {
+ for (FileDataReloader r : reloaders) {
+ r.close();
+ }
+ }
+
+ private static LoadResult resultWithHash(int... values) {
+ byte[] hash = new byte[values.length];
+ for (int i = 0; i < values.length; i++) {
+ hash[i] = (byte) values[i];
+ }
+ return new LoadResult(Collections.emptyList(), Collections.emptyList(), 0, 0, hash);
+ }
+
+ private static LoadResult resultWithoutHash() {
+ return new LoadResult(Collections.emptyList(), Collections.emptyList(), 0, 0, null);
+ }
+
+ private static FileDataException failure(String message) {
+ return new FileDataException(message, null, null);
+ }
+
+ /**
+ * A scripted loader: each call takes the next step from the queue. A step either returns a result
+ * or throws. When the queue is empty the last step repeats.
+ */
+ private static final class ScriptedLoader implements FileDataReloader.Loader {
+ final AtomicInteger calls = new AtomicInteger();
+ final LinkedBlockingQueue> steps = new LinkedBlockingQueue<>();
+ volatile Supplier last;
+ volatile CountDownLatch entered;
+ volatile CountDownLatch release;
+
+ ScriptedLoader then(LoadResult result) {
+ steps.add(() -> result);
+ return this;
+ }
+
+ ScriptedLoader thenFail(String message) {
+ steps.add(() -> {
+ throw new UncheckedFileDataException(failure(message));
+ });
+ return this;
+ }
+
+ @Override
+ public LoadResult load() throws FileDataException {
+ // Both latches are read before the entered signal, so a test that clears them after the
+ // signal cannot change what this call does.
+ CountDownLatch enteredLatch = entered;
+ CountDownLatch releaseLatch = release;
+ calls.incrementAndGet();
+ if (enteredLatch != null) {
+ enteredLatch.countDown();
+ }
+ if (releaseLatch != null) {
+ try {
+ releaseLatch.await(WAIT_MILLIS, TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ Supplier step = steps.poll();
+ if (step == null) {
+ step = last;
+ } else {
+ last = step;
+ }
+ try {
+ return step.get();
+ } catch (UncheckedFileDataException e) {
+ throw e.wrapped;
+ }
+ }
+ }
+
+ @SuppressWarnings("serial")
+ private static final class UncheckedFileDataException extends RuntimeException {
+ final FileDataException wrapped;
+
+ UncheckedFileDataException(FileDataException wrapped) {
+ this.wrapped = wrapped;
+ }
+ }
+
+ private static final class RecordingHandler implements FileDataReloader.Handler {
+ final List applied = Collections.synchronizedList(new ArrayList<>());
+ final List errors = Collections.synchronizedList(new ArrayList<>());
+
+ @Override
+ public void apply(LoadResult result) {
+ applied.add(result);
+ }
+
+ @Override
+ public void onError(FileDataException e) {
+ errors.add(e);
+ }
+ }
+
+ private FileDataReloader reloader(ScriptedLoader loader, RecordingHandler handler,
+ Duration debounce, Duration retry, boolean skipUnchanged) {
+ FileDataReloader r = new FileDataReloader(loader, handler, logger, debounce, retry, skipUnchanged);
+ reloaders.add(r);
+ return r;
+ }
+
+ private static void awaitAtLeast(AtomicInteger counter, int expected) throws InterruptedException {
+ long deadline = System.currentTimeMillis() + WAIT_MILLIS;
+ while (counter.get() < expected) {
+ if (System.currentTimeMillis() > deadline) {
+ throw new AssertionError("timed out waiting for count " + expected + ", got " + counter.get());
+ }
+ Thread.sleep(5);
+ }
+ }
+
+ private static void awaitSize(List> list, int expected) throws InterruptedException {
+ long deadline = System.currentTimeMillis() + WAIT_MILLIS;
+ while (list.size() < expected) {
+ if (System.currentTimeMillis() > deadline) {
+ throw new AssertionError("timed out waiting for size " + expected + ", got " + list.size());
+ }
+ Thread.sleep(5);
+ }
+ }
+
+ @Test
+ public void initialLoadAppliesResultSynchronously() {
+ LoadResult result = resultWithHash(1);
+ ScriptedLoader loader = new ScriptedLoader().then(result);
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, true);
+
+ r.reloadNow();
+
+ assertEquals(1, loader.calls.get());
+ assertEquals(1, handler.applied.size());
+ assertEquals(result, handler.applied.get(0));
+ assertTrue(handler.errors.isEmpty());
+ }
+
+ @Test
+ public void failedLoadReportsErrorAndAppliesNothing() {
+ ScriptedLoader loader = new ScriptedLoader().thenFail("bad file");
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, true);
+
+ r.reloadNow();
+
+ assertTrue(handler.applied.isEmpty());
+ assertEquals(1, handler.errors.size());
+ assertThat(FileDataReloader.describe(handler.errors.get(0)), equalTo("bad file"));
+ assertThat(logCapture.getMessageStrings(), hasItem("ERROR:Unable to load flags: bad file"));
+ }
+
+ @Test
+ public void reloaderThatIsNeverTriggeredStartsNoThread() {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, SHORT, SHORT, true);
+
+ r.reloadNow();
+
+ assertFalse(r.hasWorker());
+ }
+
+ @Test
+ public void triggerReloadsAfterDebounceDelay() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).then(resultWithHash(2));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, SHORT, Duration.ZERO, true);
+ r.reloadNow();
+
+ r.trigger();
+
+ awaitSize(handler.applied, 2);
+ assertTrue(r.hasWorker());
+ assertThat(logCapture.getMessageStrings(), hasItem("INFO:Reloading flag data after detecting a change"));
+ }
+
+ @Test
+ public void debounceCoalescesBurstOfTriggersIntoOneReload() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).then(resultWithHash(2));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ofMillis(150), Duration.ZERO, false);
+ r.reloadNow();
+
+ for (int i = 0; i < 10; i++) {
+ r.trigger();
+ }
+
+ awaitSize(handler.applied, 2);
+ Thread.sleep(400);
+ assertEquals(2, loader.calls.get());
+ assertEquals(2, handler.applied.size());
+ }
+
+ @Test
+ public void debounceWindowIsExtendedByEachTrigger() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).then(resultWithHash(2));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ofMillis(200), Duration.ZERO, false);
+ r.reloadNow();
+
+ // Triggers arrive every 50ms for 400ms. The window is 200ms, so no reload happens while they
+ // keep arriving.
+ for (int i = 0; i < 8; i++) {
+ r.trigger();
+ Thread.sleep(50);
+ }
+ assertEquals(1, loader.calls.get());
+
+ awaitSize(handler.applied, 2);
+ Thread.sleep(300);
+ assertEquals(2, loader.calls.get());
+ }
+
+ @Test
+ public void zeroDebounceReloadsOnEveryTrigger() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, false);
+ r.reloadNow();
+
+ r.trigger();
+ awaitAtLeast(loader.calls, 2);
+ r.trigger();
+ awaitAtLeast(loader.calls, 3);
+
+ awaitSize(handler.applied, 3);
+ }
+
+ @Test
+ public void failedInitialLoadRetriesWithoutFurtherTriggers() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().thenFail("first").thenFail("first").then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, SHORT, true);
+
+ r.reloadNow();
+ assertTrue(r.hasWorker());
+
+ awaitSize(handler.applied, 1);
+ assertThat(loader.calls.get(), greaterThanOrEqualTo(3));
+ assertThat(logCapture.getMessageStrings(), hasItem("DEBUG:Retrying flag data load after earlier failure"));
+ }
+
+ @Test
+ public void identicalRepeatedFailureIsReportedOnceUntilItChanges() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().thenFail("same").thenFail("same").thenFail("same")
+ .thenFail("different").then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, SHORT, true);
+
+ r.reloadNow();
+
+ awaitSize(handler.applied, 1);
+ assertEquals(2, handler.errors.size());
+ assertThat(FileDataReloader.describe(handler.errors.get(0)), equalTo("same"));
+ assertThat(FileDataReloader.describe(handler.errors.get(1)), equalTo("different"));
+ List messages = logCapture.getMessageStrings();
+ assertEquals(1, messages.stream().filter(m -> m.equals("ERROR:Unable to load flags: same")).count());
+ assertThat(messages, hasItem("DEBUG:Unable to load flags: same"));
+ assertThat(messages, hasItem("ERROR:Unable to load flags: different"));
+ }
+
+ @Test
+ public void successReArmsFailureReporting() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().thenFail("same").then(resultWithHash(1)).thenFail("same")
+ .then(resultWithHash(2));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, SHORT, true);
+
+ r.reloadNow(); // fails, retry applies result 1
+ awaitSize(handler.applied, 1);
+ r.trigger(); // fails again with the same message, retry applies result 2
+ awaitSize(handler.applied, 2);
+
+ assertEquals(2, handler.errors.size());
+ }
+
+ @Test
+ public void retryStopsAfterSuccess() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().thenFail("first").then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, SHORT, true);
+
+ r.reloadNow();
+ awaitSize(handler.applied, 1);
+ int callsAfterSuccess = loader.calls.get();
+
+ Thread.sleep(300);
+ assertEquals(callsAfterSuccess, loader.calls.get());
+ }
+
+ @Test
+ public void reloadWithUnchangedContentIsSkipped() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).then(resultWithHash(1)).then(resultWithHash(2));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, true);
+
+ r.reloadNow();
+ r.trigger();
+ awaitAtLeast(loader.calls, 2);
+ r.trigger();
+ awaitAtLeast(loader.calls, 3);
+
+ awaitSize(handler.applied, 2);
+ assertEquals(2, handler.applied.size());
+ }
+
+ @Test
+ public void recoveryAppliesEvenWhenContentIsUnchanged() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).thenFail("broken").then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, true);
+
+ r.reloadNow();
+ r.trigger();
+ awaitSize(handler.errors, 1);
+ r.trigger();
+
+ awaitSize(handler.applied, 2);
+ }
+
+ @Test
+ public void everyReloadIsAppliedWhenSkipUnchangedIsOff() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, false);
+
+ r.reloadNow();
+ r.trigger();
+
+ awaitSize(handler.applied, 2);
+ }
+
+ @Test
+ public void nullHashIsNeverTreatedAsUnchanged() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithoutHash());
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, true);
+
+ r.reloadNow();
+ r.trigger();
+
+ awaitSize(handler.applied, 2);
+ }
+
+ @Test
+ public void closeStopsFurtherReloads() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1));
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, false);
+ r.reloadNow();
+
+ r.close();
+ r.trigger();
+ r.reloadNow();
+ Thread.sleep(100);
+
+ assertEquals(1, loader.calls.get());
+ assertEquals(1, handler.applied.size());
+ }
+
+ @Test
+ public void closeCancelsPendingRetry() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().thenFail("broken");
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ofMillis(100), true);
+ r.reloadNow();
+
+ r.close();
+ Thread.sleep(300);
+
+ assertEquals(1, loader.calls.get());
+ }
+
+ @Test
+ public void closeIsIdempotent() {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1));
+ FileDataReloader r = reloader(loader, new RecordingHandler(), SHORT, SHORT, true);
+ r.trigger();
+ r.close();
+ r.close();
+ }
+
+ @Test
+ public void closeDoesNotWaitForInFlightReloadAndDropsItsResult() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1));
+ loader.entered = new CountDownLatch(1);
+ loader.release = new CountDownLatch(1);
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, false);
+
+ r.trigger();
+ assertTrue(loader.entered.await(WAIT_MILLIS, TimeUnit.MILLISECONDS));
+
+ long start = System.currentTimeMillis();
+ r.close();
+ assertThat(System.currentTimeMillis() - start, lessThan(1000L));
+
+ loader.release.countDown();
+ Thread.sleep(200);
+ assertTrue(handler.applied.isEmpty());
+ }
+
+ @Test
+ public void reloadsAreSerializedBetweenCallerAndWorker() throws Exception {
+ // The loader blocks on the worker thread. A synchronous reload from the caller must wait for
+ // it rather than run concurrently.
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).then(resultWithHash(2));
+ CountDownLatch entered = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ loader.entered = entered;
+ loader.release = release;
+ RecordingHandler handler = new RecordingHandler();
+ FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, false);
+
+ r.trigger();
+ assertTrue(entered.await(WAIT_MILLIS, TimeUnit.MILLISECONDS));
+ // Later calls must not block.
+ loader.entered = null;
+ loader.release = null;
+
+ AtomicReference syncReloadFinished = new AtomicReference<>(false);
+ Thread caller = new Thread(() -> {
+ r.reloadNow();
+ syncReloadFinished.set(true);
+ });
+ caller.start();
+ Thread.sleep(150);
+ assertFalse(syncReloadFinished.get());
+ assertEquals(1, loader.calls.get());
+
+ release.countDown();
+ caller.join(WAIT_MILLIS);
+ assertTrue(syncReloadFinished.get());
+ assertEquals(2, loader.calls.get());
+ assertEquals(2, handler.applied.size());
+ }
+
+ @Test
+ public void unexpectedRuntimeExceptionFromHandlerIsLoggedAndWorkerSurvives() throws Exception {
+ ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).then(resultWithHash(2));
+ AtomicInteger applies = new AtomicInteger();
+ FileDataReloader.Handler handler = new FileDataReloader.Handler() {
+ @Override
+ public void apply(LoadResult result) {
+ if (applies.incrementAndGet() == 1) {
+ throw new IllegalStateException("consumer failed");
+ }
+ }
+
+ @Override
+ public void onError(FileDataException e) {
+ }
+ };
+ FileDataReloader r = new FileDataReloader(loader, handler, logger, Duration.ZERO, Duration.ZERO, false);
+ reloaders.add(r);
+
+ r.trigger();
+ awaitAtLeast(applies, 1);
+ r.trigger();
+ awaitAtLeast(applies, 2);
+
+ assertThat(logCapture.getMessageStrings(),
+ hasItem(startsWith("ERROR:Unexpected error while reloading flag data:")));
+ }
+}
diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataWatcherTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataWatcherTest.java
new file mode 100644
index 00000000..af58b078
--- /dev/null
+++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataWatcherTest.java
@@ -0,0 +1,166 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.logging.LDLogger;
+import com.launchdarkly.testhelpers.TempDir;
+
+import org.junit.Test;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Collections;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.fail;
+
+@SuppressWarnings("javadoc")
+public class FileDataWatcherTest {
+ private static final LDLogger testLogger = LDLogger.none();
+ private static final long WAIT_MILLIS = 10000;
+ private static final long QUIET_MILLIS = 300;
+
+ private static void awaitAtLeast(AtomicInteger counter, int expected) throws InterruptedException {
+ long deadline = System.currentTimeMillis() + WAIT_MILLIS;
+ while (counter.get() < expected) {
+ if (System.currentTimeMillis() > deadline) {
+ throw new AssertionError("timed out waiting for " + expected + " change signals, got " + counter.get());
+ }
+ Thread.sleep(10);
+ }
+ }
+
+ @Test
+ public void startSignalsOnceSoThatAChangeBeforeWatchingIsNotMissed() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("a.json");
+ Files.write(file, "{}".getBytes());
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), testLogger)) {
+ watcher.start(changes::incrementAndGet);
+ awaitAtLeast(changes, 1);
+ Thread.sleep(QUIET_MILLIS);
+ assertEquals(1, changes.get());
+ }
+ }
+ }
+
+ @Test
+ public void modifiedFileSignalsChange() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("a.json");
+ Files.write(file, "{}".getBytes());
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), testLogger)) {
+ watcher.start(changes::incrementAndGet);
+ awaitAtLeast(changes, 1);
+ Thread.sleep(QUIET_MILLIS);
+
+ Files.write(file, "{\"flagValues\":{}}".getBytes());
+ awaitAtLeast(changes, 2);
+ }
+ }
+ }
+
+ @Test
+ public void fileThatDoesNotExistYetSignalsChangeWhenItAppears() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("later.json");
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), testLogger)) {
+ watcher.start(changes::incrementAndGet);
+ awaitAtLeast(changes, 1);
+ Thread.sleep(QUIET_MILLIS);
+
+ Files.write(file, "{}".getBytes());
+ awaitAtLeast(changes, 2);
+ }
+ }
+ }
+
+ @Test
+ public void deletedFileSignalsChange() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("a.json");
+ Files.write(file, "{}".getBytes());
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), testLogger)) {
+ watcher.start(changes::incrementAndGet);
+ awaitAtLeast(changes, 1);
+ Thread.sleep(QUIET_MILLIS);
+
+ Files.delete(file);
+ awaitAtLeast(changes, 2);
+ }
+ }
+ }
+
+ @Test
+ public void unrelatedFileInSameDirectoryDoesNotSignal() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("a.json");
+ Files.write(file, "{}".getBytes());
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), testLogger)) {
+ watcher.start(changes::incrementAndGet);
+ awaitAtLeast(changes, 1);
+ Thread.sleep(QUIET_MILLIS);
+
+ Files.write(dir.getPath().resolve("other.json"), "{}".getBytes());
+ Thread.sleep(QUIET_MILLIS);
+ assertEquals(1, changes.get());
+ }
+ }
+ }
+
+ @Test
+ public void relativePathIsWatchedByItsAbsoluteLocation() throws Exception {
+ // The watcher must compare the directory-relative event name against the same absolute form
+ // it registered, whatever form the caller used.
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("sub").resolve("..").resolve("a.json");
+ Files.createDirectories(dir.getPath().resolve("sub"));
+ Files.write(file, "{}".getBytes());
+ AtomicInteger changes = new AtomicInteger();
+ try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), testLogger)) {
+ watcher.start(changes::incrementAndGet);
+ awaitAtLeast(changes, 1);
+ Thread.sleep(QUIET_MILLIS);
+
+ Files.write(file, "{\"flagValues\":{}}".getBytes());
+ awaitAtLeast(changes, 2);
+ }
+ }
+ }
+
+ @Test
+ public void missingDirectoryIsAnErrorAtCreation() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("nope").resolve("a.json");
+ try {
+ FileDataWatcher.create(Collections.singletonList(file), testLogger).close();
+ fail("expected IOException");
+ } catch (IOException e) {
+ // expected
+ }
+ }
+ }
+
+ @Test
+ public void closeStopsSignalsAndIsIdempotent() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = dir.getPath().resolve("a.json");
+ Files.write(file, "{}".getBytes());
+ AtomicInteger changes = new AtomicInteger();
+ FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), testLogger);
+ watcher.start(changes::incrementAndGet);
+ awaitAtLeast(changes, 1);
+ watcher.close();
+ watcher.close();
+
+ Files.write(file, "{\"flagValues\":{}}".getBytes());
+ Thread.sleep(QUIET_MILLIS);
+ assertEquals(1, changes.get());
+ }
+ }
+}
diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileSynchronizerReloadBehaviorTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileSynchronizerReloadBehaviorTest.java
new file mode 100644
index 00000000..ab0e592c
--- /dev/null
+++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileSynchronizerReloadBehaviorTest.java
@@ -0,0 +1,204 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.logging.LDLogLevel;
+import com.launchdarkly.logging.LDLogger;
+import com.launchdarkly.logging.LogCapture;
+import com.launchdarkly.logging.Logs;
+import com.launchdarkly.sdk.fdv2.SourceResultType;
+import com.launchdarkly.sdk.fdv2.SourceSignal;
+import com.launchdarkly.sdk.server.datasources.FDv2SourceResult;
+import com.launchdarkly.sdk.server.datasources.Synchronizer;
+import com.launchdarkly.sdk.server.interfaces.DataSourceStatusProvider.ErrorKind;
+import com.launchdarkly.testhelpers.TempDir;
+import com.launchdarkly.testhelpers.TempFile;
+
+import org.junit.Test;
+
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+import static com.launchdarkly.sdk.server.integrations.FileDataSourceTestData.getResourceContents;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsString;
+import static org.hamcrest.Matchers.endsWith;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.greaterThanOrEqualTo;
+import static org.hamcrest.Matchers.not;
+import static org.hamcrest.Matchers.startsWith;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.fail;
+
+/**
+ * Pins the reload behavior of the file data source's synchronizer: it reloads on every file system
+ * event, reports every failed load, does not retry on its own, and logs a failure as the plain
+ * description at error level. The override source has different rules and its own tests.
+ */
+@SuppressWarnings("javadoc")
+public class FileSynchronizerReloadBehaviorTest {
+ private static final long CHANGE_TIMEOUT_SECONDS = 15;
+ private static final String MALFORMED = "{\"flags\"";
+
+ private final LogCapture logCapture = Logs.capture();
+ private final LDLogger logger = LDLogger.withAdapter(logCapture, "");
+
+ private Synchronizer autoUpdatingSynchronizer(TempFile file) {
+ return FileData.synchronizer()
+ .filePaths(file.getPath())
+ .autoUpdate(true)
+ .build(TestDataSourceBuildInputs.create(logger));
+ }
+
+ private static FDv2SourceResult await(CompletableFuture future) throws Exception {
+ return future.get(CHANGE_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+
+ // Consumes results until none arrives for a short quiet period, then returns the one future that is
+ // still outstanding. A single edit can produce more than one file system event, and the synchronizer
+ // reloads on each of them. Only one next() future is ever left outstanding, because the queue
+ // delivers each result to the oldest waiting future, and a forgotten one would swallow a result.
+ private static CompletableFuture settle(Synchronizer synchronizer) throws Exception {
+ CompletableFuture pending = synchronizer.next();
+ while (true) {
+ try {
+ pending.get(500, TimeUnit.MILLISECONDS);
+ pending = synchronizer.next();
+ } catch (TimeoutException e) {
+ return pending;
+ }
+ }
+ }
+
+ private void awaitErrorMessages(int atLeast) throws Exception {
+ long deadline = System.currentTimeMillis() + CHANGE_TIMEOUT_SECONDS * 1000;
+ while (errorMessages().size() < atLeast) {
+ if (System.currentTimeMillis() > deadline) {
+ fail("expected at least " + atLeast + " error logs but saw " + errorMessages());
+ }
+ Thread.sleep(50);
+ }
+ }
+
+ private List errorMessages() {
+ List out = new java.util.ArrayList<>();
+ for (LogCapture.Message m : logCapture.getMessages()) {
+ if (m.getLevel() == LDLogLevel.ERROR) {
+ out.add(m.getText());
+ }
+ }
+ return out;
+ }
+
+ @Test
+ public void rewriteWithIdenticalContentEmitsAnotherChangeSet() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ try (TempFile file = dir.tempFile(".json")) {
+ String contents = getResourceContents("flag-only.json");
+ file.setContents(contents);
+ try (Synchronizer synchronizer = autoUpdatingSynchronizer(file)) {
+ assertThat(await(synchronizer.next()).getResultType(), equalTo(SourceResultType.CHANGE_SET));
+
+ CompletableFuture next = synchronizer.next();
+ Thread.sleep(200); // let the watcher register before the change
+ file.setContents(contents); // same bytes, new modification time
+
+ // The content did not change, and the synchronizer still delivers a new change set.
+ assertThat(await(next).getResultType(), equalTo(SourceResultType.CHANGE_SET));
+ }
+ }
+ }
+ }
+
+ @Test
+ public void everyFailedReloadIsReportedAndLoggedAsThePlainDescription() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ try (TempFile file = dir.tempFile(".json")) {
+ file.setContents(getResourceContents("flag-only.json"));
+ try (Synchronizer synchronizer = autoUpdatingSynchronizer(file)) {
+ assertThat(await(synchronizer.next()).getResultType(), equalTo(SourceResultType.CHANGE_SET));
+
+ CompletableFuture next = synchronizer.next();
+ Thread.sleep(200);
+ file.setContents(MALFORMED);
+ FDv2SourceResult first = await(next);
+ assertThat(first.getResultType(), equalTo(SourceResultType.STATUS));
+ assertThat(first.getStatus().getState(), equalTo(SourceSignal.INTERRUPTED));
+ assertThat(first.getStatus().getErrorInfo().getKind(), equalTo(ErrorKind.INVALID_DATA));
+ CompletableFuture pending = settle(synchronizer);
+ int errorsAfterFirstFailure = errorMessages().size();
+
+ // The same malformed content again is the same failure. It is reported again, not deduplicated.
+ file.setContents(MALFORMED);
+ FDv2SourceResult second = await(pending);
+ assertThat(second.getResultType(), equalTo(SourceResultType.STATUS));
+ assertThat(second.getStatus().getState(), equalTo(SourceSignal.INTERRUPTED));
+ assertEquals(first.getStatus().getErrorInfo().getMessage(), second.getStatus().getErrorInfo().getMessage());
+ awaitErrorMessages(errorsAfterFirstFailure + 1);
+
+ // Each failure is logged at error level as the plain description: the parser message and
+ // the cause in brackets. There is no prefix, and the description does not name the file.
+ List errors = errorMessages();
+ assertThat(errors.size(), greaterThanOrEqualTo(2));
+ for (String message : errors) {
+ assertThat(message, startsWith("cannot parse JSON [com.google.gson.JsonSyntaxException"));
+ assertThat(message, endsWith("]"));
+ assertThat(message, not(containsString(file.getPath().toString())));
+ }
+ assertEquals(errors.get(0), first.getStatus().getErrorInfo().getMessage());
+ }
+ }
+ }
+ }
+
+ @Test
+ public void failedReloadIsNotRetriedWithoutAFileChange() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ try (TempFile file = dir.tempFile(".json")) {
+ file.setContents(getResourceContents("flag-only.json"));
+ try (Synchronizer synchronizer = autoUpdatingSynchronizer(file)) {
+ assertThat(await(synchronizer.next()).getResultType(), equalTo(SourceResultType.CHANGE_SET));
+
+ CompletableFuture next = synchronizer.next();
+ Thread.sleep(200);
+ file.setContents(MALFORMED);
+ assertThat(await(next).getStatus().getState(), equalTo(SourceSignal.INTERRUPTED));
+ CompletableFuture later = settle(synchronizer);
+ int errorsAfterFailure = errorMessages().size();
+
+ // With the file untouched, nothing reloads: no result, no further error log, no retry log.
+ try {
+ FDv2SourceResult unexpected = later.get(2, TimeUnit.SECONDS);
+ fail("unexpected reload result: " + unexpected.getResultType());
+ } catch (TimeoutException e) {
+ // expected
+ }
+ assertEquals(errorsAfterFailure, errorMessages().size());
+ for (LogCapture.Message m : logCapture.getMessages()) {
+ assertFalse(m.getText(), m.getText().contains("Retrying"));
+ }
+ assertFalse(later.isDone());
+ }
+ }
+ }
+ }
+
+ @Test
+ public void missingFileAtStartupIsAnErrorLoggedWithThePath() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ java.nio.file.Path missing = dir.getPath().resolve("missing.json");
+ try (Synchronizer synchronizer = FileData.synchronizer().filePaths(missing)
+ .build(TestDataSourceBuildInputs.create(logger))) {
+ FDv2SourceResult result = await(synchronizer.next());
+ assertThat(result.getResultType(), equalTo(SourceResultType.STATUS));
+ assertThat(result.getStatus().getState(), equalTo(SourceSignal.INTERRUPTED));
+ // The description is the cause in brackets. The path appears only inside the cause text.
+ List errors = errorMessages();
+ assertEquals(1, errors.size());
+ assertEquals("[java.nio.file.NoSuchFileException: " + missing + "]", errors.get(0));
+ assertEquals(errors.get(0), result.getStatus().getErrorInfo().getMessage());
+ }
+ }
+ }
+}
diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoaderTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoaderTest.java
new file mode 100644
index 00000000..d244316c
--- /dev/null
+++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoaderTest.java
@@ -0,0 +1,290 @@
+package com.launchdarkly.sdk.server.integrations;
+
+import com.launchdarkly.sdk.LDValue;
+import com.launchdarkly.sdk.server.integrations.FileDataSourceParsing.FileDataException;
+import com.launchdarkly.sdk.server.integrations.OverrideFileLoader.FileSummary;
+import com.launchdarkly.sdk.server.integrations.OverrideFileLoader.LoadResult;
+import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.DataKind;
+import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.ItemDescriptor;
+import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.KeyedItems;
+import com.launchdarkly.sdk.server.subsystems.SerializationException;
+import com.launchdarkly.testhelpers.TempDir;
+
+import org.junit.Test;
+
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static com.launchdarkly.sdk.server.DataModel.FEATURES;
+import static com.launchdarkly.sdk.server.DataModel.SEGMENTS;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsString;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.instanceOf;
+import static org.hamcrest.Matchers.is;
+import static org.junit.Assert.assertArrayEquals;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.fail;
+
+@SuppressWarnings("javadoc")
+public class OverrideFileLoaderTest {
+ private static Map> toMap(LoadResult result) {
+ Map> out = new HashMap<>();
+ for (Map.Entry> kind : result.getData()) {
+ Map items = new HashMap<>();
+ for (Map.Entry item : kind.getValue().getItems()) {
+ items.put(item.getKey(), item.getValue());
+ }
+ out.put(kind.getKey(), items);
+ }
+ return out;
+ }
+
+ private static LDValue json(DataKind kind, ItemDescriptor item) {
+ return LDValue.parse(kind.serialize(item));
+ }
+
+ private static Path write(TempDir dir, String name, String contents) throws Exception {
+ Path p = dir.getPath().resolve(name);
+ Files.write(p, contents.getBytes());
+ return p;
+ }
+
+ private static OverrideFileLoader loader(Path... paths) {
+ return new OverrideFileLoader(Arrays.asList(paths), FileData.DuplicateKeysHandling.FAIL);
+ }
+
+ @Test
+ public void loadsFlagsFlagValuesAndSegmentsKeepingDocumentVersions() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = write(dir, "overrides.json",
+ "{\"flags\":{\"f1\":{\"key\":\"f1\",\"version\":7,\"on\":true,\"variations\":[true],\"fallthrough\":{\"variation\":0}}},"
+ + "\"flagValues\":{\"f2\":\"x\"},"
+ + "\"segments\":{\"s1\":{\"key\":\"s1\",\"version\":9}}}");
+ LoadResult result = loader(file).load();
+
+ Map> data = toMap(result);
+ assertEquals(7, data.get(FEATURES).get("f1").getVersion());
+ assertEquals(LDValue.of(7), json(FEATURES, data.get(FEATURES).get("f1")).get("version"));
+ assertEquals(9, data.get(SEGMENTS).get("s1").getVersion());
+ assertEquals(LDValue.of(9), json(SEGMENTS, data.get(SEGMENTS).get("s1")).get("version"));
+ assertEquals(2, result.getFlagCount());
+ assertEquals(1, result.getSegmentCount());
+ assertEquals(1, result.getFiles().size());
+ assertTrue(result.getFiles().get(0).isPresent());
+ assertEquals(2, result.getFiles().get(0).getFlags());
+ assertEquals(1, result.getFiles().get(0).getSegments());
+ assertEquals(file, result.getFiles().get(0).getPath());
+ }
+ }
+
+ @Test
+ public void valueOnlyEntryBecomesOffFlagServingTheValue() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = write(dir, "overrides.json", "{\"flagValues\":{\"f2\":\"x\"}}");
+ LoadResult result = loader(file).load();
+ ItemDescriptor f2 = toMap(result).get(FEATURES).get("f2");
+ assertEquals(0, f2.getVersion());
+ LDValue flag = json(FEATURES, f2);
+ assertEquals(LDValue.of("f2"), flag.get("key"));
+ assertEquals(LDValue.of(0), flag.get("version"));
+ assertEquals(LDValue.of(false), flag.get("on"));
+ assertEquals(LDValue.of(0), flag.get("offVariation"));
+ assertEquals(LDValue.buildArray().add("x").build(), flag.get("variations"));
+ }
+ }
+
+ @Test
+ public void yamlIsAutoDetected() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = write(dir, "overrides.yaml", "flagValues:\n yaml-flag: \"override-value\"\n");
+ LoadResult result = loader(file).load();
+ assertEquals(LDValue.buildArray().add("override-value").build(),
+ json(FEATURES, toMap(result).get(FEATURES).get("yaml-flag")).get("variations"));
+ }
+ }
+
+ @Test
+ public void missingFileContributesNothing() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path present = write(dir, "present.json", "{\"flagValues\":{\"flag\":\"value\"}}");
+ Path missing = dir.getPath().resolve("missing.json");
+ Path segments = write(dir, "segments.json", "{\"segments\":{\"s1\":{\"key\":\"s1\"}}}");
+ LoadResult result = loader(present, missing, segments).load();
+
+ Map> data = toMap(result);
+ assertEquals(1, data.get(FEATURES).size());
+ assertEquals(1, data.get(SEGMENTS).size());
+ List files = result.getFiles();
+ assertEquals(3, files.size());
+ assertTrue(files.get(0).isPresent());
+ assertEquals(1, files.get(0).getFlags());
+ assertFalse(files.get(1).isPresent());
+ assertEquals(missing, files.get(1).getPath());
+ assertEquals(0, files.get(1).getFlags());
+ assertTrue(files.get(2).isPresent());
+ assertEquals(1, files.get(2).getSegments());
+ }
+ }
+
+ @Test
+ public void allFilesMissingProducesEmptyResult() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ LoadResult result = loader(dir.getPath().resolve("a.json"), dir.getPath().resolve("b.json")).load();
+ assertEquals(0, result.getFlagCount());
+ assertEquals(0, result.getSegmentCount());
+ assertFalse(result.getData().iterator().hasNext());
+ assertEquals(2, result.getFiles().size());
+ }
+ }
+
+ @Test
+ public void unreadableFileFailsTheLoad() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ // A directory exists but cannot be read as a file.
+ Path directory = dir.getPath().resolve("dir.json");
+ Files.createDirectory(directory);
+ try {
+ loader(directory).load();
+ fail("expected exception");
+ } catch (FileDataException e) {
+ assertThat(e.getMessage(), containsString(directory.toString()));
+ assertThat(e.getMessage(), containsString("unable to read file"));
+ assertThat(e.getCause(), instanceOf(java.io.IOException.class));
+ }
+ }
+ }
+
+ @Test
+ public void malformedDocumentFailsTheLoadWithFileAttribution() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path good = write(dir, "good.json", "{\"flagValues\":{\"a\":1}}");
+ Path bad = write(dir, "bad.json", "{\"flagValues\"");
+ try {
+ loader(good, bad).load();
+ fail("expected exception");
+ } catch (FileDataException e) {
+ assertThat(e.getMessage(), containsString(bad.toString()));
+ assertThat(e.getMessage(), containsString("cannot parse JSON"));
+ assertThat(e.getCause(), instanceOf(com.google.gson.JsonSyntaxException.class));
+ }
+ }
+ }
+
+ @Test
+ public void modelRejectedDocumentIsAFileDataError() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = write(dir, "bad-model.json", "{\"flags\":{\"f1\":{\"key\":\"f1\",\"rules\":5}}}");
+ try {
+ loader(file).load();
+ fail("expected exception");
+ } catch (FileDataException e) {
+ assertThat(e.getMessage(), containsString(file.toString()));
+ assertThat(e.getMessage(), containsString("cannot parse flag or segment data"));
+ assertThat(e.getCause(), instanceOf(SerializationException.class));
+ }
+ }
+ }
+
+ @Test
+ public void duplicateKeysFailByDefaultAcrossFilesAndWithinAFile() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path a = write(dir, "a.json", "{\"flagValues\":{\"flag\":\"a\"}}");
+ Path b = write(dir, "b.json", "{\"flagValues\":{\"flag\":\"b\"}}");
+ try {
+ loader(a, b).load();
+ fail("expected exception");
+ } catch (FileDataException e) {
+ assertThat(e.getMessage(), containsString("in features, key \"flag\" was already defined"));
+ assertThat(e.getMessage(), containsString(b.toString()));
+ assertNull(e.getCause());
+ }
+
+ Path both = write(dir, "both.json", "{\"flags\":{\"x\":{\"key\":\"x\"}},\"flagValues\":{\"x\":true}}");
+ try {
+ loader(both).load();
+ fail("expected exception");
+ } catch (FileDataException e) {
+ assertThat(e.getMessage(), containsString("key \"x\" was already defined"));
+ }
+
+ Path seg1 = write(dir, "seg1.json", "{\"segments\":{\"s\":{\"key\":\"s\"}}}");
+ Path seg2 = write(dir, "seg2.json", "{\"segments\":{\"s\":{\"key\":\"s\"}}}");
+ try {
+ loader(seg1, seg2).load();
+ fail("expected exception");
+ } catch (FileDataException e) {
+ assertThat(e.getMessage(), containsString("in segments, key \"s\" was already defined"));
+ }
+ }
+ }
+
+ @Test
+ public void ignoreHandlingKeepsFirstFileAndCountsOnlyKeptEntries() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path first = write(dir, "first.json", "{\"flagValues\":{\"shared\":\"first\"}}");
+ Path second = write(dir, "second.json", "{\"flagValues\":{\"shared\":\"second\",\"only-second\":\"x\"}}");
+ OverrideFileLoader loader = new OverrideFileLoader(Arrays.asList(first, second), FileData.DuplicateKeysHandling.IGNORE);
+ LoadResult result = loader.load();
+
+ Map flags = toMap(result).get(FEATURES);
+ assertEquals(LDValue.buildArray().add("first").build(), json(FEATURES, flags.get("shared")).get("variations"));
+ assertEquals(2, flags.size());
+ assertEquals(1, result.getFiles().get(0).getFlags());
+ assertEquals(1, result.getFiles().get(1).getFlags()); // the dropped duplicate is not counted
+ assertEquals(2, result.getFlagCount());
+ }
+ }
+
+ @Test
+ public void contentHashReflectsRawContentOfPresentFiles() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = write(dir, "overrides.json", "{\"flagValues\":{\"a\":1}}");
+ Path missing = dir.getPath().resolve("missing.json");
+ OverrideFileLoader loader = loader(file, missing);
+ byte[] first = loader.load().getContentHash();
+ byte[] second = loader.load().getContentHash();
+ assertArrayEquals(first, second);
+
+ Files.write(file, "{\"flagValues\":{\"a\":2}}".getBytes());
+ byte[] third = loader.load().getContentHash();
+ assertThat(Arrays.equals(first, third), is(false));
+
+ // A file that appears changes the digest.
+ Files.write(missing, "{}".getBytes());
+ byte[] fourth = loader.load().getContentHash();
+ assertThat(Arrays.equals(third, fourth), is(false));
+ }
+ }
+
+ @Test
+ public void emptyDocumentContributesNothing() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path file = write(dir, "empty.json", "{}");
+ LoadResult result = loader(file).load();
+ assertEquals(0, result.getFlagCount());
+ assertTrue(result.getFiles().get(0).isPresent());
+ assertThat(result.getFiles().get(0).getFlags(), equalTo(0));
+ }
+ }
+
+ @Test
+ public void loaderReadsFilesInConfiguredOrder() throws Exception {
+ try (TempDir dir = TempDir.create()) {
+ Path a = write(dir, "a.json", "{\"flagValues\":{\"k\":\"a\"}}");
+ Path b = write(dir, "b.json", "{\"flagValues\":{\"k\":\"b\"}}");
+ OverrideFileLoader loader = new OverrideFileLoader(Arrays.asList(b, a), FileData.DuplicateKeysHandling.IGNORE);
+ assertEquals(LDValue.buildArray().add("b").build(),
+ json(FEATURES, toMap(loader.load()).get(FEATURES).get("k")).get("variations"));
+ assertEquals(Collections.singletonList(b).get(0), loader.load().getFiles().get(0).getPath());
+ }
+ }
+}