diff --git a/src/main/java/com/flagsmith/FlagsmithClient.java b/src/main/java/com/flagsmith/FlagsmithClient.java index dd2917de..8233f53e 100644 --- a/src/main/java/com/flagsmith/FlagsmithClient.java +++ b/src/main/java/com/flagsmith/FlagsmithClient.java @@ -13,19 +13,27 @@ import com.flagsmith.interfaces.FlagsmithSdk; import com.flagsmith.mappers.EngineMappers; import com.flagsmith.models.BaseFlag; +import com.flagsmith.models.ExperimentMetadata; +import com.flagsmith.models.Flag; import com.flagsmith.models.Flags; import com.flagsmith.models.Segment; import com.flagsmith.models.SegmentMetadata; +import com.flagsmith.threads.EventProcessor; import com.flagsmith.threads.PollingManager; import com.flagsmith.utils.ModelUtils; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.CompletableFuture; import java.util.function.Function; import java.util.stream.Collectors; +import lombok.AccessLevel; import lombok.Data; +import lombok.Getter; import lombok.NonNull; +import lombok.Setter; +import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -39,6 +47,9 @@ public class FlagsmithClient { private FlagsmithSdk flagsmithSdk; private EvaluationContext evaluationContext; private PollingManager pollingManager; + @Getter(AccessLevel.PACKAGE) + @Setter(AccessLevel.NONE) + private EventProcessor eventProcessor; private FlagsmithClient() { } @@ -197,6 +208,183 @@ public List getIdentitySegments(String identifier, Map }).filter(Objects::nonNull).collect(Collectors.toList()); } + /** + * As {@link #getExperimentFlag(String, String, Map)}, with no traits. + * + * @param featureName feature name + * @param identifier identifier string + * @return the flag for the given feature + * @throws FlagsmithRuntimeError when events are not enabled + * @throws FlagsmithApiError when identity flags are unavailable and no default flag handler + * is configured + */ + public BaseFlag getExperimentFlag(String featureName, String identifier) + throws FlagsmithClientError { + return getExperimentFlag(featureName, identifier, new HashMap<>()); + } + + /** + * Get an identity's flag, recording one {@code $flag_exposure} event if the identity is enrolled + * in a running experiment on it. Only remote evaluation carries experiment metadata, so local + * evaluation and offline mode record no exposure. + * + * @param featureName feature name + * @param identifier identifier string + * @param traits a map of trait keys to trait values + * @return the flag for the given feature + * @throws FlagsmithRuntimeError when events are not enabled + * @throws FlagsmithApiError when identity flags are unavailable and no default flag handler + * is configured + */ + public BaseFlag getExperimentFlag( + String featureName, String identifier, Map traits) + throws FlagsmithClientError { + requireEventProcessor("get experiment flags"); + + Flags flags = getIdentityFlags(identifier, traits); + + if (flags == null) { + FlagsmithFlagDefaults defaults = getConfig().getFlagsmithFlagDefaults(); + if (defaults == null) { + throw new FlagsmithApiError("Failed to get feature flags."); + } + logger.info("Not recording an exposure for feature {}: identity flags are unavailable, so " + + "the default flag handler served it.", featureName); + return defaults.evaluateDefaultFlag(featureName); + } + + BaseFlag flag = flags.getFlag(featureName); + + if (!(flag instanceof Flag)) { + logger.info("Not recording an exposure for feature {}: served by the default flag handler.", + featureName); + return flag; + } + + if (!Boolean.TRUE.equals(flag.getEnabled())) { + logger.info("Not recording an exposure for feature {}: the flag is disabled.", featureName); + return flag; + } + + ExperimentMetadata experiment = ((Flag) flag).getExperiment(); + if (experiment == null || !Boolean.TRUE.equals(experiment.getInExperiment())) { + logger.info("Not recording an exposure for feature {}: the identity is not enrolled in a " + + "running experiment.", featureName); + return flag; + } + + Map metadata = new HashMap<>(); + metadata.put("experiment_id", experiment.getId()); + trackExposureEvent(featureName, identifier, ((Flag) flag).getVariant(), traits, metadata); + + return flag; + } + + /** + * Record a custom event. + * + * @param event event name + * @throws FlagsmithRuntimeError when events are not enabled + * @throws IllegalArgumentException when the event name is blank or starts with "$" + */ + public void trackEvent(String event) { + trackEvent(event, null, null, null, null); + } + + /** + * Record a custom event for an identity. + * + * @param event event name + * @param identifier identifier string + * @throws FlagsmithRuntimeError when events are not enabled + * @throws IllegalArgumentException when the event name is blank or starts with "$" + */ + public void trackEvent(String event, String identifier) { + trackEvent(event, identifier, null, null, null); + } + + /** + * Record a custom event for an identity, with a value, traits and metadata. + * + * @param event event name + * @param identifier identifier string + * @param value event value, stringified before sending + * @param traits a map of trait keys to trait values + * @param metadata a map of metadata to attach to the event + * @throws FlagsmithRuntimeError when events are not enabled + * @throws IllegalArgumentException when the event name is blank or starts with "$" + */ + public void trackEvent(String event, String identifier, Object value, + Map traits, Map metadata) { + EventProcessor processor = requireEventProcessor("track events"); + + if (StringUtils.isBlank(event)) { + throw new IllegalArgumentException("An event name is required."); + } + if (event.startsWith("$")) { + throw new IllegalArgumentException("Event names starting with \"$\" are reserved; use " + + "trackExposureEvent to record \"" + EventProcessor.FLAG_EXPOSURE_EVENT + "\"."); + } + + processor.trackEvent(event, identifier, value, traits, metadata); + } + + /** + * Record a {@code $flag_exposure} event. Skipped, with a log line, when the identifier is + * blank. + * + * @param featureName feature the identity was exposed to + * @param identifier identifier string + * @param value variant the identity was bucketed into + * @throws FlagsmithRuntimeError when events are not enabled + * @throws IllegalArgumentException when the feature name is blank + */ + public void trackExposureEvent(String featureName, String identifier, Object value) { + trackExposureEvent(featureName, identifier, value, null, null); + } + + /** + * Record a {@code $flag_exposure} event, with traits and metadata. Skipped, with a log line, + * when the identifier is blank. + * + * @param featureName feature the identity was exposed to + * @param identifier identifier string + * @param value variant the identity was bucketed into + * @param traits a map of trait keys to trait values + * @param metadata a map of metadata to attach to the event + * @throws FlagsmithRuntimeError when events are not enabled + * @throws IllegalArgumentException when the feature name is blank + */ + public void trackExposureEvent(String featureName, String identifier, Object value, + Map traits, Map metadata) { + EventProcessor processor = requireEventProcessor("track exposure events"); + + if (StringUtils.isBlank(featureName)) { + throw new IllegalArgumentException("An exposure requires a feature name."); + } + if (StringUtils.isBlank(identifier)) { + logger.info("Not sending {} for feature {}: an exposure requires an identifier.", + EventProcessor.FLAG_EXPOSURE_EVENT, featureName); + return; + } + + processor.trackExposureEvent(featureName, identifier, value, traits, metadata); + } + + /** + * Send buffered events now. + * + * @return a future completing once every event buffered so far has been sent or dropped, already + * completed when events are not enabled + */ + public CompletableFuture flushEvents() { + if (eventProcessor == null) { + return CompletableFuture.completedFuture(null); + } + + return eventProcessor.flush(); + } + /** * Should be called when terminating the client to clean up any resources that * need cleaning up. @@ -205,9 +393,23 @@ public void close() { if (pollingManager != null) { pollingManager.stopPolling(); } + + if (eventProcessor != null) { + eventProcessor.close(); + } + flagsmithSdk.close(); } + private EventProcessor requireEventProcessor(String action) { + if (eventProcessor == null) { + throw new FlagsmithRuntimeError( + "Events must be enabled to " + action + ". Use withEnableEvents(true)."); + } + + return eventProcessor; + } + private Flags getEnvironmentFlagsFromEvaluationContext() throws FlagsmithClientError { if (evaluationContext == null) { if (getConfig().getFlagsmithFlagDefaults() == null) { @@ -446,7 +648,7 @@ public Builder withApiUrl(String apiUrl) { } /** - * Add custom HTTP headers to the calls. + * Add custom HTTP headers to the Flags API calls. They are not sent to the events API. * * @param customHeaders headers. * @return the Builder @@ -501,6 +703,9 @@ public FlagsmithClient build() { if (configuration.getOfflineHandler() == null) { throw new FlagsmithRuntimeError("Offline handler must be provided to use offline mode."); } + if (configuration.getEnableEvents()) { + throw new FlagsmithRuntimeError("Events cannot be enabled in offline mode."); + } } if (this.flagsmithApiWrapper != null) { @@ -536,7 +741,30 @@ public FlagsmithClient build() { "In order to use local evaluation, please generate a server key " + "in the environment settings page."); } + } + if (configuration.getOfflineHandler() != null) { + if (configuration.getFlagsmithFlagDefaults() != null) { + throw new FlagsmithRuntimeError( + "Cannot use both default flag handler and offline handler."); + } + client.evaluationContext = EngineMappers.mapEnvironmentToContext( + configuration.getOfflineHandler().getEnvironment()); + } + + EventProcessor processor = null; + if (configuration.getEnableEvents()) { + processor = configuration.getEventProcessor() != null + ? configuration.getEventProcessor() + : new EventProcessor( + configuration.getHttpClient(), + configuration.getEventsUri(), + configuration.getEventsMaxBufferItems(), + configuration.getEventsFlushIntervalMillis()); + processor.claim(); + } + + if (configuration.getEnableLocalEvaluation()) { if (this.pollingManager != null) { client.pollingManager = pollingManager; } else { @@ -548,13 +776,14 @@ public FlagsmithClient build() { client.pollingManager.startPolling(); } - if (configuration.getOfflineHandler() != null) { - if (configuration.getFlagsmithFlagDefaults() != null) { - throw new FlagsmithRuntimeError( - "Cannot use both default flag handler and offline handler."); + if (processor != null) { + if (client.eventProcessor != null) { + client.eventProcessor.close(); } - client.evaluationContext = EngineMappers.mapEnvironmentToContext( - configuration.getOfflineHandler().getEnvironment()); + processor.setApi(client.flagsmithSdk); + processor.setLogger(client.logger); + processor.start(); + client.eventProcessor = processor; } return this.client; diff --git a/src/main/java/com/flagsmith/config/FlagsmithConfig.java b/src/main/java/com/flagsmith/config/FlagsmithConfig.java index 33d8cd1f..48adf344 100644 --- a/src/main/java/com/flagsmith/config/FlagsmithConfig.java +++ b/src/main/java/com/flagsmith/config/FlagsmithConfig.java @@ -3,6 +3,7 @@ import com.flagsmith.FlagsmithFlagDefaults; import com.flagsmith.interfaces.IOfflineHandler; import com.flagsmith.threads.AnalyticsProcessor; +import com.flagsmith.threads.EventProcessor; import java.net.Proxy; import java.util.ArrayList; import java.util.List; @@ -30,17 +31,26 @@ public final class FlagsmithConfig { private static final int DEFAULT_ENVIRONMENT_REFRESH_SECONDS = 60; private static final HttpUrl DEFAULT_BASE_URI = HttpUrl .get("https://edge.api.flagsmith.com/api/v1/"); + private static final HttpUrl DEFAULT_EVENTS_URI = HttpUrl + .get("https://events.api.flagsmith.com/"); + private static final int DEFAULT_EVENTS_MAX_BUFFER_ITEMS = 1000; + private static final int DEFAULT_EVENTS_FLUSH_INTERVAL_MILLIS = 10000; private final HttpUrl flagsUri; private final HttpUrl identitiesUri; private final HttpUrl traitsUri; private final HttpUrl environmentUri; private final OkHttpClient httpClient; private final HttpUrl baseUri; + private final HttpUrl eventsUri; + private final Boolean enableEvents; + private final int eventsMaxBufferItems; + private final int eventsFlushIntervalMillis; private final Retry retries; private Boolean enableLocalEvaluation; private Integer environmentRefreshIntervalSeconds; private AnalyticsProcessor analyticsProcessor; + private EventProcessor eventProcessor; private FlagsmithFlagDefaults flagsmithFlagDefaults = null; private Boolean raiseUpdateEnvironmentErrorsOnStartup = true; private Boolean offlineMode = false; @@ -88,6 +98,24 @@ protected FlagsmithConfig(Builder builder) { } } + this.eventsUri = builder.eventsUri; + this.enableEvents = Boolean.TRUE.equals(builder.enableEvents); + this.eventsMaxBufferItems = builder.eventsMaxBufferItems; + this.eventsFlushIntervalMillis = builder.eventsFlushIntervalMillis; + + if (enableEvents) { + if (eventsMaxBufferItems < 1) { + throw new IllegalArgumentException("maxBufferItems must be at least 1."); + } + if (eventsFlushIntervalMillis < 0) { + throw new IllegalArgumentException("flushIntervalMillis must not be negative."); + } + eventProcessor = builder.eventProcessor; + } else if (builder.eventsConfigured) { + throw new IllegalArgumentException( + "Events must be enabled with withEnableEvents(true) to configure the event processor."); + } + this.offlineMode = builder.offlineMode; this.offlineHandler = builder.offlineHandler; } @@ -120,6 +148,13 @@ public static class Builder { private Integer environmentRefreshIntervalSeconds = DEFAULT_ENVIRONMENT_REFRESH_SECONDS; private Boolean enableAnalytics = Boolean.FALSE; + private HttpUrl eventsUri = DEFAULT_EVENTS_URI; + private EventProcessor eventProcessor; + private Boolean enableEvents = Boolean.FALSE; + private Boolean eventsConfigured = Boolean.FALSE; + private int eventsMaxBufferItems = DEFAULT_EVENTS_MAX_BUFFER_ITEMS; + private int eventsFlushIntervalMillis = DEFAULT_EVENTS_FLUSH_INTERVAL_MILLIS; + private Boolean offlineMode = Boolean.FALSE; private IOfflineHandler offlineHandler; @@ -187,7 +222,8 @@ public Builder sslSocketFactory(SSLSocketFactory sslSocketFactory, } /** - * Add a custom HTTP interceptor. + * Add a custom HTTP interceptor. It runs on every request, including those to the events + * API, so an interceptor that adds credentials should check the request's host. * * @param interceptor the HTTP interceptor * @return the Builder @@ -272,6 +308,70 @@ public Builder withEnableAnalytics(Boolean enable) { return this; } + /** + * Override the events API base URL. + * + * @param eventsUri the new base URI for the events API + * @return the Builder + */ + public Builder eventsUri(String eventsUri) { + if (eventsUri != null) { + this.eventsUri = HttpUrl.get(eventsUri.endsWith("/") ? eventsUri : eventsUri + "/"); + } + return this; + } + + /** + * Enable the event processor, which records experiment exposures and custom events. + * + * @param enable boolean to enable + * @return the Builder + */ + public Builder withEnableEvents(Boolean enable) { + this.enableEvents = enable; + return this; + } + + /** + * Use a custom event processor. Also enables events. The processor can back only one client: + * building a second client from this configuration throws. + * + * @param processor the processor that buffers and sends events + * @return the Builder + */ + public Builder withEventProcessor(EventProcessor processor) { + this.eventProcessor = processor; + this.enableEvents = Boolean.TRUE; + return this; + } + + /** + * Set the number of buffered events that triggers an immediate flush. Requires events to be + * enabled; {@link #build()} throws IllegalArgumentException when it is below 1. + * + * @param items the maximum number of buffered events + * @return the Builder + */ + public Builder withEventsMaxBufferItems(int items) { + this.eventsMaxBufferItems = items; + this.eventsConfigured = Boolean.TRUE; + return this; + } + + /** + * Set the interval between timed event flushes, in milliseconds. Zero disables the timer. + * Requires events to be enabled; {@link #build()} throws IllegalArgumentException when it is + * negative. + * + * @param millis the flush interval in milliseconds + * @return the Builder + */ + public Builder withEventsFlushIntervalMillis(int millis) { + this.eventsFlushIntervalMillis = millis; + this.eventsConfigured = Boolean.TRUE; + return this; + } + /** * Enable offline mode. * diff --git a/src/main/java/com/flagsmith/config/Retry.java b/src/main/java/com/flagsmith/config/Retry.java index aedec811..eb1433cf 100644 --- a/src/main/java/com/flagsmith/config/Retry.java +++ b/src/main/java/com/flagsmith/config/Retry.java @@ -25,17 +25,40 @@ public class Retry { add(429); add(503); }}; + private Boolean statusForcelistOnly = Boolean.FALSE; public Retry(Integer total) { this.total = total; } + /** + * Create a policy with {@code statusForcelistOnly} off. + * + * @param total number of attempts before giving up + * @param attempts attempts made so far + * @param backoffFactor factor applied to the backoff between attempts + * @param backoffMax upper bound on the backoff, in seconds + * @param statusForcelist status codes that are always retried + */ + public Retry(Integer total, Integer attempts, Float backoffFactor, Float backoffMax, + Set statusForcelist) { + this(total, attempts, backoffFactor, backoffMax, statusForcelist, Boolean.FALSE); + } + /** * Should Retry or not?. * - * @param statusCode status code of last call + * @param statusCode status code of last call, or null if the call did not get a response */ public Boolean isRetry(Integer statusCode) { + if (Boolean.TRUE.equals(statusForcelistOnly)) { + if (total <= attempts) { + return Boolean.FALSE; + } + return statusCode == null + || (statusForcelist != null && statusForcelist.contains(statusCode)); + } + if (statusForcelist != null && !statusForcelist.isEmpty() && statusForcelist.contains(statusCode)) { return Boolean.TRUE; diff --git a/src/main/java/com/flagsmith/mappers/EngineMappers.java b/src/main/java/com/flagsmith/mappers/EngineMappers.java index 3beee22e..c9a23aaf 100644 --- a/src/main/java/com/flagsmith/mappers/EngineMappers.java +++ b/src/main/java/com/flagsmith/mappers/EngineMappers.java @@ -62,6 +62,7 @@ public static Flag mapFlagResultToFlag( flag.setFeatureName(flagResult.getName()); flag.setValue(flagResult.getValue()); flag.setEnabled(flagResult.getEnabled()); + flag.setReason(flagResult.getReason()); return flag; } diff --git a/src/main/java/com/flagsmith/models/ExperimentMetadata.java b/src/main/java/com/flagsmith/models/ExperimentMetadata.java new file mode 100644 index 00000000..fc234697 --- /dev/null +++ b/src/main/java/com/flagsmith/models/ExperimentMetadata.java @@ -0,0 +1,15 @@ +package com.flagsmith.models; + +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Data; + +/** + * Details of the experiment running on a feature, as returned by remote evaluation. + */ +@Data +public class ExperimentMetadata { + private Integer id; + private String name; + @JsonProperty("in_experiment") + private Boolean inExperiment = Boolean.FALSE; +} diff --git a/src/main/java/com/flagsmith/models/Flag.java b/src/main/java/com/flagsmith/models/Flag.java index 990976b9..8fc437bc 100644 --- a/src/main/java/com/flagsmith/models/Flag.java +++ b/src/main/java/com/flagsmith/models/Flag.java @@ -1,6 +1,8 @@ package com.flagsmith.models; import com.fasterxml.jackson.databind.JsonNode; +import com.flagsmith.MapperFactory; +import com.flagsmith.models.features.FeatureStateMetadata; import com.flagsmith.models.features.FeatureStateModel; import lombok.Data; @@ -8,6 +10,9 @@ public class Flag extends BaseFlag { private Integer featureId = 0; private Boolean isDefault; + private String variant; + private String reason; + private ExperimentMetadata experiment; /** * return flag from feature state model and identity id. @@ -21,6 +26,10 @@ public static Flag fromFeatureStateModel(FeatureStateModel featureState) { flag.setValue(featureState.getValue()); flag.setFeatureName(featureState.getFeature().getName()); flag.setEnabled(featureState.getEnabled()); + flag.setVariant(featureState.getVariant()); + flag.setReason(featureState.getReason()); + flag.setExperiment(featureState.getMetadata() != null + ? featureState.getMetadata().getExperiment() : null); return flag; } @@ -38,6 +47,23 @@ public static Flag fromApiFlag(JsonNode node) { flag.setFeatureName(node.get("feature").get("name").asText()); flag.setEnabled(node.get("enabled").booleanValue()); + JsonNode variant = node.get("variant"); + if (variant != null && !variant.isNull()) { + flag.setVariant(variant.asText()); + } + + JsonNode reason = node.get("reason"); + if (reason != null && !reason.isNull()) { + flag.setReason(reason.asText()); + } + + JsonNode metadata = node.get("metadata"); + if (metadata != null && !metadata.isNull()) { + flag.setExperiment(MapperFactory.getMapper() + .convertValue(metadata, FeatureStateMetadata.class) + .getExperiment()); + } + return flag; } } diff --git a/src/main/java/com/flagsmith/models/features/FeatureStateMetadata.java b/src/main/java/com/flagsmith/models/features/FeatureStateMetadata.java new file mode 100644 index 00000000..627cc32d --- /dev/null +++ b/src/main/java/com/flagsmith/models/features/FeatureStateMetadata.java @@ -0,0 +1,9 @@ +package com.flagsmith.models.features; + +import com.flagsmith.models.ExperimentMetadata; +import lombok.Data; + +@Data +public class FeatureStateMetadata { + private ExperimentMetadata experiment; +} diff --git a/src/main/java/com/flagsmith/models/features/FeatureStateModel.java b/src/main/java/com/flagsmith/models/features/FeatureStateModel.java index ccaef00f..a1c180db 100644 --- a/src/main/java/com/flagsmith/models/features/FeatureStateModel.java +++ b/src/main/java/com/flagsmith/models/features/FeatureStateModel.java @@ -20,4 +20,7 @@ public class FeatureStateModel extends BaseModel { private Object value; @JsonProperty("feature_segment") private FeatureSegmentModel featureSegment; -} \ No newline at end of file + private String variant; + private String reason; + private FeatureStateMetadata metadata; +} diff --git a/src/main/java/com/flagsmith/threads/EventProcessor.java b/src/main/java/com/flagsmith/threads/EventProcessor.java new file mode 100644 index 00000000..0da9e9d2 --- /dev/null +++ b/src/main/java/com/flagsmith/threads/EventProcessor.java @@ -0,0 +1,485 @@ +package com.flagsmith.threads; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.util.RawValue; +import com.flagsmith.FlagsmithLogger; +import com.flagsmith.MapperFactory; +import com.flagsmith.Versions; +import com.flagsmith.config.Retry; +import com.flagsmith.exceptions.FlagsmithRuntimeError; +import com.flagsmith.interfaces.FlagsmithSdk; +import com.flagsmith.models.TraitConfig; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.IdentityHashMap; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import lombok.AccessLevel; +import lombok.Getter; +import lombok.Setter; +import okhttp3.HttpUrl; +import okhttp3.MediaType; +import okhttp3.OkHttpClient; +import okhttp3.Request; +import okhttp3.RequestBody; + +/** + * Buffers experimentation events and sends them to the Flagsmith events API, on a timer, when the + * buffer fills, and on {@link #close()}. Exposures are deduplicated within a flush window. + */ +public class EventProcessor { + + /** Name of the reserved event recorded when an identity is exposed to an experiment. */ + public static final String FLAG_EXPOSURE_EVENT = "$flag_exposure"; + + private static final String EVENTS_PATH = "v1/events"; + private static final String AUTH_HEADER = "X-Environment-Key"; + private static final String USER_AGENT_HEADER = "User-Agent"; + private static final String ACCEPT_HEADER = "Accept"; + private static final String SDK_USER_AGENT_HEADER = "Flagsmith-SDK-User-Agent"; + private static final String SDK_USER_AGENT_PREFIX = "flagsmith-java-sdk/"; + private static final String SDK_VERSION_KEY = "sdk_version"; + private static final String EXPERIMENT_ID_KEY = "experiment_id"; + private static final MediaType JSON_MEDIA_TYPE = + MediaType.get("application/json; charset=utf-8"); + static final int MAX_IN_FLIGHT_BATCHES = 2; + static final int MAX_BUFFERED_EVENTS = 1_000; + static final long CLOSE_TIMEOUT_MILLIS = 25_000L; + private static final long DROP_LOG_INTERVAL_NANOS = TimeUnit.SECONDS.toNanos(10); + + @Getter + private final HttpUrl eventsEndpoint; + @Getter + private final int maxBufferItems; + @Getter + private final int flushIntervalMillis; + private final List> buffer = new ArrayList<>(); + private final Set> dedupeKeys = new HashSet<>(); + private final Map, List> dedupeKeyByEvent = new IdentityHashMap<>(); + private final Object lock = new Object(); + @Getter(AccessLevel.PACKAGE) + private final ScheduledExecutorService scheduler; + private final Set> inFlight = ConcurrentHashMap.newKeySet(); + private CompletableFuture nextBatch = new CompletableFuture<>(); + private int droppedSinceLastReport = 0; // guarded by lock + private Long lastDropReportNanos = null; // guarded by lock + @Getter(AccessLevel.PACKAGE) + private final RequestProcessor requestProcessor; + @Getter(AccessLevel.PACKAGE) + @Setter(AccessLevel.PACKAGE) + private long closeTimeoutMillis; + @Setter + private FlagsmithSdk api; + private FlagsmithLogger logger = new FlagsmithLogger(); + private final AtomicBoolean closed = new AtomicBoolean(false); + private final AtomicBoolean claimed = new AtomicBoolean(false); + private ScheduledFuture scheduledFlush; + + /** + * Create a processor that sends batches through {@code client}. + * + * @param client HTTP client; its timeouts also bound {@link #close()} + * @param eventsUri base URI of the events API, e.g. https://events.api.flagsmith.com/ + * @param maxBufferItems number of buffered events that triggers an immediate flush; at + * least 1 + * @param flushIntervalMillis interval between timed flushes; 0 disables the timer + * @throws IllegalArgumentException when maxBufferItems is below 1 or flushIntervalMillis is + * negative + */ + public EventProcessor(OkHttpClient client, HttpUrl eventsUri, int maxBufferItems, + int flushIntervalMillis) { + this(eventsUri, maxBufferItems, flushIntervalMillis, + new RequestProcessor(withCallDeadline(client), new FlagsmithLogger(), buildRetry())); + } + + EventProcessor(HttpUrl eventsUri, int maxBufferItems, int flushIntervalMillis, + RequestProcessor requestProcessor) { + if (maxBufferItems < 1) { + throw new IllegalArgumentException("maxBufferItems must be at least 1."); + } + if (flushIntervalMillis < 0) { + throw new IllegalArgumentException("flushIntervalMillis must not be negative."); + } + this.eventsEndpoint = eventsUri.newBuilder(EVENTS_PATH).build(); + this.maxBufferItems = maxBufferItems; + this.flushIntervalMillis = flushIntervalMillis; + this.requestProcessor = requestProcessor; + this.closeTimeoutMillis = worstCaseBatchMillis(requestProcessor.getClient(), buildRetry()); + this.scheduler = Executors.newSingleThreadScheduledExecutor((runnable) -> { + Thread thread = new Thread(runnable, "flagsmith-events"); + thread.setDaemon(true); + return thread; + }); + } + + private static Retry buildRetry() { + Retry retry = new Retry(2); + retry.setStatusForcelist( + IntStream.rangeClosed(500, 599).boxed().collect(Collectors.toSet())); + retry.setStatusForcelistOnly(Boolean.TRUE); + return retry; + } + + private static long worstCaseBatchMillis(OkHttpClient client, Retry retry) { + long attemptMillis = attemptMillis(client); + if (attemptMillis == 0) { + return CLOSE_TIMEOUT_MILLIS; + } + long total = retry.getTotal() * attemptMillis; + for (int retried = 1; retried < retry.getTotal(); retried++) { + total += (long) (Math.min(retry.getBackoffFactor() * 2 * retried, retry.getBackoffMax()) + * 1000); + } + return total; + } + + private static long attemptMillis(OkHttpClient client) { + if (client.callTimeoutMillis() > 0) { + return client.callTimeoutMillis(); + } + if (client.connectTimeoutMillis() > 0 && client.writeTimeoutMillis() > 0 + && client.readTimeoutMillis() > 0) { + return (long) client.connectTimeoutMillis() + client.writeTimeoutMillis() + + client.readTimeoutMillis(); + } + return 0; + } + + private static OkHttpClient withCallDeadline(OkHttpClient client) { + if (client.callTimeoutMillis() > 0) { + return client; + } + long attemptMillis = attemptMillis(client); + long callTimeout = attemptMillis > 0 ? attemptMillis : CLOSE_TIMEOUT_MILLIS; + return client.newBuilder().callTimeout(callTimeout, TimeUnit.MILLISECONDS).build(); + } + + /** + * Reserve this processor for one client; called by {@code FlagsmithClient.Builder}. + * + * @throws FlagsmithRuntimeError when another client already holds it + */ + public void claim() { + if (!claimed.compareAndSet(false, true)) { + throw new FlagsmithRuntimeError("This event processor already backs another client."); + } + } + + /** + * Set the logger, for this processor and its request processor. + * + * @param logger the client's logger, so event failures appear alongside its other output + */ + public void setLogger(FlagsmithLogger logger) { + this.logger = logger; + requestProcessor.setLogger(logger); + } + + /** + * Buffer a custom event. + * + * @param event event name + * @param identifier identity the event belongs to, may be null + * @param value event value, stringified before sending + * @param traits identity traits to attach, may be null + * @param metadata caller metadata to attach, may be null + */ + public void trackEvent(String event, String identifier, Object value, + Map traits, Map metadata) { + bufferEvent(event, null, identifier, value, traits, metadata, false); + } + + /** + * Buffer a flag exposure event. Exposures equal in feature, identifier, value and experiment + * are only sent once per flush window. + * + * @param featureName feature the identity was exposed to + * @param identifier identity the exposure belongs to + * @param value variant the identity was bucketed into + * @param traits identity traits to attach, may be null + * @param metadata caller metadata to attach, may be null + */ + public void trackExposureEvent(String featureName, String identifier, Object value, + Map traits, Map metadata) { + bufferEvent(FLAG_EXPOSURE_EVENT, featureName, identifier, value, traits, metadata, true); + } + + /** + * Send everything buffered so far. + * + * @return a future completing once every event buffered so far has been sent or dropped + */ + public CompletableFuture flush() { + List> batch = null; + CompletableFuture tracked = null; + CompletableFuture waiting = null; + + synchronized (lock) { + if (!buffer.isEmpty() && (inFlight.size() < MAX_IN_FLIGHT_BATCHES || closed.get())) { + batch = new ArrayList<>(buffer); + buffer.clear(); + dedupeKeys.clear(); + dedupeKeyByEvent.clear(); + tracked = nextBatch; + nextBatch = new CompletableFuture<>(); + inFlight.add(tracked); + } else if (!buffer.isEmpty()) { + waiting = nextBatch; + } + } + + if (batch != null) { + send(batch, tracked); + } + + CompletableFuture sent = awaitInFlight(); + return waiting == null ? sent : CompletableFuture.allOf(sent, waiting); + } + + /** + * Start the flush timer. Does nothing when the flush interval is not positive, when the timer + * is already running, or once the processor is closed. + */ + public synchronized void start() { + if (flushIntervalMillis <= 0 || scheduledFlush != null) { + return; + } + + if (closed.get()) { + logger.error("Not starting the event processor: it has been closed."); + return; + } + + scheduledFlush = scheduler.scheduleWithFixedDelay( + this::flush, flushIntervalMillis, flushIntervalMillis, TimeUnit.MILLISECONDS); + } + + /** + * Stop the timer, flush what is left and release HTTP resources. Blocks until in-flight + * batches settle, for at most one batch's worst case under the client's timeouts, or + * {@link #CLOSE_TIMEOUT_MILLIS} if a timeout is off. + */ + public void close() { + synchronized (this) { + closed.set(true); + scheduler.shutdownNow(); + } + + try { + flush().get(closeTimeoutMillis, TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + logger.error("Interrupted while flushing events on close.", e); + } catch (TimeoutException e) { + logger.error("Stopped waiting for events to be delivered after " + closeTimeoutMillis + + "ms on close."); + } catch (Exception e) { + logger.error("Failed to flush events on close.", e); + } + + requestProcessor.close(); + } + + private void bufferEvent(String event, String featureName, String identifier, Object value, + Map traits, Map metadata, boolean dedupe) { + if (closed.get()) { + logClosed(event); + return; + } + + try { + final String stringValue = value == null ? null : String.valueOf(value); + + Map eventMetadata = new HashMap<>(); + if (metadata != null) { + eventMetadata.putAll(metadata); + } + eventMetadata.put(SDK_VERSION_KEY, Versions.getVersion()); + final Object experimentId = eventMetadata.get(EXPERIMENT_ID_KEY); + + Map eventTraits = eventTraits(traits); + + Map eventPayload = new LinkedHashMap<>(); + eventPayload.put("event", event); + eventPayload.put("feature_name", featureName); + eventPayload.put("identifier", identifier); + eventPayload.put("value", stringValue); + eventPayload.put("traits", eventTraits == null ? null : toJson(eventTraits)); + eventPayload.put("metadata", toJson(eventMetadata)); + eventPayload.put("timestamp", System.currentTimeMillis()); + + boolean isFull = false; + int droppedToReport = 0; + + synchronized (lock) { + if (closed.get()) { + logClosed(event); + return; + } + if (dedupe) { + List key = dedupeKey(event, featureName, identifier, stringValue, experimentId); + if (!dedupeKeys.add(key)) { + return; + } + dedupeKeyByEvent.put(eventPayload, key); + } + buffer.add(eventPayload); + if (inFlight.size() < MAX_IN_FLIGHT_BATCHES) { + isFull = buffer.size() >= maxBufferItems; + } else if (buffer.size() > Math.max(maxBufferItems, MAX_BUFFERED_EVENTS)) { + dedupeKeys.remove(dedupeKeyByEvent.remove(buffer.remove(0))); + droppedToReport = recordDrop(1); + } + } + + logDrops(droppedToReport); + if (isFull) { + flush(); + } + } catch (JsonProcessingException | RuntimeException e) { + logger.error("Failed to buffer event " + event + ".", e); + } + } + + private static RawValue toJson(Object value) throws JsonProcessingException { + return new RawValue(MapperFactory.getMapper().writeValueAsString(value)); + } + + private void logClosed(String event) { + logger.info("Not buffering event {}: the event processor is closed.", event); + } + + private static Map eventTraits(Map traits) { + if (traits == null) { + return null; + } + + Map flattened = new LinkedHashMap<>(); + + for (Map.Entry entry : traits.entrySet()) { + TraitConfig traitConfig = TraitConfig.fromObject(entry.getValue()); + if (!traitConfig.getIsTransient()) { + flattened.put(entry.getKey(), traitConfig.getValue()); + } + } + + return flattened.isEmpty() ? null : flattened; + } + + private static List dedupeKey(String event, String featureName, String identifier, + String value, Object experimentId) { + return Arrays.asList(event, featureName, identifier, value, + experimentId == null ? null : String.valueOf(experimentId)); + } + + private void send(List> batch, CompletableFuture tracked) { + final int batchSize = batch.size(); + boolean submitted = false; + + try { + if (api == null) { + logger.error("Dropping " + batchSize + + " events: the event processor has no API wrapper."); + return; + } + + String payload = MapperFactory.getMapper() + .writeValueAsString(Collections.singletonMap("events", batch)); + + RequestBody body = RequestBody.create(payload, JSON_MEDIA_TYPE); + Request request = new Request.Builder() + .url(eventsEndpoint) + .post(body) + .header(AUTH_HEADER, api.newPostRequest(eventsEndpoint, body).headers(AUTH_HEADER).get(0)) + .header(USER_AGENT_HEADER, SDK_USER_AGENT_PREFIX + Versions.getVersion()) + .header(SDK_USER_AGENT_HEADER, SDK_USER_AGENT_PREFIX + Versions.getVersion()) + .header(ACCEPT_HEADER, "application/json") + .build(); + + requestProcessor + .submit(request, new TypeReference() {}, Boolean.FALSE, buildRetry()) + .whenComplete((response, error) -> { + try { + logRejections(response, batchSize); + } finally { + settle(tracked); + } + }); + submitted = true; + } catch (Exception e) { + logger.error("Dropping " + batchSize + " events: failed to send them.", e); + } finally { + if (!submitted) { + settle(tracked); + } + } + } + + private void logRejections(JsonNode response, int batchSize) { + JsonNode rejected = response == null ? null : response.get("rejected"); + if (rejected != null && rejected.isArray() && rejected.size() > 0) { + logger.error("The events API rejected " + rejected.size() + " of " + batchSize + + " events. First rejection: " + rejected.get(0)); + } + } + + private void settle(CompletableFuture tracked) { + boolean waiting; + synchronized (lock) { + inFlight.remove(tracked); + waiting = !buffer.isEmpty(); + } + tracked.complete(null); + if (waiting) { + flush(); + } + } + + private int recordDrop(int count) { + droppedSinceLastReport += count; + long now = System.nanoTime(); + if (lastDropReportNanos != null && now - lastDropReportNanos < DROP_LOG_INTERVAL_NANOS) { + return 0; + } + lastDropReportNanos = now; + int toReport = droppedSinceLastReport; + droppedSinceLastReport = 0; + return toReport; + } + + private void logDrops(int dropped) { + if (dropped > 0) { + logger.error("Dropped the " + dropped + " oldest events: " + MAX_IN_FLIGHT_BATCHES + + " batches are in flight to the events API and the buffer is full. Further drops are" + + " reported at most every " + TimeUnit.NANOSECONDS.toSeconds(DROP_LOG_INTERVAL_NANOS) + + "s."); + } + } + + private CompletableFuture awaitInFlight() { + return CompletableFuture.allOf(inFlight.toArray(new CompletableFuture[0])); + } + + List> bufferedEvents() { + synchronized (lock) { + return new ArrayList<>(buffer); + } + } +} diff --git a/src/main/java/com/flagsmith/threads/RequestProcessor.java b/src/main/java/com/flagsmith/threads/RequestProcessor.java index 25f4867f..36bc4bce 100644 --- a/src/main/java/com/flagsmith/threads/RequestProcessor.java +++ b/src/main/java/com/flagsmith/threads/RequestProcessor.java @@ -14,6 +14,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.Future; import lombok.Getter; +import lombok.Setter; import okhttp3.Call; import okhttp3.OkHttpClient; import okhttp3.Request; @@ -25,6 +26,7 @@ public class RequestProcessor { @Getter private OkHttpClient client; @Getter + @Setter private FlagsmithLogger logger; private Retry retries = new Retry(3); @@ -78,6 +80,22 @@ public Future executeAsync(Request request, Boolean doThrow) { */ public Future executeAsync( Request request, TypeReference clazz, Boolean doThrow, Retry retries) { + return submit(request, clazz, doThrow, retries); + } + + /** + * Execute the request in async mode, returning a future callers can compose on. + * + * @param request request to send + * @param clazz type to unmarshal the response body into + * @param doThrow whether a failed call completes the future exceptionally + * @param retries retry policy, copied for this call + * @param response type + * @return a future completed with the unmarshalled response, or null when the call failed and + * doThrow is false + */ + public CompletableFuture submit( + Request request, TypeReference clazz, Boolean doThrow, Retry retries) { CompletableFuture completableFuture = new CompletableFuture<>(); Retry localRetry = retries.toBuilder().build(); // run the execute method in a fixed thread with retries. diff --git a/src/test/java/com/flagsmith/FlagsmithClientTest.java b/src/test/java/com/flagsmith/FlagsmithClientTest.java index 5f8acc74..4e2df3ab 100644 --- a/src/test/java/com/flagsmith/FlagsmithClientTest.java +++ b/src/test/java/com/flagsmith/FlagsmithClientTest.java @@ -5,6 +5,7 @@ import static org.mockito.Mockito.*; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -16,6 +17,7 @@ import com.fasterxml.jackson.databind.JsonNode; import com.flagsmith.config.FlagsmithCacheConfig; import com.flagsmith.config.FlagsmithConfig; +import com.flagsmith.exceptions.FeatureNotFoundError; import com.flagsmith.exceptions.FlagsmithApiError; import com.flagsmith.exceptions.FlagsmithClientError; import com.flagsmith.exceptions.FlagsmithRuntimeError; @@ -26,20 +28,25 @@ import com.flagsmith.models.DefaultFlag; import com.flagsmith.models.environments.EnvironmentModel; import com.flagsmith.models.features.FeatureStateModel; +import com.flagsmith.models.Flag; import com.flagsmith.models.Flags; import com.flagsmith.models.SdkTraitModel; import com.flagsmith.models.Segment; import com.flagsmith.models.TraitConfig; import com.flagsmith.models.TraitModel; import com.flagsmith.responses.FlagsAndTraitsResponse; +import com.flagsmith.threads.EventProcessor; import com.flagsmith.threads.PollingManager; import com.flagsmith.threads.RequestProcessor; import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; import java.util.regex.Pattern; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -998,4 +1005,475 @@ public void testFlagsmithUsesOfflineHandlerIfSetAndNoAPIResponse() throws Flagsm assertTrue(environmentFlags.isFeatureEnabled("some_feature")); assertTrue(identityFlags.isFeatureEnabled("some_feature")); } + + /** A client on a mock event processor, serving the experiment identity flags or none. */ + private static FlagsmithClient experimentClient( + EventProcessor processor, boolean withDefaultHandler, boolean flagsUnavailable) { + MockInterceptor interceptor = new MockInterceptor(); + interceptor.addRule() + .post("http://bad-url/identities/") + .anyTimes() + .respond(FlagsmithTestHelper.getIdentitiesFlagsWithExperiment(), MEDIATYPE_JSON); + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .baseUri("http://bad-url") + .addHttpInterceptor(interceptor) + .withEventProcessor(processor) + .build(); + FlagsmithClient.Builder builder = FlagsmithClient.newBuilder() + .withConfiguration(config) + .setApiKey("api-key"); + if (withDefaultHandler) { + builder.setDefaultFlagValueFunction(FlagsmithClientTest::defaultHandler); + } + if (flagsUnavailable) { + FlagsmithApiWrapper mockApiWrapper = mock(FlagsmithApiWrapper.class); + when(mockApiWrapper.getConfig()).thenReturn(config); + when(mockApiWrapper.identifyUserWithTraits(any(), any(), anyBoolean(), anyBoolean())) + .thenReturn(null); + builder.withFlagsmithApiWrapper(mockApiWrapper); + } + return builder.build(); + } + + private static FlagsmithClient experimentClient(EventProcessor processor) { + return experimentClient(processor, false, false); + } + + private static Stream invalidEventsConfigs() { + return Stream.of( + Arguments.of(FlagsmithConfig.newBuilder().withEventsMaxBufferItems(10)), + Arguments.of(FlagsmithConfig.newBuilder().withEventsFlushIntervalMillis(10)), + Arguments.of(FlagsmithConfig.newBuilder() + .withEnableEvents(Boolean.TRUE).withEventsMaxBufferItems(0)), + Arguments.of(FlagsmithConfig.newBuilder() + .withEnableEvents(Boolean.TRUE).withEventsFlushIntervalMillis(-1))); + } + + @ParameterizedTest + @MethodSource("invalidEventsConfigs") + public void testInvalidEventsConfigThrowsAtBuild(FlagsmithConfig.Builder builder) { + assertThrows(IllegalArgumentException.class, builder::build); + } + + @Test + public void testEventsStayDisabledUnlessEnabled() { + FlagsmithConfig config = FlagsmithConfig.newBuilder().eventsUri("http://events-uri").build(); + + assertFalse(config.getEnableEvents()); + assertEquals("http://events-uri/", config.getEventsUri().toString()); + assertFalse(FlagsmithConfig.newBuilder().withEnableEvents(null).build().getEnableEvents()); + } + + @Test + public void testEventsSettingsReachTheProcessor() { + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .eventsUri("http://events-uri") + .withEnableEvents(Boolean.TRUE) + .withEventsMaxBufferItems(5) + .withEventsFlushIntervalMillis(0) + .build(); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(config) + .setApiKey("api-key") + .build(); + + EventProcessor processor = client.getEventProcessor(); + assertEquals("http://events-uri/v1/events", processor.getEventsEndpoint().toString()); + assertEquals(5, processor.getMaxBufferItems()); + assertEquals(0, processor.getFlushIntervalMillis()); + client.close(); + } + + @Test + public void testEventsInOfflineModeThrowsAtBuild() { + FlagsmithClient.Builder clientBuilder = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .withOfflineMode(Boolean.TRUE) + .withOfflineHandler(new DummyOfflineHandler()) + .withEventProcessor(mock(EventProcessor.class)) + .build()) + .setApiKey("api-key"); + + FlagsmithRuntimeError ex = assertThrows(FlagsmithRuntimeError.class, clientBuilder::build); + assertEquals("Events cannot be enabled in offline mode.", ex.getMessage()); + } + + @Test + public void testFailedBuildDoesNotStartTheEventProcessor() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient.Builder clientBuilder = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .withLocalEvaluation(true) + .withEventProcessor(processor) + .build()) + // Local evaluation needs a server key, so this build fails. + .setApiKey("api-key"); + + assertThrows(FlagsmithRuntimeError.class, clientBuilder::build); + verify(processor, never()).claim(); + verify(processor, never()).start(); + } + + @Test + public void testEventApisThrowWhenEventsAreDisabled() { + FlagsmithClient client = FlagsmithClient.newBuilder().setApiKey("api-key").build(); + + assertThrows(FlagsmithRuntimeError.class, + () -> client.getExperimentFlag("checkout_cta", "user-1")); + assertThrows(FlagsmithRuntimeError.class, () -> client.trackEvent("purchase")); + assertThrows(FlagsmithRuntimeError.class, + () -> client.trackExposureEvent("checkout_cta", "user-1", "treatment")); + assertTrue(client.flushEvents().isDone()); + } + + private static Stream> invalidEventCalls() { + return Stream.of( + (client) -> client.trackEvent("$flag_exposure"), + (client) -> client.trackEvent(null), + (client) -> client.trackEvent(" ", "user-1"), + (client) -> client.trackExposureEvent(null, "user-1", "treatment"), + (client) -> client.trackExposureEvent("", "user-1", "treatment")); + } + + @ParameterizedTest + @MethodSource("invalidEventCalls") + public void testEventApisRejectReservedAndBlankNames(Consumer call) { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = experimentClient(processor); + + assertThrows(IllegalArgumentException.class, () -> call.accept(client)); + verify(processor, never()).trackEvent(any(), any(), any(), any(), any()); + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testTrackExposureEventWithBlankIdentifierSendsNothing() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = experimentClient(processor); + + client.trackExposureEvent("checkout_cta", " ", "treatment"); + client.trackExposureEvent("checkout_cta", null, "treatment"); + + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testGetExperimentFlagRecordsAnExposurePerIdentityWhenEnrolled() + throws FlagsmithClientError { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = experimentClient(processor); + Map traits = new HashMap<>(); + traits.put("plan", "premium"); + traits.put("session_id", new TraitConfig("abc123", true)); + + Flag flag = (Flag) client.getExperimentFlag("checkout_cta", "user-1", traits); + client.getExperimentFlag("checkout_cta", "user-2"); + + assertEquals("treatment", flag.getVariant()); + assertEquals("SPLIT; weight=70.0", flag.getReason()); + assertEquals(42, flag.getExperiment().getId()); + assertEquals(Boolean.TRUE, flag.getExperiment().getInExperiment()); + Map metadata = Collections.singletonMap("experiment_id", 42); + verify(processor).trackExposureEvent( + "checkout_cta", "user-1", "treatment", traits, metadata); + verify(processor).trackExposureEvent( + "checkout_cta", "user-2", "treatment", new HashMap<>(), metadata); + } + + @Test + public void testGetExperimentFlagRecordsNoExposureWhenNotEnrolledOrDisabled() + throws FlagsmithClientError { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = experimentClient(processor); + + // bucketed but outside the rollout + assertEquals("control", + ((Flag) client.getExperimentFlag("pricing_page", "user-1")).getVariant()); + // no metadata at all + assertNull(((Flag) client.getExperimentFlag("some_feature", "user-1")).getExperiment()); + // enrolled, but the flag is off + assertEquals(Boolean.FALSE, + client.getExperimentFlag("disabled_feature", "user-1").getEnabled()); + + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + private static Stream experimentFlagFallbacks() { + return Stream.of( + Arguments.of("no_such_feature", false, FeatureNotFoundError.class), + Arguments.of("checkout_cta", true, FlagsmithApiError.class)); + } + + @ParameterizedTest + @MethodSource("experimentFlagFallbacks") + public void testGetExperimentFlagFallsBackToTheDefaultHandlerWithoutAnExposure( + String featureName, boolean flagsUnavailable, Class noDefaultError) + throws FlagsmithClientError { + EventProcessor processor = mock(EventProcessor.class); + + BaseFlag flag = experimentClient(processor, true, flagsUnavailable) + .getExperimentFlag(featureName, "user-1"); + assertTrue(flag instanceof DefaultFlag); + assertEquals(DEFAULT_FLAG_VALUE, flag.getValue()); + + FlagsmithClient noDefault = experimentClient(processor, false, flagsUnavailable); + assertThrows(noDefaultError, () -> noDefault.getExperimentFlag(featureName, "user-1")); + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testGetExperimentFlagRecordsNoExposureWithLocalEvaluation() + throws JsonProcessingException, FlagsmithClientError { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + EventProcessor processor = mock(EventProcessor.class); + interceptor.addRule() + .get(baseUrl + "/environment-document/") + .anyTimes() + .respond( + MapperFactory.getMapper() + .writeValueAsString(FlagsmithTestHelper.environmentModel()), + MEDIATYPE_JSON); + + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .baseUri(baseUrl) + .addHttpInterceptor(interceptor) + .withEventProcessor(processor) + .withLocalEvaluation(true) + .build()) + .setApiKey("ser.abcdefg") + .build(); + client.updateEnvironment(); + + BaseFlag flag = client.getExperimentFlag("some_feature", "user-1"); + + assertTrue(flag instanceof Flag); + assertNull(((Flag) flag).getVariant()); + assertNull(((Flag) flag).getExperiment()); + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + /** + * An events-enabled config serving the experiment identity flags, which records each events + * batch as "environment key|body". + */ + private static FlagsmithConfig.Builder recordingEventsConfig(List batches) { + MockInterceptor interceptor = new MockInterceptor(); + interceptor.addRule() + .post("http://bad-url/identities/") + .anyTimes() + .respond(FlagsmithTestHelper.getIdentitiesFlagsWithExperiment(), MEDIATYPE_JSON); + interceptor.addRule() + .post("http://events-uri/v1/events") + .headerMatches("Flagsmith-SDK-User-Agent", Pattern.compile("flagsmith-java-sdk/.*")) + .anyTimes() + .respond("{\"accepted\": 1, \"rejected\": []}", MEDIATYPE_JSON); + // Added before the MockInterceptor, which short-circuits the chain once it matches. + return FlagsmithConfig.newBuilder() + .baseUri("http://bad-url") + .addHttpInterceptor((chain) -> { + Request request = chain.request(); + if (request.url().toString().endsWith("/v1/events")) { + Buffer buffer = new Buffer(); + request.body().writeTo(buffer); + batches.add(request.header("X-Environment-Key") + "|" + buffer.readUtf8()); + } + return chain.proceed(request); + }) + .addHttpInterceptor(interceptor) + .eventsUri("http://events-uri") + .withEnableEvents(Boolean.TRUE) + .withEventsFlushIntervalMillis(0); + } + + @Test + public void testCloseFlushesBufferedEvents() throws FlagsmithClientError, IOException { + List batches = Collections.synchronizedList(new ArrayList<>()); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(recordingEventsConfig(batches).build()) + .setApiKey("api-key") + .build(); + + client.getExperimentFlag("checkout_cta", "user-1"); + assertTrue(batches.isEmpty()); + + client.close(); + + assertEquals(1, batches.size()); + String body = batches.get(0).substring(batches.get(0).indexOf('|') + 1); + JsonNode events = MapperFactory.getMapper().readTree(body).get("events"); + assertEquals(1, events.size()); + assertEquals("$flag_exposure", events.get(0).get("event").asText()); + assertEquals("checkout_cta", events.get(0).get("feature_name").asText()); + assertEquals("user-1", events.get(0).get("identifier").asText()); + assertEquals("treatment", events.get(0).get("value").asText()); + assertEquals(42, events.get(0).get("metadata").get("experiment_id").asInt()); + } + + @Test + public void testRebuildingClosesThePreviousEventProcessor() { + List batches = Collections.synchronizedList(new ArrayList<>()); + FlagsmithClient.Builder builder = FlagsmithClient.newBuilder() + .withConfiguration(recordingEventsConfig(batches).build()) + .setApiKey("api-key"); + FlagsmithClient client = builder.build(); + client.trackEvent("purchase", "user-1"); + + builder.build(); + assertEquals(1, batches.size()); + + client.trackEvent("purchase", "user-2"); + client.close(); + assertEquals(2, batches.size()); + } + + @Test + public void testEventsRequestLeavesOutCustomHeaders() { + List requests = Collections.synchronizedList(new ArrayList<>()); + MockInterceptor interceptor = new MockInterceptor(); + interceptor.addRule() + .post("http://events-uri/v1/events") + .anyTimes() + .respond("{\"accepted\": 1, \"rejected\": []}", MEDIATYPE_JSON); + HashMap customHeaders = new HashMap<>(); + customHeaders.put("Authorization", "Bearer flags-api-only"); + customHeaders.put("X-Environment-Key", "flags-api-only-key"); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .baseUri("http://bad-url") + .addHttpInterceptor((chain) -> { + requests.add(chain.request()); + return chain.proceed(chain.request()); + }) + .addHttpInterceptor(interceptor) + .eventsUri("http://events-uri") + .withEnableEvents(Boolean.TRUE) + .withEventsFlushIntervalMillis(0) + .build()) + .withCustomHttpHeaders(customHeaders) + .setApiKey("api-key") + .build(); + + client.trackEvent("purchase", "user-1"); + client.close(); + + assertEquals(1, requests.size()); + assertNull(requests.get(0).header("Authorization")); + assertEquals(Collections.singletonList("api-key"), + requests.get(0).headers("X-Environment-Key")); + assertTrue(requests.get(0).header("User-Agent").startsWith("flagsmith-java-sdk/")); + } + + @Test + public void testClientsSharingAConfigSendEventsUnderTheirOwnKeys() throws Exception { + List batches = Collections.synchronizedList(new ArrayList<>()); + FlagsmithConfig config = recordingEventsConfig(batches).build(); + FlagsmithClient clientA = FlagsmithClient.newBuilder() + .withConfiguration(config).setApiKey("key-a").build(); + FlagsmithClient clientB = FlagsmithClient.newBuilder() + .withConfiguration(config).setApiKey("key-b").build(); + + assertNotSame(clientA.getEventProcessor(), clientB.getEventProcessor()); + + clientA.trackEvent("purchase", "user-a"); + clientB.trackEvent("purchase", "user-b"); + clientA.flushEvents().get(5, TimeUnit.SECONDS); + clientB.flushEvents().get(5, TimeUnit.SECONDS); + + assertEquals(2, batches.size()); + assertTrue(batches.get(0).startsWith("key-a|")); + assertTrue(batches.get(0).contains("user-a") && !batches.get(0).contains("user-b")); + assertTrue(batches.get(1).startsWith("key-b|")); + assertTrue(batches.get(1).contains("user-b") && !batches.get(1).contains("user-a")); + clientA.close(); + clientB.close(); + } + + @Test + public void testClosingOneClientLeavesAnotherOnTheSameConfigTracking() throws Exception { + List batches = Collections.synchronizedList(new ArrayList<>()); + FlagsmithConfig config = recordingEventsConfig(batches).build(); + FlagsmithClient clientA = FlagsmithClient.newBuilder() + .withConfiguration(config).setApiKey("key-a").build(); + FlagsmithClient clientB = FlagsmithClient.newBuilder() + .withConfiguration(config).setApiKey("key-b").build(); + + clientA.close(); + clientB.trackEvent("purchase", "user-b"); + clientB.flushEvents().get(5, TimeUnit.SECONDS); + + assertEquals(1, batches.size()); + assertTrue(batches.get(0).startsWith("key-b|")); + assertTrue(batches.get(0).contains("user-b")); + clientB.close(); + } + + @Test + public void testCustomApiWrapperWithAnEventsConfigDeliversEvents() throws Exception { + List batches = Collections.synchronizedList(new ArrayList<>()); + FlagsmithApiWrapper wrapper = new FlagsmithApiWrapper( + FlagsmithConfig.newBuilder().baseUri("http://bad-url").build(), + null, new FlagsmithLogger(), "wrapper-key"); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withFlagsmithApiWrapper(wrapper) + .withConfiguration(recordingEventsConfig(batches).build()) + .setApiKey("wrapper-key") + .build(); + + client.trackEvent("purchase", "user-1"); + client.flushEvents().get(5, TimeUnit.SECONDS); + + assertEquals(1, batches.size()); + assertTrue(batches.get(0).startsWith("wrapper-key|")); + assertTrue(batches.get(0).contains("user-1")); + client.close(); + } + + @Test + public void testAnInjectedEventProcessorBacksOnlyOneClient() throws Exception { + List batches = Collections.synchronizedList(new ArrayList<>()); + FlagsmithConfig.Builder configBuilder = recordingEventsConfig(batches); + FlagsmithConfig probe = configBuilder.build(); + FlagsmithConfig config = configBuilder + .withEventProcessor(new EventProcessor( + probe.getHttpClient(), probe.getEventsUri(), 1000, 0)) + .build(); + FlagsmithClient clientA = FlagsmithClient.newBuilder() + .withConfiguration(config).setApiKey("key-a").build(); + FlagsmithClient.Builder clientBBuilder = FlagsmithClient.newBuilder() + .withConfiguration(config).setApiKey("key-b"); + + assertThrows(FlagsmithRuntimeError.class, clientBBuilder::build); + + clientA.trackEvent("purchase", "user-a"); + clientA.flushEvents().get(5, TimeUnit.SECONDS); + assertEquals(1, batches.size()); + assertTrue(batches.get(0).startsWith("key-a|")); + clientA.close(); + } + + @Test + public void testFailedClaimDoesNotStartPolling() { + FlagsmithConfig probe = FlagsmithConfig.newBuilder().build(); + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .withLocalEvaluation(true) + .withEventProcessor(new EventProcessor( + probe.getHttpClient(), probe.getEventsUri(), 1000, 0)) + .build(); + FlagsmithClient clientA = FlagsmithClient.newBuilder() + .withConfiguration(config) + .withPollingManager(mock(PollingManager.class)) + .setApiKey("ser.key-a") + .build(); + PollingManager pollingB = mock(PollingManager.class); + FlagsmithClient.Builder clientBBuilder = FlagsmithClient.newBuilder() + .withConfiguration(config) + .withPollingManager(pollingB) + .setApiKey("ser.key-b"); + + assertThrows(FlagsmithRuntimeError.class, clientBBuilder::build); + verify(pollingB, never()).startPolling(); + clientA.close(); + } } diff --git a/src/test/java/com/flagsmith/FlagsmithRetryTest.java b/src/test/java/com/flagsmith/FlagsmithRetryTest.java index cd021087..eb54ccbc 100644 --- a/src/test/java/com/flagsmith/FlagsmithRetryTest.java +++ b/src/test/java/com/flagsmith/FlagsmithRetryTest.java @@ -2,10 +2,14 @@ import com.flagsmith.config.Retry; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; import static org.junit.jupiter.api.Assertions.*; import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashSet; import java.util.List; public class FlagsmithRetryTest { @@ -66,6 +70,33 @@ public void FlagsmithRetry_validateSleep() { assertTrue(attempts.equals(3)); } + @ParameterizedTest + @CsvSource({ + "1, 503, true", "1, , true", + "1, 400, false", "1, 404, false", + "2, 503, false", "2, , false"}) + public void FlagsmithRetry_statusForcelistOnly_retriesListedStatusesWithinBudget( + int attempts, Integer status, boolean expected) { + Retry retry = new Retry(2); + retry.setStatusForcelist(new HashSet<>(Arrays.asList(500, 502, 503, 504))); + retry.setStatusForcelistOnly(Boolean.TRUE); + for (int i = 0; i < attempts; i++) { + retry.retryAttempted(); + } + + assertEquals(expected, retry.isRetry(status)); + } + + @Test + public void FlagsmithRetry_defaultPolicyIsUnchanged() { + Retry retry = new Retry(1); + + assertFalse(retry.getStatusForcelistOnly()); + retry.retryAttempted(); + assertTrue(retry.isRetry(503)); + assertFalse(retry.isRetry(401)); + } + @Test public void FlagsmithRetry_shouldNotExceedBackoffMax() { Retry retryObject = new Retry(7); diff --git a/src/test/java/com/flagsmith/FlagsmithTestHelper.java b/src/test/java/com/flagsmith/FlagsmithTestHelper.java index d1922d56..4fe1dcd0 100644 --- a/src/test/java/com/flagsmith/FlagsmithTestHelper.java +++ b/src/test/java/com/flagsmith/FlagsmithTestHelper.java @@ -431,6 +431,26 @@ public static String getIdentitiesFlags() { return featureJson; } + /** Identity flags covering every branch of the experiment gate. */ + public static String getIdentitiesFlagsWithExperiment() { + return "{\"traits\": [], \"flags\": [\n" + + " {\"id\": 1, \"feature\": {\"id\": 1, \"name\": \"checkout_cta\", \"type\": \"MULTIVARIATE\"},\n" + + " \"feature_state_value\": \"buy-now\", \"enabled\": true, \"variant\": \"treatment\",\n" + + " \"reason\": \"SPLIT; weight=70.0\", \"metadata\": {\"experiment\":\n" + + " {\"id\": 42, \"name\": \"checkout_experiment\", \"in_experiment\": true}}},\n" + + " {\"id\": 2, \"feature\": {\"id\": 2, \"name\": \"pricing_page\", \"type\": \"MULTIVARIATE\"},\n" + + " \"feature_state_value\": \"old-pricing\", \"enabled\": true, \"variant\": \"control\",\n" + + " \"reason\": \"SPLIT; weight=30.0\", \"metadata\": {\"experiment\":\n" + + " {\"id\": 43, \"name\": \"pricing_experiment\", \"in_experiment\": false}}},\n" + + " {\"id\": 3, \"feature\": {\"id\": 3, \"name\": \"some_feature\", \"type\": \"STANDARD\"},\n" + + " \"feature_state_value\": \"some-value\", \"enabled\": true},\n" + + " {\"id\": 4, \"feature\": {\"id\": 4, \"name\": \"disabled_feature\", \"type\": \"MULTIVARIATE\"},\n" + + " \"feature_state_value\": \"off\", \"enabled\": false, \"variant\": \"treatment\",\n" + + " \"reason\": \"SPLIT; weight=50.0\", \"metadata\": {\"experiment\":\n" + + " {\"id\": 44, \"name\": \"disabled_experiment\", \"in_experiment\": true}}}\n" + + "]}"; + } + public static Future futurableReturn(T response) { CompletableFuture promise = new CompletableFuture<>(); promise.complete(response); diff --git a/src/test/java/com/flagsmith/models/FeatureStateModelTest.java b/src/test/java/com/flagsmith/models/FeatureStateModelTest.java new file mode 100644 index 00000000..2a692a6d --- /dev/null +++ b/src/test/java/com/flagsmith/models/FeatureStateModelTest.java @@ -0,0 +1,82 @@ +package com.flagsmith.models; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.flagsmith.MapperFactory; +import com.flagsmith.models.features.FeatureStateModel; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; + +public class FeatureStateModelTest { + + private static FeatureStateModel parse(String json) throws JsonProcessingException { + return MapperFactory.getMapper().readValue(json, FeatureStateModel.class); + } + + @Test + public void parsesVariantReasonAndExperiment() throws JsonProcessingException { + FeatureStateModel featureState = parse( + "{\"feature\": {\"id\": 220175, \"name\": \"checkout_cta\", \"type\": \"MULTIVARIATE\"}," + + " \"enabled\": true," + + " \"feature_state_value\": \"buy-now\"," + + " \"variant\": \"treatment\"," + + " \"reason\": \"SPLIT; weight=70.0\"," + + " \"metadata\": {\"experiment\": {\"id\": 167, \"name\": \"flutter_demo_exp\"," + + " \"in_experiment\": true}}}"); + + assertEquals("treatment", featureState.getVariant()); + assertEquals("SPLIT; weight=70.0", featureState.getReason()); + + ExperimentMetadata experiment = featureState.getMetadata().getExperiment(); + assertEquals(167, experiment.getId()); + assertEquals("flutter_demo_exp", experiment.getName()); + assertEquals(Boolean.TRUE, experiment.getInExperiment()); + + Flag flag = Flag.fromFeatureStateModel(featureState); + assertEquals("treatment", flag.getVariant()); + assertEquals("SPLIT; weight=70.0", flag.getReason()); + assertSame(experiment, flag.getExperiment()); + } + + @Test + public void ignoresUnknownMetadataKeys() throws JsonProcessingException { + FeatureStateModel featureState = parse( + "{\"feature\": {\"id\": 1, \"name\": \"checkout_cta\"}," + + " \"enabled\": true," + + " \"metadata\": {\"something_else\": {\"a\": 1}, \"experiment\": {\"id\": 3," + + " \"in_experiment\": true}}}"); + + assertNotNull(featureState.getMetadata()); + assertEquals(3, featureState.getMetadata().getExperiment().getId()); + } + + @Test + public void inExperimentDefaultsToFalseWhenMissing() throws JsonProcessingException { + FeatureStateModel featureState = parse( + "{\"feature\": {\"id\": 1, \"name\": \"checkout_cta\"}," + + " \"enabled\": true," + + " \"metadata\": {\"experiment\": {\"id\": 3, \"name\": \"exp\"}}}"); + + assertEquals(Boolean.FALSE, featureState.getMetadata().getExperiment().getInExperiment()); + } + + @Test + public void parsesPayloadWithoutAnyExperimentFields() throws JsonProcessingException { + FeatureStateModel featureState = parse( + "{\"feature\": {\"id\": 1, \"name\": \"some_feature\"}," + + " \"enabled\": true, \"feature_state_value\": \"some-value\"}"); + + assertNull(featureState.getVariant()); + assertNull(featureState.getReason()); + assertNull(featureState.getMetadata()); + assertEquals("some-value", featureState.getValue()); + + Flag flag = Flag.fromFeatureStateModel(featureState); + assertNull(flag.getVariant()); + assertNull(flag.getReason()); + assertNull(flag.getExperiment()); + } +} diff --git a/src/test/java/com/flagsmith/threads/EventProcessorTest.java b/src/test/java/com/flagsmith/threads/EventProcessorTest.java new file mode 100644 index 00000000..e161a433 --- /dev/null +++ b/src/test/java/com/flagsmith/threads/EventProcessorTest.java @@ -0,0 +1,896 @@ +package com.flagsmith.threads; + +import static okhttp3.mock.MediaTypes.MEDIATYPE_JSON; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.contains; +import static org.mockito.ArgumentMatchers.startsWith; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.mockingDetails; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.util.RawValue; +import com.flagsmith.FlagsmithLogger; +import com.flagsmith.MapperFactory; +import com.flagsmith.config.FlagsmithConfig; +import com.flagsmith.interfaces.FlagsmithSdk; +import com.flagsmith.models.TraitConfig; +import java.io.IOException; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Consumer; +import java.util.regex.Pattern; +import java.util.stream.Stream; +import lombok.SneakyThrows; +import okhttp3.HttpUrl; +import okhttp3.Interceptor; +import okhttp3.MediaType; +import okhttp3.OkHttpClient; +import okhttp3.Protocol; +import okhttp3.Request; +import okhttp3.RequestBody; +import okhttp3.Response; +import okhttp3.ResponseBody; +import okhttp3.mock.MockInterceptor; +import okio.Buffer; +import org.mockito.invocation.Invocation; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.CsvSource; +import org.junit.jupiter.params.provider.MethodSource; + +public class EventProcessorTest { + + private static final String EVENTS_URI = "http://events-uri/"; + private static final String EVENTS_ENDPOINT = EVENTS_URI + "v1/events"; + private static final String ACCEPTED_BODY = "{\"accepted\": 1, \"rejected\": []}"; + private static final long WAIT_SECONDS = 10L; + + private MockInterceptor interceptor; + private RecordingInterceptor recorder; + private FlagsmithSdk api; + private EventProcessor eventProcessor; + + @BeforeEach + public void init() { + interceptor = new MockInterceptor(); + recorder = new RecordingInterceptor(); + api = mock(FlagsmithSdk.class); + when(api.newPostRequest(any(), any())).thenAnswer((invocation) -> new Request.Builder() + .url((HttpUrl) invocation.getArgument(0)) + .header("X-Environment-Key", "api-key") + .header("User-Agent", "flagsmith-java-sdk/test") + .addHeader("Accept", "application/json") + .post((RequestBody) invocation.getArgument(1)) + .build()); + } + + @AfterEach + public void tearDown() { + if (eventProcessor != null) { + // Appended last, so it only catches whatever close() flushes out of the buffer. + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + eventProcessor.close(); + eventProcessor = null; + } + } + + private EventProcessor newProcessor(int maxBufferItems, int flushIntervalMillis) { + return newProcessor(maxBufferItems, flushIntervalMillis, null); + } + + private EventProcessor newProcessor( + int maxBufferItems, int flushIntervalMillis, Interceptor extra) { + OkHttpClient.Builder clientBuilder = new OkHttpClient.Builder().addInterceptor(recorder); + if (extra != null) { + clientBuilder.addInterceptor(extra); + } + OkHttpClient client = clientBuilder.addInterceptor(interceptor).build(); + eventProcessor = new EventProcessor( + HttpUrl.get(EVENTS_URI), + maxBufferItems, + flushIntervalMillis, + new RequestProcessor(client, new FlagsmithLogger())); + eventProcessor.setApi(api); + return eventProcessor; + } + + @SneakyThrows + private void flushAndWait(EventProcessor processor) { + processor.flush().get(WAIT_SECONDS, TimeUnit.SECONDS); + } + + private static FlagsmithLogger mockLogger(EventProcessor processor) { + FlagsmithLogger logger = mock(FlagsmithLogger.class); + processor.setLogger(logger); + return logger; + } + + /** Parse a buffered event's pre-serialised traits or metadata. */ + @SneakyThrows + private static JsonNode json(Object buffered) { + return MapperFactory.getMapper().readTree(((RawValue) buffered).rawValue().toString()); + } + + @Test + public void trackEvent_buffersStringifiedValueSdkVersionAndTimestamp() { + EventProcessor processor = newProcessor(1000, 0); + + Map traits = new HashMap<>(); + traits.put("plan", "premium"); + Map metadata = new HashMap<>(); + metadata.put("source", "checkout"); + + long before = System.currentTimeMillis(); + processor.trackEvent("purchase", "user-123", 49.0, traits, metadata); + + assertEquals(1, processor.bufferedEvents().size()); + + Map event = processor.bufferedEvents().get(0); + assertEquals( + Arrays.asList("event", "feature_name", "identifier", "value", "traits", "metadata", + "timestamp"), + new ArrayList<>(event.keySet())); + assertEquals("purchase", event.get("event")); + assertNull(event.get("feature_name")); + assertEquals("user-123", event.get("identifier")); + assertEquals("49.0", event.get("value")); + assertEquals(MapperFactory.getMapper().valueToTree(traits), json(event.get("traits"))); + + JsonNode eventMetadata = json(event.get("metadata")); + assertEquals("checkout", eventMetadata.get("source").asText()); + assertNotNull(eventMetadata.get("sdk_version")); + + long timestamp = (Long) event.get("timestamp"); + assertTrue(timestamp >= before); + + processor.trackEvent("purchase", "user-123", null, null, null); + assertNull(processor.bufferedEvents().get(1).get("value")); + assertNull(processor.bufferedEvents().get(1).get("traits")); + + processor.trackEvent("purchase", "user-123", null, new HashMap<>(), null); + assertNull(processor.bufferedEvents().get(2).get("traits")); + } + + @Test + public void trackExposureEvent_dedupesOnlyEqualExposuresWithinTheFlushWindow() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + Map metadata = Collections.singletonMap("experiment_id", 42); + + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, metadata); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, metadata); + assertEquals(1, processor.bufferedEvents().size()); + + processor.trackExposureEvent("checkout_cta", "user-2", "treatment", null, metadata); + processor.trackExposureEvent("checkout_cta", "user-1", "control", null, metadata); + processor.trackExposureEvent("pricing_page", "user-1", "treatment", null, metadata); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, + Collections.singletonMap("experiment_id", 43)); + processor.trackExposureEvent("checkout_cta", "a\u0000b", "c", null, metadata); + processor.trackExposureEvent("checkout_cta", "a", "b\u0000c", null, metadata); + assertEquals(7, processor.bufferedEvents().size()); + + processor.trackEvent("purchase", "user-1", "49.00", null, null); + processor.trackEvent("purchase", "user-1", "49.00", null, null); + assertEquals(9, processor.bufferedEvents().size()); + + flushAndWait(processor); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, metadata); + assertEquals(1, processor.bufferedEvents().size()); + assertEquals(1, recorder.count()); + } + + @Test + @SneakyThrows + public void flush_postsTheBatchToTheEventsEndpoint() { + EventProcessor processor = newProcessor(1000, 0); + flushAndWait(processor); + assertEquals(0, recorder.count()); + + interceptor.addRule() + .post(EVENTS_ENDPOINT) + .headerMatches("X-Environment-Key", Pattern.compile("api-key")) + .headerMatches("Flagsmith-SDK-User-Agent", Pattern.compile("flagsmith-java-sdk/.*")) + .headerMatches("Content-Type", Pattern.compile("application/json; charset=utf-8")) + .anyTimes() + .respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, + Collections.singletonMap("experiment_id", 42)); + flushAndWait(processor); + + assertEquals(1, recorder.count()); + + JsonNode body = MapperFactory.getMapper().readTree(recorder.bodies().get(0)); + assertEquals(1, body.get("events").size()); + + JsonNode event = body.get("events").get(0); + assertEquals("$flag_exposure", event.get("event").asText()); + assertEquals("checkout_cta", event.get("feature_name").asText()); + assertEquals("user-1", event.get("identifier").asText()); + assertEquals("treatment", event.get("value").asText()); + assertEquals(42, event.get("metadata").get("experiment_id").asInt()); + assertTrue(event.get("metadata").has("sdk_version")); + assertTrue(event.get("timestamp").isNumber()); + assertTrue(processor.bufferedEvents().isEmpty()); + } + + @Test + @SneakyThrows + public void trackEvent_flushesWhenTheBufferIsFull() { + EventProcessor processor = newProcessor(2, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + assertEquals(1, processor.bufferedEvents().size()); + + processor.trackEvent("purchase", "user-2", "2", null, null); + assertTrue(processor.bufferedEvents().isEmpty()); + + flushAndWait(processor); + + assertEquals(1, recorder.count()); + JsonNode body = MapperFactory.getMapper().readTree(recorder.bodies().get(0)); + assertEquals(2, body.get("events").size()); + } + + @Test + @SneakyThrows + public void flush_completesOnlyAfterTheInFlightPostCompletes() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule() + .post(EVENTS_ENDPOINT) + .anyTimes() + .delay(600) + .respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + + long start = System.currentTimeMillis(); + flushAndWait(processor); + long elapsed = System.currentTimeMillis() - start; + + assertEquals(1, recorder.count()); + assertTrue(elapsed >= 500, "flush() returned after " + elapsed + "ms, before the POST"); + } + + /** + * The first {@code failures} attempts fail with {@code status}, or with a connection failure + * when it is empty; the rest are accepted. + */ + @ParameterizedTest + @CsvSource({ + "503, 1, 2", "501, 1, 2", "500, 2, 2", + "400, 1, 1", + ", 1, 2", ", 2, 2"}) + @SneakyThrows + public void flush_retriesOnceOnAServerErrorOrConnectionFailure( + Integer status, int failures, int expectedAttempts) { + EventProcessor processor = status == null + ? newProcessor(1000, 0, new FailingInterceptor(failures)) + : newProcessor(1000, 0); + if (status != null) { + interceptor.addRule().post(EVENTS_ENDPOINT).times(failures).respond(status); + } + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(expectedAttempts, recorder.count()); + assertEquals(recorder.bodies().get(0), recorder.bodies().get(expectedAttempts - 1)); + assertTrue(processor.bufferedEvents().isEmpty()); + } + + @Test + @SneakyThrows + public void flush_waitsForABatchAnotherThreadIsAlreadySending() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + CountDownLatch insideSend = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + doAnswer((invocation) -> { + insideSend.countDown(); + release.await(WAIT_SECONDS, TimeUnit.SECONDS); + return new Request.Builder() + .url((HttpUrl) invocation.getArgument(0)) + .header("X-Environment-Key", "api-key") + .post((RequestBody) invocation.getArgument(1)) + .build(); + }).when(api).newPostRequest(any(), any()); + + processor.trackEvent("purchase", "user-1", "1", null, null); + + Thread sender = new Thread(processor::flush, "test-sender"); + sender.start(); + assertTrue(insideSend.await(WAIT_SECONDS, TimeUnit.SECONDS)); + + // The buffer is already empty here, but the batch is still on its way out. + CompletableFuture second = processor.flush(); + assertFalse(second.isDone(), "flush() completed while another thread was mid-send"); + + release.countDown(); + second.get(WAIT_SECONDS, TimeUnit.SECONDS); + sender.join(TimeUnit.SECONDS.toMillis(WAIT_SECONDS)); + + assertEquals(1, recorder.count()); + } + + private static Stream> brokenSends() { + return Stream.of( + (test) -> doThrow(new IllegalStateException("boom")).when(test.api) + .newPostRequest(any(), any()), + (test) -> test.eventProcessor.getRequestProcessor().close(), + (test) -> test.eventProcessor.setApi(null)); + } + + @ParameterizedTest + @MethodSource("brokenSends") + public void trackEvent_dropsTheBatchWithoutThrowingWhenItCannotBeSent( + Consumer breakSend) { + EventProcessor processor = newProcessor(1, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + breakSend.accept(this); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(0, recorder.count()); + assertTrue(processor.bufferedEvents().isEmpty()); + } + + @Test + @SneakyThrows + public void trackEvent_isANoOpAfterClose() { + EventProcessor processor = newProcessor(2, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.close(); + eventProcessor = null; + int postsAfterClose = recorder.count(); + + processor.trackEvent("purchase", "user-1", "1", null, null); + processor.trackEvent("purchase", "user-2", "2", null, null); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + + assertTrue(processor.bufferedEvents().isEmpty()); + flushAndWait(processor); + assertEquals(postsAfterClose, recorder.count()); + } + + @Test + @SneakyThrows + public void trackEvent_refusesAnEventThatRacesCloseInsteadOfStrandingIt() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + // Serialising the traits happens after the first closed check and before the buffer lock, + // so a getter that blocks parks the tracking thread exactly in that window. + CountDownLatch serialising = new CountDownLatch(1); + CountDownLatch proceed = new CountDownLatch(1); + Object slowTrait = new Object() { + @SuppressWarnings("unused") + public String getValue() throws InterruptedException { + serialising.countDown(); + proceed.await(WAIT_SECONDS, TimeUnit.SECONDS); + return "slow"; + } + }; + Thread tracker = new Thread(() -> processor.trackEvent( + "purchase", "user-1", "1", Collections.singletonMap("slow", slowTrait), null)); + tracker.start(); + assertTrue(serialising.await(WAIT_SECONDS, TimeUnit.SECONDS)); + + processor.close(); + eventProcessor = null; + proceed.countDown(); + tracker.join(TimeUnit.SECONDS.toMillis(WAIT_SECONDS)); + + assertTrue(processor.bufferedEvents().isEmpty(), "an event was stranded after close()"); + assertEquals(0, recorder.count()); + } + + @Test + public void trackExposureEvent_unwrapsTraitConfigsAndDropsTransientTraits() { + EventProcessor processor = newProcessor(1000, 0); + + Map traits = new LinkedHashMap<>(); + traits.put("plan", "premium"); + traits.put("tier", new TraitConfig("gold", false)); + traits.put("session_id", new TraitConfig("abc123", true)); + + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", traits, null); + + JsonNode buffered = json(processor.bufferedEvents().get(0).get("traits")); + assertEquals(2, buffered.size()); + assertEquals("premium", buffered.get("plan").asText()); + assertEquals("gold", buffered.get("tier").asText()); + assertFalse(buffered.has("session_id"), "a transient trait reached the events API"); + } + + @Test + public void trackEvent_copiesTraitsAndMetadataDeeplyAtBufferTime() { + EventProcessor processor = newProcessor(1000, 0); + + Map traits = new LinkedHashMap<>(); + traits.put("plan", "premium"); + Map nested = new HashMap<>(); + nested.put("step", "payment"); + Map metadata = new HashMap<>(); + metadata.put("context", nested); + processor.trackEvent("purchase", "user-1", "1", traits, metadata); + + traits.put("plan", "mutated"); + traits.put("added_later", "nope"); + nested.put("step", "mutated"); + + Map event = processor.bufferedEvents().get(0); + assertEquals( + MapperFactory.getMapper().valueToTree(Collections.singletonMap("plan", "premium")), + json(event.get("traits"))); + assertEquals("payment", + json(event.get("metadata")).get("context").get("step").asText()); + } + + @Test + @SneakyThrows + public void trackEvent_dropsOnlyTheEventWhoseValuesCannotBeSerialised() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + // Jackson has no serialiser for a bean without properties. + processor.trackEvent("purchase", "user-2", "2", + Collections.singletonMap("opaque", new Object()), null); + processor.trackEvent("purchase", "user-3", "3", null, + Collections.singletonMap("opaque", new Object())); + Map cyclic = new HashMap<>(); + cyclic.put("self", cyclic); + List cyclicList = new ArrayList<>(); + cyclicList.add(cyclicList); + processor.trackEvent("purchase", "user-5", "5", cyclic, null); + processor.trackEvent("purchase", "user-6", "6", null, + Collections.singletonMap("list", cyclicList)); + processor.trackEvent("purchase", "user-4", "4", null, null); + + assertEquals(2, processor.bufferedEvents().size()); + flushAndWait(processor); + + assertEquals(1, recorder.count()); + JsonNode events = MapperFactory.getMapper().readTree(recorder.bodies().get(0)).get("events"); + assertEquals(2, events.size()); + assertEquals("user-1", events.get(0).get("identifier").asText()); + assertEquals("user-4", events.get(1).get("identifier").asText()); + } + + @Test + @SneakyThrows + public void start_flushesOnTheTimer() { + EventProcessor processor = newProcessor(1000, 100); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.start(); + processor.trackEvent("purchase", "user-1", "1", null, null); + + assertTrue(recorder.awaitFirstRequest(WAIT_SECONDS), "the timer never flushed"); + assertTrue(processor.bufferedEvents().isEmpty()); + } + + @Test + @SneakyThrows + public void close_flushesAndStopsTheSchedulerThread() { + EventProcessor processor = newProcessor(1000, 500); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.start(); + processor.trackEvent("purchase", "user-1", "1", null, null); + processor.close(); + eventProcessor = null; + + assertEquals(1, recorder.count()); + assertTrue(processor.getScheduler().awaitTermination(WAIT_SECONDS, TimeUnit.SECONDS)); + assertTrue(processor.getScheduler().isTerminated()); + } + + @Test + @SneakyThrows + public void start_isANoOpAfterClose() { + EventProcessor processor = newProcessor(1000, 100); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + processor.close(); + eventProcessor = null; + + processor.start(); + + assertTrue(processor.getScheduler().isShutdown()); + } + + private static Stream clientTimeouts() { + FlagsmithConfig longRead = FlagsmithConfig.newBuilder().readTimeout(30_000).build(); + return Stream.of( + // Two attempts at 2s connect + 5s write + 5s read, and 200ms backoff before the second. + Arguments.of(FlagsmithConfig.newBuilder().build().getHttpClient(), 2 * 12_000 + 200, + 12_000), + Arguments.of(longRead.getHttpClient(), 2 * 37_000 + 200, 37_000), + Arguments.of(new OkHttpClient.Builder().callTimeout(4, TimeUnit.SECONDS).build(), + 2 * 4_000 + 200, 4_000), + Arguments.of(new OkHttpClient.Builder().readTimeout(0, TimeUnit.SECONDS).build(), + 2 * EventProcessor.CLOSE_TIMEOUT_MILLIS + 200, EventProcessor.CLOSE_TIMEOUT_MILLIS)); + } + + @ParameterizedTest + @MethodSource("clientTimeouts") + public void close_waitsForOneBatchAtTheClientTimeouts(OkHttpClient client, long expected, + long callTimeout) { + EventProcessor processor = new EventProcessor(client, HttpUrl.get(EVENTS_URI), 1, 0); + + assertEquals(expected, processor.getCloseTimeoutMillis()); + assertEquals(callTimeout, processor.getRequestProcessor().getClient().callTimeoutMillis()); + processor.close(); + } + + @Test + @SneakyThrows + public void close_sendsTheLastBatchAlongsideTheBatchesInFlight() { + AcceptingInterceptor eventsApi = AcceptingInterceptor.blocked(); + EventProcessor processor = newProcessor(1, 0, eventsApi); + eventProcessor = null; + FlagsmithLogger logger = mockLogger(processor); + + for (int i = 0; i <= EventProcessor.MAX_IN_FLIGHT_BATCHES; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + } + assertTrue(recorder.awaitCount(EventProcessor.MAX_IN_FLIGHT_BATCHES)); + assertEquals(1, processor.bufferedEvents().size()); + + CompletableFuture closed = CompletableFuture.runAsync(processor::close); + assertTrue(recorder.awaitCount(EventProcessor.MAX_IN_FLIGHT_BATCHES + 1)); + eventsApi.release(); + closed.get(WAIT_SECONDS, TimeUnit.SECONDS); + + assertEquals(EventProcessor.MAX_IN_FLIGHT_BATCHES + 1, deliveredEvents()); + assertEquals(Collections.emptyList(), errorCalls(logger)); + } + + @Test + @SneakyThrows + public void close_stopsWaitingAtTheTimeout() { + AcceptingInterceptor eventsApi = AcceptingInterceptor.blocked(); + EventProcessor processor = newProcessor(1000, 0, eventsApi); + eventProcessor = null; + processor.setCloseTimeoutMillis(400); + FlagsmithLogger logger = mockLogger(processor); + + try { + processor.trackEvent("purchase", "user-1", "1", null, null); + long start = System.nanoTime(); + processor.close(); + long elapsedMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start); + + verify(logger).error(contains("Stopped waiting")); + assertTrue(elapsedMillis >= 400, "close() returned after " + elapsedMillis + "ms"); + assertTrue(elapsedMillis < TimeUnit.SECONDS.toMillis(WAIT_SECONDS) / 2, + "close() waited " + elapsedMillis + "ms"); + } finally { + eventsApi.release(); + } + } + + @Test + @SneakyThrows + public void settle_sendsEventsThatWaitedBehindTheInFlightLimit() { + AcceptingInterceptor eventsApi = AcceptingInterceptor.blocked(); + EventProcessor processor = newProcessor(1000, 0, eventsApi); + FlagsmithLogger logger = mockLogger(processor); + + for (int i = 0; i < EventProcessor.MAX_IN_FLIGHT_BATCHES; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + processor.flush(); + } + for (int i = 0; i < 100; i++) { + processor.trackEvent("purchase", "waiting-" + i, "1", null, null); + processor.flush(); + } + + assertTrue(recorder.awaitCount(EventProcessor.MAX_IN_FLIGHT_BATCHES)); + assertEquals(100, processor.bufferedEvents().size()); + + eventsApi.release(); + + assertTrue(awaitDelivered(EventProcessor.MAX_IN_FLIGHT_BATCHES + 100)); + assertEquals(Collections.emptyList(), errorCalls(logger)); + } + + @Test + @SneakyThrows + public void flush_waitsForTheEventsBufferedBehindTheInFlightLimit() { + CountDownLatch inFlightReleased = new CountDownLatch(1); + CountDownLatch waitingReleased = new CountDownLatch(1); + EventProcessor processor = newProcessor(1000, 0, (chain) -> { + Buffer body = new Buffer(); + chain.request().body().writeTo(body); + try { + (body.readUtf8().contains("waiting") ? waitingReleased : inFlightReleased) + .await(WAIT_SECONDS, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IOException(e); + } + return AcceptingInterceptor.accepted(chain.request()); + }); + + for (int i = 0; i < EventProcessor.MAX_IN_FLIGHT_BATCHES; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + processor.flush(); + } + processor.trackEvent("purchase", "waiting", "1", null, null); + CompletableFuture flushed = processor.flush(); + + assertTrue(recorder.awaitCount(EventProcessor.MAX_IN_FLIGHT_BATCHES)); + inFlightReleased.countDown(); + assertTrue(recorder.awaitCount(EventProcessor.MAX_IN_FLIGHT_BATCHES + 1)); + assertFalse(flushed.isDone()); + + waitingReleased.countDown(); + flushed.get(WAIT_SECONDS, TimeUnit.SECONDS); + assertEquals(EventProcessor.MAX_IN_FLIGHT_BATCHES + 1, deliveredEvents()); + } + + @Test + @SneakyThrows + public void trackEvent_dropsTheOldestEventsOnceTheBufferIsFullBehindTheLimit() { + AcceptingInterceptor eventsApi = AcceptingInterceptor.blocked(); + EventProcessor processor = newProcessor(1000, 0, eventsApi); + FlagsmithLogger logger = mockLogger(processor); + int inFlight = EventProcessor.MAX_IN_FLIGHT_BATCHES * 1000; + + for (int i = 0; i < inFlight; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + } + for (int i = 0; i < EventProcessor.MAX_BUFFERED_EVENTS + 500; i++) { + processor.trackEvent("purchase", "overflow-" + i + "-", "1", null, null); + } + + assertEquals(EventProcessor.MAX_BUFFERED_EVENTS, processor.bufferedEvents().size()); + verify(logger).error(contains("Dropped the 1 oldest events")); + verify(logger, times(1)).error(startsWith("Dropped")); + + eventsApi.release(); + + assertTrue(awaitDelivered(inFlight + EventProcessor.MAX_BUFFERED_EVENTS)); + for (String body : recorder.bodies()) { + assertTrue(MapperFactory.getMapper().readTree(body).get("events").size() <= 1000); + } + String bodies = String.join("", recorder.bodies()); + assertFalse(bodies.contains("overflow-499-")); + assertTrue(bodies.contains("overflow-500-")); + } + + @Test + @SneakyThrows + public void trackExposureEvent_buffersAnExposureAgainOnceItsCopyWasDropped() { + AcceptingInterceptor eventsApi = AcceptingInterceptor.blocked(); + EventProcessor processor = newProcessor(1000, 0, eventsApi); + mockLogger(processor); + + for (int i = 0; i < EventProcessor.MAX_IN_FLIGHT_BATCHES * 1000; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + } + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + for (int i = 0; i < EventProcessor.MAX_BUFFERED_EVENTS; i++) { + processor.trackEvent("purchase", "overflow-" + i, "1", null, null); + } + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + + List> buffered = processor.bufferedEvents(); + assertEquals(EventProcessor.MAX_BUFFERED_EVENTS, buffered.size()); + assertEquals("$flag_exposure", buffered.get(buffered.size() - 1).get("event")); + eventsApi.release(); + } + + private static Stream healthyApiLoads() { + return Stream.of( + Arguments.of(Integer.MAX_VALUE, EventProcessor.MAX_BUFFERED_EVENTS + 1), + Arguments.of(1, 2000)); + } + + @ParameterizedTest + @MethodSource("healthyApiLoads") + @SneakyThrows + public void flush_neverDropsEventsOnAHealthyApi(int maxBufferItems, int events) { + EventProcessor processor = newProcessor(maxBufferItems, 0, AcceptingInterceptor.open()); + FlagsmithLogger logger = mockLogger(processor); + + for (int i = 0; i < events; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + } + processor.flush(); + + assertTrue(awaitDelivered(events)); + assertEquals(Collections.emptyList(), errorCalls(logger)); + } + + /** + * Every error-level call on a mocked logger. Reads invocations rather than verify(...), whose + * varargs matching silently misses calls with a different argument count. + */ + private static List errorCalls(FlagsmithLogger logger) { + List calls = new ArrayList<>(); + for (Invocation invocation : mockingDetails(logger).getInvocations()) { + String method = invocation.getMethod().getName(); + if (method.equals("error") || method.equals("httpError")) { + calls.add(invocation.toString()); + } + } + return calls; + } + + /** The number of events across every request the events API received. */ + @SneakyThrows + private int deliveredEvents() { + int delivered = 0; + for (String body : recorder.bodies()) { + delivered += MapperFactory.getMapper().readTree(body).get("events").size(); + } + return delivered; + } + + @SneakyThrows + private boolean awaitDelivered(int expected) { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(WAIT_SECONDS); + while (deliveredEvents() < expected && System.nanoTime() < deadline) { + Thread.sleep(10); + } + return deliveredEvents() == expected; + } + + @Test + @SneakyThrows + public void flush_logsOnlyEventsTheApiRejects() { + EventProcessor processor = newProcessor(1000, 0); + FlagsmithLogger logger = mockLogger(processor); + interceptor.addRule().post(EVENTS_ENDPOINT).times(1).respond(ACCEPTED_BODY, MEDIATYPE_JSON); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond( + "{\"accepted\": 1, \"rejected\": [{\"index\": 1, \"error\": \"event too long\"}]}", + MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + assertEquals(1, recorder.count()); + assertEquals(Collections.emptyList(), errorCalls(logger)); + + processor.trackEvent("purchase", "user-1", "1", null, null); + processor.trackEvent("purchase", "user-2", "2", null, null); + flushAndWait(processor); + + verify(logger).error(contains("rejected 1 of 2 events")); + verify(logger).error(contains("event too long")); + } + + /** + * Accepts every request, optionally holding each until released. Answers itself because + * MockInterceptor's canned bodies share one buffer and break under concurrent calls. + */ + private static class AcceptingInterceptor implements Interceptor { + + private final CountDownLatch released; + + private AcceptingInterceptor(boolean blocked) { + this.released = new CountDownLatch(blocked ? 1 : 0); + } + + static AcceptingInterceptor blocked() { + return new AcceptingInterceptor(true); + } + + static AcceptingInterceptor open() { + return new AcceptingInterceptor(false); + } + + void release() { + released.countDown(); + } + + @Override + public Response intercept(Chain chain) throws IOException { + try { + released.await(WAIT_SECONDS, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IOException(e); + } + return accepted(chain.request()); + } + + static Response accepted(Request request) { + return new Response.Builder() + .request(request) + .protocol(Protocol.HTTP_1_1) + .code(202) + .message("Accepted") + .body(ResponseBody.create(ACCEPTED_BODY, MediaType.get("application/json"))) + .build(); + } + } + + /** Records every request that reaches the network, with its body. */ + private static class RecordingInterceptor implements Interceptor { + + private final List bodies = Collections.synchronizedList(new ArrayList<>()); + private final CountDownLatch firstRequest = new CountDownLatch(1); + + @Override + public Response intercept(Chain chain) throws IOException { + Request request = chain.request(); + Buffer buffer = new Buffer(); + if (request.body() != null) { + request.body().writeTo(buffer); + } + bodies.add(buffer.readUtf8()); + firstRequest.countDown(); + return chain.proceed(request); + } + + int count() { + return bodies.size(); + } + + List bodies() { + return new ArrayList<>(bodies); + } + + boolean awaitFirstRequest(long seconds) throws InterruptedException { + return firstRequest.await(seconds, TimeUnit.SECONDS); + } + + boolean awaitCount(int expected) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(WAIT_SECONDS); + while (count() < expected && System.nanoTime() < deadline) { + Thread.sleep(10); + } + return count() >= expected; + } + } + + /** Simulates a connection failure for the first n attempts. */ + private static class FailingInterceptor implements Interceptor { + + private final AtomicInteger remainingFailures; + + FailingInterceptor(int failures) { + this.remainingFailures = new AtomicInteger(failures); + } + + @Override + public Response intercept(Chain chain) throws IOException { + if (remainingFailures.getAndDecrement() > 0) { + throw new IOException("connection refused"); + } + return chain.proceed(chain.request()); + } + } +}