From 0b86a5f0fe445256ab18aaf0e83bc02d44d49a5e Mon Sep 17 00:00:00 2001 From: ryan echternacht Date: Wed, 30 Sep 2026 11:07:09 -0400 Subject: [PATCH 1/2] Serve replicator-mode flag checks only once the cache is ready Read the replicator health body whatever the HTTP status, so a 503 from /ready still records cache_version, and set readiness from the ready field. A failed poll sets not ready and keeps the last known cache version. Add isCacheReady() on DataStreamClient and Schematic. Single and bulk flag checks both gate on it: before the replicator reports ready they skip the cache and use the API; once ready they evaluate from the cache with the existing API fallback. isConnected() and isDatastreamConnected() are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --- README.md | 16 + .../java/com/schematic/api/Schematic.java | 45 +- .../api/datastream/DataStreamClient.java | 141 ++++-- .../datastream/ReplicatorCacheReadyTest.java | 422 ++++++++++++++++++ 4 files changed, 578 insertions(+), 46 deletions(-) create mode 100644 src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java diff --git a/README.md b/README.md index cc4b567..63599b6 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. + +Call `isCacheReady()` to see whether flag checks are currently being served from the cache: + +```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. + ## 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..be3e684 100644 --- a/src/main/java/com/schematic/api/Schematic.java +++ b/src/main/java/com/schematic/api/Schematic.java @@ -290,17 +290,33 @@ 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 can be served 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. In + * direct WebSocket mode it matches {@link #isDatastreamConnected()}. Returns false when + * datastream is not configured. + */ + 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 - * evaluation fails. + *

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

    - *
  1. DataStream evaluation (if configured and connected)
  2. + *
  3. DataStream evaluation (if configured 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. *
@@ -343,13 +360,18 @@ private RulesengineCheckFlagResult defaultFlagResult(String flagKey, String reas /** * 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 datastream cache is not ready (see + * {@link #isCacheReady()}) or evaluation failed. Callers are responsible for emitting a + * {@code flag_check} event when appropriate. Single-flag checks do, bulk checks do not. + * + *

Single ({@link #checkFlagWithEntitlement}) and bulk ({@link #checkFlags}) checks + * both go through this method, and both gate on {@link #isCacheReady()}, so they serve + * from the cache under the same condition. In replicator mode that is once the + * replicator reports ready; before that, both take the API path. */ private RulesengineCheckFlagResult tryDatastreamCheckFlag( String flagKey, Map company, Map user) { - if (dataStreamClient == null || !dataStreamClient.isConnected()) { + if (!isCacheReady()) { return null; } try { @@ -415,8 +437,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 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 +463,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 (isCacheReady() && 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..392de84 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,24 @@ public boolean isConnected() { return wsClient != null && wsClient.isReady(); } + /** + * Returns whether the datastream cache is ready to serve flag checks. + * + *

    In replicator mode this is true only once the replicator's health endpoint reports + * {@code ready: true}, meaning its cache is complete for the reported cache version. Until + * then flag checks skip the cache and go to the Schematic API. A failed health poll + * (connection error, timeout, unparseable body) also reports not ready. + * + *

    In direct WebSocket mode the SDK fills its own cache over the WebSocket, and this + * returns the same value as {@link #isConnected()}. + */ + public boolean isCacheReady() { + if (options.isReplicatorMode()) { + return replicatorReady.get(); + } + return isConnected(); + } + /** * Returns whether this client is running in replicator mode. */ @@ -631,6 +655,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 +672,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..37e8b41 --- /dev/null +++ b/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java @@ -0,0 +1,422 @@ +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(); + } + } + + // --- 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 RulesengineFlag flag(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()); + try { + return objectMapper.treeToValue(node, 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); + } +} From ed588dc2858dd269f9b9bc77ed2a0315e835225e Mon Sep 17 00:00:00 2001 From: ryan echternacht Date: Wed, 30 Sep 2026 11:11:47 -0400 Subject: [PATCH 2/2] Match the Go gate: shared useDataStreamCache helper, isCacheReady true outside replicator mode Align with SchematicHQ/schematic-go#240. isCacheReady() now returns true outside replicator mode. Single and bulk flag checks share one helper, useDataStreamCache(), which is the existing datastream check (configured and connected, so WebSocket mode is unchanged) plus isCacheReady(). Co-Authored-By: Claude Opus 5.5 (1M context) --- README.md | 4 +- .../java/com/schematic/api/Schematic.java | 53 +++++++++------- .../api/datastream/DataStreamClient.java | 18 +++--- .../datastream/ReplicatorCacheReadyTest.java | 60 ++++++++++++++++++- 4 files changed, 102 insertions(+), 33 deletions(-) diff --git a/README.md b/README.md index 63599b6..0de3443 100644 --- a/README.md +++ b/README.md @@ -506,7 +506,7 @@ The SDK serves flag checks from the replicator cache only once the replicator re 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. -Call `isCacheReady()` to see whether flag checks are currently being served from the cache: +`isCacheReady()` reports the same readiness the flag checks use: ```java if (schematic.isCacheReady()) { @@ -514,7 +514,7 @@ if (schematic.isCacheReady()) { } ``` -`isDatastreamConnected()` returns the same value in replicator mode and is kept for backward compatibility. +`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 diff --git a/src/main/java/com/schematic/api/Schematic.java b/src/main/java/com/schematic/api/Schematic.java index be3e684..7ba9113 100644 --- a/src/main/java/com/schematic/api/Schematic.java +++ b/src/main/java/com/schematic/api/Schematic.java @@ -300,12 +300,14 @@ public boolean isDatastreamConnected() { } /** - * Returns whether flag checks can be served from the datastream cache. + * 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. In - * direct WebSocket mode it matches {@link #isDatastreamConnected()}. Returns false when - * datastream is not configured. + * 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(); @@ -314,9 +316,10 @@ public boolean isCacheReady() { /** * Checks a feature flag, returning a boolean value. * - *

    If datastream is configured and its cache is ready (see {@link #isCacheReady()}), - * evaluates the flag locally using cached data and the rules engine. Falls back to the API - * if the cache is not ready or evaluation fails. + *

    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) { return checkFlagWithEntitlement(flagKey, company, user).getValue(); @@ -328,7 +331,7 @@ public boolean checkFlag(String flagKey, Map company, MapPriority order: *

      - *
    1. DataStream evaluation (if configured and the cache is ready, see + *
    2. DataStream evaluation (if configured and connected and the cache is ready, see * {@link #isCacheReady()})
    3. *
    4. API call with result caching (fallback)
    5. *
    6. Flag default value (if all else fails)
    7. @@ -358,20 +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 the datastream cache is not ready (see - * {@link #isCacheReady()}) or evaluation failed. Callers are responsible for emitting a - * {@code flag_check} event when appropriate. Single-flag checks do, bulk checks do not. - * - *

      Single ({@link #checkFlagWithEntitlement}) and bulk ({@link #checkFlags}) checks - * both go through this method, and both gate on {@link #isCacheReady()}, so they serve - * from the cache under the same condition. In replicator mode that is once the - * replicator reports ready; before that, both take the API path. + * 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 (!isCacheReady()) { + if (!useDataStreamCache()) { return null; } try { @@ -437,9 +448,9 @@ private void enqueueFlagCheckEvent( *

      Evaluation order: *

        *
      1. Offline mode → return flag defaults for the requested keys
      2. - *
      3. DataStream / replicator (if configured and the cache is ready, see - * {@link #isCacheReady()}) → 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. @@ -463,7 +474,7 @@ public List checkFlags( } // 2. DataStream/replicator path: evaluate each key; on any failure fall back to API. - if (isCacheReady() && 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 392de84..3d51e94 100644 --- a/src/main/java/com/schematic/api/datastream/DataStreamClient.java +++ b/src/main/java/com/schematic/api/datastream/DataStreamClient.java @@ -163,21 +163,23 @@ public boolean isConnected() { } /** - * Returns whether the datastream cache is ready to serve flag checks. + * Returns whether flag checks may be evaluated from the cache. * - *

        In replicator mode this is true only once the replicator's health endpoint reports - * {@code ready: true}, meaning its cache is complete for the reported cache version. Until - * then flag checks skip the cache and go to the Schematic API. A failed health poll - * (connection error, timeout, unparseable body) also reports not ready. + *

        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. * - *

        In direct WebSocket mode the SDK fills its own cache over the WebSocket, and this - * returns the same value as {@link #isConnected()}. + *

        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 isConnected(); + return true; } /** diff --git a/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java b/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java index 37e8b41..7c12583 100644 --- a/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java +++ b/src/test/java/com/schematic/api/datastream/ReplicatorCacheReadyTest.java @@ -293,6 +293,58 @@ void healthWithoutCacheVersion_keepsCacheVersion() { } } + @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) { @@ -376,7 +428,7 @@ private static CheckFlagsResponse bulkApiResponse(List keys, List