From 79250a7133b968a4b563d9cd0c5ac15202b04591 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:38:12 -0700 Subject: [PATCH 1/7] feat: Add reload infrastructure for file-based flag overrides The flag overrides feature needs a file source that reloads reliably: coalesced change notifications, retry after a failed read, retention of the last good data, sequential reloads, and support for files that do not exist yet. This adds that infrastructure to the integrations package as a purely additive change. The existing file data source keeps its current behavior. - FileDataReloader serializes reloads, debounces change signals, keeps the last good data on failure by not applying, retries a failed load after a bounded delay, reports an identical repeated failure once, and skips a reload whose raw content did not change. Close never waits for an in-flight read. - FileDataPoller detects changes by examining each file's modification time and size on a fixed interval, including a file that appears or disappears. - FileDataWatcher watches parent directories and signals once after start so a change between an initial load and the start of watching is not missed. - OverrideFileLoader reads a list of files with the override source's rules: a missing file 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, a document that the model rejects is a file data error, and the result carries a per-file summary and a content digest. Two additive FlagFactory methods keep document versions. FileDataSourceImpl, FileSynchronizer, FileInitializer, FileDataSourceBase, and FileDataException are unchanged. Two new test classes pin the file data source's reload and error behavior: a reload on every watch event with no debounce, no retry, and no skip, every failure reported and logged as the plain description at error level, a model-rejected document escaping as a SerializationException, and a description that requires a cause. --- .../server/integrations/FileDataPoller.java | 134 +++++ .../server/integrations/FileDataReloader.java | 324 +++++++++++ .../integrations/FileDataSourceParsing.java | 33 ++ .../server/integrations/FileDataWatcher.java | 155 ++++++ .../integrations/OverrideFileLoader.java | 247 +++++++++ .../FileDataLoadingBehaviorTest.java | 80 +++ .../integrations/FileDataPollerTest.java | 182 +++++++ .../integrations/FileDataReloaderTest.java | 512 ++++++++++++++++++ .../integrations/FileDataWatcherTest.java | 166 ++++++ .../FileSynchronizerReloadBehaviorTest.java | 204 +++++++ .../integrations/OverrideFileLoaderTest.java | 290 ++++++++++ 11 files changed, 2327 insertions(+) create mode 100644 lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataPoller.java create mode 100644 lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataReloader.java create mode 100644 lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/FileDataWatcher.java create mode 100644 lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoader.java create mode 100644 lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataLoadingBehaviorTest.java create mode 100644 lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataPollerTest.java create mode 100644 lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataReloaderTest.java create mode 100644 lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileDataWatcherTest.java create mode 100644 lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/FileSynchronizerReloadBehaviorTest.java create mode 100644 lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/integrations/OverrideFileLoaderTest.java 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()); + } + } +} From b9ceace01326456735c71013df1ac448c883db05 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Thu, 1 Oct 2026 23:16:17 +0000 Subject: [PATCH 2/7] fix: Expand value-only overrides to a flag served by fallthrough --- .../integrations/FileDataSourceParsing.java | 21 +++---------------- .../integrations/OverrideFileLoader.java | 2 +- .../integrations/OverrideFileLoaderTest.java | 7 ++++--- 3 files changed, 8 insertions(+), 22 deletions(-) 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 bae5309d..56fb4f22 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 @@ -202,26 +202,11 @@ 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 - * variation and setting the fallthrough variation to that. + * variation and setting the fallthrough variation to that. The file data source passes its + * load version. The override source passes version zero, because its documents supply no + * version for a value-only entry. */ static ItemDescriptor flagWithValue(String key, LDValue jsonValue, int version) { LDValue o = LDValue.buildObject() 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 index 77dca424..987056d0 100644 --- 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 @@ -187,7 +187,7 @@ LoadResult load() throws FileDataException { } if (fileContents.flagValues != null) { for (Map.Entry e : fileContents.flagValues.entrySet()) { - if (add(data, FEATURES, e.getKey(), FlagFactory.offFlagWithValue(e.getKey(), e.getValue()))) { + if (add(data, FEATURES, e.getKey(), FlagFactory.flagWithValue(e.getKey(), e.getValue(), 0))) { flagsAdded++; } } 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 index d244316c..2c00651c 100644 --- 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 @@ -87,7 +87,7 @@ public void loadsFlagsFlagValuesAndSegmentsKeepingDocumentVersions() throws Exce } @Test - public void valueOnlyEntryBecomesOffFlagServingTheValue() throws Exception { + public void valueOnlyEntryBecomesFallthroughFlagServingTheValue() throws Exception { try (TempDir dir = TempDir.create()) { Path file = write(dir, "overrides.json", "{\"flagValues\":{\"f2\":\"x\"}}"); LoadResult result = loader(file).load(); @@ -96,8 +96,9 @@ public void valueOnlyEntryBecomesOffFlagServingTheValue() throws Exception { 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.of(true), flag.get("on")); + assertEquals(LDValue.ofNull(), flag.get("offVariation")); + assertEquals(LDValue.of(0), flag.get("fallthrough").get("variation")); assertEquals(LDValue.buildArray().add("x").build(), flag.get("variations")); } } From 71971c6f3bd269b8ecbdc4688f477d948f5d2128 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:25:07 +0000 Subject: [PATCH 3/7] fix: Register a watched directory again when it is missing or deleted A directory that does not exist at start, or whose watch key becomes invalid because it was deleted, is remembered as missing and registered again on a fixed schedule from the worker thread. A recovered directory signals one change so that files written while it was not watched are picked up. The worker ends on close. --- .../server/integrations/FileDataWatcher.java | 169 +++++++++++++++--- .../integrations/FileDataWatcherTest.java | 79 +++++++- 2 files changed, 213 insertions(+), 35 deletions(-) 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 index 172f1a0f..70ceb6e8 100644 --- 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 @@ -7,13 +7,17 @@ import java.io.IOException; import java.nio.file.ClosedWatchServiceException; import java.nio.file.FileSystems; +import java.nio.file.NoSuchFileException; 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.time.Duration; import java.util.HashSet; +import java.util.Iterator; import java.util.Set; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import static java.nio.file.StandardWatchEventKinds.ENTRY_CREATE; @@ -26,9 +30,10 @@ * 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)}. + * parent directory of every file. A file that does not exist yet is reported when it appears. A + * directory that does not exist, or that is deleted while it is watched, is not an error. The + * watcher tries to register it again on a fixed schedule until it exists, and then signals one + * change so that a file written to it in the meantime is picked up. *

