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()); + } + } +}