diff --git a/README.md b/README.md index cc4b567..0de3443 100644 --- a/README.md +++ b/README.md @@ -500,6 +500,22 @@ When running in Replicator Mode, the client will: - Use cached data populated by the external replicator service - Fall back to direct API calls if the replicator is not available +### Cache Readiness + +The SDK serves flag checks from the replicator cache only once the replicator reports that its cache is ready. The SDK polls the health URL and reads `ready` and `cache_version` from the JSON body, including on the 503 the replicator returns while its cache is still loading. Until the replicator reports `ready: true`, `checkFlag`, `checkFlagWithEntitlement` and `checkFlags` skip the cache and call the Schematic API instead. If the API call fails, they return the flag default. Once the cache is ready, flag checks evaluate locally from the cache, and a flag missing from the cache still falls back to the API. Single and bulk flag checks follow the same rule. + +If the health URL can't be reached, times out, or returns a body that isn't JSON, the SDK treats the cache as not ready and keeps the last cache version it saw. + +`isCacheReady()` reports the same readiness the flag checks use: + +```java +if (schematic.isCacheReady()) { + // Flag checks are evaluated locally from the replicator cache +} +``` + +`isDatastreamConnected()` returns the same value in replicator mode and is kept for backward compatibility. Outside replicator mode `isCacheReady()` returns true whenever datastream is configured, since the SDK fills its own cache over the WebSocket. + ## Contributing While we value open-source contributions to this SDK, this library is generated programmatically. diff --git a/src/main/java/com/schematic/api/Schematic.java b/src/main/java/com/schematic/api/Schematic.java index 2c426d5..7ba9113 100644 --- a/src/main/java/com/schematic/api/Schematic.java +++ b/src/main/java/com/schematic/api/Schematic.java @@ -290,16 +290,35 @@ public boolean isReplicatorMode() { /** * Returns whether the datastream connection is active and ready. + * + *