* 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 @@ -36,23 +41,49 @@ * {@link FileDataReloader}, which coalesces the burst and skips no-op reloads. */ final class FileDataWatcher implements Closeable, Runnable { + /** + * The time between attempts to register a directory that cannot be watched because it does not + * exist. + */ + static final Duration DEFAULT_DIRECTORY_RETRY_DELAY = Duration.ofSeconds(1); + private final WatchService watchService; private final Set watchedFilePaths; + private final long directoryRetryDelayMillis; private final LDLogger logger; private final Thread thread; private final AtomicBoolean stopped = new AtomicBoolean(false); private volatile Runnable onChange; + // The directories that are not registered because they do not exist, and the time of the next + // attempt to register them. Both are written during creation and then only by the worker thread. + private final Set missingDirectories = new HashSet<>(); + private long nextRetryAtMillis; + /** * 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 + * @throws IOException if the watch service cannot be created */ static FileDataWatcher create(Iterable filePaths, LDLogger logger) throws IOException { + return create(filePaths, DEFAULT_DIRECTORY_RETRY_DELAY, logger); + } + + /** + * Creates a watcher with a specific delay between attempts to register a directory that does not + * exist. Visible for tests. + * + * @param filePaths the files to watch, absolute or relative to the working directory + * @param directoryRetryDelay the time between registration attempts for a missing directory + * @param logger the logger + * @return the watcher + * @throws IOException if the watch service cannot be created + */ + static FileDataWatcher create(Iterable filePaths, Duration directoryRetryDelay, LDLogger logger) + throws IOException { Set directoryPaths = new HashSet<>(); Set absoluteFilePaths = new HashSet<>(); for (Path p : filePaths) { @@ -61,20 +92,27 @@ static FileDataWatcher create(Iterable filePaths, LDLogger logger) throws directoryPaths.add(absolutePath.getParent()); } WatchService ws = FileSystems.getDefault().newWatchService(); + FileDataWatcher watcher = new FileDataWatcher(ws, absoluteFilePaths, directoryRetryDelay, logger); try { for (Path d : directoryPaths) { - d.register(ws, ENTRY_CREATE, ENTRY_MODIFY, ENTRY_DELETE); + watcher.registerOrRememberAsMissing(d); } - } catch (IOException | RuntimeException e) { + } catch (RuntimeException e) { ws.close(); throw e; } - return new FileDataWatcher(ws, absoluteFilePaths, logger); + return watcher; } - private FileDataWatcher(WatchService watchService, Set watchedFilePaths, LDLogger logger) { + private FileDataWatcher( + WatchService watchService, + Set watchedFilePaths, + Duration directoryRetryDelay, + LDLogger logger + ) { this.watchService = watchService; this.watchedFilePaths = watchedFilePaths; + this.directoryRetryDelayMillis = Math.max(directoryRetryDelay.toMillis(), 1); this.logger = logger; this.thread = new Thread(this, "LaunchDarkly-FileDataWatcher"); this.thread.setDaemon(true); @@ -97,33 +135,100 @@ public void run() { while (!stopped.get()) { WatchKey key; try { - key = watchService.take(); // blocks until a change is available or the thread is interrupted + if (missingDirectories.isEmpty()) { + key = watchService.take(); // blocks until a change is available or the thread is interrupted + } else { + // Wake up for the next registration attempt even when no change arrives. + long waitMillis = nextRetryAtMillis - System.currentTimeMillis(); + key = waitMillis > 0 ? watchService.poll(waitMillis, TimeUnit.MILLISECONDS) : null; + } } 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. + if (key != null) { + processKey(key); + } + if (!missingDirectories.isEmpty() && System.currentTimeMillis() >= nextRetryAtMillis) { + retryMissingDirectories(); + } + } + } + + private void processKey(WatchKey key) { + 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; } - 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(); + } + // Without the reset, the watch on this key stops working. The reset fails when the key is no + // longer valid, which means that the directory was deleted. The watch on it is gone, so the + // directory is registered again once it exists. + if (!key.reset() && !stopped.get() && key.watchable() instanceof Path) { + Path directory = (Path) key.watchable(); + logger.warn("Directory {} no longer exists. It is watched again once it appears.", directory); + rememberAsMissing(directory); + } + if (watchedFileWasChanged && !stopped.get()) { + signal(); + } + } + + // Registers a directory with the watch service. A directory that cannot be registered is + // remembered as missing, so that the worker thread tries again on the retry schedule. + private void registerOrRememberAsMissing(Path directory) { + try { + directory.register(watchService, ENTRY_CREATE, ENTRY_MODIFY, ENTRY_DELETE); + } catch (NoSuchFileException e) { + logger.warn("Directory {} does not exist. It is watched once it appears.", directory); + rememberAsMissing(directory); + } catch (IOException e) { + logger.warn("Unable to watch directory {}: {}. The attempt is repeated until it succeeds.", directory, + LogValues.exceptionSummary(e)); + rememberAsMissing(directory); + } + } + + private void rememberAsMissing(Path directory) { + if (missingDirectories.isEmpty()) { + nextRetryAtMillis = System.currentTimeMillis() + directoryRetryDelayMillis; + } + missingDirectories.add(directory); + } + + // Tries to register every missing directory. A directory that is registered leaves the missing + // set, and one change is signalled so that files written to it while it was not watched are + // picked up. A directory that still cannot be registered stays in the set without a new log + // entry, because its absence was logged when it was first found missing. + private void retryMissingDirectories() { + boolean registered = false; + for (Iterator it = missingDirectories.iterator(); it.hasNext();) { + Path directory = it.next(); + try { + directory.register(watchService, ENTRY_CREATE, ENTRY_MODIFY, ENTRY_DELETE); + } catch (IOException | ClosedWatchServiceException e) { + continue; } + logger.info("Directory {} exists. Watching it for changes.", directory); + it.remove(); + registered = true; + } + nextRetryAtMillis = System.currentTimeMillis() + directoryRetryDelayMillis; + if (registered && !stopped.get()) { + signal(); } } @@ -136,6 +241,18 @@ private void signal() { } } + /** + * Waits for the worker thread to end after {@link #close()}. Visible for tests. + * + * @param timeoutMillis how long to wait + * @return true if the worker thread has ended + * @throws InterruptedException if the wait is interrupted + */ + boolean awaitStop(long timeoutMillis) throws InterruptedException { + thread.join(timeoutMillis); + return !thread.isAlive(); + } + /** * Stops watching and releases the watch service. It does not wait for the worker thread. */ 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 index af58b078..246ccf0d 100644 --- 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 @@ -1,24 +1,28 @@ 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.testhelpers.TempDir; import org.junit.Test; -import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; +import java.time.Duration; import java.util.Collections; import java.util.concurrent.atomic.AtomicInteger; import static org.junit.Assert.assertEquals; -import static org.junit.Assert.fail; +import static org.junit.Assert.assertTrue; @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 final Duration DIRECTORY_RETRY = Duration.ofMillis(50); private static void awaitAtLeast(AtomicInteger counter, int expected) throws InterruptedException { long deadline = System.currentTimeMillis() + WAIT_MILLIS; @@ -134,18 +138,75 @@ public void relativePathIsWatchedByItsAbsoluteLocation() throws Exception { } @Test - public void missingDirectoryIsAnErrorAtCreation() throws Exception { + public void deletedAndRecreatedDirectoryIsWatchedAgain() throws Exception { + LogCapture logCapture = Logs.capture(); + LDLogger logger = LDLogger.withAdapter(logCapture, ""); 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 + Path sub = dir.getPath().resolve("sub"); + Files.createDirectory(sub); + Path file = sub.resolve("a.json"); + Files.write(file, "{}".getBytes()); + AtomicInteger changes = new AtomicInteger(); + try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), DIRECTORY_RETRY, logger)) { + watcher.start(changes::incrementAndGet); + awaitAtLeast(changes, 1); + Thread.sleep(QUIET_MILLIS); + + // Delete the file and then its directory, and let the deletion notifications settle. + Files.delete(file); + Files.delete(sub); + awaitAtLeast(changes, 2); + Thread.sleep(QUIET_MILLIS); + int afterDeletion = changes.get(); + + // A file written into the recreated directory must be reported. + Files.createDirectory(sub); + Files.write(file, "{\"flagValues\":{}}".getBytes()); + awaitAtLeast(changes, afterDeletion + 1); + + // The lost directory is reported once, not on every registration attempt. + long warnings = logCapture.getMessages().stream().filter(m -> m.getLevel() == LDLogLevel.WARN).count(); + assertEquals(1, warnings); } } } + @Test + public void directoryMissingAtStartIsWatchedOnceItExists() throws Exception { + try (TempDir dir = TempDir.create()) { + Path sub = dir.getPath().resolve("later"); + Path file = sub.resolve("a.json"); + AtomicInteger changes = new AtomicInteger(); + try (FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), DIRECTORY_RETRY, testLogger)) { + watcher.start(changes::incrementAndGet); + awaitAtLeast(changes, 1); + Thread.sleep(QUIET_MILLIS); + + // Create the directory and write the file into it. + Files.createDirectory(sub); + Files.write(file, "{}".getBytes()); + // The file is reported even though its directory did not exist when watching began. + awaitAtLeast(changes, 2); + } + } + } + + @Test + public void closeEndsTheWorkerWhileADirectoryIsMissing() throws Exception { + try (TempDir dir = TempDir.create()) { + Path file = dir.getPath().resolve("later").resolve("a.json"); + AtomicInteger changes = new AtomicInteger(); + FileDataWatcher watcher = FileDataWatcher.create(Collections.singletonList(file), DIRECTORY_RETRY, testLogger); + watcher.start(changes::incrementAndGet); + awaitAtLeast(changes, 1); + + // Close while the worker is waiting for its next registration attempt. + watcher.close(); + // The worker must end instead of continuing to retry. + assertTrue(watcher.awaitStop(WAIT_MILLIS)); + } + } + @Test public void closeStopsSignalsAndIsIdempotent() throws Exception { try (TempDir dir = TempDir.create()) { From fd26195933ea46796d720ef508c9962835baad7e Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:25:07 +0000 Subject: [PATCH 4/7] fix: Keep polling after the change callback throws A RuntimeException from the callback is logged at error level instead of ending the repeating task. --- .../server/integrations/FileDataPoller.java | 16 ++++++- .../integrations/FileDataPollerTest.java | 46 +++++++++++++++---- 2 files changed, 52 insertions(+), 10 deletions(-) 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 index 2cc5e488..954e22e6 100644 --- 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 @@ -1,5 +1,8 @@ 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.Files; @@ -31,6 +34,7 @@ final class FileDataPoller implements Closeable { private final List paths; private final Runnable onChange; + private final LDLogger logger; private final ScheduledThreadPoolExecutor executor; private final AtomicBoolean closed = new AtomicBoolean(false); private List last; @@ -42,10 +46,12 @@ final class FileDataPoller implements Closeable { * @param paths the files to examine * @param interval the time between examinations * @param onChange called when any file changed since the previous examination + * @param logger receives log output about a callback that fails */ - FileDataPoller(List paths, Duration interval, Runnable onChange) { + FileDataPoller(List paths, Duration interval, Runnable onChange, LDLogger logger) { this.paths = new ArrayList<>(paths); this.onChange = onChange; + this.logger = logger; this.last = observeAll(this.paths); ThreadFactory threadFactory = runnable -> { Thread t = new Thread(runnable, "LaunchDarkly-FileDataPoller"); @@ -80,7 +86,13 @@ private void examine() { boolean changed = !current.equals(last); last = current; if (changed && !closed.get()) { - onChange.run(); + try { + onChange.run(); + } catch (RuntimeException e) { + // The executor would otherwise cancel the repeating task and the poller would go quiet. + logger.error("Unexpected error while handling a file change: {}", LogValues.exceptionSummary(e)); + logger.debug(LogValues.exceptionTrace(e)); + } } } 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 index b243d968..cd6bc3d6 100644 --- 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 @@ -1,5 +1,9 @@ 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.testhelpers.TempDir; import org.junit.Test; @@ -22,6 +26,7 @@ @SuppressWarnings("javadoc") public class FileDataPollerTest { + private static final LDLogger testLogger = LDLogger.none(); private static final Duration INTERVAL = Duration.ofMillis(20); private static final long WAIT_MILLIS = 5000; private static final long QUIET_MILLIS = 250; @@ -53,7 +58,7 @@ 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)) { + try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet, testLogger)) { Thread.sleep(QUIET_MILLIS); assertEquals(0, changes.get()); @@ -68,7 +73,7 @@ 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)) { + try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet, testLogger)) { rewrite(file, "{\"flagValues\":{}}", 1000_000); awaitCount(changes, 1); } @@ -80,7 +85,7 @@ 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)) { + try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet, testLogger)) { rewrite(file, "{}", 2000_000); awaitCount(changes, 1); Thread.sleep(QUIET_MILLIS); @@ -94,7 +99,7 @@ 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)) { + try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet, testLogger)) { Thread.sleep(QUIET_MILLIS); assertEquals(0, changes.get()); @@ -109,7 +114,7 @@ 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)) { + try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet, testLogger)) { Files.delete(file); awaitCount(changes, 1); } @@ -123,7 +128,7 @@ public void watchesEveryConfiguredFile() throws Exception { 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)) { + try (FileDataPoller poller = new FileDataPoller(paths, INTERVAL, changes::incrementAndGet, testLogger)) { rewrite(second, "{}", 2000_000); awaitCount(changes, 1); rewrite(first, "{}", 2000_000); @@ -137,7 +142,7 @@ 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); + FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, changes::incrementAndGet, testLogger); poller.close(); poller.close(); // idempotent @@ -160,7 +165,7 @@ public void closeReturnsWhileCallbackBlocks() throws Exception { } catch (InterruptedException e) { Thread.currentThread().interrupt(); } - }); + }, testLogger); rewrite(file, "{}", 2000_000); assertTrue(entered.await(WAIT_MILLIS, TimeUnit.MILLISECONDS)); @@ -171,6 +176,31 @@ public void closeReturnsWhileCallbackBlocks() throws Exception { } } + @Test + public void pollingContinuesAfterTheCallbackThrows() throws Exception { + LogCapture logCapture = Logs.capture(); + LDLogger logger = LDLogger.withAdapter(logCapture, ""); + try (TempDir dir = TempDir.create()) { + Path file = createFile(dir, "a.json", "{}", 1000_000); + AtomicInteger changes = new AtomicInteger(); + Runnable onChange = () -> { + if (changes.incrementAndGet() == 1) { + throw new IllegalStateException("consumer failed"); + } + }; + try (FileDataPoller poller = new FileDataPoller(Collections.singletonList(file), INTERVAL, onChange, logger)) { + // The first change makes the callback throw. + rewrite(file, "{}", 2000_000); + awaitCount(changes, 1); + + // The next change must still be detected and delivered, and the failure logged as an error. + rewrite(file, "{}", 3000_000); + awaitCount(changes, 2); + assertTrue(logCapture.getMessages().stream().anyMatch(m -> m.getLevel() == LDLogLevel.ERROR)); + } + } + } + @Test public void unreadableAttributesCountAsAbsent() throws Exception { try (TempDir dir = TempDir.create()) { From f9613962b84e719fc98a96c90c365beba3e4ee01 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:25:07 +0000 Subject: [PATCH 5/7] fix: Report a failed apply as a load failure and retry it The last good hash and the failure state are recorded only after the handler accepts the result. An exception from apply goes through the same failure path as a read or parse failure, so it is reported once and retried after the retry delay. --- .../server/integrations/FileDataReloader.java | 22 +++++-- .../integrations/FileDataReloaderTest.java | 58 ++++++++++++++----- 2 files changed, 60 insertions(+), 20 deletions(-) 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 index 259ee185..97a73c07 100644 --- 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 @@ -54,16 +54,19 @@ interface Loader { */ interface Handler { /** - * Called with each successfully merged result. + * Called with each successfully merged result. An exception thrown here fails the reload: it is + * reported through {@link #onError(FileDataException)}, nothing from the load is remembered as + * the last good result, and the reload is retried. * * @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. + * Called when a reload fails, including when {@link #apply(LoadResult)} throws, 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 */ @@ -291,13 +294,20 @@ private boolean reload() { // 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; } + // Nothing is remembered until the consumer has accepted the result. A result that the + // consumer rejects must not become the baseline that skip-unchanged compares against, and + // must not count as a recovery. + try { + handler.apply(result); + } catch (RuntimeException e) { + return fail(new FileDataException("unable to apply flag data", e)); + } + lastErrorMessage = null; lastGoodHash = hash; - handler.apply(result); return true; } } 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 index 4c2c1704..017906bc 100644 --- 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 @@ -24,6 +24,7 @@ import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.greaterThanOrEqualTo; import static org.hamcrest.Matchers.hasItem; +import static org.hamcrest.Matchers.instanceOf; import static org.hamcrest.Matchers.lessThan; import static org.hamcrest.Matchers.startsWith; import static org.junit.Assert.assertEquals; @@ -85,6 +86,13 @@ ScriptedLoader thenFail(String message) { return this; } + ScriptedLoader thenThrow(RuntimeException e) { + steps.add(() -> { + throw e; + }); + return this; + } + @Override public LoadResult load() throws FileDataException { // Both latches are read before the entered signal, so a test that clears them after the @@ -125,7 +133,7 @@ private static final class UncheckedFileDataException extends RuntimeException { } } - private static final class RecordingHandler implements FileDataReloader.Handler { + private static class RecordingHandler implements FileDataReloader.Handler { final List applied = Collections.synchronizedList(new ArrayList<>()); final List errors = Collections.synchronizedList(new ArrayList<>()); @@ -483,29 +491,51 @@ public void reloadsAreSerializedBetweenCallerAndWorker() throws Exception { } @Test - public void unexpectedRuntimeExceptionFromHandlerIsLoggedAndWorkerSurvives() throws Exception { - ScriptedLoader loader = new ScriptedLoader().then(resultWithHash(1)).then(resultWithHash(2)); + public void applyThatThrowsIsReportedAsAFailureAndRetried() throws Exception { + LoadResult result = resultWithHash(1); + ScriptedLoader loader = new ScriptedLoader().then(result); AtomicInteger applies = new AtomicInteger(); - FileDataReloader.Handler handler = new FileDataReloader.Handler() { + RecordingHandler handler = new RecordingHandler() { @Override - public void apply(LoadResult result) { + public void apply(LoadResult r) { if (applies.incrementAndGet() == 1) { throw new IllegalStateException("consumer failed"); } - } - - @Override - public void onError(FileDataException e) { + super.apply(r); } }; - FileDataReloader r = new FileDataReloader(loader, handler, logger, Duration.ZERO, Duration.ZERO, false); - reloaders.add(r); + FileDataReloader r = reloader(loader, handler, Duration.ZERO, SHORT, true); + // The initial load's apply throws. The reload returns normally and reports the failure. + r.reloadNow(); + assertTrue(handler.applied.isEmpty()); + assertEquals(1, handler.errors.size()); + assertThat(handler.errors.get(0).getCause(), instanceOf(IllegalStateException.class)); + assertThat(logCapture.getMessageStrings(), hasItem(startsWith("ERROR:Unable to load flags:"))); + + // The retry applies the same content without a further trigger. + awaitSize(handler.applied, 1); + assertEquals(result, handler.applied.get(0)); + + // Only now is the content remembered: a further reload of it is skipped. r.trigger(); - awaitAtLeast(applies, 1); - r.trigger(); - awaitAtLeast(applies, 2); + awaitAtLeast(loader.calls, 3); + assertEquals(1, handler.applied.size()); + } + @Test + public void unexpectedRuntimeExceptionFromLoaderIsLoggedAndWorkerSurvives() throws Exception { + ScriptedLoader loader = new ScriptedLoader().thenThrow(new IllegalStateException("loader failed")) + .then(resultWithHash(1)); + RecordingHandler handler = new RecordingHandler(); + FileDataReloader r = reloader(loader, handler, Duration.ZERO, Duration.ZERO, false); + + // The first reload's loader throws something other than a file data failure. + r.trigger(); + awaitAtLeast(loader.calls, 1); + // The worker must still run the next reload, and the exception must have been logged. + r.trigger(); + awaitSize(handler.applied, 1); assertThat(logCapture.getMessageStrings(), hasItem(startsWith("ERROR:Unexpected error while reloading flag data:"))); } From 548e2cfeb33f95510e75c6fe9fc9205eba204f86 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:25:07 +0000 Subject: [PATCH 6/7] docs: Describe value-only overrides as flags served by fallthrough --- .../sdk/server/integrations/OverrideFileLoader.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 index 987056d0..ece5b8b7 100644 --- 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 @@ -32,9 +32,9 @@ *

* 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. + * versions that the documents specify, a value-only entry becomes a flag that is on and serves the + * value by fallthrough, 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 { /** From 8edc7dbe7b0d438919255cf355726b9700446413 Mon Sep 17 00:00:00 2001 From: Ryan Lamb <4955475+kinyoklion@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:51:05 +0000 Subject: [PATCH 7/7] fix: Schedule directory registration retries on the monotonic clock The time of the next registration attempt and the wait before it come from System.nanoTime, so a step of the system clock neither delays nor advances the retry. --- .../server/integrations/FileDataWatcher.java | 20 ++++++++++--------- 1 file changed, 11 insertions(+), 9 deletions(-) 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 index 70ceb6e8..dbbd897c 100644 --- 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 @@ -49,16 +49,18 @@ final class FileDataWatcher implements Closeable, Runnable { private final WatchService watchService; private final Set watchedFilePaths; - private final long directoryRetryDelayMillis; + private final long directoryRetryDelayNanos; private final LDLogger logger; private final Thread thread; private final AtomicBoolean stopped = new AtomicBoolean(false); private volatile Runnable onChange; // The directories that are not registered because they do not exist, and the time of the next - // attempt to register them. Both are written during creation and then only by the worker thread. + // attempt to register them. The time comes from the monotonic clock, so a change to the system + // clock does not move the schedule. Both are written during creation and then only by the + // worker thread. private final Set missingDirectories = new HashSet<>(); - private long nextRetryAtMillis; + private long nextRetryAtNanos; /** * Creates a watcher for the given files. Nothing is watched until {@link #start(Runnable)}. @@ -112,7 +114,7 @@ private FileDataWatcher( ) { this.watchService = watchService; this.watchedFilePaths = watchedFilePaths; - this.directoryRetryDelayMillis = Math.max(directoryRetryDelay.toMillis(), 1); + this.directoryRetryDelayNanos = TimeUnit.MILLISECONDS.toNanos(Math.max(directoryRetryDelay.toMillis(), 1)); this.logger = logger; this.thread = new Thread(this, "LaunchDarkly-FileDataWatcher"); this.thread.setDaemon(true); @@ -139,8 +141,8 @@ public void run() { key = watchService.take(); // blocks until a change is available or the thread is interrupted } else { // Wake up for the next registration attempt even when no change arrives. - long waitMillis = nextRetryAtMillis - System.currentTimeMillis(); - key = waitMillis > 0 ? watchService.poll(waitMillis, TimeUnit.MILLISECONDS) : null; + long waitNanos = nextRetryAtNanos - System.nanoTime(); + key = waitNanos > 0 ? watchService.poll(waitNanos, TimeUnit.NANOSECONDS) : null; } } catch (InterruptedException e) { continue; // if stopped, the loop condition ends the thread @@ -150,7 +152,7 @@ public void run() { if (key != null) { processKey(key); } - if (!missingDirectories.isEmpty() && System.currentTimeMillis() >= nextRetryAtMillis) { + if (!missingDirectories.isEmpty() && System.nanoTime() - nextRetryAtNanos >= 0) { retryMissingDirectories(); } } @@ -204,7 +206,7 @@ private void registerOrRememberAsMissing(Path directory) { private void rememberAsMissing(Path directory) { if (missingDirectories.isEmpty()) { - nextRetryAtMillis = System.currentTimeMillis() + directoryRetryDelayMillis; + nextRetryAtNanos = System.nanoTime() + directoryRetryDelayNanos; } missingDirectories.add(directory); } @@ -226,7 +228,7 @@ private void retryMissingDirectories() { it.remove(); registered = true; } - nextRetryAtMillis = System.currentTimeMillis() + directoryRetryDelayMillis; + nextRetryAtNanos = System.nanoTime() + directoryRetryDelayNanos; if (registered && !stopped.get()) { signal(); }