Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 30 additions & 10 deletions src/main/java/com/schematic/api/Schematic.java
Original file line number Diff line number Diff line change
Expand Up @@ -298,9 +298,9 @@ public boolean isDatastreamConnected() {
/**
* Checks a feature flag, returning a boolean value.
*
* <p>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.
* <p>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<String, String> company, Map<String, String> user) {
return checkFlagWithEntitlement(flagKey, company, user).getValue();
Expand All @@ -312,7 +312,8 @@ public boolean checkFlag(String flagKey, Map<String, String> company, Map<String
*
* <p>Priority order:
* <ol>
* <li>DataStream evaluation (if configured and connected)</li>
* <li>DataStream evaluation (if configured; in direct WebSocket mode it must also be
* connected, while replicator mode evaluates from the cache regardless of readiness)</li>
* <li>API call with result caching (fallback)</li>
* <li>Flag default value (if all else fails)</li>
* </ol>
Expand Down Expand Up @@ -341,15 +342,33 @@ private RulesengineCheckFlagResult defaultFlagResult(String flagKey, String reas
.build();
}

/**
* Whether flag checks should be evaluated through the datastream client.
*
* <p>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<String, String> company, Map<String, String> user) {
if (dataStreamClient == null || !dataStreamClient.isConnected()) {
if (!canEvaluateViaDatastream()) {
return null;
}
try {
Expand Down Expand Up @@ -415,7 +434,8 @@ private void enqueueFlagCheckEvent(
* <p>Evaluation order:
* <ol>
* <li>Offline mode → return flag defaults for the requested keys</li>
* <li>DataStream / replicator (if configured and connected) → evaluate each key
* <li>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</li>
* <li>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
Expand All @@ -440,7 +460,7 @@ public List<RulesengineCheckFlagResult> 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<RulesengineCheckFlagResult> dsResults = new ArrayList<>(flagKeys.size());
boolean dsOk = true;
for (String key : flagKeys) {
Expand Down
80 changes: 50 additions & 30 deletions src/main/java/com/schematic/api/datastream/DataStreamClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, String> 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<RulesengineCheckFlagResult> 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<RulesengineFlag> 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);
}
}
Loading