In replicator mode this reports replicator readiness, the same value + * {@link #isCacheReady()} returns. Use {@link #isCacheReady()} to ask whether flag checks + * are being served from the replicator cache. */ public boolean isDatastreamConnected() { return this.dataStreamClient != null && this.dataStreamClient.isConnected(); } + /** + * Returns whether flag checks may be evaluated from the datastream cache. + * + *

In replicator mode this is true only once the replicator reports its cache is ready. + * Until then, single and bulk flag checks skip the cache and use the Schematic API. + * Outside replicator mode there is nothing to wait for and this returns true whenever + * datastream is configured (flag checks there still require + * {@link #isDatastreamConnected()}, as before). Returns false when datastream is not + * configured. See {@link DataStreamClient#isCacheReady()}. + */ + public boolean isCacheReady() { + return this.dataStreamClient != null && this.dataStreamClient.isCacheReady(); + } + /** * Checks a feature flag, returning a boolean value. * - *

If datastream is configured and connected, evaluates the flag locally using cached - * data and the rules engine. Falls back to the API if datastream is unavailable or + *

If datastream is configured and connected and its cache is ready (in replicator mode, + * once the replicator reports ready; see {@link #isCacheReady()}), evaluates the flag + * locally using cached data and the rules engine. Falls back to the API otherwise, or if * evaluation fails. */ public boolean checkFlag(String flagKey, Map company, Map user) { @@ -312,7 +331,8 @@ public boolean checkFlag(String flagKey, Map company, MapPriority order: *

    - *
  1. DataStream evaluation (if configured and connected)
  2. + *
  3. DataStream evaluation (if configured and connected and the cache is ready, see + * {@link #isCacheReady()})
  4. *
  5. API call with result caching (fallback)
  6. *
  7. Flag default value (if all else fails)
  8. *
@@ -341,15 +361,28 @@ private RulesengineCheckFlagResult defaultFlagResult(String flagKey, String reas .build(); } + /** + * Whether a flag check should be evaluated from the datastream cache. Single + * ({@link #checkFlagWithEntitlement}) and bulk ({@link #checkFlags}) flag checks both ask + * this, so they cannot drift apart. It is the existing datastream check (configured and + * connected, which leaves WebSocket mode unchanged) plus {@link #isCacheReady()}, so in + * replicator mode the cache is read only once the replicator reports it ready; until then + * flag checks take the API path. + */ + private boolean useDataStreamCache() { + return dataStreamClient != null && dataStreamClient.isConnected() && dataStreamClient.isCacheReady(); + } + /** * Attempts to evaluate a flag via the datastream client. Returns the result on - * success, or {@code null} if datastream is not configured/connected or evaluation - * failed. Callers are responsible for emitting a {@code flag_check} event when - * appropriate — single-flag checks do, bulk checks do not. + * success, or {@code null} if the cache should not be used (see + * {@link #useDataStreamCache()}) or evaluation failed. Callers are responsible for + * emitting a {@code flag_check} event when appropriate. Single-flag checks do, bulk + * checks do not. */ private RulesengineCheckFlagResult tryDatastreamCheckFlag( String flagKey, Map company, Map user) { - if (dataStreamClient == null || !dataStreamClient.isConnected()) { + if (!useDataStreamCache()) { return null; } try { @@ -415,8 +448,9 @@ private void enqueueFlagCheckEvent( *

Evaluation order: *

    *
  1. Offline mode → return flag defaults for the requested keys
  2. - *
  3. DataStream / replicator (if configured and connected) → evaluate each key - * locally; falls back to the API if any key fails
  4. + *
  5. DataStream / replicator (if configured and connected and the cache is ready, + * see {@link #isCacheReady()}) → evaluate each key locally; falls back to the API + * if any key fails
  6. *
  7. Otherwise → look up each requested key in the result cache; if any are * missing, issue a single bulk {@code features.checkFlags} API call to fetch * fresh values, refresh the cache, and merge the results
  8. @@ -440,7 +474,7 @@ public List checkFlags( } // 2. DataStream/replicator path: evaluate each key; on any failure fall back to API. - if (dataStreamClient != null && dataStreamClient.isConnected() && flagKeys != null && !flagKeys.isEmpty()) { + if (useDataStreamCache() && flagKeys != null && !flagKeys.isEmpty()) { List dsResults = new ArrayList<>(flagKeys.size()); boolean dsOk = true; for (String key : flagKeys) { diff --git a/src/main/java/com/schematic/api/datastream/DataStreamClient.java b/src/main/java/com/schematic/api/datastream/DataStreamClient.java index 1390106..3d51e94 100644 --- a/src/main/java/com/schematic/api/datastream/DataStreamClient.java +++ b/src/main/java/com/schematic/api/datastream/DataStreamClient.java @@ -148,6 +148,12 @@ public void start() { /** * Returns whether the datastream is connected and ready for flag checks. + * + *

    In direct WebSocket mode this reports whether the WebSocket is connected and + * initialized. In replicator mode it reports replicator readiness (the {@code ready} + * field of the replicator's health response), which is the same value + * {@link #isCacheReady()} returns. Prefer {@link #isCacheReady()} when the question is + * whether flag checks can be served from the replicator cache. */ public boolean isConnected() { if (options.isReplicatorMode()) { @@ -156,6 +162,26 @@ public boolean isConnected() { return wsClient != null && wsClient.isReady(); } + /** + * Returns whether flag checks may be evaluated from the cache. + * + *

    In replicator mode this is the replicator's readiness from its last health poll: true + * once the replicator reports {@code ready: true}, meaning its cache is complete for the + * current cache version, and false before that or after a failed poll (connection error, + * timeout, unparseable body). While it is false, flag checks do not read the cache and + * are answered by the Schematic API. + * + *

    Outside replicator mode the SDK fills its own cache over the WebSocket and fetches + * what it lacks on demand, so there is nothing to wait for and this returns true. Flag + * checks in that mode still require {@link #isConnected()}, as before. + */ + public boolean isCacheReady() { + if (options.isReplicatorMode()) { + return replicatorReady.get(); + } + return true; + } + /** * Returns whether this client is running in replicator mode. */ @@ -631,6 +657,15 @@ private void startReplicatorMode() { this::checkReplicatorHealth, 0, intervalMs, TimeUnit.MILLISECONDS); } + /** + * Polls the replicator health endpoint once. + * + *

    The JSON body is read whatever the HTTP status: while its cache is still loading the + * replicator answers 503 with {@code ready: false} and a {@code cache_version}. Readiness + * comes from the {@code ready} field, and the cache version is recorded from any response + * that carries a non-empty one. A failed poll (connection error, timeout, unparseable + * body) sets not ready and keeps the last known cache version. + */ void checkReplicatorHealth() { try { Request request = new Request.Builder() @@ -639,46 +674,84 @@ void checkReplicatorHealth() { .build(); try (Response response = httpClient.newCall(request).execute()) { - if (response.isSuccessful() && response.body() != null) { - JsonNode body = objectMapper.readTree(response.body().string()); - boolean ready = body.has("ready") && body.get("ready").asBoolean(false); - boolean wasReady = replicatorReady.getAndSet(ready); - - String newCacheVersion = null; - if (body.has("cache_version")) { - newCacheVersion = body.get("cache_version").asText(); - } else if (body.has("cacheVersion")) { - newCacheVersion = body.get("cacheVersion").asText(); - } - if (newCacheVersion != null && !newCacheVersion.equals(replicatorCacheVersion)) { - String oldVersion = replicatorCacheVersion; - replicatorCacheVersion = newCacheVersion; - log( - "info", - "Replicator cache version changed from " - + (oldVersion == null ? "(null)" : oldVersion) + " to " - + newCacheVersion); - } + JsonNode body = readHealthBody(response); + if (body == null) { + setReplicatorNotReady("unparseable health response (status " + response.code() + ")"); + return; + } - if (ready && !wasReady) { - log("info", "Replicator is now ready"); - } else if (!ready && wasReady) { - log("warn", "Replicator is no longer ready"); - } - } else { - boolean wasReady = replicatorReady.getAndSet(false); - if (wasReady) { - log("warn", "Replicator health check failed with status: " + response.code()); - } + updateReplicatorCacheVersion(body); + + JsonNode readyNode = body.path("ready"); + boolean ready = readyNode.isBoolean() && readyNode.booleanValue(); + boolean wasReady = replicatorReady.getAndSet(ready); + + if (ready && !wasReady) { + log("info", "Replicator is now ready (cache_version: " + replicatorCacheVersion + ")"); + } else if (!ready && wasReady) { + log("warn", "Replicator is no longer ready (status: " + response.code() + ")"); } } + } catch (Exception e) { + // Catch everything, not just IOException: an exception escaping this method would + // cancel the scheduled health check and leave readiness stuck at its last value. + setReplicatorNotReady(e.getMessage()); + } + } + + /** + * Reads the health response body as a JSON object, or returns null if there is no body + * or it is not a JSON object. + */ + private JsonNode readHealthBody(Response response) { + if (response.body() == null) { + return null; + } + try { + JsonNode body = objectMapper.readTree(response.body().string()); + return body != null && body.isObject() ? body : null; } catch (IOException e) { - boolean wasReady = replicatorReady.getAndSet(false); - if (wasReady) { - log("warn", "Replicator health check failed: " + e.getMessage()); - } - log("debug", "Replicator health check error: " + e.getMessage()); + log("debug", "Failed to parse replicator health response: " + e.getMessage()); + return null; + } + } + + /** + * Records the cache version from a health response body. Keeps the last known version + * when the body reports none. + */ + private void updateReplicatorCacheVersion(JsonNode body) { + String newCacheVersion = null; + if (body.hasNonNull("cache_version")) { + newCacheVersion = body.get("cache_version").asText(); + } else if (body.hasNonNull("cacheVersion")) { + newCacheVersion = body.get("cacheVersion").asText(); + } + if (newCacheVersion != null && !newCacheVersion.isEmpty() && !newCacheVersion.equals(replicatorCacheVersion)) { + String oldVersion = replicatorCacheVersion; + replicatorCacheVersion = newCacheVersion; + log( + "info", + "Replicator cache version changed from " + + (oldVersion == null ? "(null)" : oldVersion) + " to " + + newCacheVersion); + } + } + + private void setReplicatorNotReady(String reason) { + boolean wasReady = replicatorReady.getAndSet(false); + if (wasReady) { + log("warn", "Replicator health check failed: " + reason); } + log("debug", "Replicator health check error: " + reason); + } + + /** + * Returns the last cache version reported by the replicator, or null if none has been + * reported. Package-private for tests. + */ + String getReplicatorCacheVersion() { + return replicatorCacheVersion; } // --- Message handling --- diff --git a/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java b/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java new file mode 100644 index 0000000..7c12583 --- /dev/null +++ b/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java @@ -0,0 +1,478 @@ +package com.schematic.api.datastream; + +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.*; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.schematic.api.Schematic; +import com.schematic.api.cache.LocalCache; +import com.schematic.api.datastream.DataStreamMessages.DataStreamResp; +import com.schematic.api.datastream.DataStreamMessages.EntityType; +import com.schematic.api.datastream.DataStreamMessages.MessageType; +import com.schematic.api.logger.SchematicLogger; +import com.schematic.api.resources.features.FeaturesClient; +import com.schematic.api.resources.features.types.CheckFlagResponse; +import com.schematic.api.resources.features.types.CheckFlagsResponse; +import com.schematic.api.types.CheckFlagRequestBody; +import com.schematic.api.types.CheckFlagResponseData; +import com.schematic.api.types.CheckFlagsResponseData; +import com.schematic.api.types.RulesengineCheckFlagResult; +import com.schematic.api.types.RulesengineFlag; +import com.sun.net.httpserver.HttpServer; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +/** + * Replicator mode serves flag checks from the cache only once the replicator reports ready. + * Before that, single and bulk flag checks both skip the cache and use the API. + */ +@ExtendWith(MockitoExtension.class) +class ReplicatorCacheReadyTest { + + private static final String CACHE_VERSION = "v1"; + private static final Map COMPANY = Collections.singletonMap("customer_id", "cust-1"); + private static final List KEYS = Arrays.asList("flag-on", "flag-off"); + + @Mock + private SchematicLogger logger; + + private final ObjectMapper objectMapper = new ObjectMapper(); + private final AtomicInteger healthStatus = new AtomicInteger(200); + private final AtomicReference healthBody = new AtomicReference<>("{}"); + private HttpServer server; + private String healthUrl; + private LocalCache flagCache; + private Schematic schematic; + + @BeforeEach + void setUp() throws Exception { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/ready", exchange -> { + byte[] bytes = healthBody.get().getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().add("Content-Type", "application/json"); + exchange.sendResponseHeaders(healthStatus.get(), bytes.length); + try (OutputStream os = exchange.getResponseBody()) { + os.write(bytes); + } + }); + server.start(); + healthUrl = "http://127.0.0.1:" + server.getAddress().getPort() + "/ready"; + flagCache = new LocalCache<>(); + } + + @AfterEach + void tearDown() { + if (schematic != null) { + schematic.close(); + } + if (server != null) { + server.stop(0); + } + } + + // --- Flag checks, cache not ready --- + + @Test + void notReady_singleAndBulkSkipCacheAndUseApi() { + setHealth(503, false, CACHE_VERSION); + Schematic spySchematic = spy(buildSchematic(Collections.emptyMap())); + assertFalse(spySchematic.isCacheReady()); + assertEquals(CACHE_VERSION, schematic.getDataStreamClient().getReplicatorCacheVersion()); + seedReplicatorCache(schematic.getDataStreamClient()); + + // The API answers the opposite of the cache, so each result shows which one was read. + FeaturesClient featuresClient = mock(FeaturesClient.class); + when(spySchematic.features()).thenReturn(featuresClient); + when(featuresClient.checkFlag(eq("flag-on"), any(CheckFlagRequestBody.class))) + .thenReturn(singleApiResponse("flag-on", false)); + when(featuresClient.checkFlag(eq("flag-off"), any(CheckFlagRequestBody.class))) + .thenReturn(singleApiResponse("flag-off", true)); + when(featuresClient.checkFlags(any(CheckFlagRequestBody.class))) + .thenReturn(bulkApiResponse(Arrays.asList("flag-on", "flag-off"), Arrays.asList(false, true))); + + RulesengineCheckFlagResult singleOn = spySchematic.checkFlagWithEntitlement("flag-on", COMPANY, null); + RulesengineCheckFlagResult singleOff = spySchematic.checkFlagWithEntitlement("flag-off", COMPANY, null); + List bulk = spySchematic.checkFlags(KEYS, COMPANY, null); + + assertFalse(singleOn.getValue()); + assertTrue(singleOff.getValue()); + assertEquals("api", singleOn.getReason()); + assertEquals(2, bulk.size()); + assertEquals("flag-on", bulk.get(0).getFlagKey()); + assertFalse(bulk.get(0).getValue()); + assertEquals("flag-off", bulk.get(1).getFlagKey()); + assertTrue(bulk.get(1).getValue()); + assertEquals("api", bulk.get(0).getReason()); + + verify(featuresClient).checkFlag(eq("flag-on"), any(CheckFlagRequestBody.class)); + verify(featuresClient).checkFlag(eq("flag-off"), any(CheckFlagRequestBody.class)); + verify(featuresClient).checkFlags(any(CheckFlagRequestBody.class)); + } + + @Test + void notReady_apiFails_singleAndBulkReturnFlagDefaults() { + setHealth(503, false, CACHE_VERSION); + // SDK flag defaults are the opposite of the cached flags' default values, so a + // result read from the cache would not match. + Map flagDefaults = new HashMap<>(); + flagDefaults.put("flag-on", false); + flagDefaults.put("flag-off", true); + Schematic spySchematic = spy(buildSchematic(flagDefaults)); + assertFalse(spySchematic.isCacheReady()); + seedReplicatorCache(schematic.getDataStreamClient()); + + FeaturesClient featuresClient = mock(FeaturesClient.class); + when(spySchematic.features()).thenReturn(featuresClient); + when(featuresClient.checkFlag(any(String.class), any(CheckFlagRequestBody.class))) + .thenThrow(new RuntimeException("API unavailable")); + when(featuresClient.checkFlags(any(CheckFlagRequestBody.class))) + .thenThrow(new RuntimeException("API unavailable")); + + assertFalse(spySchematic.checkFlag("flag-on", COMPANY, null)); + assertTrue(spySchematic.checkFlag("flag-off", COMPANY, null)); + + List bulk = spySchematic.checkFlags(KEYS, COMPANY, null); + assertEquals(2, bulk.size()); + assertFalse(bulk.get(0).getValue()); + assertTrue(bulk.get(1).getValue()); + } + + // --- Flag checks, cache ready --- + + @Test + void ready_singleAndBulkEvaluateFromCacheWithoutApi() { + setHealth(200, true, CACHE_VERSION); + Schematic spySchematic = spy(buildSchematic(Collections.emptyMap())); + assertTrue(spySchematic.isCacheReady()); + seedReplicatorCache(schematic.getDataStreamClient()); + + FeaturesClient featuresClient = mock(FeaturesClient.class); + lenient().when(spySchematic.features()).thenReturn(featuresClient); + + RulesengineCheckFlagResult singleOn = spySchematic.checkFlagWithEntitlement("flag-on", COMPANY, null); + RulesengineCheckFlagResult singleOff = spySchematic.checkFlagWithEntitlement("flag-off", COMPANY, null); + List bulk = spySchematic.checkFlags(KEYS, COMPANY, null); + + assertTrue(singleOn.getValue()); + assertFalse(singleOff.getValue()); + assertEquals("flag_flag-on", singleOn.getFlagId().orElse(null)); + assertEquals("comp-1", singleOn.getCompanyId().orElse(null)); + + assertEquals(2, bulk.size()); + assertSameEvaluation(singleOn, bulk.get(0)); + assertSameEvaluation(singleOff, bulk.get(1)); + + verifyNoInteractions(featuresClient); + } + + @Test + void ready_flagMissingFromCache_singleAndBulkFallBackToApi() { + setHealth(200, true, CACHE_VERSION); + Schematic spySchematic = spy(buildSchematic(Collections.emptyMap())); + assertTrue(spySchematic.isCacheReady()); + seedReplicatorCache(schematic.getDataStreamClient()); + + FeaturesClient featuresClient = mock(FeaturesClient.class); + when(spySchematic.features()).thenReturn(featuresClient); + when(featuresClient.checkFlag(eq("uncached"), any(CheckFlagRequestBody.class))) + .thenReturn(singleApiResponse("uncached", true)); + when(featuresClient.checkFlags(any(CheckFlagRequestBody.class))) + .thenReturn(bulkApiResponse(Arrays.asList("flag-on", "uncached"), Arrays.asList(true, true))); + + assertTrue(spySchematic.checkFlag("uncached", COMPANY, null)); + List bulk = + spySchematic.checkFlags(Arrays.asList("flag-on", "uncached"), COMPANY, null); + assertEquals(2, bulk.size()); + assertTrue(bulk.get(1).getValue()); + assertEquals("api", bulk.get(1).getReason()); + + verify(featuresClient).checkFlag(eq("uncached"), any(CheckFlagRequestBody.class)); + verify(featuresClient).checkFlags(any(CheckFlagRequestBody.class)); + } + + // --- Health polling --- + + @Test + void health503_setsNotReadyAndRecordsCacheVersion() { + DataStreamClient client = newReplicatorClient(); + try { + // A replicator that is still loading when the SDK starts. + setHealth(503, false, "vX"); + client.checkReplicatorHealth(); + assertFalse(client.isCacheReady()); + assertFalse(client.isConnected()); + assertEquals("vX", client.getReplicatorCacheVersion()); + + // Ready, then back to loading under a new cache version. + setHealth(200, true, "vX"); + client.checkReplicatorHealth(); + assertTrue(client.isCacheReady()); + + setHealth(503, false, "vY"); + client.checkReplicatorHealth(); + assertFalse(client.isCacheReady()); + assertEquals("vY", client.getReplicatorCacheVersion()); + } finally { + client.close(); + } + } + + @Test + void healthUnreachable_setsNotReadyAndKeepsCacheVersion() { + DataStreamClient client = newReplicatorClient(); + try { + setHealth(200, true, CACHE_VERSION); + client.checkReplicatorHealth(); + assertTrue(client.isCacheReady()); + assertEquals(CACHE_VERSION, client.getReplicatorCacheVersion()); + + server.stop(0); + server = null; + client.checkReplicatorHealth(); + + assertFalse(client.isCacheReady()); + assertEquals(CACHE_VERSION, client.getReplicatorCacheVersion()); + } finally { + client.close(); + } + } + + @Test + void healthUnparseableBody_setsNotReadyAndKeepsCacheVersion() { + DataStreamClient client = newReplicatorClient(); + try { + setHealth(200, true, CACHE_VERSION); + client.checkReplicatorHealth(); + assertTrue(client.isCacheReady()); + + healthStatus.set(200); + healthBody.set("not json"); + client.checkReplicatorHealth(); + + assertFalse(client.isCacheReady()); + assertEquals(CACHE_VERSION, client.getReplicatorCacheVersion()); + } finally { + client.close(); + } + } + + @Test + void healthWithoutCacheVersion_keepsCacheVersion() { + DataStreamClient client = newReplicatorClient(); + try { + setHealth(200, true, CACHE_VERSION); + client.checkReplicatorHealth(); + + healthStatus.set(200); + healthBody.set("{\"ready\":true,\"cache_version\":\"\"}"); + client.checkReplicatorHealth(); + + assertTrue(client.isCacheReady()); + assertEquals(CACHE_VERSION, client.getReplicatorCacheVersion()); + } finally { + client.close(); + } + } + + @Test + void isCacheReady_trueOutsideReplicatorMode() { + DataStreamClient client = new DataStreamClient( + DatastreamOptions.builder().build(), "test-key", "https://api.schematichq.com", logger); + try { + // WebSocket mode keeps its connection gate; the cache itself has nothing to wait for. + assertTrue(client.isCacheReady()); + assertFalse(client.isConnected()); + } finally { + client.close(); + } + } + + @Test + void webSocketModeNotConnected_singleAndBulkStillUseApi() { + schematic = Schematic.builder() + .apiKey("test_api_key") + .logger(logger) + .basePath("http://127.0.0.1:1") + .eventCaptureBaseUrl("http://127.0.0.1:1") + .cacheProviders(Collections.emptyList()) + .datastreamOptions( + DatastreamOptions.builder().flagCacheProvider(flagCache).build()) + .build(); + Schematic spySchematic = spy(schematic); + assertTrue(spySchematic.isCacheReady()); + assertFalse(spySchematic.isDatastreamConnected()); + // Cached under whatever version key the SDK uses in this mode, so an open gate would + // find it. + schematic + .getDataStreamClient() + .handleMessage(buildResp(EntityType.FLAG.getValue(), null, flagNode("flag-on", true))); + assertNotNull(schematic.getDataStreamClient().getCachedFlag("flag-on")); + + FeaturesClient featuresClient = mock(FeaturesClient.class); + when(spySchematic.features()).thenReturn(featuresClient); + when(featuresClient.checkFlag(eq("flag-on"), any(CheckFlagRequestBody.class))) + .thenReturn(singleApiResponse("flag-on", false)); + when(featuresClient.checkFlags(any(CheckFlagRequestBody.class))) + .thenReturn(bulkApiResponse(Collections.singletonList("flag-on"), Collections.singletonList(false))); + + assertEquals( + "api", + spySchematic.checkFlagWithEntitlement("flag-on", null, null).getReason()); + assertEquals( + "api", + spySchematic + .checkFlags(Collections.singletonList("flag-on"), null, null) + .get(0) + .getReason()); + } + + // --- Helpers --- + + private Schematic buildSchematic(Map flagDefaults) { + schematic = Schematic.builder() + .apiKey("test_api_key") + .logger(logger) + .basePath("http://127.0.0.1:1") + .eventCaptureBaseUrl("http://127.0.0.1:1") + .flagDefaults(new HashMap<>(flagDefaults)) + // No flag-check result cache, so every API-path check reaches the API mock. + .cacheProviders(Collections.emptyList()) + .datastreamOptions(DatastreamOptions.builder() + .withReplicatorMode(healthUrl) + .replicatorHealthCheckInterval(Duration.ofHours(1)) + .flagCacheProvider(flagCache) + .build()) + .build(); + // The client polls once on start in the background; poll again here so the state is + // settled before the test continues. Both polls see the same response. + schematic.getDataStreamClient().checkReplicatorHealth(); + return schematic; + } + + private DataStreamClient newReplicatorClient() { + DatastreamOptions options = DatastreamOptions.builder() + .withReplicatorMode(healthUrl) + .flagCacheProvider(flagCache) + .build(); + return new DataStreamClient(options, "test-key", "https://api.schematichq.com", logger); + } + + private void setHealth(int status, boolean ready, String cacheVersion) { + healthStatus.set(status); + healthBody.set("{\"ready\":" + ready + ",\"cache_version\":\"" + cacheVersion + "\"}"); + } + + /** + * Seeds the cache the way the replicator lays it out: flags at + * {@code flags:{cache_version}:{key}}, and the company at + * {@code company:{cache_version}:{id}} with a key lookup at + * {@code company:{cache_version}:{key}:{value}}. + */ + private void seedReplicatorCache(DataStreamClient client) { + flagCache.set("flags:" + CACHE_VERSION + ":flag-on", flag("flag-on", true)); + flagCache.set("flags:" + CACHE_VERSION + ":flag-off", flag("flag-off", false)); + client.handleMessage(companyResp("comp-1", "customer_id", "cust-1")); + assertNotNull(client.getCachedCompany(COMPANY)); + } + + private static void assertSameEvaluation(RulesengineCheckFlagResult expected, RulesengineCheckFlagResult actual) { + assertEquals(expected.getFlagKey(), actual.getFlagKey()); + assertEquals(expected.getValue(), actual.getValue()); + assertEquals(expected.getReason(), actual.getReason()); + assertEquals(expected.getFlagId(), actual.getFlagId()); + assertEquals(expected.getCompanyId(), actual.getCompanyId()); + } + + private static CheckFlagResponse singleApiResponse(String key, boolean value) { + return CheckFlagResponse.builder() + .data(CheckFlagResponseData.builder() + .flag(key) + .reason("api") + .value(value) + .build()) + .build(); + } + + private static CheckFlagsResponse bulkApiResponse(List keys, List values) { + CheckFlagResponseData[] flags = new CheckFlagResponseData[keys.size()]; + for (int i = 0; i < keys.size(); i++) { + flags[i] = CheckFlagResponseData.builder() + .flag(keys.get(i)) + .reason("api") + .value(values.get(i)) + .build(); + } + return CheckFlagsResponse.builder() + .data(CheckFlagsResponseData.builder() + .flags(Arrays.asList(flags)) + .build()) + .build(); + } + + private ObjectNode flagNode(String key, boolean defaultValue) { + ObjectNode node = objectMapper.createObjectNode(); + node.put("key", key); + node.put("id", "flag_" + key); + node.put("account_id", "acc_1"); + node.put("environment_id", "env_1"); + node.put("default_value", defaultValue); + node.set("rules", objectMapper.createArrayNode()); + return node; + } + + private RulesengineFlag flag(String key, boolean defaultValue) { + try { + return objectMapper.treeToValue(flagNode(key, defaultValue), RulesengineFlag.class); + } catch (Exception e) { + throw new IllegalStateException(e); + } + } + + private DataStreamResp companyResp(String id, String keyName, String keyValue) { + ObjectNode node = objectMapper.createObjectNode(); + node.put("id", id); + node.put("account_id", "acc_1"); + node.put("environment_id", "env_1"); + ObjectNode keys = objectMapper.createObjectNode(); + keys.put(keyName, keyValue); + node.set("keys", keys); + node.set("traits", objectMapper.createArrayNode()); + node.set("metrics", objectMapper.createArrayNode()); + node.set("rules", objectMapper.createArrayNode()); + node.set("billing_product_ids", objectMapper.createArrayNode()); + node.set("credit_balances", objectMapper.createObjectNode()); + node.set("plan_ids", objectMapper.createArrayNode()); + node.set("plan_version_ids", objectMapper.createArrayNode()); + return buildResp(EntityType.COMPANY.getValue(), id, node); + } + + private DataStreamResp buildResp(String entityType, String entityId, JsonNode data) { + ObjectNode respNode = objectMapper.createObjectNode(); + respNode.put("entity_type", entityType); + respNode.put("message_type", MessageType.FULL.getValue()); + if (entityId != null) { + respNode.put("entity_id", entityId); + } + respNode.set("data", data); + return objectMapper.convertValue(respNode, DataStreamResp.class); + } +}