diff --git a/src/main/java/com/flagsmith/FlagsmithClient.java b/src/main/java/com/flagsmith/FlagsmithClient.java index dd2917de..644f9f46 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,186 @@ 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) { + // The API wrapper returns null, not throws, on a timed-out or interrupted request. + 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"); + + // A missing feature name is a bug in the caller, and the events API rejects the exposure. A + // missing identifier is ordinary at runtime (an anonymous visitor), so it is logged instead. + 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 in-flight batch is done, 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 +396,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) { @@ -501,6 +706,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) { @@ -557,6 +765,23 @@ public FlagsmithClient build() { configuration.getOfflineHandler().getEnvironment()); } + // Last, once nothing else can throw: starting the processor starts its flush timer, which a + // failed build would otherwise leave running with no client to close it. + if (configuration.getEnableEvents()) { + EventProcessor processor = configuration.getEventProcessor() != null + ? configuration.getEventProcessor() + : new EventProcessor( + configuration.getHttpClient(), + configuration.getEventsUri(), + configuration.getEventsMaxBufferItems(), + configuration.getEventsFlushIntervalMillis()); + processor.claim(); + 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..c65b0ddd 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,27 @@ 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; + /** The processor from withEventProcessor, or null; each client otherwise builds its own. */ + private EventProcessor eventProcessor; private FlagsmithFlagDefaults flagsmithFlagDefaults = null; private Boolean raiseUpdateEnvironmentErrorsOnStartup = true; private Boolean offlineMode = false; @@ -88,6 +99,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 +149,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; @@ -272,6 +308,71 @@ public Builder withEnableAnalytics(Boolean enable) { return this; } + /** + * Override the events API base URL. Allowed with events disabled, so a shared configuration + * can carry it. + * + * @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..4a411fa8 100644 --- a/src/main/java/com/flagsmith/config/Retry.java +++ b/src/main/java/com/flagsmith/config/Retry.java @@ -25,17 +25,45 @@ public class Retry { add(429); add(503); }}; + /** + * When true, only force-listed statuses and connection failures (null status) are retried, and + * never past {@link #total} attempts. False keeps the historical policy: any status retries + * within the budget, and a force-listed one regardless of it. + */ + 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..12d17303 --- /dev/null +++ b/src/main/java/com/flagsmith/models/ExperimentMetadata.java @@ -0,0 +1,19 @@ +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; + /** + * Whether this identity is enrolled in the experiment. The variant alone cannot tell: an + * identity outside the rollout is still bucketed into a variant. + */ + @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..625f4c9d 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,18 @@ public class Flag extends BaseFlag { private Integer featureId = 0; private Boolean isDefault; + /** + * Variant key the identity was bucketed into. Set by remote evaluation only; null otherwise. + */ + private String variant; + /** + * Evaluation reason, e.g. "DEFAULT", "SPLIT; weight=70.0", "TARGETING_MATCH; segment=...". + */ + private String reason; + /** + * Experiment running on this feature. Set by remote identity evaluation only; null otherwise. + */ + private ExperimentMetadata experiment; /** * return flag from feature state model and identity id. @@ -21,6 +35,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 +56,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..007a4e3c --- /dev/null +++ b/src/main/java/com/flagsmith/models/features/FeatureStateMetadata.java @@ -0,0 +1,13 @@ +package com.flagsmith.models.features; + +import com.flagsmith.models.ExperimentMetadata; +import lombok.Data; + +/** + * The {@code metadata} object returned alongside a remotely evaluated feature state. Keys other + * than {@code experiment} are ignored. + */ +@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..86bdc680 --- /dev/null +++ b/src/main/java/com/flagsmith/threads/EventProcessor.java @@ -0,0 +1,522 @@ +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.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 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 ships them to the Flagsmith events API. + * + *

Events are flushed on a timer, when the buffer fills, and on {@link #close()}. Exposures are + * deduplicated within a flush window. Exceptions are logged, never thrown to the caller; + * {@link Error}s are not caught. + */ +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 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 String KEY_SEPARATOR = "\u0000"; + private static final MediaType JSON_MEDIA_TYPE = + MediaType.get("application/json; charset=utf-8"); + /** + * Cap on events awaiting the events API. Past it a flush drops its batch rather than queue it, + * so an outage cannot grow memory with the host's traffic. It counts events, not batches, so a + * small buffer is not throttled on a healthy API. The true bound is this plus one buffer. + */ + static final int MAX_IN_FLIGHT_EVENTS = 10_000; + /** The close timeout when the HTTP client's timeouts do not bound a request. */ + static final long UNBOUNDED = -1L; + /** The least time between two error lines reporting dropped events. */ + private static final long DROP_LOG_INTERVAL_NANOS = TimeUnit.SECONDS.toNanos(10); + + /** The URL batches are POSTed to. */ + @Getter + private final HttpUrl eventsEndpoint; + /** The number of buffered events that triggers an immediate flush. */ + @Getter + private final int maxBufferItems; + /** The interval between timed flushes; 0 means there is no timer. */ + @Getter + private final int flushIntervalMillis; + // No public getters below: each one would become supported API. + private final List> buffer = new ArrayList<>(); + private final Set dedupeKeys = new HashSet<>(); + private final Object lock = new Object(); + @Getter(AccessLevel.PACKAGE) + private final ScheduledExecutorService scheduler; + private final Set> inFlight = ConcurrentHashMap.newKeySet(); + private int inFlightEvents = 0; // guarded by lock + private int droppedSinceLastReport = 0; // guarded by lock + private Long lastDropReportNanos = null; // guarded by lock + @Getter(AccessLevel.PACKAGE) + private final RequestProcessor requestProcessor; + /** + * How long {@link #close()} waits for in-flight batches: the worst case of one batch under the + * HTTP client's timeouts and the retry policy, or {@link #UNBOUNDED} when a timeout is off. + */ + @Getter(AccessLevel.PACKAGE) + private final long closeTimeoutMillis; + /** The API wrapper used to build requests; injected by {@code FlagsmithClient.Builder}. */ + @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(client, new FlagsmithLogger(), buildRetry())); + } + + /** For tests: sends through the given request processor. */ + EventProcessor(HttpUrl eventsUri, int maxBufferItems, int flushIntervalMillis, + RequestProcessor requestProcessor) { + // The buffer limit is the only bound on the buffer between timed flushes, or with no timer. + 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; + }); + } + + /** + * The retry policy for an event batch: at most one retry, on a connection failure or a 500, + * 502, 503 or 504, and never on any other status. + */ + private static Retry buildRetry() { + Retry retry = new Retry(2); + retry.setStatusForcelist(new HashSet<>(Arrays.asList(500, 502, 503, 504))); + retry.setStatusForcelistOnly(Boolean.TRUE); + return retry; + } + + /** + * The longest one batch can take: every attempt the retry policy allows, each at the client's + * timeouts, plus backoff. An attempt is bounded by the call timeout if set, else by connect, + * write and read in turn. + * + * @return the worst case in milliseconds, or {@link #UNBOUNDED} + */ + static long worstCaseBatchMillis(OkHttpClient client, Retry retry) { + long attemptMillis; + if (client.callTimeoutMillis() > 0) { + attemptMillis = client.callTimeoutMillis(); + } else if (client.connectTimeoutMillis() > 0 && client.writeTimeoutMillis() > 0 + && client.readTimeoutMillis() > 0) { + attemptMillis = (long) client.connectTimeoutMillis() + client.writeTimeoutMillis() + + client.readTimeoutMillis(); + } else { + return UNBOUNDED; + } + + // Walk a copy of the policy the way RequestProcessor does: back off, then attempt. + Retry walk = retry.toBuilder().build(); + long total = 0; + for (int attempt = 0; attempt < walk.getTotal(); attempt++) { + total += walk.calculateSleepTime() + attemptMillis; + walk.retryAttempted(); + } + return total; + } + + /** + * 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 in-flight batch has been delivered or dropped + */ + public CompletableFuture flush() { + List> batch = null; + CompletableFuture tracked = null; + int droppedToReport = 0; + + // The batch is registered as in-flight under the same lock that empties the buffer, so a + // concurrent flush() can never observe both an empty buffer and an unregistered batch. + synchronized (lock) { + if (!buffer.isEmpty()) { + if (inFlightEvents >= MAX_IN_FLIGHT_EVENTS) { + droppedToReport = recordDrop(buffer.size()); + } else { + batch = new ArrayList<>(buffer); + tracked = new CompletableFuture<>(); + inFlight.add(tracked); + inFlightEvents += batch.size(); + } + buffer.clear(); + } + dedupeKeys.clear(); + } + + if (droppedToReport > 0) { + logger.error("Dropped " + droppedToReport + " events: at least " + MAX_IN_FLIGHT_EVENTS + + " earlier events are still waiting on the events API. Further drops are reported at" + + " most every " + TimeUnit.NANOSECONDS.toSeconds(DROP_LOG_INTERVAL_NANOS) + "s."); + } + + if (batch != null) { + send(batch, tracked); + } + + return awaitInFlight(); + } + + /** + * 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()) { + // Scheduling on the shut-down scheduler would throw. Reached when an injected processor + // was closed before its client was built. + 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 the worst case of one batch; unbounded if a client timeout is off. + */ + public void close() { + synchronized (this) { + closed.set(true); + scheduler.shutdownNow(); + } + + try { + CompletableFuture remaining = flush(); + if (closeTimeoutMillis == UNBOUNDED) { + remaining.get(); + } else { + remaining.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; batches still in flight carry on in the background."); + } catch (Exception e) { + logger.error("Failed to flush events on close.", e); + } + + // A graceful shutdown: interrupting a POST mid-flight would lose its batch. + 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); + + // Serialised now, not at flush: an unserialisable value drops only this event, not its + // batch, and the JSON text is a deep copy that later caller mutation cannot reach. + 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; + + synchronized (lock) { + // Re-checked under the lock: close() sets the flag before its final flush takes the + // lock, so an event either makes that flush or is refused here, never stranded. + if (closed.get()) { + logClosed(event); + return; + } + if (dedupe && !dedupeKeys.add( + dedupeKey(event, featureName, identifier, stringValue, experimentId))) { + return; + } + buffer.add(eventPayload); + isFull = buffer.size() >= maxBufferItems; + } + + if (isFull) { + flush(); + } + } catch (JsonProcessingException | RuntimeException e) { + logger.error("Failed to buffer event " + event + ".", e); + } + } + + /** + * JSON text written verbatim into the batch. Uses writeValueAsString, not valueToTree: on a + * self-containing map valueToTree throws a raw StackOverflowError instead of an exception. + */ + 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); + } + + /** + * Unwrap {@link TraitConfig} values and drop transient traits: transient means "do not + * persist", and an event store keeps what it is sent. + */ + 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; + } + + private static String dedupeKey(String event, String featureName, String identifier, + String value, Object experimentId) { + return String.join(KEY_SEPARATOR, + nullToEmpty(event), + nullToEmpty(featureName), + nullToEmpty(identifier), + nullToEmpty(value), + experimentId == null ? "" : String.valueOf(experimentId)); + } + + private static String nullToEmpty(String value) { + return value == null ? "" : value; + } + + /** + * Hand a batch to the request processor. {@code tracked} is settled on every path out: left + * pending, it would wedge every later {@link #flush()}. + */ + private void send(List> batch, CompletableFuture tracked) { + // The callback lives as long as the POST, so it holds the size and lets the list be freed. + 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)); + + Request request = api + .newPostRequest(eventsEndpoint, RequestBody.create(payload, JSON_MEDIA_TYPE)) + .newBuilder() + .header(SDK_USER_AGENT_HEADER, SDK_USER_AGENT_PREFIX + Versions.getVersion()) + .build(); + + requestProcessor + .submit(request, new TypeReference() {}, Boolean.FALSE, buildRetry()) + .whenComplete((response, error) -> { + try { + logRejections(response, batchSize); + } finally { + settle(tracked, batchSize); + } + }); + submitted = true; + } catch (Exception e) { + logger.error("Dropping " + batchSize + " events: failed to send them.", e); + } finally { + if (!submitted) { + settle(tracked, batchSize); + } + } + } + + /** + * The events API accepts a batch with a 202 even when it rejects some of its events, listing + * those under {@code rejected}. Without this they would vanish without a trace. + */ + 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, int batchSize) { + // Idempotent: only the call that actually removes the batch gives its events back. + if (inFlight.remove(tracked)) { + synchronized (lock) { + inFlightEvents -= batchSize; + } + } + tracked.complete(null); + } + + /** + * Count dropped events. The first drop is reported at once, later ones at most once per + * {@link #DROP_LOG_INTERVAL_NANOS}, so an outage does not log per flush. Called under the lock. + * + * @return the number of drops to report now, or 0 to stay quiet + */ + 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 CompletableFuture awaitInFlight() { + return CompletableFuture.allOf(inFlight.toArray(new CompletableFuture[0])); + } + + /** A snapshot of the buffer for tests, taken under the lock. */ + 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..04188cf0 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,24 @@ 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.regex.Pattern; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -998,4 +1004,572 @@ public void testFlagsmithUsesOfflineHandlerIfSetAndNoAPIResponse() throws Flagsm assertTrue(environmentFlags.isFeatureEnabled("some_feature")); assertTrue(identityFlags.isFeatureEnabled("some_feature")); } + + private static FlagsmithConfig.Builder eventsConfigBuilder( + String baseUrl, MockInterceptor interceptor) { + return FlagsmithConfig.newBuilder() + .baseUri(baseUrl) + .addHttpInterceptor(interceptor) + .eventsUri("http://events-uri") + .withEnableEvents(Boolean.TRUE); + } + + private static void respondWithExperimentFlags(String baseUrl, MockInterceptor interceptor) { + interceptor.addRule() + .post(baseUrl + "/identities/") + .anyTimes() + .respond(FlagsmithTestHelper.getIdentitiesFlagsWithExperiment(), MEDIATYPE_JSON); + } + + @Test + public void testEventsConfigWithoutEnableEventsThrows() { + assertThrows(IllegalArgumentException.class, + () -> FlagsmithConfig.newBuilder().withEventsMaxBufferItems(10).build()); + assertThrows(IllegalArgumentException.class, + () -> FlagsmithConfig.newBuilder().withEventsFlushIntervalMillis(10).build()); + } + + @Test + public void testEventsUriWithoutEnableEventsIsHarmless() { + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .eventsUri("http://events-uri") + .build(); + + assertFalse(config.getEnableEvents()); + assertEquals("http://events-uri/", config.getEventsUri().toString()); + } + + @Test + public void testGetExperimentFlagPassesTraitsThroughToTheExposure() + throws FlagsmithClientError { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(eventsConfigBuilder(baseUrl, interceptor) + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + respondWithExperimentFlags(baseUrl, interceptor); + + Map traits = new HashMap<>(); + traits.put("plan", "premium"); + traits.put("session_id", new TraitConfig("abc123", true)); + + client.getExperimentFlag("checkout_cta", "user-1", traits); + verify(processor, times(1)).trackExposureEvent( + eq("checkout_cta"), eq("user-1"), eq("treatment"), eq(traits), any()); + + // The two-argument overload passes an empty map rather than null. + client.getExperimentFlag("checkout_cta", "user-2"); + verify(processor, times(1)).trackExposureEvent( + eq("checkout_cta"), eq("user-2"), eq("treatment"), eq(new HashMap<>()), any()); + } + + @Test + public void testEventsInOfflineModeThrowsAtBuild() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .withOfflineMode(Boolean.TRUE) + .withOfflineHandler(new DummyOfflineHandler()) + .withEventProcessor(processor) + .build(); + + FlagsmithClient.Builder clientBuilder = FlagsmithClient.newBuilder() + .withConfiguration(config) + .setApiKey("api-key"); + + FlagsmithRuntimeError ex = assertThrows(FlagsmithRuntimeError.class, clientBuilder::build); + assertEquals("Events cannot be enabled in offline mode.", ex.getMessage()); + } + + /** + * A client whose API wrapper returns null identity flags, as FlagsmithApiWrapper does when + * the identities request times out or is interrupted. + */ + private static FlagsmithClient clientWithUnavailableIdentityFlags( + FlagsmithConfig config) { + FlagsmithApiWrapper mockApiWrapper = mock(FlagsmithApiWrapper.class); + when(mockApiWrapper.getConfig()).thenReturn(config); + when(mockApiWrapper.identifyUserWithTraits(any(), any(), anyBoolean(), anyBoolean())) + .thenReturn(null); + + return FlagsmithClient.newBuilder() + .withFlagsmithApiWrapper(mockApiWrapper) + .withConfiguration(config) + .setApiKey("api-key") + .build(); + } + + @Test + public void testGetExperimentFlagServesTheDefaultWhenIdentityFlagsAreUnavailable() + throws FlagsmithClientError { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .withEventProcessor(processor) + .build(); + FlagsmithFlagDefaults defaults = new FlagsmithFlagDefaults(); + defaults.setDefaultFlagValueFunc(FlagsmithClientTest::defaultHandler); + config.setFlagsmithFlagDefaults(defaults); + FlagsmithClient client = clientWithUnavailableIdentityFlags(config); + + BaseFlag flag = client.getExperimentFlag("checkout_cta", "user-1"); + + assertTrue(flag instanceof DefaultFlag); + assertEquals(DEFAULT_FLAG_VALUE, flag.getValue()); + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testGetExperimentFlagThrowsWhenIdentityFlagsAreUnavailableWithoutADefault() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .withEventProcessor(processor) + .build(); + FlagsmithClient client = clientWithUnavailableIdentityFlags(config); + + assertThrows(FlagsmithApiError.class, + () -> client.getExperimentFlag("checkout_cta", "user-1")); + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testFailedBuildDoesNotStartTheEventProcessor() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .withLocalEvaluation(true) + .withEventProcessor(processor) + .build(); + + FlagsmithClient.Builder clientBuilder = FlagsmithClient.newBuilder() + .withConfiguration(config) + // 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()); + } + + @Test + public void testTrackEventRejectsReservedEventNames() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + + assertThrows(IllegalArgumentException.class, () -> client.trackEvent("$flag_exposure")); + verify(processor, never()).trackEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testTrackEventRejectsBlankEventNames() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + + assertThrows(IllegalArgumentException.class, () -> client.trackEvent(null)); + assertThrows(IllegalArgumentException.class, () -> client.trackEvent(" ", "user-1")); + verify(processor, never()).trackEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testTrackExposureEventRejectsBlankFeatureNames() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + + assertThrows(IllegalArgumentException.class, + () -> client.trackExposureEvent(null, "user-1", "treatment")); + assertThrows(IllegalArgumentException.class, + () -> client.trackExposureEvent("", "user-1", "treatment")); + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testInvalidEventProcessorSettingsThrowAtBuild() { + assertThrows(IllegalArgumentException.class, () -> FlagsmithConfig.newBuilder() + .withEnableEvents(Boolean.TRUE) + .withEventsMaxBufferItems(0) + .build()); + assertThrows(IllegalArgumentException.class, () -> FlagsmithConfig.newBuilder() + .withEnableEvents(Boolean.TRUE) + .withEventsFlushIntervalMillis(-1) + .build()); + } + + @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 testNullEnableEventsLeavesEventsDisabled() { + FlagsmithConfig config = FlagsmithConfig.newBuilder().withEnableEvents(null).build(); + + assertFalse(config.getEnableEvents()); + } + + @Test + public void testTrackExposureEventWithBlankIdentifierSendsNothing() { + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + + client.trackExposureEvent("checkout_cta", " ", "treatment"); + client.trackExposureEvent("checkout_cta", null, "treatment"); + + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testGetExperimentFlagRecordsAnExposureWhenEnrolled() throws FlagsmithClientError { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(eventsConfigBuilder(baseUrl, interceptor) + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + respondWithExperimentFlags(baseUrl, interceptor); + + BaseFlag flag = client.getExperimentFlag("checkout_cta", "user-1"); + + assertTrue(flag instanceof Flag); + assertEquals("treatment", ((Flag) flag).getVariant()); + assertEquals("SPLIT; weight=70.0", ((Flag) flag).getReason()); + assertEquals(42, ((Flag) flag).getExperiment().getId()); + assertEquals(Boolean.TRUE, ((Flag) flag).getExperiment().getInExperiment()); + + Map expectedMetadata = new HashMap<>(); + expectedMetadata.put("experiment_id", 42); + verify(processor, times(1)).trackExposureEvent( + eq("checkout_cta"), eq("user-1"), eq("treatment"), any(), eq(expectedMetadata)); + } + + @Test + public void testGetExperimentFlagRecordsNoExposureWhenNotEnrolledOrDisabled() + throws FlagsmithClientError { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(eventsConfigBuilder(baseUrl, interceptor) + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + respondWithExperimentFlags(baseUrl, interceptor); + + // 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()); + } + + @Test + public void testGetExperimentFlagUsesTheDefaultHandlerForAMissingFeature() + throws FlagsmithClientError { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(eventsConfigBuilder(baseUrl, interceptor) + .withEventProcessor(processor) + .build()) + .setDefaultFlagValueFunction(FlagsmithClientTest::defaultHandler) + .setApiKey("api-key") + .build(); + respondWithExperimentFlags(baseUrl, interceptor); + + BaseFlag flag = client.getExperimentFlag("no_such_feature", "user-1"); + + assertTrue(flag instanceof DefaultFlag); + assertEquals(DEFAULT_FLAG_VALUE, flag.getValue()); + verify(processor, never()).trackExposureEvent(any(), any(), any(), any(), any()); + } + + @Test + public void testGetExperimentFlagThrowsForAMissingFeatureWithoutADefaultHandler() { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(eventsConfigBuilder(baseUrl, interceptor) + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + respondWithExperimentFlags(baseUrl, interceptor); + + assertThrows(FeatureNotFoundError.class, + () -> client.getExperimentFlag("no_such_feature", "user-1")); + } + + @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(eventsConfigBuilder(baseUrl, 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()); + } + + @Test + public void testGetExperimentFlagRecordsOneExposurePerIdentity() throws FlagsmithClientError { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + EventProcessor processor = mock(EventProcessor.class); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(eventsConfigBuilder(baseUrl, interceptor) + .withEventProcessor(processor) + .build()) + .setApiKey("api-key") + .build(); + respondWithExperimentFlags(baseUrl, interceptor); + + client.getExperimentFlag("checkout_cta", "user-1"); + client.getExperimentFlag("checkout_cta", "user-2"); + + verify(processor, times(1)).trackExposureEvent( + eq("checkout_cta"), eq("user-1"), eq("treatment"), any(), any()); + verify(processor, times(1)).trackExposureEvent( + eq("checkout_cta"), eq("user-2"), eq("treatment"), any(), any()); + } + + @Test + public void testCloseFlushesBufferedEvents() throws FlagsmithClientError, IOException { + String baseUrl = "http://bad-url"; + MockInterceptor interceptor = new MockInterceptor(); + List eventBodies = new ArrayList<>(); + // Added before the MockInterceptor, which short-circuits the chain once it matches. + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .baseUri(baseUrl) + .addHttpInterceptor((chain) -> { + Request request = chain.request(); + if (request.url().toString().endsWith("/v1/events")) { + Buffer buffer = new Buffer(); + request.body().writeTo(buffer); + eventBodies.add(buffer.readUtf8()); + } + return chain.proceed(request); + }) + .addHttpInterceptor(interceptor) + .eventsUri("http://events-uri") + .withEnableEvents(Boolean.TRUE) + .withEventsFlushIntervalMillis(0) + .build()) + .setApiKey("api-key") + .build(); + respondWithExperimentFlags(baseUrl, interceptor); + 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); + + client.getExperimentFlag("checkout_cta", "user-1"); + assertTrue(eventBodies.isEmpty()); + + client.close(); + + assertEquals(1, eventBodies.size()); + JsonNode events = MapperFactory.getMapper().readTree(eventBodies.get(0)).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()); + } + + /** An events-enabled config that records each events batch as "environment key|body". */ + private static FlagsmithConfig recordingEventsConfig(List batches) { + return recordingEventsConfigBuilder(batches).build(); + } + + private static FlagsmithConfig.Builder recordingEventsConfigBuilder(List batches) { + MockInterceptor interceptor = new MockInterceptor(); + interceptor.addRule() + .post("http://events-uri/v1/events") + .anyTimes() + .respond("{\"accepted\": 1, \"rejected\": []}", MEDIATYPE_JSON); + 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 testClientsSharingAConfigSendEventsUnderTheirOwnKeys() throws Exception { + List batches = Collections.synchronizedList(new ArrayList<>()); + FlagsmithConfig config = recordingEventsConfig(batches); + 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); + 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)) + .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 = recordingEventsConfigBuilder(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(); + } } diff --git a/src/test/java/com/flagsmith/FlagsmithRetryTest.java b/src/test/java/com/flagsmith/FlagsmithRetryTest.java index cd021087..89e0306c 100644 --- a/src/test/java/com/flagsmith/FlagsmithRetryTest.java +++ b/src/test/java/com/flagsmith/FlagsmithRetryTest.java @@ -6,6 +6,8 @@ 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 +68,60 @@ public void FlagsmithRetry_validateSleep() { assertTrue(attempts.equals(3)); } + private static Retry oneRetryOnServerErrors() { + Retry retry = new Retry(2); + retry.setStatusForcelist(new HashSet<>(Arrays.asList(500, 502, 503, 504))); + retry.setStatusForcelistOnly(Boolean.TRUE); + return retry; + } + + @Test + public void FlagsmithRetry_statusForcelistOnly_retriesForcedStatusWithinBudget() { + Retry retry = oneRetryOnServerErrors(); + + retry.retryAttempted(); + assertTrue(retry.isRetry(503), "a force-listed status should retry while attempts remain"); + } + + @Test + public void FlagsmithRetry_statusForcelistOnly_stopsAtTheAttemptsBudget() { + Retry retry = oneRetryOnServerErrors(); + + retry.retryAttempted(); + retry.retryAttempted(); + // Without this bound a permanently failing endpoint retries forever. + assertFalse(retry.isRetry(503), "a force-listed status must not retry past the budget"); + assertFalse(retry.isRetry(null), "a connection failure must not retry past the budget"); + } + + @Test + public void FlagsmithRetry_statusForcelistOnly_doesNotRetryUnlistedStatus() { + Retry retry = oneRetryOnServerErrors(); + + retry.retryAttempted(); + assertFalse(retry.isRetry(400), "a 4xx must not be retried"); + assertFalse(retry.isRetry(404), "a 4xx must not be retried"); + } + + @Test + public void FlagsmithRetry_statusForcelistOnly_retriesConnectionFailuresWithinBudget() { + Retry retry = oneRetryOnServerErrors(); + + retry.retryAttempted(); + assertTrue(retry.isRetry(null), "a connection failure should retry while attempts remain"); + } + + @Test + public void FlagsmithRetry_defaultPolicyIsUnchanged() { + Retry retry = new Retry(1); + + assertFalse(retry.getStatusForcelistOnly()); + retry.retryAttempted(); + // Historical behaviour: a force-listed status retries regardless of the budget. + 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..03aeb887 100644 --- a/src/test/java/com/flagsmith/FlagsmithTestHelper.java +++ b/src/test/java/com/flagsmith/FlagsmithTestHelper.java @@ -431,6 +431,90 @@ public static String getIdentitiesFlags() { return featureJson; } + /** + * An identity flags payload exercising every branch of the experiment gate: an enrolled flag, + * a bucketed but unenrolled flag, a flag with no metadata at all, and a disabled flag whose + * experiment is running. + */ + public static String getIdentitiesFlagsWithExperiment() { + return "{\n" + + " \"traits\": [],\n" + + " \"flags\": [\n" + + " {\n" + + " \"id\": 1,\n" + + " \"feature\": {\n" + + " \"id\": 1,\n" + + " \"name\": \"checkout_cta\",\n" + + " \"type\": \"MULTIVARIATE\",\n" + + " \"project\": 1\n" + + " },\n" + + " \"feature_state_value\": \"buy-now\",\n" + + " \"enabled\": true,\n" + + " \"variant\": \"treatment\",\n" + + " \"reason\": \"SPLIT; weight=70.0\",\n" + + " \"metadata\": {\n" + + " \"experiment\": {\n" + + " \"id\": 42,\n" + + " \"name\": \"checkout_experiment\",\n" + + " \"in_experiment\": true\n" + + " }\n" + + " }\n" + + " },\n" + + " {\n" + + " \"id\": 2,\n" + + " \"feature\": {\n" + + " \"id\": 2,\n" + + " \"name\": \"pricing_page\",\n" + + " \"type\": \"MULTIVARIATE\",\n" + + " \"project\": 1\n" + + " },\n" + + " \"feature_state_value\": \"old-pricing\",\n" + + " \"enabled\": true,\n" + + " \"variant\": \"control\",\n" + + " \"reason\": \"SPLIT; weight=30.0\",\n" + + " \"metadata\": {\n" + + " \"experiment\": {\n" + + " \"id\": 43,\n" + + " \"name\": \"pricing_experiment\",\n" + + " \"in_experiment\": false\n" + + " }\n" + + " }\n" + + " },\n" + + " {\n" + + " \"id\": 3,\n" + + " \"feature\": {\n" + + " \"id\": 3,\n" + + " \"name\": \"some_feature\",\n" + + " \"type\": \"STANDARD\",\n" + + " \"project\": 1\n" + + " },\n" + + " \"feature_state_value\": \"some-value\",\n" + + " \"enabled\": true\n" + + " },\n" + + " {\n" + + " \"id\": 4,\n" + + " \"feature\": {\n" + + " \"id\": 4,\n" + + " \"name\": \"disabled_feature\",\n" + + " \"type\": \"MULTIVARIATE\",\n" + + " \"project\": 1\n" + + " },\n" + + " \"feature_state_value\": \"off\",\n" + + " \"enabled\": false,\n" + + " \"variant\": \"treatment\",\n" + + " \"reason\": \"SPLIT; weight=50.0\",\n" + + " \"metadata\": {\n" + + " \"experiment\": {\n" + + " \"id\": 44,\n" + + " \"name\": \"disabled_experiment\",\n" + + " \"in_experiment\": true\n" + + " }\n" + + " }\n" + + " }\n" + + " ]\n" + + "}"; + } + public static Future futurableReturn(T response) { CompletableFuture promise = new CompletableFuture<>(); promise.complete(response); diff --git a/src/test/java/com/flagsmith/flagengine/models/FlagTest.java b/src/test/java/com/flagsmith/flagengine/models/FlagTest.java index c5aebcee..757a6157 100644 --- a/src/test/java/com/flagsmith/flagengine/models/FlagTest.java +++ b/src/test/java/com/flagsmith/flagengine/models/FlagTest.java @@ -1,9 +1,15 @@ package com.flagsmith.flagengine.models; +import com.flagsmith.MapperFactory; +import com.flagsmith.models.ExperimentMetadata; import com.flagsmith.models.Flag; -import org.junit.Test; +import com.flagsmith.models.features.FeatureStateModel; +import com.fasterxml.jackson.core.JsonProcessingException; +import org.junit.jupiter.api.Test; -import static org.junit.Assert.assertEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; public class FlagTest { @Test @@ -15,13 +21,56 @@ public void testToString() { flag.setValue("foo"); flag.setFeatureId(1); - String expected = String.join("Flag(super=", - "BaseFlag(enabled=true, ", - "value=foo, ", - "featureName=my_feature), ", - "featureId=1, ", - "isDefault=false)"); + // BaseModel has no toString(), so the "super=" part carries an identity hash. Only the + // fields Lombok renders are asserted. + String expected = ", enabled=true, " + + "value=foo, " + + "featureName=my_feature), " + + "featureId=1, " + + "isDefault=false, " + + "variant=null, " + + "reason=null, " + + "experiment=null)"; - assertEquals(expected, flag.toString()); + assertTrue(flag.toString().startsWith("Flag(super=BaseFlag(super="), flag.toString()); + assertTrue(flag.toString().endsWith(expected), flag.toString()); + } + + @Test + public void fromFeatureStateModel_copiesExperimentMetadata() throws JsonProcessingException { + FeatureStateModel featureState = MapperFactory.getMapper().readValue( + "{\"feature\": {\"id\": 1, \"name\": \"checkout_cta\"}," + + " \"enabled\": true," + + " \"feature_state_value\": \"buy-now\"," + + " \"variant\": \"treatment\"," + + " \"reason\": \"SPLIT; weight=70.0\"," + + " \"metadata\": {\"experiment\": {\"id\": 42, \"name\": \"exp\"," + + " \"in_experiment\": true}}}", + FeatureStateModel.class); + + Flag flag = Flag.fromFeatureStateModel(featureState); + + assertEquals("treatment", flag.getVariant()); + assertEquals("SPLIT; weight=70.0", flag.getReason()); + + ExperimentMetadata experiment = flag.getExperiment(); + assertEquals(42, experiment.getId()); + assertEquals("exp", experiment.getName()); + assertEquals(Boolean.TRUE, experiment.getInExperiment()); + } + + @Test + public void fromFeatureStateModel_leavesExperimentMetadataNullWhenAbsent() + throws JsonProcessingException { + FeatureStateModel featureState = MapperFactory.getMapper().readValue( + "{\"feature\": {\"id\": 1, \"name\": \"some_feature\"}," + + " \"enabled\": true, \"feature_state_value\": \"some-value\"}", + FeatureStateModel.class); + + Flag flag = Flag.fromFeatureStateModel(featureState); + + assertNull(flag.getVariant()); + assertNull(flag.getReason()); + assertNull(flag.getExperiment()); } } 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..1e4e120f --- /dev/null +++ b/src/test/java/com/flagsmith/models/FeatureStateModelTest.java @@ -0,0 +1,71 @@ +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; + +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()); + } + + @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()); + } +} 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..3f98efbb --- /dev/null +++ b/src/test/java/com/flagsmith/threads/EventProcessorTest.java @@ -0,0 +1,910 @@ +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.config.Retry; +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.regex.Pattern; +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; + +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\": []}"; + /** Every wait in this test is bounded so a broken retry loop fails instead of hanging. */ + 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); + } + + /** 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); + } + + @Test + public void trackEvent_buffersNullValueAsNull() { + EventProcessor processor = newProcessor(1000, 0); + + processor.trackEvent("purchase", "user-123", null, null, null); + + Map event = processor.bufferedEvents().get(0); + assertNull(event.get("value")); + assertNull(event.get("traits")); + } + + @Test + public void trackExposureEvent_dedupesIdenticalExposuresWithinTheFlushWindow() { + EventProcessor processor = newProcessor(1000, 0); + 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()); + } + + @Test + public void trackExposureEvent_doesNotDedupeWhenAnyKeyPartDiffers() { + EventProcessor processor = newProcessor(1000, 0); + Map metadata = Collections.singletonMap("experiment_id", 42); + + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, metadata); + 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)); + + assertEquals(5, processor.bufferedEvents().size()); + } + + @Test + public void trackExposureEvent_buffersAgainAfterAFlush() { + 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); + flushAndWait(processor); + + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, metadata); + + assertEquals(1, processor.bufferedEvents().size()); + assertEquals(1, recorder.count()); + } + + @Test + public void trackEvent_neverDedupesCustomEvents() { + EventProcessor processor = newProcessor(1000, 0); + + processor.trackEvent("purchase", "user-1", "49.00", null, null); + processor.trackEvent("purchase", "user-1", "49.00", null, null); + + assertEquals(2, processor.bufferedEvents().size()); + } + + @Test + @SneakyThrows + public void flush_postsTheBatchToTheEventsEndpoint() { + EventProcessor processor = newProcessor(1000, 0); + 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 flush_doesNotPostWhenTheBufferIsEmpty() { + EventProcessor processor = newProcessor(1000, 0); + + flushAndWait(processor); + + assertEquals(0, recorder.count()); + } + + @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"); + } + + @Test + @SneakyThrows + public void flush_retriesOnceOnServerErrorThenDelivers() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).times(1).respond(503); + interceptor.addRule().post(EVENTS_ENDPOINT).times(1).respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(2, recorder.count()); + assertEquals(recorder.bodies().get(0), recorder.bodies().get(1)); + } + + @Test + @SneakyThrows + public void flush_dropsTheBatchAfterTwoServerErrors() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(500); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(2, recorder.count()); + assertTrue(processor.bufferedEvents().isEmpty()); + } + + @Test + @SneakyThrows + public void flush_doesNotRetryOnClientError() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(400); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(1, recorder.count()); + assertTrue(processor.bufferedEvents().isEmpty()); + } + + @Test + @SneakyThrows + public void flush_retriesOnceOnAConnectionFailureThenDelivers() { + EventProcessor processor = newProcessor(1000, 0, new FailingInterceptor(1)); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(2, recorder.count()); + assertEquals(recorder.bodies().get(0), recorder.bodies().get(1)); + } + + @Test + @SneakyThrows + public void flush_dropsTheBatchAfterTwoConnectionFailures() { + EventProcessor processor = newProcessor(1000, 0, new FailingInterceptor(Integer.MAX_VALUE)); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(2, recorder.count()); + 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()); + } + + @Test + @SneakyThrows + public void flush_completesWhenBuildingTheRequestThrows() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + doThrow(new IllegalStateException("boom")).when(api).newPostRequest(any(), any()); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(0, recorder.count()); + // A later flush must not inherit a batch that was never settled. + flushAndWait(processor); + } + + @Test + @SneakyThrows + public void flush_completesWhenTheRequestProcessorIsAlreadyShutDown() { + EventProcessor processor = newProcessor(1000, 0); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + processor.getRequestProcessor().close(); + + flushAndWait(processor); + + assertEquals(0, recorder.count()); + } + + @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)); + + // close() runs its final flush while the event is still on its way into the buffer. + 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())); + 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 + public void trackEvent_dropsAnEventWhoseTraitsOrMetadataContainThemselves() { + EventProcessor processor = newProcessor(1000, 0); + Map cyclic = new HashMap<>(); + cyclic.put("self", cyclic); + List cyclicList = new ArrayList<>(); + cyclicList.add(cyclicList); + + // Serialising these recurses without end; it must surface as a dropped event, not an Error. + processor.trackEvent("purchase", "user-1", "1", cyclic, null); + processor.trackEvent("purchase", "user-2", "2", null, + Collections.singletonMap("list", cyclicList)); + processor.trackEvent("purchase", "user-3", "3", null, null); + + assertEquals(1, processor.bufferedEvents().size()); + assertEquals("user-3", processor.bufferedEvents().get(0).get("identifier")); + } + + @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; + + // Scheduling on the shut-down scheduler would throw RejectedExecutionException. + processor.start(); + + assertTrue(processor.getScheduler().isShutdown()); + } + + @Test + public void worstCaseBatchMillis_coversEveryAttemptAtTheClientTimeoutsPlusBackoff() { + OkHttpClient client = new OkHttpClient.Builder() + .connectTimeout(1000, TimeUnit.MILLISECONDS) + .writeTimeout(2000, TimeUnit.MILLISECONDS) + .readTimeout(3000, TimeUnit.MILLISECONDS) + .build(); + Retry retry = new Retry(2); + + // Two attempts of connect + write + read, and the 200ms backoff before the second. + assertEquals(2 * 6000 + 200, EventProcessor.worstCaseBatchMillis(client, retry)); + } + + @Test + public void worstCaseBatchMillis_prefersTheCallTimeout() { + OkHttpClient client = new OkHttpClient.Builder() + .callTimeout(4000, TimeUnit.MILLISECONDS) + .build(); + + assertEquals(2 * 4000 + 200, + EventProcessor.worstCaseBatchMillis(client, new Retry(2))); + } + + @Test + public void worstCaseBatchMillis_isUnboundedWhenATimeoutIsOff() { + OkHttpClient client = new OkHttpClient.Builder() + .readTimeout(0, TimeUnit.MILLISECONDS) + .build(); + + assertEquals(EventProcessor.UNBOUNDED, + EventProcessor.worstCaseBatchMillis(client, new Retry(2))); + } + + @Test + @SneakyThrows + public void close_stopsWaitingAtTheWorstCaseBound() { + // A call timeout of 100ms bounds a batch at 2 x 100ms + 200ms backoff. The hung API below + // ignores the cancellation, as a stuck interceptor or proxy would. + AcceptingInterceptor eventsApi = AcceptingInterceptor.blocked(); + OkHttpClient client = new OkHttpClient.Builder() + .callTimeout(100, TimeUnit.MILLISECONDS) + .addInterceptor(eventsApi) + .build(); + EventProcessor processor = new EventProcessor( + HttpUrl.get(EVENTS_URI), 1000, 0, new RequestProcessor(client, new FlagsmithLogger())); + processor.setApi(api); + FlagsmithLogger logger = mock(FlagsmithLogger.class); + processor.setLogger(logger); + assertEquals(400, processor.getCloseTimeoutMillis()); + + 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 < TimeUnit.SECONDS.toMillis(WAIT_SECONDS) / 2, + "close() waited " + elapsedMillis + "ms"); + } finally { + eventsApi.release(); + } + } + + @Test + public void closeTimeout_followsTheConfiguredTimeouts() { + FlagsmithConfig config = FlagsmithConfig.newBuilder() + .connectTimeout(1000) + .writeTimeout(2000) + .readTimeout(30000) + .build(); + EventProcessor processor = + new EventProcessor(config.getHttpClient(), config.getEventsUri(), 1, 0); + + // The read timeout the caller configured, not the SDK default. + assertEquals(2 * (1000 + 2000 + 30000) + 200, processor.getCloseTimeoutMillis()); + processor.close(); + } + + @Test + public void trackEvent_neverThrowsWhenTheApiIsMissing() { + EventProcessor processor = newProcessor(1, 0); + processor.setApi(null); + + processor.trackEvent("purchase", "user-1", "1", null, null); + + assertEquals(0, recorder.count()); + assertTrue(processor.bufferedEvents().isEmpty()); + } + + @Test + @SneakyThrows + public void flush_dropsEventsBeyondTheInFlightLimitInsteadOfQueueingThem() { + AcceptingInterceptor eventsApi = AcceptingInterceptor.blocked(); + EventProcessor processor = newProcessor(1000, 0, eventsApi); + FlagsmithLogger logger = mock(FlagsmithLogger.class); + processor.setLogger(logger); + + // The events API hangs, so every batch stays in flight until it is released. Each full + // buffer flushes itself, which puts exactly the limit in flight. + for (int i = 0; i < EventProcessor.MAX_IN_FLIGHT_EVENTS; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + } + processor.trackEvent("purchase", "one-too-many", "1", null, null); + CompletableFuture all = processor.flush(); + + assertTrue(processor.bufferedEvents().isEmpty(), "the dropped batch stayed in the buffer"); + verify(logger).error(contains("Dropped 1 events")); + assertFalse(all.isDone()); + + // A caller saturating the processor does not get an error line per flush. + for (int i = 0; i < 100; i++) { + processor.trackEvent("purchase", "also-dropped-" + i, "1", null, null); + processor.flush(); + } + verify(logger, times(1)).error(startsWith("Dropped")); + + eventsApi.release(); + all.get(WAIT_SECONDS, TimeUnit.SECONDS); + + assertEquals(EventProcessor.MAX_IN_FLIGHT_EVENTS, deliveredEvents()); + for (String body : recorder.bodies()) { + assertFalse(body.contains("one-too-many")); + assertFalse(body.contains("also-dropped")); + } + + // Once the backlog clears, batches flow again. + processor.trackEvent("purchase", "after", "1", null, null); + flushAndWait(processor); + assertEquals(EventProcessor.MAX_IN_FLIGHT_EVENTS + 1, deliveredEvents()); + } + + @Test + @SneakyThrows + public void flush_neverDropsEventsOnAHealthyApiWithASmallBuffer() { + // A one-event buffer sends a batch per event: the case a cap on batches would throttle. + EventProcessor processor = newProcessor(1, 0, AcceptingInterceptor.open()); + FlagsmithLogger logger = mock(FlagsmithLogger.class); + processor.setLogger(logger); + int events = 2000; + + for (int i = 0; i < events; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + } + flushAndWait(processor); + + assertEquals(events, deliveredEvents()); + 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; + } + + @Test + @SneakyThrows + public void flush_logsEventsTheApiRejects() { + EventProcessor processor = newProcessor(1000, 0); + FlagsmithLogger logger = mock(FlagsmithLogger.class); + processor.setLogger(logger); + 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); + 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")); + } + + @Test + @SneakyThrows + public void flush_logsNothingWhenEveryEventIsAccepted() { + EventProcessor processor = newProcessor(1000, 0); + FlagsmithLogger logger = mock(FlagsmithLogger.class); + processor.setLogger(logger); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(1, recorder.count()); + assertEquals(Collections.emptyList(), errorCalls(logger)); + } + + /** + * 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 new Response.Builder() + .request(chain.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); + } + } + + /** 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()); + } + } +}