diff --git a/providers/go-feature-flag/README.md b/providers/go-feature-flag/README.md
index cc1d6cb0fa..a936e52706 100644
--- a/providers/go-feature-flag/README.md
+++ b/providers/go-feature-flag/README.md
@@ -68,21 +68,21 @@ The `targetingKey` is mandatory for GO Feature Flag in order to evaluate the fea
You can configure the provider with several options to customize its behavior. The following options are available:
-| name | mandatory | Description |
-|-----------------------------------|-----------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
-| **`endpoint`** | `true` | endpoint contains the DNS of your GO Feature Flag relay proxy _(ex: https://mydomain.com/gofeatureflagproxy/)_ |
-| **`evaluationType`** | `false` | evaluationType is the type of evaluation you want to use.
If you want to have a local evaluation, you should use IN_PROCESS.
If you want to have an evaluation on the relay-proxy directly, you should use REMOTE.
Default: IN_PROCESS |
-| **`timeout`** | `false` | timeout in millisecond we are waiting when calling the relay proxy API. _(default: `10000`)_ |
-| **`maxIdleConnections`** | `false` | maxIdleConnections is the maximum number of connections in the connection pool. _(default: `1000`)_ |
-| **`keepAliveDuration`** | `false` | keepAliveDuration is the time in millisecond we keep the connection open. _(default: `7200000` (2 hours))_ |
-| **`apiKey`** | `false` | If the relay proxy is configured to authenticate the requests, you should provide an API Key to the provider. Please ask the administrator of the relay proxy to provide an API Key. (This feature is available only if you are using GO Feature Flag relay proxy v1.7.0 or above). _(default: null)_ |
-| **`flushIntervalMs`** | `false` | interval time we publish statistics collection data to the proxy. The parameter is used only if the cache is enabled, otherwise the collection of the data is done directly when calling the evaluation API. default: `1000` ms |
-| **`maxPendingEvents`** | `false` | max pending events aggregated before publishing for collection data to the proxy. When event is added while events collection is full, event is omitted. _(default: `10000`)_ |
-| **`disableDataCollection`** | `false` | set to true if you don't want to collect the usage of flags retrieved in the cache. _(default: `false`)_ |
-| **`exporterMetadata`** | `false` | exporterMetadata is the metadata we send to the GO Feature Flag relay proxy when we report the evaluation data usage. |
-| **`evaluationFlagList`** | `false` | If you are using in process evaluation, by default we will load in memory all the flags available in the relay proxy. If you want to limit the number of flags loaded in memory, you can use this parameter. By setting this parameter, you will only load the flags available in the list.
If null or empty, all the flags available in the relay proxy will be loaded.
|
-| **`flagChangePollingIntervalMs`** | `false` | interval time we poll the proxy to check if the configuration has changed. It is used for the in process evaluation to check if we should refresh our internal cache. default: `120000` |
-| **`wasmEvaluatorPoolSize`** | `false` | _(IN_PROCESS only)_ Number of WASM instances kept in the evaluation pool. Each instance owns independent memory, allowing fully concurrent flag evaluations without serialisation. Must be `>= 1`. _(default: number of available CPU cores)_ |
+| name | mandatory | Description |
+|-----------------------------------|-----------|-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
+| **`endpoint`** | `true` | endpoint contains the DNS of your GO Feature Flag relay proxy _(ex: https://mydomain.com/gofeatureflagproxy/)_ |
+| **`evaluationType`** | `false` | evaluationType is the type of evaluation you want to use.
If you want to have a local evaluation, you should use IN_PROCESS.
If you want to have an evaluation on the relay-proxy directly, you should use REMOTE.
Default: IN_PROCESS |
+| **`timeout`** | `false` | timeout in millisecond we are waiting when calling the relay proxy API. _(default: `10000`)_ |
+| **`apiKey`** | `false` | If the relay proxy is configured to authenticate the requests, you should provide an API Key to the provider. Please ask the administrator of the relay proxy to provide an API Key. (This feature is available only if you are using GO Feature Flag relay proxy v1.7.0 or above). _(default: null)_ |
+| **`dataCollectorBaseUrl`** | `false` | base URL used to publish the evaluation and tracking events, when the data collector is not served by the relay proxy itself. It replaces the whole base of the collector route, scheme, host, port and path prefix included; the flag configuration and the evaluations keep using `endpoint`. _(default: `endpoint`)_ |
+| **`customHeaders`** | `false` | Extra HTTP headers added to every request the provider makes (flag configuration, evaluation and data collection, including `dataCollectorBaseUrl`), for deployments behind a gateway that needs its own authentication. A configured `apiKey` always wins over a custom `X-API-Key`. `Content-Type` and `If-None-Match` are set by the provider and refused here, as are the headers the Java HTTP client restricts (`Host`, `Connection`, `Content-Length`, `Expect`, `Upgrade`). _(default: none)_ |
+| **`flushIntervalMs`** | `false` | interval time in millisecond we publish the collected evaluation and tracking events to the data collector. _(default: `60000` (1 minute))_ |
+| **`maxPendingEvents`** | `false` | max pending events aggregated before publishing for collection data to the proxy. Once that many events are pending they are published without waiting for `flushIntervalMs`. If they cannot be published, at most twice that many are kept and the oldest are dropped. _(default: `10000`)_ |
+| **`disableDataCollection`** | `false` | set to true if you don't want to send the evaluation and tracking events to the data collector. _(default: `false`)_ |
+| **`exporterMetadata`** | `false` | exporterMetadata is the metadata we send to the GO Feature Flag relay proxy when we report the evaluation data usage. |
+| **`evaluationFlagList`** | `false` | If you are using in process evaluation, by default we will load in memory all the flags available in the relay proxy. If you want to limit the number of flags loaded in memory, you can use this parameter. By setting this parameter, you will only load the flags available in the list.
If null or empty, all the flags available in the relay proxy will be loaded.
|
+| **`flagChangePollingIntervalMs`** | `false` | interval time we poll the proxy to check if the configuration has changed. It is used for the in process evaluation to check if the flag configuration it holds should be refreshed. Each poll is randomly shortened or lengthened by up to 10%, so that providers started together do not poll in lockstep. default: `120000` |
+| **`wasmEvaluatorPoolSize`** | `false` | _(IN_PROCESS only)_ Number of WASM instances kept in the evaluation pool. Each instance owns independent memory, allowing fully concurrent flag evaluations without serialisation. Each instance holds about 2.3 MiB of memory once warm. Must be `>= 1`. _(default: number of available CPU cores)_ |
### Evaluate a feature flag
The OpenFeature client is used to retrieve values for the current `EvaluationContext`. For example, retrieving a boolean value for the flag **"my-flag"**:
@@ -116,11 +116,23 @@ client.getObjectDetails("my-flag",Value.objectToValue(new MutableStructure().add
When the provider is configured to use in process evaluation, it will fetch the flag configuration from the GO Feature Flag relay-proxy API and evaluate the flags directly in the provider.
The evaluation is done inside the provider using a webassembly module that is compiled from the GO Feature Flag source code.
-The `wasm` module is used to evaluate the flags and the source code is available in the [thomaspoignant/go-feature-flag](https://github.com/thomaspoignant/go-feature-flag/tree/main/wasm) repository.
+The `wasm` module is used to evaluate the flags and the source code is available in the [thomaspoignant/go-feature-flag](https://github.com/thomaspoignant/go-feature-flag/tree/main/cmd/wasm) repository.
The provider will call the GO Feature Flag relay-proxy API to fetch the flag configuration and then evaluate the flags using the `wasm` module.
+The `wasm` module is compiled into Java bytecode when the provider is built, so it ships as classes inside the provider jar. There is no `.wasm` file to locate at runtime, and repackaging the provider (shaded or fat jars) keeps the engine with it. The engine version is pinned by the provider release.
+
+#### Performance
+Some of the engine's functions compile to Java methods larger than the JVM's default JIT limit (`HugeMethodLimit`, 8000 bytes of bytecode), and HotSpot leaves such methods to the bytecode interpreter. This includes the evaluation path itself, so the engine runs about 3x slower than it could: roughly 350 µs instead of 110 µs per evaluation on an Apple M4 Pro with JDK 21.
+
+If this matters to you, start the JVM with `-XX:-DontCompileHugeMethods`. The flag is process-wide and lifts the limit for every class, not only the engine's.
+
### Remote evaluation
When the provider is configured to use remote evaluation, it will call the GO Feature Flag relay-proxy for each flag evaluation.
It will perform an HTTP request to the GO Feature Flag relay-proxy API with the flag name and the evaluation context for each flag evaluation.
+
+### Logging
+The provider logs through [SLF4J](https://www.slf4j.org/), so its diagnostics go wherever your SLF4J backend sends them. Every logger it uses sits under `dev.openfeature.contrib.providers.gofeatureflag`, and remote evaluation also logs under `dev.openfeature.contrib.providers.ofrep`.
+
+What the evaluation engine prints, such as the panic it reports before a trap, is logged too, at error level by `dev.openfeature.contrib.providers.gofeatureflag.wasm.EvaluationWasm`, rather than written to the process's standard streams.
diff --git a/providers/go-feature-flag/pom.xml b/providers/go-feature-flag/pom.xml
index 78359cf5e1..693f49eea2 100644
--- a/providers/go-feature-flag/pom.xml
+++ b/providers/go-feature-flag/pom.xml
@@ -20,7 +20,9 @@
${groupId}.gofeatureflag
- 0.2.3
+ 0.2.4
+
+ **/e2e/*.java
@@ -72,7 +74,7 @@
org.apache.logging.log4jlog4j-slf4j2-impl
- 2.25.2
+ 2.26.1test
@@ -103,12 +105,38 @@
test
+
+ org.testcontainers
+ testcontainers
+ 2.0.4
+ test
+
+
+
+ org.testcontainers
+ testcontainers-junit-jupiter
+ 2.0.4
+ test
+
+
com.github.spotbugsspotbugs-annotations4.9.8provided
+
+
+ dev.openfeature.contrib.providers
+ ofrep
+ 0.0.2
+
+
+
+ com.google.guava
+ guava
+ 33.4.8-jre
+
@@ -131,4 +159,14 @@
+
+
+
+
+ e2e
+
+
+
+
+
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProvider.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProvider.java
index 75dba19ec8..709f692f03 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProvider.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProvider.java
@@ -11,7 +11,6 @@
import dev.openfeature.contrib.providers.gofeatureflag.hook.DataCollectorHook;
import dev.openfeature.contrib.providers.gofeatureflag.hook.DataCollectorHookOptions;
import dev.openfeature.contrib.providers.gofeatureflag.hook.EnrichEvaluationContextHook;
-import dev.openfeature.contrib.providers.gofeatureflag.service.EvaluationService;
import dev.openfeature.contrib.providers.gofeatureflag.service.EventsPublisher;
import dev.openfeature.contrib.providers.gofeatureflag.util.Const;
import dev.openfeature.contrib.providers.gofeatureflag.util.EvaluationContextUtil;
@@ -20,6 +19,7 @@
import dev.openfeature.sdk.Hook;
import dev.openfeature.sdk.Metadata;
import dev.openfeature.sdk.ProviderEvaluation;
+import dev.openfeature.sdk.ProviderEvent;
import dev.openfeature.sdk.ProviderEventDetails;
import dev.openfeature.sdk.Tracking;
import dev.openfeature.sdk.TrackingEventDetails;
@@ -29,6 +29,7 @@
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.function.BiConsumer;
import java.util.function.Consumer;
import lombok.extern.slf4j.Slf4j;
import lombok.val;
@@ -41,7 +42,7 @@ public final class GoFeatureFlagProvider extends EventProvider implements Tracki
/** Options to configure the provider. */
private final GoFeatureFlagProviderOptions options;
/** Service to evaluate the flags. */
- private final EvaluationService evalService;
+ private final IEvaluator evaluator;
/** List of the hooks used by the provider. */
private final List hooks = new ArrayList<>();
/** API layer to contact GO Feature Flag. */
@@ -50,8 +51,6 @@ public final class GoFeatureFlagProvider extends EventProvider implements Tracki
private final EventsPublisher eventsPublisher;
/** exporter metadata contains the metadata that we want to send to the exporter. */
private final Map exporterMetadata;
- /** DataCollectorHook is the hook to send usage of the flags. */
- private DataCollectorHook dataCollectorHook;
/**
* Constructor of the provider.
@@ -66,24 +65,16 @@ public GoFeatureFlagProvider(final GoFeatureFlagProviderOptions options) throws
options.validate();
this.options = options;
this.api = GoFeatureFlagApi.builder().options(options).build();
- this.evalService = new EvaluationService(getEvaluator(this.api));
+ this.evaluator = getEvaluator();
- long flushIntervalMs =
- (options.getFlushIntervalMs() == null) ? Const.DEFAULT_FLUSH_INTERVAL_MS : options.getFlushIntervalMs();
- int maxPendingEvents = (options.getMaxPendingEvents() == null)
- ? Const.DEFAULT_MAX_PENDING_EVENTS
- : options.getMaxPendingEvents();
Consumer> publisher = this::publishEvents;
- this.eventsPublisher = new EventsPublisher<>(publisher, flushIntervalMs, maxPendingEvents);
-
- if (options.getExporterMetadata() == null) {
- this.exporterMetadata = new HashMap<>();
- } else {
- val exp = new HashMap<>(options.getExporterMetadata());
- exp.put("provider", "java");
- exp.put("openfeature", true);
- this.exporterMetadata = exp;
- }
+ this.eventsPublisher =
+ new EventsPublisher<>(publisher, options.getFlushIntervalMs(), options.getMaxPendingEvents());
+
+ val exp = new HashMap<>(options.getExporterMetadata());
+ exp.put(Const.METADATA_PROVIDER, "java");
+ exp.put(Const.METADATA_OPENFEATURE, true);
+ this.exporterMetadata = exp;
}
@Override
@@ -99,48 +90,54 @@ public List getProviderHooks() {
@Override
public ProviderEvaluation getBooleanEvaluation(
String key, Boolean defaultValue, EvaluationContext evaluationContext) {
- return this.evalService.getEvaluation(key, defaultValue, evaluationContext, Boolean.class);
+ return this.evaluator.getBooleanEvaluation(key, defaultValue, evaluationContext);
}
@Override
public ProviderEvaluation getStringEvaluation(
String key, String defaultValue, EvaluationContext evaluationContext) {
- return this.evalService.getEvaluation(key, defaultValue, evaluationContext, String.class);
+ return this.evaluator.getStringEvaluation(key, defaultValue, evaluationContext);
}
@Override
public ProviderEvaluation getIntegerEvaluation(
String key, Integer defaultValue, EvaluationContext evaluationContext) {
- return this.evalService.getEvaluation(key, defaultValue, evaluationContext, Integer.class);
+ return this.evaluator.getIntegerEvaluation(key, defaultValue, evaluationContext);
}
@Override
public ProviderEvaluation getDoubleEvaluation(
String key, Double defaultValue, EvaluationContext evaluationContext) {
- return this.evalService.getEvaluation(key, defaultValue, evaluationContext, Double.class);
+ return this.evaluator.getDoubleEvaluation(key, defaultValue, evaluationContext);
}
@Override
public ProviderEvaluation getObjectEvaluation(
String key, Value defaultValue, EvaluationContext evaluationContext) {
- return this.evalService.getEvaluation(key, defaultValue, evaluationContext, Value.class);
+ return this.evaluator.getObjectEvaluation(key, defaultValue, evaluationContext);
}
@Override
public void initialize(EvaluationContext evaluationContext) throws Exception {
+ this.initialize(evaluationContext, "");
+ }
+
+ @Override
+ public void initialize(EvaluationContext evaluationContext, String domain) throws Exception {
super.initialize(evaluationContext);
- this.evalService.init();
- this.hooks.add(new EnrichEvaluationContextHook(this.options.getExporterMetadata()));
+ // re-initialization must reset the publisher: its shutdown flag and its scheduler are both
+ // one-shot, so without this a provider that is shut down and initialized again never flushes.
+ this.eventsPublisher.start();
+ this.evaluator.initialize(evaluationContext, domain);
+ this.hooks.clear();
+ this.hooks.add(new EnrichEvaluationContextHook(this.exporterMetadata));
// In case of remote evaluation, we don't need to send the data to the collector
// because the relay-proxy will collect events directly server side.
if (!this.options.isDisableDataCollection() && this.options.getEvaluationType() != EvaluationType.REMOTE) {
- this.dataCollectorHook = new DataCollectorHook(DataCollectorHookOptions.builder()
+ this.hooks.add(new DataCollectorHook(DataCollectorHookOptions.builder()
.eventsPublisher(this.eventsPublisher)
- .collectUnCachedEvaluation(true)
- .evalService(this.evalService)
- .build());
-
- this.hooks.add(this.dataCollectorHook);
+ .evaluator(this.evaluator)
+ .build()));
}
log.info("finishing initializing provider");
}
@@ -148,10 +145,8 @@ public void initialize(EvaluationContext evaluationContext) throws Exception {
@Override
public void shutdown() {
super.shutdown();
- this.evalService.destroy();
- if (this.dataCollectorHook != null) {
- this.dataCollectorHook.shutdown();
- }
+ this.evaluator.shutdown();
+ this.eventsPublisher.shutdown();
}
@Override
@@ -171,10 +166,14 @@ public void track(final String eventName, final TrackingEventDetails trackingEve
@Override
public void track(final String eventName, final EvaluationContext context, final TrackingEventDetails details) {
+ if (this.options.isDisableDataCollection()) {
+ return;
+ }
+
val trackingEvent = TrackingEvent.builder()
.evaluationContext((context != null) ? context.asObjectMap() : Collections.emptyMap())
- .userKey(context != null ? context.getTargetingKey() : "undefined-targetingKey")
- .contextKind(EvaluationContextUtil.isAnonymousUser(context) ? "anonymousUser" : "user")
+ .userKey(EvaluationContextUtil.userKey(context))
+ .contextKind(EvaluationContextUtil.contextKind(context))
.kind("tracking")
.key(eventName)
.trackingEventDetails(details != null ? details.asObjectMap() : Collections.emptyMap())
@@ -185,17 +184,16 @@ public void track(final String eventName, final EvaluationContext context, final
/**
* Get the evaluator based on the evaluation type.
- * It will initialize the evaluator based on the evaluation type.
*
* @return the evaluator
*/
- private IEvaluator getEvaluator(GoFeatureFlagApi api) {
+ private IEvaluator getEvaluator() {
// Select the evaluator based on the evaluation type
+ BiConsumer emitter = this::emit;
if (options.getEvaluationType() == null || options.getEvaluationType() == EvaluationType.IN_PROCESS) {
- Consumer emitProviderConfigurationChanged = this::emitProviderConfigurationChanged;
- return new InProcessEvaluator(api, this.options, emitProviderConfigurationChanged);
+ return new InProcessEvaluator(this.api, this.options, emitter);
}
- return new RemoteEvaluator(api);
+ return new RemoteEvaluator(this.options, emitter);
}
/**
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProviderOptions.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProviderOptions.java
index 9cd620553b..187515bfb6 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProviderOptions.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/GoFeatureFlagProviderOptions.java
@@ -4,20 +4,31 @@
import dev.openfeature.contrib.providers.gofeatureflag.exception.InvalidEndpoint;
import dev.openfeature.contrib.providers.gofeatureflag.exception.InvalidExporterMetadata;
import dev.openfeature.contrib.providers.gofeatureflag.exception.InvalidOptions;
+import dev.openfeature.contrib.providers.gofeatureflag.util.Const;
import java.net.MalformedURLException;
import java.net.URL;
+import java.net.http.HttpRequest;
+import java.util.Collections;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import lombok.Builder;
import lombok.Getter;
import lombok.val;
/**
* GoFeatureFlagProviderOptions contains the options to initialise the provider.
+ *
+ *
Every optional field is unset by default, and its getter resolves the documented default value,
+ * so the provider and the evaluators can read the options without repeating the fallbacks.
*/
@Builder
@Getter
public class GoFeatureFlagProviderOptions {
+ /** Default timeout in millisecond when calling the GO Feature Flag relay proxy API. */
+ private static final int DEFAULT_TIMEOUT_MS = 10000;
+
/**
* evaluationType is the type of evaluation you want to use.
* - If you want to have a local evaluation, you should use IN_PROCESS.
@@ -30,21 +41,18 @@ public class GoFeatureFlagProviderOptions {
* https://mydomain.com/gofeatureflagproxy/
*/
private String endpoint;
+ /**
+ * (optional) dataCollectorBaseUrl is the base URL used to publish the evaluation data, when the
+ * data collector is not served by the relay proxy itself. It replaces the whole base of the
+ * collector route, scheme, host, port and path prefix included, and applies to that route only:
+ * the flag configuration and the evaluations keep using endpoint. Default: endpoint
+ */
+ private String dataCollectorBaseUrl;
/**
* (optional) timeout in millisecond we are waiting when calling the go-feature-flag relay proxy
* API. Default: 10000 ms
*/
private int timeout;
- /**
- * (optional) maxIdleConnections is the maximum number of connexions in the connexion pool.
- * Default: 1000
- */
- private int maxIdleConnections;
- /**
- * (optional) keepAliveDuration is the time in millisecond we keep the connexion open. Default:
- * 7200000 (2 hours)
- */
- private Long keepAliveDuration;
/**
* (optional) If the relay proxy is configured to authenticate the requests, you should provide an
* API Key to the provider. Please ask the administrator of the relay proxy to provide an API Key.
@@ -53,25 +61,33 @@ public class GoFeatureFlagProviderOptions {
*/
private String apiKey;
/**
- * (optional) interval time we publish statistics collection data to the proxy. The parameter is
- * used only if the cache is enabled, otherwise the collection of the data is done directly when
- * calling the evaluation API. default: 1000 ms
+ * (optional) customHeaders are extra HTTP headers added to every request the provider makes, to the
+ * relay proxy and to dataCollectorBaseUrl, for deployments behind a gateway that needs its own
+ * authentication. A configured apiKey always wins over a custom X-API-Key. Content-Type and
+ * If-None-Match are set by the provider, and they are refused here, as are the headers the Java HTTP
+ * client restricts (Host, Connection, Content-Length, Expect, Upgrade). Default: none
+ */
+ private Map customHeaders;
+ /**
+ * (optional) interval time in millisecond we publish the collected evaluation and tracking events
+ * to the data collector. default: 60000 ms (1 minute)
*/
private Long flushIntervalMs;
/**
* (optional) max pending events aggregated before publishing for collection data to the proxy.
- * When an event is added while an events collection is full, the event is omitted. default: 10000
+ * Once that many events are pending they are published without waiting for flushIntervalMs. If
+ * they cannot be published, at most twice that many are kept and the oldest are dropped. default: 10000
*/
private Integer maxPendingEvents;
/**
- * (optional) disableDataCollection set to true if you don't want to collect the usage of flags
- * retrieved in the cache. default: false
+ * (optional) disableDataCollection set to true if you don't want to send the evaluation and
+ * tracking events to the data collector. default: false
*/
private boolean disableDataCollection;
/**
* (optional) exporterMetadata is the metadata we send to the GO Feature Flag relay proxy when we report the
- * evaluation data usage.
+ * evaluation data usage. default: empty
*/
private Map exporterMetadata;
@@ -85,9 +101,9 @@ public class GoFeatureFlagProviderOptions {
private List evaluationFlagList;
/**
- * (optional) interval time we poll the proxy to check if the configuration has changed. If the
- * cache is enabled, we will poll the relay-proxy every X milliseconds to check if the
- * configuration has changed. default: 120000
+ * (optional) interval time in millisecond we poll the relay proxy to check if the flag
+ * configuration has changed, for in process evaluation. Each poll is randomly shortened or
+ * lengthened by up to 10%. default: 120000
*/
private Long flagChangePollingIntervalMs;
@@ -99,29 +115,133 @@ public class GoFeatureFlagProviderOptions {
*/
private Integer wasmEvaluatorPoolSize;
+ /**
+ * Get the type of evaluation to use.
+ *
+ * @return the configured evaluation type, IN_PROCESS if none was set
+ */
+ public EvaluationType getEvaluationType() {
+ return evaluationType == null ? EvaluationType.IN_PROCESS : evaluationType;
+ }
+
+ /**
+ * Get the base URL used to publish the evaluation data.
+ *
+ * @return the configured data collector base URL, the endpoint if none was set
+ */
+ public String getDataCollectorBaseUrl() {
+ return dataCollectorBaseUrl == null || dataCollectorBaseUrl.isEmpty() ? endpoint : dataCollectorBaseUrl;
+ }
+
+ /**
+ * Get the timeout in millisecond when calling the GO Feature Flag relay proxy API.
+ *
+ * @return the configured timeout, 10000 ms if none was set
+ */
+ public int getTimeout() {
+ return timeout == 0 ? DEFAULT_TIMEOUT_MS : timeout;
+ }
+
+ /**
+ * Get the extra HTTP headers added to every request to the relay proxy.
+ *
+ * @return the configured headers, an empty map if none was set
+ */
+ public Map getCustomHeaders() {
+ return customHeaders == null ? Collections.emptyMap() : customHeaders;
+ }
+
+ /**
+ * Get the interval time we publish the collected events to the data collector.
+ *
+ * @return the configured interval, 60000 ms if none was set
+ */
+ public Long getFlushIntervalMs() {
+ return Objects.requireNonNullElse(flushIntervalMs, Const.DEFAULT_FLUSH_INTERVAL_MS);
+ }
+
+ /**
+ * Get the maximum number of events aggregated before publishing them to the proxy.
+ *
+ * @return the configured maximum, 10000 if none was set
+ */
+ public Integer getMaxPendingEvents() {
+ return Objects.requireNonNullElse(maxPendingEvents, Const.DEFAULT_MAX_PENDING_EVENTS);
+ }
+
+ /**
+ * Get the metadata sent to the relay proxy when reporting the evaluation data usage.
+ *
+ * @return the configured metadata, an empty map if none was set
+ */
+ public Map getExporterMetadata() {
+ return exporterMetadata == null ? Collections.emptyMap() : exporterMetadata;
+ }
+
+ /**
+ * Get the list of flags to load for in process evaluation.
+ *
+ * @return the configured list, an empty list if none was set, meaning all the flags are loaded
+ */
+ public List getEvaluationFlagList() {
+ return evaluationFlagList == null ? Collections.emptyList() : evaluationFlagList;
+ }
+
+ /**
+ * Get the interval time we poll the proxy to check if the configuration has changed.
+ *
+ * @return the configured interval, 120000 ms if none was set
+ */
+ public Long getFlagChangePollingIntervalMs() {
+ return Objects.requireNonNullElse(
+ flagChangePollingIntervalMs, Const.DEFAULT_POLLING_CONFIG_FLAG_CHANGE_INTERVAL_MS);
+ }
+
+ /**
+ * Get the number of WASM instances kept in the evaluation pool.
+ *
+ * @return the configured pool size, the number of available CPU cores if none was set
+ */
+ public Integer getWasmEvaluatorPoolSize() {
+ return Objects.requireNonNullElse(wasmEvaluatorPoolSize, Const.DEFAULT_WASM_EVALUATOR_POOL_SIZE);
+ }
+
/**
* Validate the options provided to the provider.
*
+ *
Validation reads the fields directly and not the getters, to check what the caller has
+ * really set and not the resolved default values.
+ *
* @throws InvalidOptions - if options are invalid
*/
public void validate() throws InvalidOptions {
- if (getEndpoint() == null || getEndpoint().isEmpty()) {
+ if (endpoint == null || endpoint.isEmpty()) {
throw new InvalidEndpoint("endpoint is a mandatory field when initializing the provider");
}
try {
- new URL(getEndpoint());
+ new URL(endpoint);
} catch (MalformedURLException e) {
- throw new InvalidEndpoint("malformed endpoint: " + getEndpoint());
+ throw new InvalidEndpoint("malformed endpoint: " + endpoint);
}
- if (getWasmEvaluatorPoolSize() != null && getWasmEvaluatorPoolSize() < 1) {
+ if (dataCollectorBaseUrl != null && !dataCollectorBaseUrl.isEmpty()) {
+ try {
+ new URL(dataCollectorBaseUrl);
+ } catch (MalformedURLException e) {
+ throw new InvalidEndpoint("malformed dataCollectorBaseUrl: " + dataCollectorBaseUrl);
+ }
+ }
+
+ validateCustomHeaders(customHeaders);
+
+ if (wasmEvaluatorPoolSize != null && wasmEvaluatorPoolSize < 1) {
throw new InvalidOptions("wasmEvaluatorPoolSize must be at least 1");
}
- if (getExporterMetadata() != null) {
+ if (exporterMetadata != null) {
val acceptableExporterMetadataTypes = List.of("String", "Boolean", "Integer", "Double");
- for (Map.Entry entry : getExporterMetadata().entrySet()) {
+ for (Map.Entry entry : exporterMetadata.entrySet()) {
if (!acceptableExporterMetadataTypes.contains(
entry.getValue().getClass().getSimpleName())) {
throw new InvalidExporterMetadata(
@@ -130,4 +250,42 @@ public void validate() throws InvalidOptions {
}
}
}
+
+ /**
+ * validateCustomHeaders rejects, at construction, a custom header the provider would otherwise
+ * send wrongly or fail on at every request. The messages name the header but never carry its
+ * value, which is typically a gateway credential.
+ */
+ private static void validateCustomHeaders(final Map headers) throws InvalidOptions {
+ if (headers == null || headers.isEmpty()) {
+ return;
+ }
+
+ val probe = HttpRequest.newBuilder();
+ val headersSet = new HashSet();
+ for (Map.Entry header : headers.entrySet()) {
+ val name = header.getKey();
+ if (name == null) {
+ throw new InvalidOptions("customHeaders cannot contain a null header name");
+ }
+ if (Const.HTTP_HEADER_CONTENT_TYPE.equalsIgnoreCase(name)
+ || Const.HTTP_HEADER_IF_NONE_MATCH.equalsIgnoreCase(name)) {
+ throw new InvalidOptions("custom header " + name + " is set by the provider itself");
+ }
+ try {
+ probe.header(name, "value");
+ } catch (IllegalArgumentException e) {
+ throw new InvalidOptions("invalid custom header name, " + e.getMessage());
+ }
+
+ if (header.getValue() == null) {
+ throw new InvalidOptions("null value for header: " + name);
+ }
+
+ val setSuccess = headersSet.add(name.toLowerCase());
+ if (!setSuccess) {
+ throw new InvalidOptions("more than one header configured with the name, " + name);
+ }
+ }
+ }
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/GoFeatureFlagApi.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/GoFeatureFlagApi.java
index 591a05338e..1ee692f25f 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/GoFeatureFlagApi.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/GoFeatureFlagApi.java
@@ -1,28 +1,23 @@
package dev.openfeature.contrib.providers.gofeatureflag.api;
import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationFeature;
import com.fasterxml.jackson.databind.SerializationFeature;
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import dev.openfeature.contrib.providers.gofeatureflag.GoFeatureFlagProviderOptions;
import dev.openfeature.contrib.providers.gofeatureflag.api.bean.ExporterRequest;
import dev.openfeature.contrib.providers.gofeatureflag.api.bean.FlagConfigApiRequest;
import dev.openfeature.contrib.providers.gofeatureflag.api.bean.FlagConfigApiResponse;
-import dev.openfeature.contrib.providers.gofeatureflag.api.bean.OfrepRequest;
-import dev.openfeature.contrib.providers.gofeatureflag.api.bean.OfrepResponse;
import dev.openfeature.contrib.providers.gofeatureflag.bean.FlagConfigResponse;
-import dev.openfeature.contrib.providers.gofeatureflag.bean.GoFeatureFlagResponse;
import dev.openfeature.contrib.providers.gofeatureflag.bean.IEvent;
+import dev.openfeature.contrib.providers.gofeatureflag.exception.AuthenticationFailure;
import dev.openfeature.contrib.providers.gofeatureflag.exception.FlagConfigurationEndpointNotFound;
import dev.openfeature.contrib.providers.gofeatureflag.exception.ImpossibleToRetrieveConfiguration;
import dev.openfeature.contrib.providers.gofeatureflag.exception.ImpossibleToSendEventsException;
import dev.openfeature.contrib.providers.gofeatureflag.exception.InvalidEndpoint;
import dev.openfeature.contrib.providers.gofeatureflag.exception.InvalidOptions;
import dev.openfeature.contrib.providers.gofeatureflag.util.Const;
-import dev.openfeature.sdk.EvaluationContext;
-import dev.openfeature.sdk.exceptions.FlagNotFoundError;
import dev.openfeature.sdk.exceptions.GeneralError;
-import dev.openfeature.sdk.exceptions.InvalidContextError;
-import dev.openfeature.sdk.exceptions.OpenFeatureError;
import java.io.IOException;
import java.net.HttpURLConnection;
import java.net.URI;
@@ -37,6 +32,7 @@
import java.util.List;
import java.util.Locale;
import java.util.Map;
+import java.util.Optional;
import lombok.Builder;
import lombok.extern.slf4j.Slf4j;
import lombok.val;
@@ -55,9 +51,15 @@ public final class GoFeatureFlagApi {
/** endpoint is the endpoint of the GO Feature Flag relay proxy. */
private final URI endpoint;
+ /** dataCollectorBaseUrl is the base of the data collector route, the endpoint unless overridden. */
+ private final URI dataCollectorBaseUrl;
+
/** timeout is the timeout in milliseconds for the HTTP requests. */
private final int timeout;
+ /** customHeaders are the extra headers added to every request, before the provider's own. */
+ private final Map customHeaders;
+
/**
* GoFeatureFlagController is the constructor of the controller to contact the GO Feature Flag
* relay proxy.
@@ -72,9 +74,11 @@ private GoFeatureFlagApi(final GoFeatureFlagProviderOptions options) throws Inva
}
options.validate();
this.apiKey = options.getApiKey();
+ this.customHeaders = Map.copyOf(options.getCustomHeaders());
try {
- this.endpoint = new URI(options.getEndpoint());
+ this.endpoint = asBaseUri(options.getEndpoint());
+ this.dataCollectorBaseUrl = asBaseUri(options.getDataCollectorBaseUrl());
} catch (URISyntaxException e) {
throw new InvalidEndpoint(e);
}
@@ -90,110 +94,39 @@ private GoFeatureFlagApi(final GoFeatureFlagProviderOptions options) throws Inva
.build();
}
- /**
- * evaluateFlag is calling the GO Feature Flag relay proxy to evaluate the feature flag.
- *
- * @param key - name of the flag
- * @param evaluationContext - context of the evaluation
- * @return EvaluationResponse with the evaluation of the flag
- * @throws OpenFeatureError - if an error occurred while evaluating the flag
- */
- public GoFeatureFlagResponse evaluateFlag(final String key, final EvaluationContext evaluationContext)
- throws OpenFeatureError {
- return this.evaluateFlag(key, evaluationContext, 0);
- }
-
- /**
- * evaluateFlag is calling the GO Feature Flag relay proxy to evaluate the feature flag.\
- * It will retry once if the relay proxy is unavailable.
- *
- * @param key - name of the flag
- * @param evaluationContext - context of the evaluation
- * @param retryCount - number of retries already done
- * @return EvaluationResponse with the evaluation of the flag
- * @throws OpenFeatureError - if an error occurred while evaluating the flag
- */
- private GoFeatureFlagResponse evaluateFlag(
- final String key, final EvaluationContext evaluationContext, final int retryCount) throws OpenFeatureError {
- try {
- URI url = this.endpoint.resolve("/ofrep/v1/evaluate/flags/" + key);
-
- val requestBody = OfrepRequest.builder()
- .context(evaluationContext.asObjectMap())
- .build();
-
- HttpRequest request = prepareHttpRequest(url, requestBody);
-
- HttpResponse response = this.httpClient.send(request, HttpResponse.BodyHandlers.ofString());
- String body = response.body();
-
- switch (response.statusCode()) {
- case HttpURLConnection.HTTP_OK:
- val goffResp = Const.DESERIALIZE_OBJECT_MAPPER.readValue(body, OfrepResponse.class);
- return goffResp.toGoFeatureFlagResponse();
- case HttpURLConnection.HTTP_UNAUTHORIZED:
- case HttpURLConnection.HTTP_FORBIDDEN:
- throw new GeneralError("authentication/authorization error");
- case HttpURLConnection.HTTP_BAD_REQUEST:
- throw new InvalidContextError("Invalid context: " + body);
- case HttpURLConnection.HTTP_UNAVAILABLE:
- // If the relay proxy is unavailable, we can retry once.
- if (retryCount < 1) {
- log.warn("GO Feature Flag relay proxy is unavailable, retrying evaluation for flag: {}", key);
- return this.evaluateFlag(key, evaluationContext, retryCount + 1);
- }
- throw new GeneralError("Service Unavailable: " + body);
- case HttpURLConnection.HTTP_NOT_FOUND:
- throw new FlagNotFoundError("Flag " + key + " not found");
- default:
- throw new GeneralError("Unknown error while retrieving flag " + body);
- }
- } catch (IOException | InterruptedException e) {
- if (e instanceof InterruptedException) {
- Thread.currentThread().interrupt();
- }
- throw new GeneralError("unknown error while retrieving flag " + key, e);
- }
- }
-
/**
* retrieveFlagConfiguration is calling the GO Feature Flag relay proxy to retrieve the flags'
* configuration.
*
- * @param etag - etag of the request
- * @return FlagConfigResponse with the flag configuration
+ *
A {@code 304 Not Modified} response is reported as an empty Optional rather than as an
+ * empty configuration object, so that the not-modified branch is structurally incapable of
+ * carrying a configuration and cannot be mistaken for one downstream.
+ *
+ * @param etag - etag of the request
+ * @param flags - flags to retrieve, empty for all of them
+ * @return the flag configuration, or empty if the configuration has not been modified
*/
- public FlagConfigResponse retrieveFlagConfiguration(final String etag, final List flags) {
+ public Optional retrieveFlagConfiguration(final String etag, final List flags) {
try {
val request = new FlagConfigApiRequest(flags == null ? Collections.emptyList() : flags);
- final URI url = this.endpoint.resolve("/v1/flag/configuration");
-
- HttpRequest.Builder reqBuilder =
- HttpRequest.newBuilder().uri(url).header(Const.HTTP_HEADER_CONTENT_TYPE, Const.APPLICATION_JSON);
+ final URI url = route(this.endpoint, Const.PATH_FLAG_CONFIGURATION);
- if (this.apiKey != null && !this.apiKey.isEmpty()) {
- reqBuilder.header(Const.HTTP_HEADER_AUTHORIZATION, Const.BEARER_TOKEN + this.apiKey);
- }
+ final HttpRequest httpRequest = etag != null && !etag.isEmpty()
+ ? prepareHttpRequest(url, request, Const.HTTP_HEADER_IF_NONE_MATCH, etag)
+ : prepareHttpRequest(url, request);
- if (etag != null && !etag.isEmpty()) {
- reqBuilder.header(Const.HTTP_HEADER_IF_NONE_MATCH, etag);
- }
-
- reqBuilder.POST(
- HttpRequest.BodyPublishers.ofByteArray(Const.SERIALIZE_OBJECT_MAPPER.writeValueAsBytes(request)));
-
- HttpResponse response =
- this.httpClient.send(reqBuilder.build(), HttpResponse.BodyHandlers.ofString());
+ HttpResponse response = this.httpClient.send(httpRequest, HttpResponse.BodyHandlers.ofString());
String body = response.body();
switch (response.statusCode()) {
case HttpURLConnection.HTTP_OK:
+ return Optional.of(handleFlagConfigurationSuccess(response, body));
case HttpURLConnection.HTTP_NOT_MODIFIED:
- return handleFlagConfigurationSuccess(response, body);
+ return Optional.empty();
case HttpURLConnection.HTTP_NOT_FOUND:
throw new FlagConfigurationEndpointNotFound();
case HttpURLConnection.HTTP_UNAUTHORIZED:
case HttpURLConnection.HTTP_FORBIDDEN:
- throw new ImpossibleToRetrieveConfiguration(
+ throw new AuthenticationFailure(
"retrieve flag configuration error: authentication/authorization error");
case HttpURLConnection.HTTP_BAD_REQUEST:
throw new ImpossibleToRetrieveConfiguration(
@@ -221,7 +154,7 @@ public FlagConfigResponse retrieveFlagConfiguration(final String etag, final Lis
public void sendEventToDataCollector(final List eventsList, final Map exporterMetadata) {
try {
ExporterRequest requestBody = new ExporterRequest(eventsList, exporterMetadata);
- URI url = this.endpoint.resolve("/v1/data/collector");
+ URI url = route(this.dataCollectorBaseUrl, Const.PATH_DATA_COLLECTOR);
HttpRequest request = prepareHttpRequest(url, requestBody);
@@ -250,28 +183,39 @@ public void sendEventToDataCollector(final List eventsList, final Map response, final String body)
throws JsonProcessingException {
- var result = FlagConfigResponse.builder()
+ // without FAIL_ON_TRAILING_TOKENS, a body cut or garbled after a complete JSON value would
+ // be read as that value, and {"flags":{}}} would wipe every flag and advance the ETag.
+ final FlagConfigApiResponse goffResp = Const.DESERIALIZE_OBJECT_MAPPER
+ .readerFor(FlagConfigApiResponse.class)
+ .with(DeserializationFeature.FAIL_ON_TRAILING_TOKENS)
+ .readValue(body);
+
+ // A 200 that decodes to no flag map is a failed refresh, not an empty configuration:
+ // accepting it would wipe every flag and advance the ETag, making the empty state permanent.
+ // A null evaluationContextEnrichment is NOT the same case - the relay proxy builds that field
+ // from a Go map and a nil map marshals to null - so it is accepted as "no enrichment".
+ if (goffResp == null || goffResp.getFlags() == null) {
+ throw new ImpossibleToRetrieveConfiguration(
+ "retrieve flag configuration error: the response contains no flag map");
+ }
+
+ return FlagConfigResponse.builder()
.etag(response.headers().firstValue(Const.HTTP_HEADER_ETAG).orElse(null))
.lastUpdated(extractLastUpdatedFromHeaders(response))
+ .flags(goffResp.getFlags())
+ .evaluationContextEnrichment(goffResp.getEvaluationContextEnrichment())
.build();
-
- if (response.statusCode() == HttpURLConnection.HTTP_OK) {
- val goffResp = Const.DESERIALIZE_OBJECT_MAPPER.readValue(body, FlagConfigApiResponse.class);
- result.setFlags(goffResp.getFlags());
- result.setEvaluationContextEnrichment(goffResp.getEvaluationContextEnrichment());
- }
-
- return result;
}
/**
@@ -294,6 +238,31 @@ private Date extractLastUpdatedFromHeaders(final HttpResponse response)
}
}
+ /**
+ * route builds the URL of an API route from an arbitrary base, so that the data collector can be
+ * addressed somewhere other than the relay proxy.
+ *
+ * @param base - base URL of the route, already normalised by {@link #asBaseUri(String)}
+ * @param path - route path, relative to the base and without a leading slash
+ * @return the URL to call
+ */
+ private static URI route(final URI base, final String path) {
+ return base.resolve(path.startsWith("/") ? path.substring(1) : path);
+ }
+
+ /**
+ * asBaseUri normalises a configured URL into a base other paths can be resolved against. The
+ * trailing slash makes it directory-like, so resolving a relative path appends to any prefix it
+ * carries instead of replacing it.
+ *
+ * @param url - the configured URL
+ * @return the URL as a base
+ * @throws URISyntaxException - if the URL is not a valid URI
+ */
+ private static URI asBaseUri(final String url) throws URISyntaxException {
+ return new URI(url.endsWith("/") ? url : url + "/");
+ }
+
/**
* prepareHttpRequest is preparing the request to be sent to the GO Feature Flag relay proxy.
*
@@ -302,18 +271,22 @@ private Date extractLastUpdatedFromHeaders(final HttpResponse response)
* @return HttpRequest ready to be sent
* @throws JsonProcessingException - if an error occurred while processing the json
*/
- private HttpRequest prepareHttpRequest(final URI url, final T requestBody) throws JsonProcessingException {
+ private HttpRequest prepareHttpRequest(final URI url, final T requestBody, final String... requestHeaders)
+ throws JsonProcessingException {
HttpRequest.Builder reqBuilder = HttpRequest.newBuilder()
.uri(url)
.timeout(Duration.ofMillis(timeout))
- .header(Const.HTTP_HEADER_CONTENT_TYPE, Const.APPLICATION_JSON)
.POST(HttpRequest.BodyPublishers.ofByteArray(
Const.SERIALIZE_OBJECT_MAPPER.writeValueAsBytes(requestBody)));
+ this.customHeaders.forEach(reqBuilder::header);
+ reqBuilder.setHeader(Const.HTTP_HEADER_CONTENT_TYPE, Const.APPLICATION_JSON);
+ for (int i = 0; i + 1 < requestHeaders.length; i += 2) {
+ reqBuilder.setHeader(requestHeaders[i], requestHeaders[i + 1]);
+ }
if (this.apiKey != null && !this.apiKey.isEmpty()) {
- reqBuilder.header(Const.HTTP_HEADER_AUTHORIZATION, Const.BEARER_TOKEN + this.apiKey);
+ reqBuilder.setHeader(Const.HTTP_HEADER_API_KEY, this.apiKey);
}
-
return reqBuilder.build();
}
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/FlagConfigApiResponse.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/FlagConfigApiResponse.java
index 4d5c9d6bd9..5125779093 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/FlagConfigApiResponse.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/FlagConfigApiResponse.java
@@ -1,17 +1,21 @@
package dev.openfeature.contrib.providers.gofeatureflag.api.bean;
import com.fasterxml.jackson.annotation.JsonProperty;
-import dev.openfeature.contrib.providers.gofeatureflag.bean.Flag;
+import com.fasterxml.jackson.databind.JsonNode;
import java.util.Map;
import lombok.Data;
/**
* Represents the response body for the flag configuration API.
+ *
+ *
Flags are kept as raw JSON: the evaluation engine owns the flag schema, so deserialising it
+ * into a typed model here would silently drop any field a newer engine adds and hand a truncated
+ * flag to evaluation.
*/
@Data
public class FlagConfigApiResponse {
@JsonProperty("flags")
- private Map flags;
+ private Map flags;
@JsonProperty("evaluationContextEnrichment")
private Map evaluationContextEnrichment;
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/OfrepRequest.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/OfrepRequest.java
deleted file mode 100644
index 479e5f79e8..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/OfrepRequest.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.api.bean;
-
-import java.util.Map;
-import lombok.Builder;
-import lombok.Data;
-
-/**
- * Represents the request body for the OFREP API request.
- */
-@Data
-@Builder
-public class OfrepRequest {
- private Map context;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/OfrepResponse.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/OfrepResponse.java
deleted file mode 100644
index e1246aa26b..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/api/bean/OfrepResponse.java
+++ /dev/null
@@ -1,53 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.api.bean;
-
-import dev.openfeature.contrib.providers.gofeatureflag.bean.GoFeatureFlagResponse;
-import java.util.Map;
-import lombok.Data;
-import lombok.val;
-
-/**
- * This class represents the response from an OFREP response.
- */
-@Data
-public class OfrepResponse {
- private Object value;
- private String key;
- private String variant;
- private String reason;
- private boolean cacheable;
- private Map metadata;
-
- private String errorCode;
- private String errorDetails;
-
- /**
- * Converts the OFREP response to a GO Feature Flag response.
- *
- * @return the converted GO Feature Flag response
- */
- public GoFeatureFlagResponse toGoFeatureFlagResponse() {
- val goff = new GoFeatureFlagResponse();
- goff.setValue(value);
- goff.setVariationType(variant);
- goff.setReason(reason);
- goff.setErrorCode(errorCode);
- goff.setErrorDetails(errorDetails);
- goff.setFailed(errorCode != null);
-
- if (metadata != null) {
- val cacheable = metadata.get("gofeatureflag_cacheable");
- if (cacheable instanceof Boolean) {
- goff.setCacheable((Boolean) cacheable);
- metadata.remove("gofeatureflag_cacheable");
- }
-
- val version = metadata.get("gofeatureflag_version");
- if (version instanceof String) {
- goff.setVersion((String) version);
- metadata.remove("gofeatureflag_version");
- }
- goff.setMetadata(metadata);
- }
- return goff;
- }
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ExperimentationRollout.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ExperimentationRollout.java
deleted file mode 100644
index 8e6b3052a9..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ExperimentationRollout.java
+++ /dev/null
@@ -1,13 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.bean;
-
-import java.util.Date;
-import lombok.Data;
-
-/**
- * This class represents the rollout of an experimentation.
- */
-@Data
-public class ExperimentationRollout {
- private Date start;
- private Date end;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FeatureEvent.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FeatureEvent.java
index ea8c48224d..a3fa01b2ed 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FeatureEvent.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FeatureEvent.java
@@ -18,6 +18,7 @@ public class FeatureEvent implements IEvent {
private String key;
private String kind;
+ private String source;
private String userKey;
private Object value;
private String variation;
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/Flag.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/Flag.java
deleted file mode 100644
index e6ed6126fd..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/Flag.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.bean;
-
-import java.util.List;
-import lombok.Data;
-import lombok.EqualsAndHashCode;
-
-/**
- * Flag is a class that represents a feature flag for GO Feature Flag.
- */
-@EqualsAndHashCode(callSuper = true)
-@Data
-public class Flag extends FlagBase {
- private List scheduledRollout;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FlagBase.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FlagBase.java
deleted file mode 100644
index ad301eadd8..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FlagBase.java
+++ /dev/null
@@ -1,21 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.bean;
-
-import java.util.List;
-import java.util.Map;
-import lombok.Data;
-
-/**
- * FlagBase is a class that represents the base structure of a feature flag for GO Feature Flag.
- */
-@Data
-public abstract class FlagBase {
- private Map variations;
- private List targeting;
- private String bucketingKey;
- private Rule defaultRule;
- private ExperimentationRollout experimentation;
- private Boolean trackEvents;
- private Boolean disable;
- private String version;
- private Map metadata;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FlagConfigResponse.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FlagConfigResponse.java
index a373135aa1..4143f6e6cc 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FlagConfigResponse.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/FlagConfigResponse.java
@@ -1,5 +1,6 @@
package dev.openfeature.contrib.providers.gofeatureflag.bean;
+import com.fasterxml.jackson.databind.JsonNode;
import java.util.Date;
import java.util.Map;
import lombok.Builder;
@@ -11,7 +12,7 @@
@Data
@Builder
public class FlagConfigResponse {
- private Map flags;
+ private Map flags;
private Map evaluationContextEnrichment;
private String etag;
private Date lastUpdated;
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ProgressiveRollout.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ProgressiveRollout.java
deleted file mode 100644
index 6c39471e13..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ProgressiveRollout.java
+++ /dev/null
@@ -1,12 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.bean;
-
-import lombok.Data;
-
-/**
- * ProgressiveRollout is a class that represents the progressive rollout of a feature flag.
- */
-@Data
-public class ProgressiveRollout {
- private ProgressiveRolloutStep initial;
- private ProgressiveRolloutStep end;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ProgressiveRolloutStep.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ProgressiveRolloutStep.java
deleted file mode 100644
index 843becfd75..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ProgressiveRolloutStep.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.bean;
-
-import java.util.Date;
-import lombok.Data;
-
-/**
- * ProgressiveRolloutStep is a class that represents a step in the progressive rollout of a feature flag.
- */
-@Data
-public class ProgressiveRolloutStep {
- private String variation;
- private Float percentage;
- private Date date;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/Rule.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/Rule.java
deleted file mode 100644
index dc7a59ad28..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/Rule.java
+++ /dev/null
@@ -1,17 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.bean;
-
-import java.util.Map;
-import lombok.Data;
-
-/**
- * This class represents a rule in the GO Feature Flag system.
- */
-@Data
-public class Rule {
- private String name;
- private String query;
- private String variation;
- private Map percentage;
- private Boolean disable;
- private ProgressiveRollout progressiveRollout;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ScheduledStep.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ScheduledStep.java
deleted file mode 100644
index ca9a5d219b..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/ScheduledStep.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.bean;
-
-import java.util.Date;
-import lombok.Data;
-import lombok.EqualsAndHashCode;
-
-/**
- * ScheduledStep is a class that represents a scheduled step in the rollout of a feature flag.
- */
-@EqualsAndHashCode(callSuper = true)
-@Data
-public class ScheduledStep extends FlagBase {
- private Date date;
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/TrackingEvent.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/TrackingEvent.java
index b03fb02bb8..49299681d2 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/TrackingEvent.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/bean/TrackingEvent.java
@@ -33,7 +33,7 @@ public class TrackingEvent implements IEvent {
private String userKey;
/**
- * CreationDate When the feature flag was requested at Unix epoch time in milliseconds.
+ * CreationDate When the event happened, at Unix epoch time in seconds.
*/
private Long creationDate;
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/IEvaluator.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/IEvaluator.java
index 78f1d7aed7..aeabe6fd01 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/IEvaluator.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/IEvaluator.java
@@ -1,7 +1,8 @@
package dev.openfeature.contrib.providers.gofeatureflag.evaluator;
-import dev.openfeature.contrib.providers.gofeatureflag.bean.GoFeatureFlagResponse;
import dev.openfeature.sdk.EvaluationContext;
+import dev.openfeature.sdk.ProviderEvaluation;
+import dev.openfeature.sdk.Value;
/**
* IEvaluator is an interface that represents the evaluation of a feature flag.
@@ -9,27 +10,79 @@
*/
public interface IEvaluator {
/**
- * Initialize the evaluator.
+ * initialize is called when the provider is initialized for a specific domain.
+ *
+ * @param ctx - evaluation context
+ * @param domain - domain the provider is bound to
+ * @throws Exception - if the evaluator cannot be initialized
+ */
+ void initialize(final EvaluationContext ctx, final String domain) throws Exception;
+
+ /**
+ * initialize is called when the provider is initialized.
+ *
+ * @param ctx - evaluation context
+ * @throws Exception - if the evaluator cannot be initialized
*/
- void init();
+ void initialize(final EvaluationContext ctx) throws Exception;
/**
- * Destroy the evaluator.
+ * shutdown releases everything the evaluator holds, so that it stops doing background work.
+ */
+ void shutdown();
+
+ /**
+ * getBooleanEvaluation resolves the value of a boolean flag.
+ *
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param ctx - evaluation context
+ * @return the evaluation result for this flag
+ */
+ ProviderEvaluation getBooleanEvaluation(String key, Boolean defaultValue, EvaluationContext ctx);
+
+ /**
+ * getStringEvaluation resolves the value of a string flag.
+ *
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param ctx - evaluation context
+ * @return the evaluation result for this flag
+ */
+ ProviderEvaluation getStringEvaluation(String key, String defaultValue, EvaluationContext ctx);
+
+ /**
+ * getIntegerEvaluation resolves the value of an integer flag.
+ *
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param ctx - evaluation context
+ * @return the evaluation result for this flag
+ */
+ ProviderEvaluation getIntegerEvaluation(String key, Integer defaultValue, EvaluationContext ctx);
+
+ /**
+ * getDoubleEvaluation resolves the value of a float flag.
+ *
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param ctx - evaluation context
+ * @return the evaluation result for this flag
*/
- void destroy();
+ ProviderEvaluation getDoubleEvaluation(String key, Double defaultValue, EvaluationContext ctx);
/**
- * Evaluate the flag.
+ * getObjectEvaluation resolves the value of a flag holding a structure.
*
- * @param key - name of the flag
- * @param defaultValue - default value
- * @param evaluationContext - evaluation context
- * @return the evaluation response
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param ctx - evaluation context
+ * @return the evaluation result for this flag
*/
- GoFeatureFlagResponse evaluate(String key, Object defaultValue, EvaluationContext evaluationContext);
+ ProviderEvaluation getObjectEvaluation(String key, Value defaultValue, EvaluationContext ctx);
/**
- * Check if the flag is trackable or not.
+ * isFlagTrackable returns true if we should collect the usage of this flag.
*
* @param flagKey - name of the flag
* @return true if the flag is trackable, false otherwise
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/InProcessEvaluator.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/InProcessEvaluator.java
index 4a63cb3adf..f81132f0ee 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/InProcessEvaluator.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/InProcessEvaluator.java
@@ -1,28 +1,45 @@
package dev.openfeature.contrib.providers.gofeatureflag.evaluator;
+import static dev.openfeature.sdk.Value.objectToValue;
+
+import com.fasterxml.jackson.databind.JsonNode;
import dev.openfeature.contrib.providers.gofeatureflag.GoFeatureFlagProviderOptions;
import dev.openfeature.contrib.providers.gofeatureflag.api.GoFeatureFlagApi;
-import dev.openfeature.contrib.providers.gofeatureflag.bean.Flag;
import dev.openfeature.contrib.providers.gofeatureflag.bean.FlagConfigResponse;
import dev.openfeature.contrib.providers.gofeatureflag.bean.GoFeatureFlagResponse;
import dev.openfeature.contrib.providers.gofeatureflag.util.Const;
+import dev.openfeature.contrib.providers.gofeatureflag.util.JsonValueUtil;
+import dev.openfeature.contrib.providers.gofeatureflag.util.MetadataUtil;
import dev.openfeature.contrib.providers.gofeatureflag.wasm.WasmEvaluatorPool;
import dev.openfeature.contrib.providers.gofeatureflag.wasm.bean.FlagContext;
import dev.openfeature.contrib.providers.gofeatureflag.wasm.bean.WasmInput;
import dev.openfeature.sdk.ErrorCode;
import dev.openfeature.sdk.EvaluationContext;
+import dev.openfeature.sdk.ProviderEvaluation;
+import dev.openfeature.sdk.ProviderEvent;
import dev.openfeature.sdk.ProviderEventDetails;
import dev.openfeature.sdk.Reason;
+import dev.openfeature.sdk.Value;
+import dev.openfeature.sdk.exceptions.ExceptionUtils;
+import dev.openfeature.sdk.exceptions.TypeMismatchError;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import io.reactivex.rxjava3.core.Observable;
import io.reactivex.rxjava3.disposables.Disposable;
import io.reactivex.rxjava3.schedulers.Schedulers;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Date;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.TimeUnit;
-import java.util.function.Consumer;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.BiConsumer;
import lombok.extern.slf4j.Slf4j;
import lombok.val;
@@ -32,61 +49,105 @@
*/
@Slf4j
public class InProcessEvaluator implements IEvaluator {
+ /**
+ * Raw engine error codes that make the engine, rather than the flag, the suspect: the relay proxy
+ * runs the same engine against the same flag and may well answer correctly. FLAG_CONFIG is
+ * deliberately absent, being a misconfiguration the relay proxy would reproduce identically.
+ */
+ private static final Set FALLBACK_TRIGGERS = Set.of(ErrorCode.PARSE_ERROR.name(), ErrorCode.GENERAL.name());
+
/** API to contact GO Feature Flag. */
private final GoFeatureFlagApi api;
/** Pool of WASM evaluation engine instances for thread-safe concurrent evaluation. */
- private final WasmEvaluatorPool evaluationPool;
+ private volatile WasmEvaluatorPool evaluationPool;
/** Options to configure the provider. */
private final GoFeatureFlagProviderOptions options;
- /** Method to call when we have a configuration change. */
- private final Consumer emitProviderConfigurationChanged;
+ /** Method to call to emit a provider event to the SDK. */
+ private final BiConsumer emitter;
/** Immutable snapshot of all flag configuration state; updated atomically by the polling daemon. */
private volatile EvaluatorState state;
/** disposable which manage the polling of the flag configurations. */
private Disposable configurationDisposable;
+ /** Number of refreshes that have failed in a row, reset by any successful one. */
+ private final AtomicInteger consecutiveRefreshFailures = new AtomicInteger();
+ /** true between the announcement of a stale configuration and the announcement of its end. */
+ private final AtomicBoolean staleAnnounced = new AtomicBoolean();
+ /** Builds the evaluator a failed local evaluation falls back to. */
+ private final IEvaluator fallbackEvaluator;
private static final class EvaluatorState {
- final Map flags;
+ final Map flags;
final Map evaluationContextEnrichment;
final String etag;
final Date lastUpdate;
+ /**
+ * true once a configuration has been retrieved and applied. A configuration carrying an empty
+ * flag map counts as loaded: an empty configuration is a valid one, not a missing one.
+ * The marker lives in the snapshot rather than beside it, so that readiness and the flag map
+ * it describes can never be read out of step.
+ */
+ final boolean configurationLoaded;
+
+ /** the state held until a configuration has been applied for the first time. */
+ static EvaluatorState notLoaded() {
+ return new EvaluatorState(Collections.emptyMap(), null, "", new Date(0), false);
+ }
EvaluatorState(
- Map flags,
+ Map flags,
Map evaluationContextEnrichment,
String etag,
Date lastUpdate) {
+ // only reachable from a retrieved configuration, so it is loaded by construction
+ this(flags, evaluationContextEnrichment, etag, lastUpdate, true);
+ }
+
+ private EvaluatorState(
+ Map flags,
+ Map evaluationContextEnrichment,
+ String etag,
+ Date lastUpdate,
+ boolean configurationLoaded) {
this.flags = flags;
this.evaluationContextEnrichment = evaluationContextEnrichment;
this.etag = etag;
this.lastUpdate = lastUpdate;
+ this.configurationLoaded = configurationLoaded;
}
}
/**
* Constructor of the InProcessEvaluator.
*
- * @param api - API to contact GO Feature Flag
- * @param options - options to configure the provider
- * @param emitProviderConfigurationChanged - method to call when we have a configuration change
+ * @param api - API to contact GO Feature Flag
+ * @param options - options to configure the provider
+ * @param emitter - method to call to emit a provider event to the SDK
*/
public InProcessEvaluator(
GoFeatureFlagApi api,
GoFeatureFlagProviderOptions options,
- Consumer emitProviderConfigurationChanged) {
+ BiConsumer emitter) {
this.api = api;
this.options = options;
- this.emitProviderConfigurationChanged = emitProviderConfigurationChanged;
- this.state = new EvaluatorState(Collections.emptyMap(), null, "", new Date(0));
- int poolSize = options.getWasmEvaluatorPoolSize() != null
- ? options.getWasmEvaluatorPoolSize()
- : Const.DEFAULT_WASM_EVALUATOR_POOL_SIZE;
- this.evaluationPool = new WasmEvaluatorPool(poolSize);
+ this.emitter = emitter;
+ this.fallbackEvaluator = new RemoteEvaluator(options, emitter);
+ this.state = EvaluatorState.notLoaded();
+ this.evaluationPool = new WasmEvaluatorPool(options.getWasmEvaluatorPoolSize());
}
- @Override
- public GoFeatureFlagResponse evaluate(String key, Object defaultValue, EvaluationContext evaluationContext) {
+ private GoFeatureFlagResponse evaluate(String key, Object defaultValue, EvaluationContext evaluationContext) {
EvaluatorState current = this.state;
+ // With no configuration ever loaded every key is absent, and answering FLAG_NOT_FOUND would
+ // blame the caller's flag key for an infrastructure failure.
+ if (!current.configurationLoaded) {
+ val notReady = new GoFeatureFlagResponse();
+ notReady.setReason(Reason.ERROR.name());
+ notReady.setErrorCode(ErrorCode.PROVIDER_NOT_READY.name());
+ notReady.setErrorDetails(
+ "impossible to evaluate flag " + key + ": no flag configuration has been loaded yet");
+ notReady.setValue(defaultValue);
+ return notReady;
+ }
if (current.flags.get(key) == null) {
val err = new GoFeatureFlagResponse();
err.setReason(Reason.ERROR.name());
@@ -108,28 +169,300 @@ public GoFeatureFlagResponse evaluate(String key, Object defaultValue, Evaluatio
}
@Override
- public boolean isFlagTrackable(final String flagKey) {
- Flag flag = this.state.flags.get(flagKey);
- return flag != null && (flag.getTrackEvents() == null || flag.getTrackEvents());
+ public void initialize(EvaluationContext ctx, String domain) throws Exception {
+ this.initialize(ctx);
}
@Override
- public void init() {
- val configFlags = api.retrieveFlagConfiguration(this.state.etag, options.getEvaluationFlagList());
- this.state = new EvaluatorState(
- configFlags.getFlags(),
- configFlags.getEvaluationContextEnrichment(),
- configFlags.getEtag(),
- configFlags.getLastUpdated());
+ public void initialize(EvaluationContext ctx) throws Exception {
+ // We ensure that no polling is happening before starting the initialization.
+ stopPolling();
+
+ // shutdown() closes the pool and the fallback, and the provider reuses this evaluator on re-init
+ if (this.evaluationPool.isClosed()) {
+ this.evaluationPool = new WasmEvaluatorPool(options.getWasmEvaluatorPoolSize());
+ }
+ this.fallbackEvaluator.initialize(ctx);
+
+ // an empty response means the configuration has not been modified, so the state we already
+ // hold is still current and must not be overwritten.
+ api.retrieveFlagConfiguration(this.state.etag, options.getEvaluationFlagList())
+ .ifPresent(configFlags -> this.state = new EvaluatorState(
+ configFlags.getFlags(),
+ configFlags.getEvaluationContextEnrichment(),
+ configFlags.getEtag(),
+ configFlags.getLastUpdated()));
+ // this fetch is a successful refresh too, and it bypasses the polling chain entirely. The
+ // SDK announces readiness itself once initialize() returns, so nothing is emitted here.
+ this.consecutiveRefreshFailures.set(0);
+ this.staleAnnounced.set(false);
// start the polling of the flag configuration
this.configurationDisposable = startCheckFlagConfigurationChangesDaemon();
}
@Override
- public void destroy() {
+ public void shutdown() {
+ stopPolling();
+ this.fallbackEvaluator.shutdown();
+ this.evaluationPool.close();
+ }
+
+ @Override
+ public ProviderEvaluation getBooleanEvaluation(String key, Boolean defaultValue, EvaluationContext ctx) {
+ return genericEvaluation(key, defaultValue, ctx, Boolean.class, IEvaluator::getBooleanEvaluation);
+ }
+
+ @Override
+ public ProviderEvaluation getStringEvaluation(String key, String defaultValue, EvaluationContext ctx) {
+ return genericEvaluation(key, defaultValue, ctx, String.class, IEvaluator::getStringEvaluation);
+ }
+
+ @Override
+ public ProviderEvaluation getIntegerEvaluation(String key, Integer defaultValue, EvaluationContext ctx) {
+ return genericEvaluation(key, defaultValue, ctx, Integer.class, IEvaluator::getIntegerEvaluation);
+ }
+
+ @Override
+ public ProviderEvaluation getDoubleEvaluation(String key, Double defaultValue, EvaluationContext ctx) {
+ return genericEvaluation(key, defaultValue, ctx, Double.class, IEvaluator::getDoubleEvaluation);
+ }
+
+ @Override
+ public ProviderEvaluation getObjectEvaluation(String key, Value defaultValue, EvaluationContext ctx) {
+ return genericEvaluation(key, defaultValue, ctx, Value.class, IEvaluator::getObjectEvaluation);
+ }
+
+ /**
+ * genericEvaluation evaluates the flag and converts the engine response into the resolution
+ * structure expected by the OpenFeature SDK.
+ *
+ *
When the engine reports a failure that points at this provider rather than at the flag, the
+ * relay proxy is asked instead and its answer is the one the caller receives.
+ *
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param ctx - evaluation context
+ * @param expectedType - type the resolver called by the SDK is contracted to return
+ * @param remoteResolver - resolver of the fallback evaluator matching the type asked for
+ * @param - type of the flag value
+ * @return the evaluation result for this flag
+ */
+ private ProviderEvaluation genericEvaluation(
+ final String key,
+ final T defaultValue,
+ final EvaluationContext ctx,
+ final Class> expectedType,
+ final RemoteResolver remoteResolver) {
+ val response = this.evaluate(key, defaultValue, ctx);
+ // the raw code the engine emitted, read before toProviderEvaluation maps it onto the SDK
+ // enumeration, where GO Feature Flag's own codes are folded into GENERAL and stop being
+ // distinguishable from an engine failure.
+ if (FALLBACK_TRIGGERS.contains(response.getErrorCode())) {
+ log.warn(
+ "the engine could not evaluate flag {} ({}: {}), asking the relay proxy instead",
+ key,
+ response.getErrorCode(),
+ response.getErrorDetails());
+ val remote = evaluateRemotely(key, defaultValue, ctx, remoteResolver);
+ if (remote.isPresent()) {
+ return remote.get();
+ }
+ }
+ return toProviderEvaluation(key, defaultValue, response, expectedType);
+ }
+
+ /**
+ * evaluateRemotely asks the relay proxy about a flag the engine could not evaluate.
+ *
+ *
A TYPE_MISMATCH is kept: the relay proxy did evaluate the flag, and its value does not fit the
+ * type asked for. An empty result means the relay proxy could not answer either, in which case the
+ * caller is owed the engine's error rather than the proxy's: the engine failing is the root cause,
+ * and the proxy merely failed to make up for it. The remote failure is logged here because it is about to
+ * disappear from the answer entirely.
+ *
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param ctx - evaluation context
+ * @param remoteResolver - resolver of the fallback evaluator matching the type asked for
+ * @param - type of the flag value
+ * @return the relay proxy's answer, or empty if it could not give one
+ */
+ private Optional> evaluateRemotely(
+ final String key,
+ final T defaultValue,
+ final EvaluationContext ctx,
+ final RemoteResolver remoteResolver) {
+ try {
+ val remote = remoteResolver.resolve(this.fallbackEvaluator, key, defaultValue, ctx);
+ if (remote.getErrorCode() == null || remote.getErrorCode() == ErrorCode.TYPE_MISMATCH) {
+ return Optional.of(markEvaluatedRemotely(remote));
+ }
+ log.error(
+ "the relay proxy could not evaluate flag {} either: {} {}",
+ key,
+ remote.getErrorCode(),
+ remote.getErrorMessage());
+ } catch (Exception e) {
+ // the OFREP client raises on responses it cannot read at all, and an exception escaping
+ // here would replace the engine's error with one about the recovery attempt.
+ log.error("the relay proxy could not be asked about flag {}", key, e);
+ }
+ return Optional.empty();
+ }
+
+ /**
+ * markEvaluatedRemotely records in the flag metadata that the relay proxy produced this result.
+ *
+ *
The relay proxy's own metadata keys are kept: they describe the evaluation, which is the
+ * proxy's, and only the marker is this provider's to add.
+ *
+ * @param remote - answer the relay proxy gave
+ * @param - type of the flag value
+ * @return the same evaluation, carrying the marker
+ */
+ private static ProviderEvaluation markEvaluatedRemotely(final ProviderEvaluation remote) {
+ val metadata = new LinkedHashMap();
+ if (remote.getFlagMetadata() != null) {
+ metadata.putAll(remote.getFlagMetadata().asUnmodifiableMap());
+ }
+ metadata.put(Const.METADATA_EVALUATED_REMOTELY, true);
+ remote.setFlagMetadata(MetadataUtil.convertFlagMetadata(metadata));
+ return remote;
+ }
+
+ @FunctionalInterface
+ private interface RemoteResolver {
+ ProviderEvaluation resolve(IEvaluator remote, String key, T defaultValue, EvaluationContext ctx);
+ }
+
+ /**
+ * toProviderEvaluation converts an engine response into the resolution structure expected by
+ * the OpenFeature SDK.
+ *
+ *
It is separate from genericEvaluation so that the conversion can be specified against a
+ * response directly, independently of what the WASM engine can be made to emit.
+ *
+ * @param key - name of the flag
+ * @param defaultValue - default value provided by the caller
+ * @param response - response returned by the evaluation engine
+ * @param expectedType - type the resolver called by the SDK is contracted to return
+ * @param - type of the flag value
+ * @return the evaluation result for this flag
+ */
+ static ProviderEvaluation toProviderEvaluation(
+ final String key, final T defaultValue, final GoFeatureFlagResponse response, final Class> expectedType) {
+ if (response.getErrorCode() != null && !response.getErrorCode().isEmpty()) {
+ throw ExceptionUtils.instantiateErrorByErrorCode(
+ mapErrorCode(response.getErrorCode()), response.getErrorDetails());
+ }
+
+ if (Reason.DISABLED.name().equalsIgnoreCase(response.getReason())) {
+ // we don't set a variant since we are using the default value,
+ // and we are not able to know which variant it is.
+ return ProviderEvaluation.builder()
+ .value(defaultValue)
+ .variant(response.getVariationType())
+ .reason(Reason.DISABLED.name())
+ .flagMetadata(MetadataUtil.convertFlagMetadata(response.getMetadata()))
+ .build();
+ }
+
+ if (response.getValue() == null) {
+ return ProviderEvaluation.builder()
+ .value(defaultValue)
+ .reason(response.getReason())
+ .variant(response.getVariationType())
+ .flagMetadata(MetadataUtil.convertFlagMetadata(response.getMetadata()))
+ .build();
+ }
+
+ T flagValue = convertValue(response.getValue(), expectedType);
+ if (flagValue.getClass() != expectedType) {
+ throw new TypeMismatchError(String.format(
+ "Flag value %s had unexpected type %s, expected %s.", key, flagValue.getClass(), expectedType));
+ }
+
+ return ProviderEvaluation.builder()
+ .reason(response.getReason())
+ .value(flagValue)
+ .variant(response.getVariationType())
+ .flagMetadata(MetadataUtil.convertFlagMetadata(response.getMetadata()))
+ .build();
+ }
+
+ /**
+ * convertValue is converting the value returned by the evaluation engine in the right type.
+ *
+ * @param value - the value we have received
+ * @param expectedType - the type we expect for this value
+ * @param - the type we want to convert to
+ * @return a converted object
+ */
+ // The cast to T is unchecked on purpose: toProviderEvaluation compares the runtime class with
+ // expectedType right after and raises a TypeMismatchError if the engine returned another type.
+ @SuppressWarnings("unchecked")
+ private static T convertValue(final Object value, final Class> expectedType) {
+ boolean isPrimitive = expectedType == Boolean.class
+ || expectedType == String.class
+ || expectedType == Integer.class
+ || expectedType == Double.class;
+
+ if (isPrimitive) {
+ // JSON does not distinguish 100 from 100.0, and a number too large for an int decodes to
+ // Long, so the float resolver accepts any number rather than only Integer.
+ if (expectedType == Double.class && value instanceof Number) {
+ return (T) Double.valueOf(((Number) value).doubleValue());
+ }
+ return (T) value;
+ }
+ return (T) objectToValue(JsonValueUtil.widenBigIntegers(value));
+ }
+
+ /**
+ * mapErrorCode is mapping the error code in string received from the evaluation engine to the
+ * SDK ErrorCode enum.
+ *
+ * @param errorCode - string of the error code received from the evaluation engine
+ * @return an item from the enum, null if the engine reported no error
+ */
+ private static ErrorCode mapErrorCode(final String errorCode) {
+ if (errorCode == null || errorCode.isEmpty()) {
+ return null;
+ }
+
+ try {
+ return ErrorCode.valueOf(errorCode);
+ } catch (IllegalArgumentException e) {
+ // an error the SDK does not know about, such as GO Feature Flag's own FLAG_CONFIG, is
+ // still an error: reporting it as GENERAL keeps the evaluation from looking successful.
+ return ErrorCode.GENERAL;
+ }
+ }
+
+ @Override
+ public boolean isFlagTrackable(final String flagKey) {
+ // trackEvents is the only field of the flag configuration a provider may read: everything
+ // else belongs to the evaluation engine and is passed through untouched.
+ JsonNode flag = this.state.flags.get(flagKey);
+ if (flag == null) {
+ return true;
+ }
+ JsonNode trackEvents = flag.get(Const.FIELD_TRACK_EVENTS);
+ return trackEvents == null || trackEvents.isNull() || trackEvents.asBoolean(true);
+ }
+
+ /**
+ * stopPolling cancels the configuration polling task, if one is running.
+ * Once dispose() returns, the subscription is guaranteed to deliver no further emission to the
+ * refresh consumer, so a request still in flight can no longer reach the configuration state. The
+ * field is cleared so that the disposed subscription cannot be disposed, or mistaken for a live
+ * one, a second time.
+ */
+ private synchronized void stopPolling() {
if (this.configurationDisposable != null) {
this.configurationDisposable.dispose();
+ this.configurationDisposable = null;
}
}
@@ -143,45 +476,181 @@ private Disposable startCheckFlagConfigurationChangesDaemon() {
? options.getFlagChangePollingIntervalMs()
: Const.DEFAULT_POLLING_CONFIG_FLAG_CHANGE_INTERVAL_MS;
- Observable intervalObservable =
- Observable.interval(pollingIntervalMs, TimeUnit.MILLISECONDS, Schedulers.io());
- Observable apiCallObservable = intervalObservable
- .flatMap(tick -> Observable.fromCallable(() ->
- this.api.retrieveFlagConfiguration(this.state.etag, options.getEvaluationFlagList()))
+ Observable pollObservable = Observable.defer(() ->
+ Observable.timer(nextPollDelayMs(pollingIntervalMs), TimeUnit.MILLISECONDS, Schedulers.io()))
+ .repeat();
+ Observable apiCallObservable = pollObservable
+ .flatMap(tick -> Observable.fromCallable(this::refreshFlagConfiguration)
.onErrorResumeNext(e -> {
log.error("error while calling flag configuration API", e);
- return Observable.empty();
+ emitRefreshFailureEvent();
+ return Observable.>empty();
}))
+ // a 304 emits an empty Optional: drop it here so the refresh consumer below is
+ // structurally unable to write state for a response that carries no configuration.
+ .filter(Optional::isPresent)
+ .map(Optional::get)
.subscribeOn(Schedulers.io());
return apiCallObservable.subscribe(
response -> {
- EvaluatorState current = this.state;
- if (response.getEtag().equals(current.etag)) {
- log.debug("flag configuration has not changed: {}", response);
- return;
+ try {
+ applyFlagConfiguration(response);
+ } catch (Exception e) {
+ // an exception escaping the subscriber disposes it, which would stop polling
+ // for the lifetime of the provider rather than for this one refresh.
+ log.error("error while applying the flag configuration", e);
}
+ },
+ throwable -> log.error("flag configuration polling has stopped and will not resume", throwable));
+ }
- if (response.getLastUpdated().before(current.lastUpdate)) {
- log.info("configuration received is older than the current one");
- return;
- }
+ /**
+ * refreshFlagConfiguration fetches the flag configuration once and counts it as a successful refresh.
+ *
+ * @return the configuration, or an empty Optional when the relay proxy answered not modified
+ */
+ private Optional refreshFlagConfiguration() {
+ val configuration = this.api.retrieveFlagConfiguration(this.state.etag, options.getEvaluationFlagList());
+ emitRefreshSuccessEvent();
+ return configuration;
+ }
- log.info("flag configuration has changed");
- val flagChanges = findFlagConfigurationChanges(current.flags, response.getFlags());
- this.state = new EvaluatorState(
- response.getFlags(),
- response.getEvaluationContextEnrichment(),
- response.getEtag(),
- response.getLastUpdated());
- val changeDetails = ProviderEventDetails.builder()
- .flagsChanged(flagChanges)
- .message("flag configuration has changed")
- .build();
- this.emitProviderConfigurationChanged.accept(changeDetails);
- },
- throwable ->
- log.error("error while calling flag configuration API, error: {}", throwable.getMessage()));
+ /**
+ * nextPollDelayMs is the polling interval with jitter applied, so that a fleet restarted together
+ * does not poll the relay proxy in lockstep for as long as it stays up.
+ *
+ * @param pollingIntervalMs - the configured polling interval
+ * @return the interval, randomly shortened or lengthened by up to {@link Const#POLLING_JITTER_RATIO}
+ */
+ @SuppressFBWarnings(value = "PREDICTABLE_RANDOM", justification = "the poll jitter is not security-relevant")
+ static long nextPollDelayMs(final long pollingIntervalMs) {
+ return (long) (pollingIntervalMs
+ * ThreadLocalRandom.current()
+ .nextDouble(1 - Const.POLLING_JITTER_RATIO, 1 + Const.POLLING_JITTER_RATIO));
+ }
+
+ /**
+ * emitRefreshFailureEvent counts a failed refresh and marks the configuration stale once enough of
+ * them have happened in a row.
+ */
+ private void emitRefreshFailureEvent() {
+ if (this.consecutiveRefreshFailures.incrementAndGet() != Const.STALE_AFTER_CONSECUTIVE_FAILURES) {
+ return;
+ }
+
+ log.warn(
+ "{} consecutive failed refreshes, still serving the last known good configuration",
+ Const.STALE_AFTER_CONSECUTIVE_FAILURES);
+ this.staleAnnounced.set(true);
+ try {
+ this.emitter.accept(
+ ProviderEvent.PROVIDER_STALE,
+ ProviderEventDetails.builder()
+ .message("the flag configuration could not be refreshed "
+ + Const.STALE_AFTER_CONSECUTIVE_FAILURES + " times in a row")
+ .build());
+ } catch (Exception e) {
+ log.error("error while emitting the {} event", ProviderEvent.PROVIDER_STALE, e);
+ }
+ }
+
+ /**
+ * emitRefreshSuccessEvent counts a refresh that worked and, if the configuration had been announced
+ * as stale, announces that it no longer is.
+ */
+ private void emitRefreshSuccessEvent() {
+ this.consecutiveRefreshFailures.set(0);
+ if (!this.staleAnnounced.compareAndSet(true, false)) {
+ return;
+ }
+
+ log.info("the flag configuration could be refreshed again");
+ try {
+ this.emitter.accept(
+ ProviderEvent.PROVIDER_READY,
+ ProviderEventDetails.builder()
+ .message("the flag configuration could be refreshed again")
+ .build());
+ } catch (Exception e) {
+ log.error("error while emitting the {} event", ProviderEvent.PROVIDER_READY, e);
+ }
+ }
+
+ /**
+ * applyFlagConfiguration replaces the configuration state with a newly retrieved one, unless the
+ * response describes the configuration already held or an older one.
+ *
+ *
Both validators are optional: a relay proxy behind a cache or a reverse proxy may answer
+ * without an {@code ETag} or with a {@code Last-Modified} this provider cannot parse, so a null
+ * on either side means "cannot rule this response out" rather than a comparison.
+ *
+ * @param response - configuration returned by the last successful refresh
+ */
+ private void applyFlagConfiguration(final FlagConfigResponse response) {
+ EvaluatorState current = this.state;
+ if (response.getEtag() != null && response.getEtag().equals(current.etag)) {
+ log.debug("flag configuration has not changed: {}", response);
+ return;
+ }
+
+ if (response.getLastUpdated() != null
+ && current.lastUpdate != null
+ && response.getLastUpdated().before(current.lastUpdate)) {
+ log.info("configuration received is older than the current one");
+ return;
+ }
+
+ val flagChanges =
+ enrichmentChanged(current.evaluationContextEnrichment, response.getEvaluationContextEnrichment())
+ ? allFlagKeys(current.flags, response.getFlags())
+ : findFlagConfigurationChanges(current.flags, response.getFlags());
+ this.state = new EvaluatorState(
+ response.getFlags(),
+ response.getEvaluationContextEnrichment(),
+ response.getEtag(),
+ response.getLastUpdated());
+
+ if (flagChanges.isEmpty()) {
+ log.debug("flag configuration has not changed: {}", response);
+ return;
+ }
+
+ log.info("flag configuration has changed");
+ val changeDetails = ProviderEventDetails.builder()
+ .flagsChanged(flagChanges)
+ .message("flag configuration has changed")
+ .build();
+ this.emitter.accept(ProviderEvent.PROVIDER_CONFIGURATION_CHANGED, changeDetails);
+ }
+
+ /**
+ * enrichmentChanged reports whether the evaluation context enrichment differs between two
+ * configurations. A null enrichment is the same as an empty one.
+ *
+ * @param original - enrichment currently in use
+ * @param updated - enrichment of the new configuration
+ * @return true if the enrichment has changed
+ */
+ private static boolean enrichmentChanged(final Map original, final Map updated) {
+ return !Optional.ofNullable(original)
+ .orElse(Collections.emptyMap())
+ .equals(Optional.ofNullable(updated).orElse(Collections.emptyMap()));
+ }
+
+ /**
+ * allFlagKeys lists every flag of either configuration: a new enrichment can change the result of
+ * any of them.
+ *
+ * @param originalFlags - list of original flags
+ * @param newFlags - list of new flags
+ * @return the keys of every flag in either configuration
+ */
+ private static List allFlagKeys(
+ final Map originalFlags, final Map newFlags) {
+ Set keys = new LinkedHashSet<>(newFlags.keySet());
+ keys.addAll(originalFlags.keySet());
+ return new ArrayList<>(keys);
}
/**
@@ -192,16 +661,16 @@ private Disposable startCheckFlagConfigurationChangesDaemon() {
* @return - list of flags that have changed
*/
private List findFlagConfigurationChanges(
- final Map originalFlags, final Map newFlags) {
+ final Map originalFlags, final Map newFlags) {
// this function should return a list of flags that have changed between the two maps
// it should contain all updated, added and removed flags
List changedFlags = new ArrayList<>();
// Find added or updated flags
- for (Map.Entry entry : newFlags.entrySet()) {
+ for (Map.Entry entry : newFlags.entrySet()) {
String key = entry.getKey();
- Flag newFlag = entry.getValue();
- Flag originalFlag = originalFlags.get(key);
+ JsonNode newFlag = entry.getValue();
+ JsonNode originalFlag = originalFlags.get(key);
if (originalFlag == null || !originalFlag.equals(newFlag)) {
changedFlags.add(key);
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/RemoteEvaluator.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/RemoteEvaluator.java
index 15b60c5f85..931afdf0b0 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/RemoteEvaluator.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/evaluator/RemoteEvaluator.java
@@ -1,9 +1,23 @@
package dev.openfeature.contrib.providers.gofeatureflag.evaluator;
-import dev.openfeature.contrib.providers.gofeatureflag.api.GoFeatureFlagApi;
-import dev.openfeature.contrib.providers.gofeatureflag.bean.GoFeatureFlagResponse;
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableMap;
+import dev.openfeature.contrib.providers.gofeatureflag.GoFeatureFlagProviderOptions;
+import dev.openfeature.contrib.providers.gofeatureflag.util.Const;
+import dev.openfeature.contrib.providers.ofrep.OfrepProvider;
+import dev.openfeature.contrib.providers.ofrep.OfrepProviderOptions;
+import dev.openfeature.sdk.ErrorCode;
import dev.openfeature.sdk.EvaluationContext;
+import dev.openfeature.sdk.ProviderEvaluation;
+import dev.openfeature.sdk.ProviderEvent;
+import dev.openfeature.sdk.ProviderEventDetails;
+import dev.openfeature.sdk.Value;
+import java.time.Duration;
+import java.util.HashMap;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.function.BiConsumer;
import lombok.extern.slf4j.Slf4j;
+import lombok.val;
/**
* RemoteEvaluator is an implementation of the IEvaluator interface.
@@ -11,21 +25,56 @@
*/
@Slf4j
public class RemoteEvaluator implements IEvaluator {
- /** API to contact GO Feature Flag. */
- public final GoFeatureFlagApi api;
+ /**
+ * Start of the error message the OFREP client reports for a 401 or 403. That client folds every
+ * status into a GENERAL error code, so the message is the only part of the result that still
+ * distinguishes rejected credentials from an ordinary failure.
+ */
+ private static final String OFREP_AUTHENTICATION_ERROR = "authentication/authorization error for flag:";
+
+ /** Options to configure the provider, kept to rebuild the OFREP provider after a shutdown. */
+ private final GoFeatureFlagProviderOptions options;
+ /** OFREP provider doing the actual remote evaluation. */
+ private volatile OfrepProvider ofrep;
+ /** true once shutdown() has stopped the OFREP provider, which cannot be restarted. */
+ private volatile boolean ofrepShutDown;
+ /** Method to call to emit a provider event to the SDK. */
+ private final BiConsumer emitter;
+ /** Guards against re-reporting the same authentication failure on every later evaluation. */
+ private final AtomicBoolean authenticationFailureReported = new AtomicBoolean();
/**
* Constructor of the evaluator.
*
- * @param api - api service to evaluate the flags
+ * @param opts - options to configure the provider
+ * @param emitter - method to call to emit a provider event to the SDK
*/
- public RemoteEvaluator(GoFeatureFlagApi api) {
- this.api = api;
+ public RemoteEvaluator(GoFeatureFlagProviderOptions opts, BiConsumer emitter) {
+ this.options = opts;
+ this.emitter = emitter;
+ this.ofrep = newOfrepProvider(opts);
}
- @Override
- public GoFeatureFlagResponse evaluate(String key, Object defaultValue, EvaluationContext evaluationContext) {
- return this.api.evaluateFlag(key, evaluationContext);
+ /**
+ * newOfrepProvider builds the OFREP provider from the provider options.
+ *
+ * @param opts - options to configure the provider
+ * @return a new OFREP provider
+ */
+ private static OfrepProvider newOfrepProvider(final GoFeatureFlagProviderOptions opts) {
+ val headers = new HashMap>();
+ opts.getCustomHeaders().forEach((name, value) -> headers.put(name, ImmutableList.of(value)));
+ if (opts.getApiKey() != null && !opts.getApiKey().isEmpty()) {
+ headers.keySet().removeIf(Const.HTTP_HEADER_API_KEY::equalsIgnoreCase);
+ headers.put(Const.HTTP_HEADER_API_KEY, ImmutableList.of(opts.getApiKey()));
+ }
+
+ return OfrepProvider.constructProvider(OfrepProviderOptions.builder()
+ .baseUrl(opts.getEndpoint().replaceAll("/+$", ""))
+ .connectTimeout(Duration.ofMillis(opts.getTimeout()))
+ .requestTimeout(Duration.ofMillis(opts.getTimeout()))
+ .headers(ImmutableMap.copyOf(headers))
+ .build());
}
@Override
@@ -34,12 +83,89 @@ public boolean isFlagTrackable(String flagKey) {
}
@Override
- public void init() {
- // do nothing
+ public void initialize(final EvaluationContext ctx, final String domain) throws Exception {
+ restartAfterShutdown();
+ this.ofrep.initialize(ctx, domain);
}
@Override
- public void destroy() {
- // do nothing
+ public void initialize(final EvaluationContext ctx) throws Exception {
+ restartAfterShutdown();
+ this.ofrep.initialize(ctx);
+ }
+
+ /**
+ * restartAfterShutdown prepares the evaluator for a new initialization. The OFREP provider's
+ * shutdown terminates the executor its HTTP client runs on, so a shut down one is replaced.
+ */
+ private void restartAfterShutdown() {
+ this.authenticationFailureReported.set(false);
+ if (this.ofrepShutDown) {
+ this.ofrep = newOfrepProvider(this.options);
+ this.ofrepShutDown = false;
+ }
+ }
+
+ @Override
+ public void shutdown() {
+ this.ofrepShutDown = true;
+ this.ofrep.shutdown();
+ }
+
+ @Override
+ public ProviderEvaluation getBooleanEvaluation(String key, Boolean defaultValue, EvaluationContext ctx) {
+ return reportAuthenticationFailure(this.ofrep.getBooleanEvaluation(key, defaultValue, ctx));
+ }
+
+ @Override
+ public ProviderEvaluation getStringEvaluation(String key, String defaultValue, EvaluationContext ctx) {
+ return reportAuthenticationFailure(this.ofrep.getStringEvaluation(key, defaultValue, ctx));
+ }
+
+ @Override
+ public ProviderEvaluation getIntegerEvaluation(String key, Integer defaultValue, EvaluationContext ctx) {
+ return reportAuthenticationFailure(this.ofrep.getIntegerEvaluation(key, defaultValue, ctx));
+ }
+
+ @Override
+ public ProviderEvaluation getDoubleEvaluation(String key, Double defaultValue, EvaluationContext ctx) {
+ return reportAuthenticationFailure(this.ofrep.getDoubleEvaluation(key, defaultValue, ctx));
+ }
+
+ @Override
+ public ProviderEvaluation getObjectEvaluation(String key, Value defaultValue, EvaluationContext ctx) {
+ return reportAuthenticationFailure(this.ofrep.getObjectEvaluation(key, defaultValue, ctx));
+ }
+
+ /**
+ * reportAuthenticationFailure moves the provider to the fatal state when the relay proxy has
+ * rejected our credentials.
+ *
+ *
Remote evaluation holds no configuration, so initialization has nothing to fetch and cannot
+ * discover that the API key is wrong. The first evaluation is therefore the earliest point at
+ * which the provider can learn it, and rejected credentials cannot be repaired by retrying, so
+ * the SDK must be told rather than left reporting a per-call error forever.
+ *
+ * @param evaluation - result returned by the OFREP client
+ * @param - type of the flag value
+ * @return the evaluation, unchanged
+ */
+ private ProviderEvaluation reportAuthenticationFailure(final ProviderEvaluation evaluation) {
+ if (evaluation.getErrorCode() != ErrorCode.GENERAL
+ || evaluation.getErrorMessage() == null
+ || !evaluation.getErrorMessage().startsWith(OFREP_AUTHENTICATION_ERROR)) {
+ return evaluation;
+ }
+
+ if (this.authenticationFailureReported.compareAndSet(false, true)) {
+ log.error("the relay proxy rejected our credentials, the provider cannot recover by retrying");
+ this.emitter.accept(
+ ProviderEvent.PROVIDER_ERROR,
+ ProviderEventDetails.builder()
+ .errorCode(ErrorCode.PROVIDER_FATAL)
+ .message("authentication/authorization error while evaluating a flag remotely")
+ .build());
+ }
+ return evaluation;
}
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/exception/AuthenticationFailure.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/exception/AuthenticationFailure.java
new file mode 100644
index 0000000000..ccd3b682e0
--- /dev/null
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/exception/AuthenticationFailure.java
@@ -0,0 +1,14 @@
+package dev.openfeature.contrib.providers.gofeatureflag.exception;
+
+import dev.openfeature.sdk.exceptions.FatalError;
+import lombok.experimental.StandardException;
+
+/**
+ * Thrown when the relay proxy rejects our credentials (HTTP 401 or 403).
+ *
+ *
It extends the SDK's FatalError, whose error code is PROVIDER_FATAL, because that is what the
+ * SDK inspects to move the provider to the FATAL state: credentials cannot be repaired by retrying,
+ * so a provider that keeps retrying them would never recover and would hide the real problem.
+ */
+@StandardException
+public class AuthenticationFailure extends FatalError {}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHook.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHook.java
index 4ea6101e40..c04766e390 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHook.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHook.java
@@ -2,14 +2,14 @@
import dev.openfeature.contrib.providers.gofeatureflag.bean.FeatureEvent;
import dev.openfeature.contrib.providers.gofeatureflag.bean.IEvent;
+import dev.openfeature.contrib.providers.gofeatureflag.evaluator.IEvaluator;
import dev.openfeature.contrib.providers.gofeatureflag.exception.InvalidOptions;
-import dev.openfeature.contrib.providers.gofeatureflag.service.EvaluationService;
import dev.openfeature.contrib.providers.gofeatureflag.service.EventsPublisher;
+import dev.openfeature.contrib.providers.gofeatureflag.util.Const;
import dev.openfeature.contrib.providers.gofeatureflag.util.EvaluationContextUtil;
import dev.openfeature.sdk.FlagEvaluationDetails;
import dev.openfeature.sdk.Hook;
import dev.openfeature.sdk.HookContext;
-import dev.openfeature.sdk.Reason;
import java.util.Map;
import lombok.extern.slf4j.Slf4j;
@@ -19,12 +19,10 @@
*/
@Slf4j
public final class DataCollectorHook implements Hook> {
- /** options contains all the options of this hook. */
- private final DataCollectorHookOptions options;
/** eventsPublisher is the system collecting all the information to send to GO Feature Flag. */
private final EventsPublisher eventsPublisher;
- /** evalService is the service to evaluate the flags. */
- private final EvaluationService evalService;
+ /** evaluator is the service to evaluate the flags. */
+ private final IEvaluator evaluator;
/**
* Constructor of the hook.
@@ -38,46 +36,69 @@ public DataCollectorHook(final DataCollectorHookOptions options) throws InvalidO
}
options.validate();
eventsPublisher = options.getEventsPublisher();
- evalService = options.getEvalService();
- this.options = options;
+ evaluator = options.getEvaluator();
}
@Override
public void after(HookContext ctx, FlagEvaluationDetails details, Map hints) {
- if (!this.evalService.isFlagTrackable(ctx.getFlagKey())
- || (!Boolean.TRUE.equals(this.options.getCollectUnCachedEvaluation())
- && !Reason.CACHED.name().equals(details.getReason()))) {
+ // the relay proxy evaluated this flag itself and recorded it server side as it did so
+ if (wasEvaluatedRemotely(details) || !this.evaluator.isFlagTrackable(ctx.getFlagKey())) {
return;
}
IEvent event = FeatureEvent.builder()
.key(ctx.getFlagKey())
.kind("feature")
- .contextKind(EvaluationContextUtil.isAnonymousUser(ctx.getCtx()) ? "anonymousUser" : "user")
+ .source("INPROCESS")
+ .contextKind(EvaluationContextUtil.contextKind(ctx.getCtx()))
.defaultValue(false)
- .variation(details.getVariant())
+ .variation(details.getVariant() != null ? details.getVariant() : "SdkDefault")
.value(details.getValue())
- .userKey(ctx.getCtx().getTargetingKey())
+ .userKey(EvaluationContextUtil.userKey(ctx.getCtx()))
.creationDate(System.currentTimeMillis() / 1000L)
.build();
eventsPublisher.add(event);
}
+ /**
+ * finallyAfter records a failed evaluation. It stands in for the error stage, which only sees the
+ * exception and so cannot tell that the relay proxy produced the result.
+ */
@Override
- public void error(HookContext ctx, Exception error, Map hints) {
+ public void finallyAfter(HookContext ctx, FlagEvaluationDetails details, Map hints) {
+ if (details.getErrorCode() == null
+ || wasEvaluatedRemotely(details)
+ || !this.evaluator.isFlagTrackable(ctx.getFlagKey())) {
+ return;
+ }
+
IEvent event = FeatureEvent.builder()
.key(ctx.getFlagKey())
.kind("feature")
- .contextKind(EvaluationContextUtil.isAnonymousUser(ctx.getCtx()) ? "anonymousUser" : "user")
+ .source("INPROCESS")
+ .contextKind(EvaluationContextUtil.contextKind(ctx.getCtx()))
.creationDate(System.currentTimeMillis() / 1000L)
.defaultValue(true)
.variation("SdkDefault")
.value(ctx.getDefaultValue())
- .userKey(ctx.getCtx().getTargetingKey())
+ .userKey(EvaluationContextUtil.userKey(ctx.getCtx()))
.build();
eventsPublisher.add(event);
}
+ /**
+ * wasEvaluatedRemotely reports whether this result came from the relay proxy rather than the
+ * local engine. The relay proxy recorded such a result itself, so recording it here too would
+ * count it twice.
+ *
+ * @param details - result the SDK is about to hand to the caller
+ * @return true if the relay proxy produced this result
+ */
+ private static boolean wasEvaluatedRemotely(final FlagEvaluationDetails> details) {
+ return details.getFlagMetadata() != null
+ && Boolean.TRUE.equals(details.getFlagMetadata().getBoolean(Const.METADATA_EVALUATED_REMOTELY));
+ }
+
/** shutdown should be called when we stop the hook, it will publish the remaining event. */
public void shutdown() {
// eventsPublisher is required so no need to check if it is null
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHookOptions.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHookOptions.java
index ffb67400ae..608fed7c76 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHookOptions.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/DataCollectorHookOptions.java
@@ -1,8 +1,8 @@
package dev.openfeature.contrib.providers.gofeatureflag.hook;
import dev.openfeature.contrib.providers.gofeatureflag.bean.IEvent;
+import dev.openfeature.contrib.providers.gofeatureflag.evaluator.IEvaluator;
import dev.openfeature.contrib.providers.gofeatureflag.exception.InvalidOptions;
-import dev.openfeature.contrib.providers.gofeatureflag.service.EvaluationService;
import dev.openfeature.contrib.providers.gofeatureflag.service.EventsPublisher;
import lombok.Builder;
import lombok.Getter;
@@ -14,20 +14,15 @@
@Builder
@Getter
public class DataCollectorHookOptions {
- /**
- * collectUnCachedEvent (optional) set to true if you want to send all events not only the cached
- * evaluations.
- */
- private Boolean collectUnCachedEvaluation;
/**
* eventsPublisher is the system collecting all the information to send to GO Feature Flag.
*/
private EventsPublisher eventsPublisher;
/**
- * evalService is the service to evaluate the flags.
+ * evaluator is used to know whether the usage of a flag should be collected.
*/
- private EvaluationService evalService;
+ private IEvaluator evaluator;
/**
* Validate the options provided to the data collector hook.
@@ -38,5 +33,8 @@ public void validate() throws InvalidOptions {
if (getEventsPublisher() == null) {
throw new InvalidOptions("No events publisher provided");
}
+ if (getEvaluator() == null) {
+ throw new InvalidOptions("No evaluator provided");
+ }
}
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/EnrichEvaluationContextHook.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/EnrichEvaluationContextHook.java
index 68bfa2d93e..9300f6b162 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/EnrichEvaluationContextHook.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/hook/EnrichEvaluationContextHook.java
@@ -9,11 +9,14 @@
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
+import lombok.val;
/**
* EnrichEvaluationContextHook is an OpenFeature Hook in charge of enriching the evaluation context.
*/
public class EnrichEvaluationContextHook implements Hook {
+ private static final String GOFEATUREFLAG_KEY = "gofeatureflag";
+ private static final String EXPORTER_METADATA_KEY = "exporterMetadata";
private final Map exporterMetadata;
public EnrichEvaluationContextHook(Map exporterMetadata) {
@@ -29,10 +32,10 @@ public Optional before(HookContext ctx, Map entry : exporterMetadata.entrySet()) {
switch (entry.getValue().getClass().getSimpleName()) {
case "String":
@@ -53,9 +56,13 @@ public Optional before(HookContext ctx, Map expMetadata = new HashMap<>();
- expMetadata.put("exporterMetadata", new Value(metadata));
- mutableContext.add("gofeatureflag", new MutableStructure(expMetadata));
+ Map goffNamespace = new HashMap<>();
+ val existing = ctx.getCtx().getValue(GOFEATUREFLAG_KEY);
+ if (existing != null && existing.isStructure()) {
+ goffNamespace.putAll(existing.asStructure().asMap());
+ }
+ goffNamespace.put(EXPORTER_METADATA_KEY, new Value(metadata));
+ mutableContext.add(GOFEATUREFLAG_KEY, new MutableStructure(goffNamespace));
return Optional.of(mutableContext);
}
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/service/EvaluationService.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/service/EvaluationService.java
deleted file mode 100644
index da452d39f7..0000000000
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/service/EvaluationService.java
+++ /dev/null
@@ -1,155 +0,0 @@
-package dev.openfeature.contrib.providers.gofeatureflag.service;
-
-import static dev.openfeature.sdk.Value.objectToValue;
-
-import dev.openfeature.contrib.providers.gofeatureflag.evaluator.IEvaluator;
-import dev.openfeature.contrib.providers.gofeatureflag.util.MetadataUtil;
-import dev.openfeature.sdk.ErrorCode;
-import dev.openfeature.sdk.EvaluationContext;
-import dev.openfeature.sdk.ProviderEvaluation;
-import dev.openfeature.sdk.Reason;
-import dev.openfeature.sdk.exceptions.FlagNotFoundError;
-import dev.openfeature.sdk.exceptions.TargetingKeyMissingError;
-import dev.openfeature.sdk.exceptions.TypeMismatchError;
-import lombok.AllArgsConstructor;
-import lombok.val;
-
-/**
- * EvaluationService is responsible for evaluating feature flags using the provided evaluator.
- * It can use different evaluators based on the configuration and context.
- */
-@AllArgsConstructor
-public class EvaluationService {
- /**
- * The evaluator used to evaluate the flags.
- */
- private IEvaluator evaluator;
-
- /**
- * Return true if we should track the usage of the flag.
- *
- * @param flagKey - name of the flag
- * @return true if the flag is trackable, false otherwise
- */
- public boolean isFlagTrackable(final String flagKey) {
- return this.evaluator.isFlagTrackable(flagKey);
- }
-
- /**
- * Init the evaluator.
- */
- public void init() {
- this.evaluator.init();
- }
-
- /**
- * Destroy the evaluator.
- */
- public void destroy() {
- this.evaluator.destroy();
- }
-
- /**
- * Get the evaluation response from the evaluator.
- *
- * @param flagKey - name of the flag
- * @param defaultValue - default value
- * @param evaluationContext - evaluation context
- * @param expectedType - expected type of the value
- * @param - type of the value
- * @return the evaluation response
- */
- public ProviderEvaluation getEvaluation(
- String flagKey, T defaultValue, EvaluationContext evaluationContext, Class> expectedType) {
-
- if (evaluationContext.getTargetingKey() == null) {
- throw new TargetingKeyMissingError("GO Feature Flag requires a targeting key");
- }
-
- val goffResp = evaluator.evaluate(flagKey, defaultValue, evaluationContext);
-
- // Check for FLAG_NOT_FOUND error first, before general error handling
- if (goffResp.getErrorCode() != null
- && ErrorCode.FLAG_NOT_FOUND.name().equalsIgnoreCase(goffResp.getErrorCode())) {
- throw new FlagNotFoundError("Flag " + flagKey + " was not found in your configuration");
- }
-
- // If we have an error code, we return the error directly.
- if (goffResp.getErrorCode() != null && !goffResp.getErrorCode().isEmpty()) {
- return ProviderEvaluation.builder()
- .errorCode(mapErrorCode(goffResp.getErrorCode()))
- .errorMessage(goffResp.getErrorDetails())
- .reason(Reason.ERROR.name())
- .value(defaultValue)
- .build();
- }
-
- if (Reason.DISABLED.name().equalsIgnoreCase(goffResp.getReason())) {
- // we don't set a variant since we are using the default value,
- // and we are not able to know which variant it is.
- return ProviderEvaluation.builder()
- .value(defaultValue)
- .variant(goffResp.getVariationType())
- .reason(Reason.DISABLED.name())
- .build();
- }
-
- // Convert the value received from the API.
- T flagValue = convertValue(goffResp.getValue(), expectedType);
-
- if (flagValue.getClass() != expectedType) {
- throw new TypeMismatchError(String.format(
- "Flag value %s had unexpected type %s, expected %s.", flagKey, flagValue.getClass(), expectedType));
- }
-
- return ProviderEvaluation.builder()
- .errorCode(mapErrorCode(goffResp.getErrorCode()))
- .reason(goffResp.getReason())
- .value(flagValue)
- .variant(goffResp.getVariationType())
- .flagMetadata(MetadataUtil.convertFlagMetadata(goffResp.getMetadata()))
- .build();
- }
-
- /**
- * convertValue is converting the object return by the proxy response in the right type.
- *
- * @param value - The value we have received
- * @param expectedType - the type we expect for this value
- * @param the type we want to convert to.
- * @return A converted object
- */
- private T convertValue(Object value, Class> expectedType) {
- boolean isPrimitive = expectedType == Boolean.class
- || expectedType == String.class
- || expectedType == Integer.class
- || expectedType == Double.class;
-
- if (isPrimitive) {
- if (value.getClass() == Integer.class && expectedType == Double.class) {
- return (T) Double.valueOf((Integer) value);
- }
- return (T) value;
- }
- return (T) objectToValue(value);
- }
-
- /**
- * mapErrorCode is mapping the errorCode in string received by the API to our internal SDK
- * ErrorCode enum.
- *
- * @param errorCode - string of the errorCode received from the API
- * @return an item from the enum
- */
- private ErrorCode mapErrorCode(String errorCode) {
- if (errorCode == null || errorCode.isEmpty()) {
- return null;
- }
-
- try {
- return ErrorCode.valueOf(errorCode);
- } catch (IllegalArgumentException e) {
- return null;
- }
- }
-}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/service/EventsPublisher.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/service/EventsPublisher.java
index 32e1272f7a..42b1967968 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/service/EventsPublisher.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/service/EventsPublisher.java
@@ -7,6 +7,7 @@
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.Executors;
+import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -27,13 +28,18 @@ public final class EventsPublisher {
public final AtomicBoolean isShutdown = new AtomicBoolean(false);
private final int maxPendingEvents;
private final Consumer> publisher;
+ /** true while a batch is being posted, so that a second publish skips rather than overlapping it. */
+ private final AtomicBoolean publishing = new AtomicBoolean(false);
+ /** true while a flush of a full buffer is queued on the scheduler and has not started yet. */
+ private final AtomicBoolean flushRequested = new AtomicBoolean(false);
private final ReadWriteLock readWriteLock = new ReentrantReadWriteLock();
private final Lock readLock = readWriteLock.readLock();
private final Lock writeLock = readWriteLock.writeLock();
- private final ScheduledExecutorService scheduledExecutorService = Executors.newScheduledThreadPool(1);
+ private final long flushIntervalMs;
private final List eventsList;
+ private volatile ScheduledExecutorService scheduledExecutorService;
/**
* Constructor.
@@ -47,6 +53,21 @@ public EventsPublisher(Consumer> publisher, long flushIntervalMs, int ma
eventsList = new CopyOnWriteArrayList<>();
this.publisher = publisher;
this.maxPendingEvents = maxPendingEvents;
+ this.flushIntervalMs = flushIntervalMs;
+ start();
+ }
+
+ /**
+ * start schedules the periodic flush.
+ * Calling start() on a running publisher does nothing, so it is safe to call from both the
+ * constructor and provider initialization.
+ */
+ public synchronized void start() {
+ if (scheduledExecutorService != null && !scheduledExecutorService.isShutdown()) {
+ return;
+ }
+ isShutdown.set(false);
+ scheduledExecutorService = Executors.newScheduledThreadPool(1);
log.debug("Scheduling events publishing at fixed rate of {} milliseconds", flushIntervalMs);
scheduledExecutorService.scheduleAtFixedRate(
this::publish, flushIntervalMs, flushIntervalMs, TimeUnit.MILLISECONDS);
@@ -73,50 +94,121 @@ public void add(T event) {
}
if (shouldPublish) {
- log.warn("events collection is full. Publishing before adding new events.");
- publish();
+ requestFlush();
}
try {
writeLock.lock();
if (eventsList != null) {
eventsList.add(event);
+ discardOverflow();
}
} finally {
writeLock.unlock();
}
}
+ /**
+ * requestFlush has the scheduler's thread publish the buffer, so the thread adding an event, often
+ * one evaluating a flag, never waits for the data collector. Requests made while one is already
+ * queued are merged into it.
+ */
+ private void requestFlush() {
+ if (!flushRequested.compareAndSet(false, true)) {
+ return;
+ }
+ log.warn("events collection is full, publishing it");
+ try {
+ scheduledExecutorService.execute(() -> {
+ flushRequested.set(false);
+ publish();
+ });
+ } catch (RejectedExecutionException e) {
+ flushRequested.set(false);
+ log.debug("the publisher is shutting down, the final drain will publish the buffer");
+ }
+ }
+
+ /**
+ * discardOverflow keeps the buffer within twice maxPendingEvents, dropping the oldest events
+ * first. Without a cap a data collector outage is an unbounded memory leak, and the oldest
+ * events are the least useful to keep.
+ *
+ *
Callers must hold {@link #writeLock}.
+ */
+ private void discardOverflow() {
+ long overflow = eventsList.size() - (2L * maxPendingEvents);
+ if (overflow > 0) {
+ log.warn("events buffer is full, discarding the {} oldest events", overflow);
+ eventsList.subList(0, (int) overflow).clear();
+ }
+ }
+
/**
* publish events.
*
* @return count of publish events
*/
public int publish() {
- int publishedEvents = 0;
+ if (!publishing.compareAndSet(false, true)) {
+ log.debug("a publish is already in progress, skipping this one");
+ return 0;
+ }
+ try {
+ return drainAndPost();
+ } finally {
+ publishing.set(false);
+ }
+ }
+
+ /**
+ * drainAndPost swaps the buffer out under the lock, releases it, and only then posts, so the
+ * data collector's availability cannot hold up an evaluation. A batch that fails to publish goes
+ * back to the head of the buffer, keeping the events in chronological order.
+ *
+ *
Callers must have set {@link #publishing}, or have stopped the scheduler.
+ */
+ private int drainAndPost() {
+ List batch;
writeLock.lock();
try {
if (eventsList.isEmpty()) {
log.debug("Not publishing, no events");
- return publishedEvents;
+ return 0;
}
- log.info("publishing {} events", eventsList.size());
- publisher.accept(new ArrayList<>(eventsList));
- publishedEvents = eventsList.size();
+ batch = new ArrayList<>(eventsList);
eventsList.clear();
- return publishedEvents;
+ } finally {
+ writeLock.unlock();
+ }
+
+ try {
+ log.info("publishing {} events", batch.size());
+ publisher.accept(batch);
+ return batch.size();
} catch (Exception e) {
log.error("Error publishing events", e);
+ writeLock.lock();
+ try {
+ eventsList.addAll(0, batch);
+ discardOverflow();
+ } finally {
+ writeLock.unlock();
+ }
return 0;
- } finally {
- writeLock.unlock();
}
}
- /** Shutdown. */
- public void shutdown() {
+ /**
+ * Shutdown: stop accepting events, stop the scheduler, letting a publish in progress finish, then
+ * drain what is buffered.
+ */
+ public synchronized void shutdown() {
log.info("shutdown, draining remaining events");
- publish();
- ConcurrentUtil.shutdownAndAwaitTermination(scheduledExecutorService, 10);
+ isShutdown.set(true);
+ if (scheduledExecutorService != null) {
+ ConcurrentUtil.shutdownAndAwaitTermination(scheduledExecutorService, 10);
+ }
+ drainAndPost();
}
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/Const.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/Const.java
index 6b3c6a326c..8c3bccb238 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/Const.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/Const.java
@@ -10,24 +10,43 @@
*/
public class Const {
// HTTP
- public static final String BEARER_TOKEN = "Bearer ";
public static final String APPLICATION_JSON = "application/json; charset=utf-8";
public static final String HTTP_HEADER_CONTENT_TYPE = "Content-Type";
- public static final String HTTP_HEADER_AUTHORIZATION = "Authorization";
+ public static final String HTTP_HEADER_API_KEY = "X-API-Key";
public static final String HTTP_HEADER_ETAG = "ETag";
public static final String HTTP_HEADER_IF_NONE_MATCH = "If-None-Match";
public static final String HTTP_HEADER_LAST_MODIFIED = "Last-Modified";
+ // API ROUTES (relative to the configured endpoint, so that a path prefix on the endpoint survives)
+ public static final String PATH_FLAG_CONFIGURATION = "v1/flag/configuration";
+ public static final String PATH_DATA_COLLECTOR = "v1/data/collector";
+ // FLAG CONFIGURATION FIELDS
+ // the only field of a flag configuration a provider may read, everything else is the engine's
+ public static final String FIELD_TRACK_EVENTS = "trackEvents";
+ // FLAG METADATA KEYS
+ public static final String METADATA_EVALUATED_REMOTELY = "gofeatureflag_evaluated_remotely";
+ // EXPORTER METADATA KEYS
+ // the collector groups by provider, so the value must not change between releases
+ public static final String METADATA_PROVIDER = "provider";
+ public static final String METADATA_OPENFEATURE = "openfeature";
+ // EVENT FIELDS
+ public static final String UNDEFINED_TARGETING_KEY = "undefined-targetingKey";
// DEFAULT VALUES
public static final long DEFAULT_POLLING_CONFIG_FLAG_CHANGE_INTERVAL_MS = 2L * 60L * 1000L;
public static final long DEFAULT_FLUSH_INTERVAL_MS = Duration.ofMinutes(1).toMillis();
public static final int DEFAULT_MAX_PENDING_EVENTS = 10000;
public static final int DEFAULT_WASM_EVALUATOR_POOL_SIZE =
Runtime.getRuntime().availableProcessors();
+ /** consecutive failed refreshes after which the configuration is announced as stale. */
+ public static final int STALE_AFTER_CONSECUTIVE_FAILURES = 3;
+ /** fraction by which each poll interval is randomly shortened or lengthened. */
+ public static final double POLLING_JITTER_RATIO = 0.1;
// MAPPERS
public static final ObjectMapper DESERIALIZE_OBJECT_MAPPER =
new ObjectMapper().configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false);
public static final ObjectMapper SERIALIZE_OBJECT_MAPPER = new ObjectMapper();
public static final ObjectMapper SERIALIZE_WASM_MAPPER = new ObjectMapper()
- .setSerializationInclusion(com.fasterxml.jackson.annotation.JsonInclude.Include.NON_NULL)
+ .setDefaultPropertyInclusion(com.fasterxml.jackson.annotation.JsonInclude.Value.construct(
+ com.fasterxml.jackson.annotation.JsonInclude.Include.NON_NULL,
+ com.fasterxml.jackson.annotation.JsonInclude.Include.ALWAYS))
.setDateFormat(new StdDateFormat().withColonInTimeZone(true));
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/EvaluationContextUtil.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/EvaluationContextUtil.java
index fc085bad84..16218d9ed3 100644
--- a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/EvaluationContextUtil.java
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/EvaluationContextUtil.java
@@ -14,6 +14,9 @@ public class EvaluationContextUtil {
*/
private static final String anonymousFieldName = "anonymous";
+ private static final String ANONYMOUS_USER_CONTEXT_KIND = "anonymousUser";
+ private static final String USER_CONTEXT_KIND = "user";
+
/**
* isAnonymousUser is checking if the user in the evaluationContext is anonymous.
*
@@ -25,6 +28,31 @@ public static boolean isAnonymousUser(final EvaluationContext ctx) {
return true;
}
Value value = ctx.getValue(anonymousFieldName);
- return value != null && value.asBoolean();
+ return value != null && value.isBoolean() && Boolean.TRUE.equals(value.asBoolean());
+ }
+
+ /**
+ * contextKind is the bucket an event is counted under.
+ *
+ * @param ctx - EvaluationContext from open-feature
+ * @return the bucket this evaluation belongs to
+ */
+ public static String contextKind(final EvaluationContext ctx) {
+ return isAnonymousUser(ctx) ? ANONYMOUS_USER_CONTEXT_KIND : USER_CONTEXT_KIND;
+ }
+
+ /**
+ * userKey is the key an event is attributed to.
+ *
+ * @param ctx - EvaluationContext from open-feature
+ * @return the targeting key, or a placeholder when there is none
+ */
+ public static String userKey(final EvaluationContext ctx) {
+ if (ctx == null
+ || ctx.getTargetingKey() == null
+ || ctx.getTargetingKey().isEmpty()) {
+ return Const.UNDEFINED_TARGETING_KEY;
+ }
+ return ctx.getTargetingKey();
}
}
diff --git a/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/JsonValueUtil.java b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/JsonValueUtil.java
new file mode 100644
index 0000000000..64bc9ebb4a
--- /dev/null
+++ b/providers/go-feature-flag/src/main/java/dev/openfeature/contrib/providers/gofeatureflag/util/JsonValueUtil.java
@@ -0,0 +1,46 @@
+package dev.openfeature.contrib.providers.gofeatureflag.util;
+
+import java.math.BigInteger;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import lombok.AccessLevel;
+import lombok.NoArgsConstructor;
+
+/**
+ * JsonValueUtil is a utility class to prepare values decoded from JSON before they are converted to
+ * an Open Feature Value.
+ */
+@NoArgsConstructor(access = AccessLevel.PRIVATE)
+public class JsonValueUtil {
+ /**
+ * widenBigIntegers replaces every BigInteger in a decoded JSON value with its Double equivalent.
+ *
+ *
Jackson decodes an integer beyond the long range as BigInteger, which Value.objectToValue
+ * rejects. The engine writes those numbers from a float64, so a Double holds them exactly.
+ *
+ * @param value - a value decoded from the engine's JSON output
+ * @return the same value, with BigInteger replaced by Double at any depth
+ */
+ public static Object widenBigIntegers(final Object value) {
+ if (value instanceof BigInteger) {
+ return ((BigInteger) value).doubleValue();
+ }
+ if (value instanceof Map) {
+ Map