From 0908d2f9d27559ea7605228e8c4d2a41c6bcc0fb Mon Sep 17 00:00:00 2001 From: ryan echternacht Date: Tue, 29 Sep 2026 17:01:37 -0400 Subject: [PATCH] Evaluate flags from the replicator cache when the replicator is not ready In replicator mode, checkFlag, checkFlagWithEntitlement, and checkFlags skipped the datastream whenever the replicator's health endpoint reported ready: false, and went to the API instead. The replicator keeps its Redis cache when it loses its connection to Schematic, so those checks now evaluate from the cache regardless of readiness, matching the Go SDK. A flag missing from the cache still falls back to the API. Direct WebSocket mode is unchanged and still requires an active connection. The replicator answers 503 with a JSON body while not ready, and the health check only read the body on 2xx. An SDK that started while the replicator was not ready never learned cache_version and built cache keys under the wrong version. The health check now reads the body regardless of status. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../java/com/schematic/api/Schematic.java | 40 +++- .../api/datastream/DataStreamClient.java | 80 ++++--- .../datastream/ReplicatorNotReadyTest.java | 196 ++++++++++++++++++ 3 files changed, 276 insertions(+), 40 deletions(-) create mode 100644 src/test/java/com/schematic/api/datastream/ReplicatorNotReadyTest.java diff --git a/src/main/java/com/schematic/api/Schematic.java b/src/main/java/com/schematic/api/Schematic.java index 2c426d5..e80b67e 100644 --- a/src/main/java/com/schematic/api/Schematic.java +++ b/src/main/java/com/schematic/api/Schematic.java @@ -298,9 +298,9 @@ public boolean isDatastreamConnected() { /** * 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, evaluates the flag locally using cached data and the + * rules engine (in direct WebSocket mode this requires an active connection; in replicator + * mode it does not). Falls back to the API if datastream is unavailable or evaluation fails. */ public boolean checkFlag(String flagKey, Map company, Map user) { return checkFlagWithEntitlement(flagKey, company, user).getValue(); @@ -312,7 +312,8 @@ public boolean checkFlag(String flagKey, Map company, MapPriority order: *

    - *
  1. DataStream evaluation (if configured and connected)
  2. + *
  3. DataStream evaluation (if configured; in direct WebSocket mode it must also be + * connected, while replicator mode evaluates from the cache regardless of readiness)
  4. *
  5. API call with result caching (fallback)
  6. *
  7. Flag default value (if all else fails)
  8. *
