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:
*
- * - DataStream evaluation (if configured and connected)
+ * - DataStream evaluation (if configured; in direct WebSocket mode it must also be
+ * connected, while replicator mode evaluates from the cache regardless of readiness)
* - API call with result caching (fallback)
* - Flag default value (if all else fails)
*
@@ -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:
*
* - Offline mode → return flag defaults for the requested keys
- * - DataStream / replicator (if configured and connected) → evaluate each key
+ *
- 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
* - 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);
+ }
+}