Skip to content
Open
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
16 changes: 16 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
54 changes: 44 additions & 10 deletions src/main/java/com/schematic/api/Schematic.java
Original file line number Diff line number Diff line change
Expand Up @@ -290,16 +290,35 @@ public boolean isReplicatorMode() {

/**
* Returns whether the datastream connection is active and ready.
*
* <p>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.
*
* <p>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.
*
* <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
* <p>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<String, String> company, Map<String, String> user) {
Expand All @@ -312,7 +331,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 and connected and the cache is ready, see
* {@link #isCacheReady()})</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 +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<String, String> company, Map<String, String> user) {
if (dataStreamClient == null || !dataStreamClient.isConnected()) {
if (!useDataStreamCache()) {
return null;
}
try {
Expand Down Expand Up @@ -415,8 +448,9 @@ 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
* locally; falls back to the API if any key fails</li>
* <li>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</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
* fresh values, refresh the cache, and merge the results</li>
Expand All @@ -440,7 +474,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 (useDataStreamCache() && flagKeys != null && !flagKeys.isEmpty()) {
List<RulesengineCheckFlagResult> dsResults = new ArrayList<>(flagKeys.size());
boolean dsOk = true;
for (String key : flagKeys) {
Expand Down
143 changes: 108 additions & 35 deletions src/main/java/com/schematic/api/datastream/DataStreamClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,12 @@ public void start() {

/**
* Returns whether the datastream is connected and ready for flag checks.
*
* <p>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()) {
Expand All @@ -156,6 +162,26 @@ public boolean isConnected() {
return wsClient != null && wsClient.isReady();
}

/**
* Returns whether flag checks may be evaluated from the cache.
*
* <p>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.
*
* <p>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.
*/
Expand Down Expand Up @@ -631,6 +657,15 @@ private void startReplicatorMode() {
this::checkReplicatorHealth, 0, intervalMs, TimeUnit.MILLISECONDS);
}

/**
* Polls the replicator health endpoint once.
*
* <p>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()
Expand All @@ -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 ---
Expand Down
Loading
Loading