@@ -341,15 +342,33 @@ private RulesengineCheckFlagResult defaultFlagResult(String flagKey, String reas .build(); } + /** + * Whether flag checks should be evaluated through the datastream client. + * + *

In replicator mode the external replicator owns the cache and keeps it when it + * loses its upstream connection to Schematic (it reports {@code ready: false} but the + * cached flags, companies, and users are still there). Readiness therefore does not + * gate evaluation: flags are evaluated from whatever the cache holds, and a flag + * missing from the cache still falls back to the API. This matches the Go SDK. In + * direct WebSocket mode evaluation still requires an active connection. + */ + private boolean canEvaluateViaDatastream() { + if (dataStreamClient == null) { + return false; + } + return dataStreamClient.isReplicatorMode() || dataStreamClient.isConnected(); + } + /** * 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 datastream is not usable (see + * {@link #canEvaluateViaDatastream()}) 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 (!canEvaluateViaDatastream()) { return null; } try { @@ -415,7 +434,8 @@ 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 + *
  4. DataStream / replicator (if configured; direct WebSocket mode must also be + * connected, replicator mode is not gated on readiness) → evaluate each key * locally; falls back to the API if any key fails
  5. *
  6. 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 @@ -440,7 +460,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 (canEvaluateViaDatastream() && 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..5718335 100644 --- a/src/main/java/com/schematic/api/datastream/DataStreamClient.java +++ b/src/main/java/com/schematic/api/datastream/DataStreamClient.java @@ -639,37 +639,25 @@ 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); - } + // The replicator answers 503 with the same JSON body while it is not ready + // (for example when it has lost its connection to Schematic but still holds + // its cache). Read the body regardless of status so the cache version stays + // current; flag checks in replicator mode evaluate from that cache even when + // the replicator is not ready, and need the version to build cache keys. + JsonNode body = readHealthBody(response); + if (body != null) { + updateReplicatorCacheVersion(body); + } - 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()); - } + boolean ready = response.isSuccessful() + && body != null + && body.path("ready").asBoolean(false); + boolean wasReady = replicatorReady.getAndSet(ready); + + if (ready && !wasReady) { + log("info", "Replicator is now ready"); + } else if (!ready && wasReady) { + log("warn", "Replicator is no longer ready (status: " + response.code() + ")"); } } } catch (IOException e) { @@ -681,6 +669,38 @@ void checkReplicatorHealth() { } } + 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) { + log("debug", "Failed to parse replicator health response: " + e.getMessage()); + return null; + } + } + + 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(); + } + // Keep the last known version when the replicator reports none. + 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); + } + } + // --- Message handling --- void handleMessage(DataStreamResp message) { diff --git a/src/test/java/com/schematic/api/datastream/ReplicatorNotReadyTest.java b/src/test/java/com/schematic/api/datastream/ReplicatorNotReadyTest.java new file mode 100644 index 0000000..ae59dbd --- /dev/null +++ b/src/test/java/com/schematic/api/datastream/ReplicatorNotReadyTest.java @@ -0,0 +1,196 @@ +package com.schematic.api.datastream; + +import static org.junit.jupiter.api.Assertions.*; +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.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.List; +import java.util.Map; +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 when the replicator reports {@code ready: false} (for example when it has + * lost its connection to Schematic but still holds its cache). Flag checks should evaluate + * from the cache rather than falling back to the API. + */ +@ExtendWith(MockitoExtension.class) +class ReplicatorNotReadyTest { + + // Nothing listens here, so every health check fails and the replicator stays not ready. + private static final String UNREACHABLE_HEALTH_URL = "http://127.0.0.1:1/ready"; + + @Mock + private SchematicLogger logger; + + private final ObjectMapper objectMapper = new ObjectMapper(); + private Schematic schematic; + private HttpServer server; + + @BeforeEach + void setUp() { + schematic = Schematic.builder() + .apiKey("test_api_key") + .logger(logger) + .basePath("http://127.0.0.1:1") + .eventCaptureBaseUrl("http://127.0.0.1:1") + .datastreamOptions(DatastreamOptions.builder() + .withReplicatorMode(UNREACHABLE_HEALTH_URL) + .replicatorHealthCheckInterval(Duration.ofHours(1)) + .build()) + .build(); + } + + @AfterEach + void tearDown() { + if (schematic != null) { + schematic.close(); + } + if (server != null) { + server.stop(0); + } + } + + @Test + void checkFlag_replicatorNotReady_evaluatesFromCache() { + DataStreamClient ds = schematic.getDataStreamClient(); + ds.handleMessage(flagResp("cached-flag", true)); + ds.handleMessage(companyResp("comp-1", "customer_id", "cust-1")); + + assertTrue(schematic.isReplicatorMode()); + assertFalse(schematic.isDatastreamConnected()); + + Schematic spySchematic = spy(schematic); + FeaturesClient featuresClient = mock(FeaturesClient.class); + lenient().when(spySchematic.features()).thenReturn(featuresClient); + + Map company = Collections.singletonMap("customer_id", "cust-1"); + RulesengineCheckFlagResult result = spySchematic.checkFlagWithEntitlement("cached-flag", company, null); + + assertEquals("cached-flag", result.getFlagKey()); + assertTrue(result.getValue()); + assertEquals("comp-1", result.getCompanyId().orElse(null)); + assertTrue(spySchematic.checkFlag("cached-flag", company, null)); + verifyNoInteractions(featuresClient); + } + + @Test + void checkFlags_replicatorNotReady_evaluatesFromCache() { + DataStreamClient ds = schematic.getDataStreamClient(); + ds.handleMessage(flagResp("flag-on", true)); + ds.handleMessage(flagResp("flag-off", false)); + ds.handleMessage(companyResp("comp-1", "customer_id", "cust-1")); + + assertFalse(schematic.isDatastreamConnected()); + + Schematic spySchematic = spy(schematic); + FeaturesClient featuresClient = mock(FeaturesClient.class); + lenient().when(spySchematic.features()).thenReturn(featuresClient); + + List results = spySchematic.checkFlags( + Arrays.asList("flag-on", "flag-off"), Collections.singletonMap("customer_id", "cust-1"), null); + + assertEquals(2, results.size()); + assertEquals("flag-on", results.get(0).getFlagKey()); + assertTrue(results.get(0).getValue()); + assertEquals("flag-off", results.get(1).getFlagKey()); + assertFalse(results.get(1).getValue()); + verifyNoInteractions(featuresClient); + } + + @Test + void checkReplicatorHealth_notReady503_stillAdoptsCacheVersion() throws Exception { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/ready", exchange -> { + byte[] bytes = + "{\"ready\":false,\"connected\":false,\"cache_version\":\"v42\"}".getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().add("Content-Type", "application/json"); + exchange.sendResponseHeaders(503, bytes.length); + try (OutputStream os = exchange.getResponseBody()) { + os.write(bytes); + } + }); + server.start(); + + LocalCache flagCache = new LocalCache<>(); + DatastreamOptions options = DatastreamOptions.builder() + .withReplicatorMode("http://127.0.0.1:" + server.getAddress().getPort() + "/ready") + .flagCacheProvider(flagCache) + .build(); + DataStreamClient client = new DataStreamClient(options, "test-key", "https://api.schematichq.com", logger); + try { + client.checkReplicatorHealth(); + assertFalse(client.isConnected()); + + // Entries are written and read under the replicator's cache version, even though + // the replicator reported not ready. + client.handleMessage(flagResp("versioned-flag", true)); + assertNotNull(flagCache.get("flags:v42:versioned-flag")); + assertTrue(client.checkFlag("versioned-flag", null, null).getValue()); + } finally { + client.close(); + } + } + + private DataStreamResp flagResp(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 buildResp(EntityType.FLAG.getValue(), null, node); + } + + 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); + } +}