diff --git a/lib/sdk/server/contract-tests/service/src/main/java/sdktest/Representations.java b/lib/sdk/server/contract-tests/service/src/main/java/sdktest/Representations.java index 99112355..cd7593e7 100644 --- a/lib/sdk/server/contract-tests/service/src/main/java/sdktest/Representations.java +++ b/lib/sdk/server/contract-tests/service/src/main/java/sdktest/Representations.java @@ -255,6 +255,7 @@ public static class EvaluationSeriesContextParam { LDContext context; LDValue defaultValue; String method; + String environmentId; } public static class IdentifyEventParams { diff --git a/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestHook.java b/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestHook.java index ac979d07..bb16db72 100644 --- a/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestHook.java +++ b/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestHook.java @@ -44,6 +44,7 @@ public Map beforeEvaluation(EvaluationSeriesContext seriesContex seriesContextParam.context = seriesContext.context; seriesContextParam.defaultValue = seriesContext.defaultValue; seriesContextParam.method = seriesContext.method; + seriesContextParam.environmentId = seriesContext.environmentId; params.evaluationSeriesContext = seriesContextParam; params.evaluationSeriesData = data; @@ -72,6 +73,7 @@ public Map afterEvaluation(EvaluationSeriesContext seriesContext seriesContextParam.context = seriesContext.context; seriesContextParam.defaultValue = seriesContext.defaultValue; seriesContextParam.method = seriesContext.method; + seriesContextParam.environmentId = seriesContext.environmentId; params.evaluationSeriesContext = seriesContextParam; params.evaluationSeriesData = data; diff --git a/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestService.java b/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestService.java index 77d48cda..45eeead7 100644 --- a/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestService.java +++ b/lib/sdk/server/contract-tests/service/src/main/java/sdktest/TestService.java @@ -34,6 +34,7 @@ public class TestService { "event-gzip", "event-sampling", "filtering", + "hook-environment-id", "inline-context-all", "migrations", "optional-event-gzip", diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataModelDependencies.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataModelDependencies.java index 57a4a16c..17ca660e 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataModelDependencies.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataModelDependencies.java @@ -131,7 +131,7 @@ public static FullDataSet sortAllCollections(FullDataSet(builder.build().entrySet(), allData.shouldPersist()); + return new FullDataSet<>(builder.build().entrySet(), allData.shouldPersist(), allData.getEnvironmentId()); } /** diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapter.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapter.java index 3281088e..ab98fde5 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapter.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapter.java @@ -131,6 +131,8 @@ public void close() { */ private static class ConvertingUpdateSink implements DataSourceUpdateSink { private final IterableAsyncQueue resultQueue; + // Only a full data set carries an environment ID; it is retained for subsequent partial updates. + private volatile String environmentId = null; public ConvertingUpdateSink(IterableAsyncQueue resultQueue) { this.resultQueue = resultQueue; @@ -138,13 +140,16 @@ public ConvertingUpdateSink(IterableAsyncQueue resultQueue) { @Override public boolean init(DataStoreTypes.FullDataSet allData) { + if (allData.getEnvironmentId() != null && !allData.getEnvironmentId().isEmpty()) { + environmentId = allData.getEnvironmentId(); + } // Convert the full data set into a ChangeSet and emit it ChangeSet>>> changeSet = new ChangeSet<>( ChangeSetType.Full, Selector.EMPTY, allData.getData(), - null, + environmentId, allData.shouldPersist() ); resultQueue.put(FDv2SourceResult.changeSet(changeSet, false)); @@ -166,7 +171,7 @@ public boolean upsert(DataKind kind, String key, ItemDescriptor item) { ChangeSetType.Partial, Selector.EMPTY, data, - null, + environmentId, true // default to true as this adapter is used for adapting FDv1 data sources which are always persistent ); resultQueue.put(FDv2SourceResult.changeSet(changeSet, false)); diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceUpdatesImpl.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceUpdatesImpl.java index 5b8e5095..2f945f64 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceUpdatesImpl.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSourceUpdatesImpl.java @@ -429,7 +429,8 @@ private boolean applyToLegacyStore(ChangeSet>>> unsortedChangeset) { // Convert ChangeSet to FullDataSet for legacy init path, preserving shouldPersist flag - return init(new FullDataSet<>(unsortedChangeset.getData(), unsortedChangeset.shouldPersist())); + return init(new FullDataSet<>(unsortedChangeset.getData(), unsortedChangeset.shouldPersist(), + unsortedChangeset.getEnvironmentId())); } private boolean applyPartialChangeSetToLegacyStore(ChangeSet>>> changeSet) { diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSystem.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSystem.java index 7b5df387..6558cc45 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSystem.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DataSystem.java @@ -55,6 +55,13 @@ interface DataSystem { * @return the data store status provider */ DataStoreStatusProvider getDataStoreStatusProvider(); + + /** + * Returns the ID of the LaunchDarkly environment the data came from, or null if it is not known. + * + * @return the environment ID, or null + */ + String getEnvironmentId(); } /** diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DefaultFeatureRequestor.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DefaultFeatureRequestor.java index 6f01f596..abe46048 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DefaultFeatureRequestor.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/DefaultFeatureRequestor.java @@ -128,7 +128,8 @@ public FullDataSet getAllData(boolean returnDataEvenIfCached) JsonReader jr = new JsonReader(response.body().charStream()); // Polling data from LaunchDarkly should be persisted - return new FullDataSet<>(parseFullDataSet(jr), true); + return new FullDataSet<>(parseFullDataSet(jr), true, + response.header(HeaderConstants.ENVIRONMENT_ID.getHeaderName())); } } } diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/EvaluatorWithHooks.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/EvaluatorWithHooks.java index aada825e..fd69420d 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/EvaluatorWithHooks.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/EvaluatorWithHooks.java @@ -1,6 +1,7 @@ package com.launchdarkly.sdk.server; import com.launchdarkly.logging.LDLogger; +import com.launchdarkly.logging.LogValues; import com.launchdarkly.sdk.LDContext; import com.launchdarkly.sdk.LDValue; import com.launchdarkly.sdk.LDValueType; @@ -11,6 +12,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.function.Supplier; /** * An {@link EvaluatorInterface} that will invoke the evaluation series methods of the provided {@link Hook} when @@ -21,13 +23,17 @@ class EvaluatorWithHooks implements EvaluatorInterface { private final EvaluatorInterface underlyingEvaluator; private final List hooks; private final LDLogger logger; + private final Supplier environmentIdSupplier; /** - * @param underlyingEvaluator that will do the actual flag evaluation - * @param hooks that will be invoked at various stages of the evaluation series - * @param hooksLogger that will be used to log + * @param underlyingEvaluator that will do the actual flag evaluation + * @param hooks that will be invoked at various stages of the evaluation series + * @param hooksLogger that will be used to log + * @param environmentIdSupplier provides the environment ID reported by LaunchDarkly, if known */ - EvaluatorWithHooks(EvaluatorInterface underlyingEvaluator, List hooks, LDLogger hooksLogger) { + EvaluatorWithHooks(EvaluatorInterface underlyingEvaluator, List hooks, LDLogger hooksLogger, + Supplier environmentIdSupplier) { + this.environmentIdSupplier = environmentIdSupplier; this.underlyingEvaluator = underlyingEvaluator; this.hooks = hooks; this.logger = hooksLogger; @@ -40,7 +46,8 @@ public EvalResultAndFlag evalAndFlag(String method, String featureKey, LDContext int size = hooks.size(); List seriesDataList = new ArrayList<>(size); - EvaluationSeriesContext seriesContext = new EvaluationSeriesContext(method, featureKey, context, defaultValue); + EvaluationSeriesContext seriesContext = new EvaluationSeriesContext(method, featureKey, context, defaultValue, + getEnvironmentId(featureKey)); Map emptyMap = Collections.emptyMap(); for (int i = 0; i < size; i++) { Hook currentHook = hooks.get(i); @@ -68,6 +75,20 @@ public EvalResultAndFlag evalAndFlag(String method, String featureKey, LDContext return result; } + /** + * Gets the current environment ID from the supplier. The supplier ultimately reads from the data store, + * which may be a customer-provided implementation, so a failure here must not prevent the evaluation. + */ + private String getEnvironmentId(String featureKey) { + try { + return environmentIdSupplier.get(); + } catch (Exception e) { + logger.error("During evaluation of flag \"{}\". Unable to determine the environment ID for hooks: {}", featureKey, + LogValues.exceptionSummary(e)); + return null; + } + } + @Override public FeatureFlagsState allFlagsState(LDContext context, FlagsStateOption... options) { // We do not support hooks for when all flags are evaluated. Perhaps in the future that will be added. diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv1DataSystem.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv1DataSystem.java index ef1d340c..a03da864 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv1DataSystem.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv1DataSystem.java @@ -155,6 +155,11 @@ public DataStoreStatusProvider getDataStoreStatusProvider() { return dataStoreStatusProvider; } + @Override + public String getEnvironmentId() { + return dataStore.getEnvironmentId(); + } + @Override public void close() throws IOException { if (disposed) { diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv2DataSystem.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv2DataSystem.java index 1fd870ba..5053625b 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv2DataSystem.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/FDv2DataSystem.java @@ -226,6 +226,11 @@ public DataStoreStatusProvider getDataStoreStatusProvider() { return dataStoreStatusProvider; } + @Override + public String getEnvironmentId() { + return store.getEnvironmentId(); + } + @Override public void close() throws IOException { if (disposed) { diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/InMemoryDataStore.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/InMemoryDataStore.java index 2e394990..97d52c95 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/InMemoryDataStore.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/InMemoryDataStore.java @@ -33,10 +33,11 @@ class InMemoryDataStore implements DataStore, TransactionalDataStore, CacheExpor private final Object selectorLock = new Object(); private volatile Selector selector = Selector.EMPTY; private volatile boolean shouldPersist = false; + private volatile String environmentId = null; @Override public void init(FullDataSet allData) { - applyFullPayload(allData.getData(), null, Selector.EMPTY, allData.shouldPersist()); + applyFullPayload(allData.getData(), allData.getEnvironmentId(), Selector.EMPTY, allData.shouldPersist()); } @Override @@ -128,7 +129,8 @@ public void apply(ChangeSet>> data, - Selector selector, boolean shouldPersist) { + String environmentId, Selector selector, boolean shouldPersist) { synchronized (this.writeLock) { // Build the complete updated dictionary before assigning to Items for transactional update ImmutableMap.Builder> itemsBuilder = ImmutableMap.builder(); @@ -192,6 +205,7 @@ private void applyPartialData(Iterable exportAll() { } // Preserve the shouldPersist value that was set when data was provided to this store - return new FullDataSet<>(builder.build(), this.shouldPersist); + return new FullDataSet<>(builder.build(), this.shouldPersist, this.environmentId); } } } diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/LDClient.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/LDClient.java index b39b5cbe..939fb51f 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/LDClient.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/LDClient.java @@ -233,8 +233,10 @@ public LDClient(String sdkKey, LDConfig config) { this.evaluator = evaluator; this.migrationEvaluator = new MigrationStageEnforcingEvaluator(evaluator, evaluationLogger); } else { - this.evaluator = new EvaluatorWithHooks(evaluator, allHooks, this.baseLogger.subLogger(Loggers.HOOKS_LOGGER_NAME)); - this.migrationEvaluator = new EvaluatorWithHooks(new MigrationStageEnforcingEvaluator(evaluator, evaluationLogger), allHooks, this.baseLogger.subLogger(Loggers.HOOKS_LOGGER_NAME)); + this.evaluator = new EvaluatorWithHooks(evaluator, allHooks, this.baseLogger.subLogger(Loggers.HOOKS_LOGGER_NAME), + this.dataSystem::getEnvironmentId); + this.migrationEvaluator = new EvaluatorWithHooks(new MigrationStageEnforcingEvaluator(evaluator, evaluationLogger), allHooks, + this.baseLogger.subLogger(Loggers.HOOKS_LOGGER_NAME), this.dataSystem::getEnvironmentId); } // Create FlagTracker using the dataSystem's flag change notifier diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreConverter.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreConverter.java index a47bfba3..54d47d62 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreConverter.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreConverter.java @@ -40,7 +40,7 @@ static FullDataSet toSerializedFormat( } // Preserve shouldPersist flag when converting formats - return new FullDataSet<>(builder.build(), inMemoryData.shouldPersist()); + return new FullDataSet<>(builder.build(), inMemoryData.shouldPersist(), inMemoryData.getEnvironmentId()); } /** diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapper.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapper.java index ef0e5e8f..0851cf6e 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapper.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapper.java @@ -63,6 +63,7 @@ final class PersistentDataStoreWrapper implements DataStore, SettableCache, Disa // themselves remain alive until GC reclaims them; the LoadingCache loaders // are short-circuited because every touch site checks this flag first. private volatile boolean cacheDisabled; + private volatile String environmentId = null; PersistentDataStoreWrapper( final PersistentDataStore core, @@ -204,7 +205,8 @@ public void init(FullDataSet allData) { KeyedItems items = PersistentDataStoreConverter.serializeAll(kind, e0.getValue()); allBuilder.add(new AbstractMap.SimpleEntry<>(kind, items)); } - RuntimeException failure = initCore(new FullDataSet<>(allBuilder.build(), allData.shouldPersist())); + RuntimeException failure = initCore(new FullDataSet<>(allBuilder.build(), allData.shouldPersist(), + allData.getEnvironmentId())); if (itemCache != null && allCache != null && !cacheDisabled) { itemCache.invalidateAll(); allCache.invalidateAll(); @@ -226,6 +228,10 @@ public void init(FullDataSet allData) { } if (failure == null || cacheIndefinitely) { inited.set(true); + String newEnvironmentId = allData.getEnvironmentId(); + if (newEnvironmentId != null && !newEnvironmentId.isEmpty()) { + environmentId = newEnvironmentId; + } } if (failure != null) { throw failure; @@ -372,6 +378,11 @@ public CacheStats getCacheStats() { itemStats.evictionCount() + allStats.evictionCount()); } + @Override + public String getEnvironmentId() { + return environmentId; + } + private ItemDescriptor getAndDeserializeItem(DataKind kind, String key) { SerializedItemDescriptor maybeSerializedItem = core.get(kind, key); return maybeSerializedItem == null ? null : PersistentDataStoreConverter.deserialize(kind, maybeSerializedItem); diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java index c7cda4a6..97df4e29 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java @@ -34,6 +34,7 @@ import com.launchdarkly.sdk.server.interfaces.DataStoreStatusProvider; import com.launchdarkly.sdk.server.subsystems.DataSource; import com.launchdarkly.sdk.server.subsystems.DataSourceUpdateSink; +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.FullDataSet; import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.ItemDescriptor; import com.launchdarkly.sdk.server.subsystems.SerializationException; @@ -304,7 +305,7 @@ private void handleMessage(MessageEvent event, CompletableFuture initFutur try { switch (event.getEventName()) { case PUT: - handlePut(event.getDataReader(), initFuture); + handlePut(event.getDataReader(), environmentIdOf(event), initFuture); break; case PATCH: @@ -366,12 +367,19 @@ private static boolean exceptionHasCause(Throwable e, Class c) { return e.getCause() != null && exceptionHasCause(e.getCause(), c); } - private void handlePut(Reader eventData, CompletableFuture initFuture) + private static String environmentIdOf(MessageEvent event) { + return event.getHeaders() == null ? null + : event.getHeaders().value(HeaderConstants.ENVIRONMENT_ID.getHeaderName()); + } + + private void handlePut(Reader eventData, String environmentId, CompletableFuture initFuture) throws StreamInputException, StreamStoreException { recordStreamInit(false); esStarted = 0; PutData putData = parseStreamJson(StreamProcessorEvents::parsePutData, eventData); - if (!dataSourceUpdates.init(putData.data)) { + FullDataSet data = new FullDataSet<>(putData.data.getData(), putData.data.shouldPersist(), + environmentId); + if (!dataSourceUpdates.init(data)) { throw new StreamStoreException(); } if (!initialized.getAndSet(true)) { diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/WriteThroughStore.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/WriteThroughStore.java index aa895228..4ef80ef3 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/WriteThroughStore.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/WriteThroughStore.java @@ -131,6 +131,11 @@ public Selector getSelector() { return txMemoryStore.getSelector(); } + @Override + public String getEnvironmentId() { + return memoryStore.getEnvironmentId(); + } + @Override public void close() throws IOException { memoryStore.close(); @@ -182,7 +187,8 @@ private boolean applyToLegacyPersistence(ChangeSet>>> sortedChangeSet) { // Preserve shouldPersist flag when converting ChangeSet to FullDataSet - persistentStore.init(new FullDataSet<>(sortedChangeSet.getData(), sortedChangeSet.shouldPersist())); + persistentStore.init(new FullDataSet<>(sortedChangeSet.getData(), sortedChangeSet.shouldPersist(), + sortedChangeSet.getEnvironmentId())); } /** diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/EvaluationSeriesContext.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/EvaluationSeriesContext.java index 7961f28a..2e90d503 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/EvaluationSeriesContext.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/integrations/EvaluationSeriesContext.java @@ -32,6 +32,13 @@ public class EvaluationSeriesContext { */ public final LDValue defaultValue; + /** + * The ID of the LaunchDarkly environment the evaluated data came from, or null if it is not + * known. It is not known before the SDK has received a successful response from LaunchDarkly, + * or when the data came from a source other than LaunchDarkly. + */ + public final String environmentId; + /** * @param method the variation method that was used to invoke the evaluation. * @param key the key of the feature flag being evaluated. @@ -39,9 +46,22 @@ public class EvaluationSeriesContext { * @param defaultValue the user-provided default value for the evaluation. */ public EvaluationSeriesContext(String method, String key, LDContext context, LDValue defaultValue) { + this(method, key, context, defaultValue, null); + } + + /** + * @param method the variation method that was used to invoke the evaluation. + * @param key the key of the feature flag being evaluated. + * @param context the context the evaluation was for. + * @param defaultValue the user-provided default value for the evaluation. + * @param environmentId the ID of the LaunchDarkly environment, or null if it is not known. + */ + public EvaluationSeriesContext(String method, String key, LDContext context, LDValue defaultValue, + String environmentId) { this.flagKey = key; this.context = context; this.defaultValue = defaultValue; this.method = method; + this.environmentId = environmentId; } } diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStore.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStore.java index e2cd3632..83a707e6 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStore.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStore.java @@ -104,4 +104,16 @@ public interface DataStore extends Closeable { * @return a cache statistics object, or null if not applicable */ CacheStats getCacheStats(); + + /** + * Returns the ID of the LaunchDarkly environment that the currently stored data came from. + *

+ * This is the environment ID that was provided with the data, if any. The SDK makes it available + * to hook implementations. + * + * @return the environment ID, or null if it is unknown + */ + default String getEnvironmentId() { + return null; + } } diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStoreTypes.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStoreTypes.java index 0667bf4e..0978a883 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStoreTypes.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/subsystems/DataStoreTypes.java @@ -263,6 +263,7 @@ public String toString() { public static final class FullDataSet { private final Iterable>> data; private final boolean shouldPersist; + private final String environmentId; /** * Returns the wrapped data set. @@ -286,14 +287,22 @@ public boolean shouldPersist() { return shouldPersist; } + /** + * Returns the ID of the LaunchDarkly environment that this data came from, if known. + * + * @return the environment ID, or null if it was not reported with the data + */ + public String getEnvironmentId() { + return environmentId; + } + /** * Constructs a new instance. Will default to shouldPersist = true and can be stored in a persistent store. * * @param data the data set */ public FullDataSet(Iterable>> data) { - this.data = data == null ? ImmutableList.of(): data; - this.shouldPersist = true; // default to true if not specified for backwards compatibility + this(data, true); // default to true if not specified for backwards compatibility } /** @@ -303,22 +312,36 @@ public FullDataSet(Iterable>> data) * @param shouldPersist true if the data should be persisted to persistent stores, false otherwise */ public FullDataSet(Iterable>> data, boolean shouldPersist) { + this(data, shouldPersist, null); + } + + /** + * Constructs a new instance. + * + * @param data the data set + * @param shouldPersist true if the data should be persisted to persistent stores, false otherwise + * @param environmentId the ID of the LaunchDarkly environment the data came from, or null if unknown + */ + public FullDataSet(Iterable>> data, boolean shouldPersist, + String environmentId) { this.data = data == null ? ImmutableList.of(): data; this.shouldPersist = shouldPersist; + this.environmentId = environmentId; } @Override public boolean equals(Object o) { if (o instanceof FullDataSet) { FullDataSet other = (FullDataSet)o; - return data.equals(other.data) && shouldPersist == other.shouldPersist; + return data.equals(other.data) && shouldPersist == other.shouldPersist + && Objects.equals(environmentId, other.environmentId); } return false; } @Override public int hashCode() { - return Objects.hash(data, shouldPersist); + return Objects.hash(data, shouldPersist, environmentId); } } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapterTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapterTest.java index 7d88f621..ce81c226 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapterTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceSynchronizerAdapterTest.java @@ -1,10 +1,14 @@ package com.launchdarkly.sdk.server; +import com.launchdarkly.sdk.fdv2.ChangeSetType; import com.launchdarkly.sdk.fdv2.SourceResultType; import com.launchdarkly.sdk.fdv2.SourceSignal; import com.launchdarkly.sdk.server.datasources.FDv2SourceResult; import com.launchdarkly.sdk.server.subsystems.DataSource; +import com.launchdarkly.sdk.server.DataStoreTestTypes.DataBuilder; import com.launchdarkly.sdk.server.subsystems.DataSourceUpdateSink; +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.FullDataSet; +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.ItemDescriptor; import org.junit.After; import org.junit.Test; @@ -17,7 +21,10 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicReference; +import static com.launchdarkly.sdk.server.DataModel.FEATURES; +import static com.launchdarkly.sdk.server.ModelBuilders.flagBuilder; import static org.junit.Assert.*; @SuppressWarnings("javadoc") @@ -295,4 +302,63 @@ public boolean isDone() { return delegate.isDone(); } } + + /** + * Test that an environment ID reported through the FDv1 update sink is carried on the + * change sets produced for the FDv2 data system. + */ + @Test + public void environmentIdIsForwardedToChangeSets() throws Exception { + AtomicReference capturedSink = new AtomicReference<>(); + + DataSourceSynchronizerAdapter adapter = new DataSourceSynchronizerAdapter(sink -> { + capturedSink.set(sink); + return new MockDataSource(new CountDownLatch(1), null); + }); + resourcesToClose.add(adapter); + + CompletableFuture nextFuture = adapter.next(); + + capturedSink.get().init(new FullDataSet<>(DataBuilder.forStandardTypes().build().getData(), true, + "env-from-fdv1")); + + FDv2SourceResult result = nextFuture.get(2, TimeUnit.SECONDS); + assertEquals(SourceResultType.CHANGE_SET, result.getResultType()); + assertEquals("env-from-fdv1", result.getChangeSet().getEnvironmentId()); + + adapter.close(); + } + + /** + * Test that the environment ID from a full data set is retained and carried on the partial change + * sets produced for subsequent upserts, and on later full data sets that do not report one. + */ + @Test + public void environmentIdIsRetainedForSubsequentChangeSets() throws Exception { + AtomicReference capturedSink = new AtomicReference<>(); + + DataSourceSynchronizerAdapter adapter = new DataSourceSynchronizerAdapter(sink -> { + capturedSink.set(sink); + return new MockDataSource(new CountDownLatch(1), null); + }); + resourcesToClose.add(adapter); + + CompletableFuture initFuture = adapter.next(); + capturedSink.get().init(new FullDataSet<>(DataBuilder.forStandardTypes().build().getData(), true, + "env-from-fdv1")); + assertEquals("env-from-fdv1", initFuture.get(2, TimeUnit.SECONDS).getChangeSet().getEnvironmentId()); + + CompletableFuture upsertFuture = adapter.next(); + capturedSink.get().upsert(FEATURES, "flag1", new ItemDescriptor(1, flagBuilder("flag1").version(1).build())); + FDv2SourceResult upsertResult = upsertFuture.get(2, TimeUnit.SECONDS); + assertEquals(SourceResultType.CHANGE_SET, upsertResult.getResultType()); + assertEquals(ChangeSetType.Partial, upsertResult.getChangeSet().getType()); + assertEquals("env-from-fdv1", upsertResult.getChangeSet().getEnvironmentId()); + + CompletableFuture reinitFuture = adapter.next(); + capturedSink.get().init(new FullDataSet<>(DataBuilder.forStandardTypes().build().getData(), true, null)); + assertEquals("env-from-fdv1", reinitFuture.get(2, TimeUnit.SECONDS).getChangeSet().getEnvironmentId()); + + adapter.close(); + } } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceUpdatesImplTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceUpdatesImplTest.java index 5df085f7..fee4df7d 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceUpdatesImplTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/DataSourceUpdatesImplTest.java @@ -57,6 +57,8 @@ import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.greaterThanOrEqualTo; import static org.hamcrest.Matchers.is; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; @SuppressWarnings("javadoc") public class DataSourceUpdatesImplTest { @@ -163,10 +165,12 @@ private static ChangeSet private static class LegacyDataStore implements DataStore { private final Map> data = new HashMap<>(); + private String environmentId; @Override public void init(FullDataSet allData) { data.clear(); + environmentId = allData.getEnvironmentId(); // recorded verbatim; retention is a store's responsibility for (Map.Entry> kindEntry : allData.getData()) { DataKind kind = kindEntry.getKey(); Map items = new HashMap<>(); @@ -177,6 +181,11 @@ public void init(FullDataSet allData) { } } + @Override + public String getEnvironmentId() { + return environmentId; + } + @Override public boolean upsert(DataKind kind, String key, ItemDescriptor item) { Map items = data.get(kind); @@ -1034,9 +1043,14 @@ public void applyFullChangeSetToLegacyStoreWithEnvironmentId() throws Exception ); updates.apply(changeSet); - // Note: Java SDK doesn't have InitMetadata/EnvironmentId support in the same way as C#, - // so this test just verifies the changeset is applied without error ItemDescriptor retrievedFlag1 = legacyStore.get(FEATURES, flag1.getKey()); assertThat(retrievedFlag1, is(org.hamcrest.Matchers.notNullValue())); + assertEquals("test-env-id", legacyStore.getEnvironmentId()); + + // The environment ID is passed through exactly as it appears on the change set. Whether an absent + // value clears a previously retained ID is decided by the store (see InMemoryDataStoreTest and + // PersistentDataStoreWrapperTest), not here. + updates.apply(new ChangeSet<>(ChangeSetType.Full, Selector.make(2, "state2"), data, null, true)); + assertNull(legacyStore.getEnvironmentId()); } } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/EvaluatorWithHookTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/EvaluatorWithHookTest.java index e88625f5..6e2b5a15 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/EvaluatorWithHookTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/EvaluatorWithHookTest.java @@ -6,6 +6,7 @@ import com.launchdarkly.sdk.LDContext; import com.launchdarkly.sdk.LDValue; import com.launchdarkly.sdk.LDValueType; +import com.launchdarkly.sdk.server.integrations.EvaluationSeriesContext; import com.launchdarkly.sdk.server.integrations.Hook; import com.launchdarkly.sdk.server.integrations.HookMetadata; import org.junit.Test; @@ -20,6 +21,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; @@ -47,7 +49,7 @@ public void beforeIsExecutedBeforeAfter() { assertTrue(beforeCalled.get()); return Collections.emptyMap(); }); - EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none()); + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none(), () -> null); evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); } @@ -61,7 +63,7 @@ public void evaluationResultIsPassedToAfter() { when(mockHook.beforeEvaluation(any(), any())).thenReturn(Collections.emptyMap()); when(mockHook.afterEvaluation(any(), any(), any())).thenReturn(Collections.emptyMap()); - EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none()); + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none(), () -> null); evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); verify(mockHook).afterEvaluation(any(), any(), eq(EvaluationDetail.fromValue(LDValue.of("aValue"), 0, EvaluationReason.fallthrough()))); @@ -95,7 +97,7 @@ public void afterExecutesInReverseOrder() { return Collections.emptyMap(); }); - EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Arrays.asList(mockHookA, mockHookB), LDLogger.none()); + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Arrays.asList(mockHookA, mockHookB), LDLogger.none(), () -> null); evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); assertEquals(calls, Arrays.asList("hookABefore", "hookBBefore", "hookBAfter", "hookAAfter")); } @@ -110,7 +112,7 @@ public void beforeIsGivenEmptySeriesData() { when(mockHook.beforeEvaluation(any(), any())).thenReturn(Collections.emptyMap()); when(mockHook.afterEvaluation(any(), any(), eq(evalResult.getResult().getAnyType()))).thenReturn(Collections.emptyMap()); - EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none()); + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none(), () -> null); evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); verify(mockHook).beforeEvaluation(any(), eq(Collections.emptyMap())); @@ -128,7 +130,7 @@ public void seriesDataFromBeforeIsPassedToAfter() { when(mockHook.beforeEvaluation(any(), any())).thenReturn(mockData); when(mockHook.afterEvaluation(any(), any(), any())).thenReturn(Collections.emptyMap()); - EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none()); + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none(), () -> null); evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); verify(mockHook).afterEvaluation(any(), eq(mockData), any()); @@ -147,7 +149,7 @@ public void beforeThrowingErrorLeadsToEmptySeriesDataPassedToAfter() { when(mockHook.getMetadata()).thenReturn(new HookMetadata("mockHookName") {}); when(mockHook.afterEvaluation(any(), any(), any())).thenReturn(Collections.emptyMap()); - EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none()); + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), LDLogger.none(), () -> null); evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); verify(mockHook, times(1)).getMetadata(); @@ -186,11 +188,61 @@ public void oneHookThrowingErrorDoesNotAffectOtherHooks() { return Collections.emptyMap(); }); - EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Arrays.asList(mockHookA, mockHookB), LDLogger.none()); + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Arrays.asList(mockHookA, mockHookB), LDLogger.none(), () -> null); evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); assertEquals(calls, Arrays.asList("hookABefore", "hookBBefore", "hookBAfter", "hookAAfter")); verify(mockHookA).afterEvaluation(any(), eq(Collections.emptyMap()), any()); verify(mockHookB).afterEvaluation(any(), eq(mockData), any()); } + @Test + public void environmentIdIsProvidedToSeriesContext() { + EvalResultAndFlag evalResult = new EvalResultAndFlag(EvalResult.of(LDValue.of("aValue"), 0, EvaluationReason.fallthrough()), null); + EvaluatorInterface mockEvaluator = mock(EvaluatorInterface.class); + when(mockEvaluator.evalAndFlag(any(), any(), any(), any(), any(), any())).thenReturn(evalResult); + + Hook mockHook = mock(Hook.class); + List environmentIds = new ArrayList<>(); + when(mockHook.beforeEvaluation(any(), any())).thenAnswer((Answer>) invocation -> { + environmentIds.add(((EvaluationSeriesContext)invocation.getArgument(0)).environmentId); + return Collections.emptyMap(); + }); + when(mockHook.afterEvaluation(any(), any(), any())).thenAnswer((Answer>) invocation -> { + environmentIds.add(((EvaluationSeriesContext)invocation.getArgument(0)).environmentId); + return Collections.emptyMap(); + }); + + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), + LDLogger.none(), () -> "the-environment-id"); + evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); + + assertEquals(Arrays.asList("the-environment-id", "the-environment-id"), environmentIds); + } + + @Test + public void environmentIdSupplierExceptionDoesNotPreventEvaluation() { + EvalResultAndFlag evalResult = new EvalResultAndFlag(EvalResult.of(LDValue.of("aValue"), 0, EvaluationReason.fallthrough()), null); + EvaluatorInterface mockEvaluator = mock(EvaluatorInterface.class); + when(mockEvaluator.evalAndFlag(any(), any(), any(), any(), any(), any())).thenReturn(evalResult); + + Hook mockHook = mock(Hook.class); + when(mockHook.getMetadata()).thenReturn(new HookMetadata("mockHookName") {}); + List environmentIds = new ArrayList<>(); + when(mockHook.beforeEvaluation(any(), any())).thenAnswer((Answer>) invocation -> { + environmentIds.add(((EvaluationSeriesContext)invocation.getArgument(0)).environmentId); + return Collections.emptyMap(); + }); + when(mockHook.afterEvaluation(any(), any(), any())).thenAnswer((Answer>) invocation -> { + environmentIds.add(((EvaluationSeriesContext)invocation.getArgument(0)).environmentId); + return Collections.emptyMap(); + }); + + EvaluatorWithHooks evaluatorUnderTest = new EvaluatorWithHooks(mockEvaluator, Collections.singletonList(mockHook), + LDLogger.none(), () -> { throw new IllegalStateException("environment ID unavailable"); }); + EvalResultAndFlag result = evaluatorUnderTest.evalAndFlag("aMethod", "aKey", LDContext.create("aKey"), LDValue.of("aDefault"), LDValueType.STRING, EvaluationOptions.NO_EVENTS); + + // The evaluation still completes, and both stages of the hook run with an unknown environment ID. + assertSame(evalResult, result); + assertEquals(Arrays.asList(new String[] { null, null }), environmentIds); + } } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/InMemoryDataStoreTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/InMemoryDataStoreTest.java index 68b9dd95..f7b6e40b 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/InMemoryDataStoreTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/InMemoryDataStoreTest.java @@ -104,6 +104,44 @@ public void applyWithFullChangeSetReplacesAllData() { assertEquals(item3, result.getItem()); } + @Test + public void environmentIdIsRetainedFromDataAndNotClearedByLaterData() { + assertNull(typedStore().getEnvironmentId()); + + typedStore().init(new FullDataSet<>(ImmutableList.of(), true, "env-id")); + assertEquals("env-id", typedStore().getEnvironmentId()); + + typedStore().init(new FullDataSet<>(ImmutableList.of(), true, null)); + assertEquals("env-id", typedStore().getEnvironmentId()); + + typedStore().apply(new ChangeSet<>(ChangeSetType.Full, Selector.make(1, "state1"), ImmutableList.of(), "", true)); + assertEquals("env-id", typedStore().getEnvironmentId()); + + typedStore().apply(new ChangeSet<>(ChangeSetType.Full, Selector.make(2, "state2"), ImmutableList.of(), + "other-env-id", true)); + assertEquals("other-env-id", typedStore().getEnvironmentId()); + } + + @Test + public void environmentIdIsUpdatedByPartialChangeSet() { + typedStore().init(new FullDataSet<>(ImmutableList.of(), true, "env-id")); + + typedStore().apply(new ChangeSet<>(ChangeSetType.Partial, Selector.make(1, "state1"), ImmutableList.of(), + "other-env-id", true)); + assertEquals("other-env-id", typedStore().getEnvironmentId()); + + typedStore().apply(new ChangeSet<>(ChangeSetType.Partial, Selector.make(2, "state2"), ImmutableList.of(), + null, true)); + assertEquals("other-env-id", typedStore().getEnvironmentId()); + } + + @Test + public void exportAllIncludesEnvironmentId() { + typedStore().init(new FullDataSet<>(ImmutableList.of(), true, "env-id")); + + assertEquals("env-id", typedStore().exportAll().getEnvironmentId()); + } + @Test public void applyWithFullChangeSetSetsSelector() { Selector selector = Selector.make(42, "test-state"); diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/LDClientHooksEnvironmentIdTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/LDClientHooksEnvironmentIdTest.java new file mode 100644 index 00000000..cf74a352 --- /dev/null +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/LDClientHooksEnvironmentIdTest.java @@ -0,0 +1,343 @@ +package com.launchdarkly.sdk.server; + +import com.google.common.collect.ImmutableList; +import com.launchdarkly.sdk.EvaluationDetail; +import com.launchdarkly.sdk.LDContext; +import com.launchdarkly.sdk.LDValue; +import com.launchdarkly.sdk.server.DataStoreTestTypes.DataBuilder; +import com.launchdarkly.sdk.server.integrations.EvaluationSeriesContext; +import com.launchdarkly.sdk.server.integrations.Hook; +import com.launchdarkly.sdk.server.integrations.MockPersistentDataStore; +import com.launchdarkly.sdk.server.interfaces.DataStoreStatusProvider.CacheStats; +import com.launchdarkly.sdk.server.subsystems.DataStore; +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.DataKind; +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.FullDataSet; +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.ItemDescriptor; +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.KeyedItems; +import com.launchdarkly.testhelpers.httptest.Handler; +import com.launchdarkly.testhelpers.httptest.Handlers; +import com.launchdarkly.testhelpers.httptest.HttpServer; + +import org.junit.Test; + +import java.io.IOException; +import java.time.Duration; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CopyOnWriteArrayList; + +import static com.launchdarkly.sdk.server.DataModel.FEATURES; +import static com.launchdarkly.sdk.server.DataModel.SEGMENTS; +import static com.launchdarkly.sdk.server.ModelBuilders.flagBuilder; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +/** + * End-to-end tests verifying that the environment ID reported by LaunchDarkly reaches the + * {@link EvaluationSeriesContext} given to hooks, for each way the SDK can receive flag data. + */ +@SuppressWarnings("javadoc") +public class LDClientHooksEnvironmentIdTest extends BaseTest { + private static final String SDK_KEY = "sdk-key"; + private static final String ENV_ID_HEADER = "x-ld-envid"; + private static final String FLAG_KEY = "flag1"; + private static final LDContext CONTEXT = LDContext.create("user-key"); + private static final DataModel.FeatureFlag FLAG = flagBuilder(FLAG_KEY) + .version(1).on(false).offVariation(0).variations(LDValue.of(true), LDValue.of(false)) + .build(); + + /** + * Records the environment ID seen at each stage of the evaluation series. + */ + private static class RecordingHook extends Hook { + final List beforeEnvironmentIds = new CopyOnWriteArrayList<>(); + final List afterEnvironmentIds = new CopyOnWriteArrayList<>(); + + RecordingHook() { + super("recording"); + } + + @Override + public Map beforeEvaluation(EvaluationSeriesContext seriesContext, Map data) { + beforeEnvironmentIds.add(String.valueOf(seriesContext.environmentId)); + return data; + } + + @Override + public Map afterEvaluation(EvaluationSeriesContext seriesContext, Map data, + EvaluationDetail evaluationDetail) { + afterEnvironmentIds.add(String.valueOf(seriesContext.environmentId)); + return data; + } + } + + private static String fdv1DataJson() { + return new DataBuilder().addAny(FEATURES, FLAG).addAny(SEGMENTS).buildJson().toJsonString(); + } + + private static Handler fdv1StreamHandler(String environmentId) { + return Handlers.all( + Handlers.header(ENV_ID_HEADER, environmentId), + Handlers.SSE.start(), + Handlers.SSE.event("event: put\ndata: {\"data\":" + fdv1DataJson() + "}"), + Handlers.SSE.leaveOpen()); + } + + private static Handler fdv1PollHandler(String environmentId) { + return Handlers.all(Handlers.header(ENV_ID_HEADER, environmentId), Handlers.bodyJson(fdv1DataJson())); + } + + private static String fdv2PutObjectJson() { + return "{\"kind\":\"flag\",\"key\":\"" + FLAG_KEY + "\",\"version\":1,\"object\":" + JsonHelpers.serialize(FLAG) + "}"; + } + + private static Handler fdv2StreamHandler(String environmentId) { + return Handlers.all( + Handlers.header(ENV_ID_HEADER, environmentId), + Handlers.SSE.start(), + Handlers.SSE.event("event: server-intent\ndata: {\"payloads\":[{\"id\":\"payload-1\",\"target\":100," + + "\"intentCode\":\"xfer-full\",\"reason\":\"payload-missing\"}]}"), + Handlers.SSE.event("event: put-object\ndata: " + fdv2PutObjectJson()), + Handlers.SSE.event("event: payload-transferred\ndata: {\"state\":\"(p:payload-1:100)\",\"version\":100}"), + Handlers.SSE.leaveOpen()); + } + + private static Handler fdv2PollHandler(String environmentId) { + String body = "{\"events\":[" + + "{\"event\":\"server-intent\",\"data\":{\"payloads\":[{\"id\":\"payload-1\",\"target\":100," + + "\"intentCode\":\"xfer-full\",\"reason\":\"payload-missing\"}]}}," + + "{\"event\":\"put-object\",\"data\":" + fdv2PutObjectJson() + "}," + + "{\"event\":\"payload-transferred\",\"data\":{\"state\":\"(p:payload-1:100)\",\"version\":100}}" + + "]}"; + return Handlers.all(Handlers.header(ENV_ID_HEADER, environmentId), Handlers.bodyJson(body)); + } + + private LDConfig.Builder configWithHook(RecordingHook hook) { + return baseConfig() + .hooks(Components.hooks().setHooks(Collections.singletonList((Hook) hook))) + .startWait(Duration.ofSeconds(10)); + } + + private static void assertHookSawEnvironmentId(RecordingHook hook, String expected) { + assertEquals(ImmutableList.of(expected), hook.beforeEnvironmentIds); + assertEquals(ImmutableList.of(expected), hook.afterEnvironmentIds); + } + + @Test + public void hookReceivesEnvironmentIdFromFDv1Stream() throws Exception { + try (HttpServer server = HttpServer.start(fdv1StreamHandler("env-from-stream"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSource(Components.streamingDataSource()) + .serviceEndpoints(Components.serviceEndpoints().streaming(server.getUri())) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertTrue(client.isInitialized()); + client.boolVariation(FLAG_KEY, CONTEXT, false); + assertHookSawEnvironmentId(hook, "env-from-stream"); + } + } + } + + @Test + public void hookReceivesEnvironmentIdFromFDv1Polling() throws Exception { + try (HttpServer server = HttpServer.start(fdv1PollHandler("env-from-poll"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSource(Components.pollingDataSource().pollInterval(Duration.ofSeconds(300))) + .serviceEndpoints(Components.serviceEndpoints().polling(server.getUri())) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertTrue(client.isInitialized()); + client.boolVariation(FLAG_KEY, CONTEXT, false); + assertHookSawEnvironmentId(hook, "env-from-poll"); + } + } + } + + @Test + public void hookReceivesEnvironmentIdFromFDv2StreamingSynchronizer() throws Exception { + try (HttpServer server = HttpServer.start(fdv2StreamHandler("env-from-fdv2-stream"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSystem(Components.dataSystem().custom() + .synchronizers(DataSystemComponents.streamingSynchronizer() + .serviceEndpointsOverride(Components.serviceEndpoints().streaming(server.getUri())))) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertTrue(client.isInitialized()); + client.boolVariation(FLAG_KEY, CONTEXT, false); + assertHookSawEnvironmentId(hook, "env-from-fdv2-stream"); + } + } + } + + @Test + public void hookReceivesEnvironmentIdFromFDv2PollingSynchronizer() throws Exception { + try (HttpServer server = HttpServer.start(fdv2PollHandler("env-from-fdv2-poll"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSystem(Components.dataSystem().custom() + .synchronizers(DataSystemComponents.pollingSynchronizer() + .pollInterval(Duration.ofSeconds(300)) + .serviceEndpointsOverride(Components.serviceEndpoints().polling(server.getUri())))) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertTrue(client.isInitialized()); + client.boolVariation(FLAG_KEY, CONTEXT, false); + assertHookSawEnvironmentId(hook, "env-from-fdv2-poll"); + } + } + } + + @Test + public void hookReceivesEnvironmentIdFromFDv2PollingInitializer() throws Exception { + try (HttpServer server = HttpServer.start(fdv2PollHandler("env-from-fdv2-init"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSystem(Components.dataSystem().custom() + .initializers(DataSystemComponents.pollingInitializer() + .serviceEndpointsOverride(Components.serviceEndpoints().polling(server.getUri()))) + .synchronizers(DataSystemComponents.pollingSynchronizer() + .pollInterval(Duration.ofSeconds(300)) + .serviceEndpointsOverride(Components.serviceEndpoints().polling(server.getUri())))) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertTrue(client.isInitialized()); + client.boolVariation(FLAG_KEY, CONTEXT, false); + assertHookSawEnvironmentId(hook, "env-from-fdv2-init"); + } + } + } + + @Test + public void hookDoesNotReceiveEnvironmentIdWhenPersistentStoreInitFails() throws Exception { + // The environment ID is only known once data has been applied. With a finite cache TTL, a failed + // write to the persistent store means nothing was applied. + MockPersistentDataStore core = new MockPersistentDataStore(); + core.fakeError = new RuntimeException("store unavailable"); + try (HttpServer server = HttpServer.start(fdv1PollHandler("env-should-not-be-reported"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSource(Components.pollingDataSource().pollInterval(Duration.ofSeconds(300))) + .dataStore(Components.persistentDataStore(ctx -> core).cacheTime(Duration.ofSeconds(30))) + .serviceEndpoints(Components.serviceEndpoints().polling(server.getUri())) + .startWait(Duration.ofMillis(500)) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertFalse(client.isInitialized()); + client.boolVariation(FLAG_KEY, CONTEXT, false); + assertHookSawEnvironmentId(hook, "null"); + } + } + } + + @Test + public void hookDoesNotReceiveEnvironmentIdFromDataStoreThatDoesNotProvideOne() throws Exception { + // A custom DataStore that does not override getEnvironmentId() reports no environment ID, even + // though the SDK received one and passed it to the store with the data. + SimpleDataStore store = new SimpleDataStore(); + try (HttpServer server = HttpServer.start(fdv1PollHandler("env-from-poll"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSource(Components.pollingDataSource().pollInterval(Duration.ofSeconds(300))) + .dataStore(ctx -> store) + .serviceEndpoints(Components.serviceEndpoints().polling(server.getUri())) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertTrue(client.isInitialized()); + client.boolVariation(FLAG_KEY, CONTEXT, false); + assertEquals("env-from-poll", store.lastInit.getEnvironmentId()); + assertHookSawEnvironmentId(hook, "null"); + } + } + } + + @Test + public void evaluationSucceedsWhenDataStoreEnvironmentIdThrows() throws Exception { + // A failure while looking up the environment ID must not turn into an exception from the + // variation methods, and must not prevent the hooks from running. + SimpleDataStore store = new SimpleDataStore() { + @Override + public String getEnvironmentId() { + throw new IllegalStateException("environment ID unavailable"); + } + }; + try (HttpServer server = HttpServer.start(fdv1PollHandler("env-from-poll"))) { + RecordingHook hook = new RecordingHook(); + LDConfig config = configWithHook(hook) + .dataSource(Components.pollingDataSource().pollInterval(Duration.ofSeconds(300))) + .dataStore(ctx -> store) + .serviceEndpoints(Components.serviceEndpoints().polling(server.getUri())) + .build(); + try (LDClient client = new LDClient(SDK_KEY, config)) { + assertTrue(client.isInitialized()); + assertTrue(client.boolVariation(FLAG_KEY, CONTEXT, false)); + assertHookSawEnvironmentId(hook, "null"); + } + } + } + + /** + * A minimal DataStore implementation that relies on the default getEnvironmentId(). + */ + private static class SimpleDataStore implements DataStore { + private final Map> data = new HashMap<>(); + volatile FullDataSet lastInit; + volatile boolean initialized; + + @Override + public void init(FullDataSet allData) { + lastInit = allData; + data.clear(); + for (Map.Entry> kindEntry : allData.getData()) { + Map items = new HashMap<>(); + for (Map.Entry itemEntry : kindEntry.getValue().getItems()) { + items.put(itemEntry.getKey(), itemEntry.getValue()); + } + data.put(kindEntry.getKey(), items); + } + initialized = true; + } + + @Override + public ItemDescriptor get(DataKind kind, String key) { + Map items = data.get(kind); + return items == null ? null : items.get(key); + } + + @Override + public KeyedItems getAll(DataKind kind) { + Map items = data.get(kind); + return items == null ? new KeyedItems<>(null) : new KeyedItems<>(ImmutableList.copyOf(items.entrySet())); + } + + @Override + public boolean upsert(DataKind kind, String key, ItemDescriptor item) { + data.computeIfAbsent(kind, k -> new HashMap<>()).put(key, item); + return true; + } + + @Override + public boolean isInitialized() { + return initialized; + } + + @Override + public boolean isStatusMonitoringEnabled() { + return false; + } + + @Override + public CacheStats getCacheStats() { + return null; + } + + @Override + public void close() throws IOException { + } + } +} diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreConverterTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreConverterTest.java index f9529b41..714fbcc7 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreConverterTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreConverterTest.java @@ -468,4 +468,18 @@ public void toSerializedFormatWithMixedDeletedAndRegularItems() { assertEquals(2, deletedCount); assertEquals(2, regularCount); } + + @Test + public void toSerializedFormatPreservesEnvironmentId() { + FullDataSet data = new FullDataSet<>( + ImmutableList.of(new AbstractMap.SimpleEntry<>(TEST_DATA_KIND, + new KeyedItems<>(ImmutableList.of(new AbstractMap.SimpleEntry<>("key1", + new ItemDescriptor(1, new TestItem("item1"))))))), + true, "env-id"); + + FullDataSet result = PersistentDataStoreConverter.toSerializedFormat(data); + + assertEquals("env-id", result.getEnvironmentId()); + assertTrue(result.shouldPersist()); + } } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapperTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapperTest.java index 1c5850e6..58d635aa 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapperTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PersistentDataStoreWrapperTest.java @@ -838,4 +838,63 @@ private void makeStoreAvailable(MockPersistentDataStore core) { core.fakeError = null; core.unavailable = false; } + + @Test + public void environmentIdIsRetainedFromSuccessfulInit() { + assertNull(wrapper.getEnvironmentId()); + + wrapper.init(makeDataSetWithEnvironmentId("env-1")); + assertThat(wrapper.getEnvironmentId(), equalTo("env-1")); + + // absent or empty values never clear a retained ID + wrapper.init(makeDataSetWithEnvironmentId(null)); + assertThat(wrapper.getEnvironmentId(), equalTo("env-1")); + wrapper.init(makeDataSetWithEnvironmentId("")); + assertThat(wrapper.getEnvironmentId(), equalTo("env-1")); + + wrapper.init(makeDataSetWithEnvironmentId("env-2")); + assertThat(wrapper.getEnvironmentId(), equalTo("env-2")); + } + + @Test + public void environmentIdIsNotRetainedWhenInitFailsAndCacheIsNotIndefinite() { + assumeThat(testMode.isCachedIndefinitely(), is(false)); + + core.fakeError = FAKE_ERROR; + try { + wrapper.init(makeDataSetWithEnvironmentId("env-1")); + fail("expected exception"); + } catch (RuntimeException e) { + assertThat(e, is(FAKE_ERROR)); + } + core.fakeError = null; + + // The data that carried this ID was never applied, so hooks must not be told about it. + assertNull(wrapper.getEnvironmentId()); + assertThat(wrapper.isInitialized(), is(false)); + } + + @Test + public void environmentIdIsRetainedWhenInitFailsButCacheIsIndefinite() { + assumeThat(testMode.isCachedIndefinitely(), is(true)); + + core.fakeError = FAKE_ERROR; + try { + wrapper.init(makeDataSetWithEnvironmentId("env-1")); + fail("expected exception"); + } catch (RuntimeException e) { + assertThat(e, is(FAKE_ERROR)); + } + core.fakeError = null; + + // With an infinite cache TTL the new data is served from the cache even though the underlying store + // failed, so the environment ID that came with it is retained too. + assertThat(wrapper.getEnvironmentId(), equalTo("env-1")); + assertThat(wrapper.isInitialized(), is(true)); + } + + private static FullDataSet makeDataSetWithEnvironmentId(String environmentId) { + return new FullDataSet<>(new DataBuilder().add(TEST_ITEMS, new TestItem("key", 1)).build().getData(), + true, environmentId); + } } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PollingProcessorTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PollingProcessorTest.java index e74f3c25..d9c87cbc 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PollingProcessorTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/PollingProcessorTest.java @@ -37,6 +37,7 @@ import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import static com.launchdarkly.sdk.server.TestComponents.clientContext; import static com.launchdarkly.sdk.server.TestComponents.dataStoreThatThrowsException; @@ -53,6 +54,7 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; @@ -65,11 +67,12 @@ public class PollingProcessorTest extends BaseTest { private static final long EXTENDED_GAP_FLOOR_MILLIS = 150L; private MockDataSourceUpdates dataSourceUpdates; + private InMemoryDataStore dataStore; @Before public void setup() { - DataStore store = new InMemoryDataStore(); - dataSourceUpdates = TestComponents.dataSourceUpdates(store, new MockDataStoreStatusProvider()); + dataStore = new InMemoryDataStore(); + dataSourceUpdates = TestComponents.dataSourceUpdates(dataStore, new MockDataStoreStatusProvider()); } private PollingProcessor makeProcessor(URI baseUri, Duration pollInterval) { @@ -486,4 +489,78 @@ private void withStatusQueue(ActionCanThrowAnyException { + int err = errorStatus.get(); + if (err == 0) { + Handlers.header("x-ld-envid", "env-from-poll").apply(ctx); + successHandler.apply(ctx); + } else { + Handlers.header("x-ld-envid", "env-from-error").apply(ctx); + ctx.setStatus(err); + } + }; + + withStatusQueue(statuses -> { + try (HttpServer server = HttpServer.start(pollingHandler)) { + try (PollingProcessor pollingProcessor = makeProcessor(server.getUri(), BRIEF_INTERVAL)) { + pollingProcessor.start(); + + // the first poll fails: nothing is taken from the error response + Status status0 = requireDataSourceStatus(statuses, State.INITIALIZING); + assertEquals(ErrorKind.ERROR_RESPONSE, status0.getLastError().getKind()); + assertNull(dataStore.getEnvironmentId()); + + // now a poll succeeds + errorStatus.set(0); + requireDataSourceStatusEventually(statuses, State.VALID, State.INITIALIZING); + assertEquals("env-from-poll", dataStore.getEnvironmentId()); + + // subsequent failures do not change the retained value + errorStatus.set(503); + requireDataSourceStatus(statuses, State.INTERRUPTED); + assertEquals("env-from-poll", dataStore.getEnvironmentId()); + } + } + }); + } } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java index b81a50b3..1db21b28 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java @@ -74,6 +74,7 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; @@ -1150,6 +1151,39 @@ private void assertFeatureInStore(DataModel.FeatureFlag feature) { assertEquals(feature.getVersion(), dataStore.get(FEATURES, feature.getKey()).getVersion()); } + @Test + public void environmentIdIsCapturedFromStreamResponseHeader() throws Exception { + Handler streamHandler = Handlers.all( + Handlers.header("x-ld-envid", "env-from-stream"), + Handlers.SSE.start(), + Handlers.SSE.event(EMPTY_DATA_EVENT), + Handlers.SSE.leaveOpen() + ); + + try (HttpServer server = HttpServer.start(streamHandler)) { + try (StreamProcessor sp = createStreamProcessor(null, server.getUri())) { + assertNull(dataStore.getEnvironmentId()); + + sp.start(); + dataSourceUpdates.awaitInit(); + + assertEquals("env-from-stream", dataStore.getEnvironmentId()); + } + } + } + + @Test + public void environmentIdIsNullWhenStreamResponseHasNoHeader() throws Exception { + try (HttpServer server = HttpServer.start(streamResponse(EMPTY_DATA_EVENT))) { + try (StreamProcessor sp = createStreamProcessor(null, server.getUri())) { + sp.start(); + dataSourceUpdates.awaitInit(); + + assertNull(dataStore.getEnvironmentId()); + } + } + } + private void assertSegmentInStore(DataModel.Segment segment) { assertEquals(segment.getVersion(), dataStore.get(SEGMENTS, segment.getKey()).getVersion()); } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/WriteThroughStoreTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/WriteThroughStoreTest.java index 1a89c4db..33e9696c 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/WriteThroughStoreTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/WriteThroughStoreTest.java @@ -54,6 +54,10 @@ private FullDataSet createTestDataSet() { } private ChangeSet>>> createFullChangeSet() { + return createFullChangeSet(null); + } + + private ChangeSet>>> createFullChangeSet(String environmentId) { Map> changeSetData = ImmutableMap.of( TEST_ITEMS, new KeyedItems<>(ImmutableList.of( @@ -66,7 +70,7 @@ private ChangeSet>>> cre ChangeSetType.Full, Selector.make(1, "state1"), changeSetData.entrySet(), - null, + environmentId, true ); } @@ -957,6 +961,72 @@ public void applySwitchesToMemoryStoreEvenWhenTransactionalStoreApplyFails() thr // Mock Stores + // Environment ID Tests + + @Test + public void environmentIdIsNullBeforeInitializingPayloadEvenIfPersistentStoreIsInitialized() { + MockPersistentStore persistentStore = new MockPersistentStore(); + persistentStore.setData(TEST_ITEMS, "key1", new ItemDescriptor(10, item1)); + persistentStore.setInitialized(true); + store = new WriteThroughStore(new InMemoryDataStore(), persistentStore, DataStoreMode.READ_WRITE); + + // Reads are served by the persistent store, which cannot know the environment ID. + assertTrue(store.isInitialized()); + assertNotNull(store.get(TEST_ITEMS, "key1")); + assertNull(store.getEnvironmentId()); + } + + @Test + public void environmentIdFromInitIsExposedAndForwardedToPersistentStore() { + MockPersistentStore persistentStore = new MockPersistentStore(); + store = new WriteThroughStore(new InMemoryDataStore(), persistentStore, DataStoreMode.READ_WRITE); + + store.init(new FullDataSet<>(createTestDataSet().getData(), true, "env-from-init")); + + assertEquals("env-from-init", store.getEnvironmentId()); + assertEquals("env-from-init", persistentStore.lastInit.getEnvironmentId()); + } + + @Test + public void environmentIdFromFullChangeSetIsExposedAndForwardedToLegacyPersistentStore() { + MockPersistentStore persistentStore = new MockPersistentStore(); + store = new WriteThroughStore(new InMemoryDataStore(), persistentStore, DataStoreMode.READ_WRITE); + + store.apply(createFullChangeSet("env-from-change-set")); + + assertEquals("env-from-change-set", store.getEnvironmentId()); + assertEquals("env-from-change-set", persistentStore.lastInit.getEnvironmentId()); + } + + @Test + public void environmentIdIsExposedButNotWrittenToPersistentStoreInReadOnlyMode() { + MockPersistentStore persistentStore = new MockPersistentStore(); + store = new WriteThroughStore(new InMemoryDataStore(), persistentStore, DataStoreMode.READ_ONLY); + + store.apply(createFullChangeSet("env-from-change-set")); + + assertEquals("env-from-change-set", store.getEnvironmentId()); + assertFalse(persistentStore.wasInitCalled); + } + + @Test + public void environmentIdMatchesMemoryDataWhenPersistentWriteFails() { + MockPersistentStore persistentStore = new MockPersistentStore(); + persistentStore.throwOnInit = true; + store = new WriteThroughStore(new InMemoryDataStore(), persistentStore, DataStoreMode.READ_WRITE); + + try { + store.apply(createFullChangeSet("env-from-change-set")); + fail("expected exception"); + } catch (RuntimeException e) { + // expected: the persistent write failed + } + + // Reads now come from the memory store, which holds the data that carried this ID. + assertNotNull(store.get(TEST_ITEMS, "key1")); + assertEquals("env-from-change-set", store.getEnvironmentId()); + } + private static class MockPersistentStore implements DataStore { private final Map> data = new HashMap<>(); private final Set keysToFailOn = new HashSet<>(); @@ -970,6 +1040,7 @@ private static class MockPersistentStore implements DataStore { public boolean failUpsert; public boolean throwOnInit; public boolean statusMonitoringEnabledValue; + public FullDataSet lastInit; public void setUpsertFailureForKey(String key) { keysToFailOn.add(key); @@ -994,6 +1065,7 @@ public void setInitialized(boolean value) { @Override public void init(FullDataSet allData) { wasInitCalled = true; + lastInit = allData; if (throwOnInit) { throw new RuntimeException("Init failed"); } diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/interfaces/DataStoreTypesTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/interfaces/DataStoreTypesTest.java index 7257c68d..9333e97d 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/interfaces/DataStoreTypesTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/interfaces/DataStoreTypesTest.java @@ -153,10 +153,14 @@ public void fullDataSetEquality() { List>> allPermutations = new ArrayList<>(); for (DataKind kind: new DataKind[] { DataModel.FEATURES, DataModel.SEGMENTS }) { for (int version: new int[] { 1, 2 }) { - allPermutations.add(() -> new FullDataSet<>( - ImmutableMap.of(kind, - new KeyedItems<>(ImmutableMap.of("key", new ItemDescriptor(version, "a")).entrySet()) - ).entrySet(), true)); + for (boolean shouldPersist: new boolean[] { true, false }) { + for (String environmentId: new String[] { null, "env-1", "env-2" }) { + allPermutations.add(() -> new FullDataSet<>( + ImmutableMap.of(kind, + new KeyedItems<>(ImmutableMap.of("key", new ItemDescriptor(version, "a")).entrySet()) + ).entrySet(), shouldPersist, environmentId)); + } + } } } TypeBehavior.checkEqualsAndHashCode(allPermutations);