diff --git a/README.md b/README.md index ba48788..d8af245 100644 --- a/README.md +++ b/README.md @@ -10,6 +10,37 @@ This SDK enables Android and Java applications to integrate with [Flagsmith](htt For full documentation visit [https://docs.flagsmith.com/clients/server-side](https://docs.flagsmith.com/clients/server-side). +## Experimentation events + +Enable events on the configuration, then record exposures and custom events: + +```java +FlagsmithClient flagsmith = FlagsmithClient.newBuilder() + .setApiKey(System.getenv("FLAGSMITH_ENVIRONMENT_KEY")) + .withConfiguration(FlagsmithConfig.newBuilder() + .withEnableEvents(true) + .build()) + .build(); + +// Records one $flag_exposure event when the identity is enrolled in a running experiment. +BaseFlag flag = flagsmith.getExperimentFlag("checkout_cta", "user-123"); + +flagsmith.trackEvent("purchase", "user-123"); +``` + +Events are buffered and sent in batches, on a timer and when the buffer fills. A batch that fails +on a retryable error (408, 429, 502, 503, 504 or a network error) is tried up to 3 times in +total, with backoff, and then kept for the next timed flush. Any other error drops it, and a 401 +or 403 stops sending until the client is re-created. + +`flushEvents()` sends what is buffered now. A short-lived process, such as a serverless function +or a CLI command, must call `close()` before it exits: it sends the remaining events and waits +for them, within a bound derived from the HTTP client's timeouts. Otherwise buffered events are +lost. + +`getDroppedEventCount()` returns how many events were dropped: when the buffer overflowed, on a +non-retryable error, when the events API rejected them, after a 401 or 403, or on close. + ## Contributing Please read [CONTRIBUTING.md](https://gist.github.com/kyle-ssg/c36a03aebe492e45cbd3eefb21cb0486) for details on our code of conduct, and the process for submitting pull requests diff --git a/src/main/java/com/flagsmith/FlagsmithClient.java b/src/main/java/com/flagsmith/FlagsmithClient.java index 8233f53..2c465fc 100644 --- a/src/main/java/com/flagsmith/FlagsmithClient.java +++ b/src/main/java/com/flagsmith/FlagsmithClient.java @@ -372,10 +372,11 @@ public void trackExposureEvent(String featureName, String identifier, Object val } /** - * Send buffered events now. + * Send buffered events now. A batch that still fails on a retryable error goes back in the + * buffer for the next flush, so a short-lived process should call {@link #close()} instead. * - * @return a future completing once every event buffered so far has been sent or dropped, already - * completed when events are not enabled + * @return a future completing once every event buffered so far has been sent, dropped or put + * back in the buffer; already completed when events are not enabled */ public CompletableFuture flushEvents() { if (eventProcessor == null) { @@ -385,9 +386,22 @@ public CompletableFuture flushEvents() { return eventProcessor.flush(); } + /** + * The number of events dropped since the client was built. It never decreases. + * + * @return the dropped event count, 0 when events are not enabled + */ + public long getDroppedEventCount() { + return eventProcessor == null ? 0 : eventProcessor.getDroppedEventCount(); + } + /** * Should be called when terminating the client to clean up any resources that * need cleaning up. + * + *

With events enabled this sends the buffered events and waits for them, within a bound + * derived from the HTTP client's timeouts; a batch failing here is dropped. Call it before a + * short-lived process, such as a serverless function or a CLI command, exits. **/ public void close() { if (pollingManager != null) { diff --git a/src/main/java/com/flagsmith/config/FlagsmithConfig.java b/src/main/java/com/flagsmith/config/FlagsmithConfig.java index 48adf34..a721b7f 100644 --- a/src/main/java/com/flagsmith/config/FlagsmithConfig.java +++ b/src/main/java/com/flagsmith/config/FlagsmithConfig.java @@ -346,8 +346,9 @@ public Builder withEventProcessor(EventProcessor processor) { } /** - * 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. + * Set the number of buffered events that triggers an immediate flush, and the most the buffer + * holds while batches are in flight: past it the oldest events are dropped. 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 diff --git a/src/main/java/com/flagsmith/threads/EventProcessor.java b/src/main/java/com/flagsmith/threads/EventProcessor.java index 0da9e9d..c5e63c3 100644 --- a/src/main/java/com/flagsmith/threads/EventProcessor.java +++ b/src/main/java/com/flagsmith/threads/EventProcessor.java @@ -1,16 +1,15 @@ 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.io.IOException; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -26,11 +25,11 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.stream.Collectors; -import java.util.stream.IntStream; +import java.util.concurrent.atomic.AtomicLong; import lombok.AccessLevel; import lombok.Getter; import lombok.Setter; @@ -39,10 +38,15 @@ import okhttp3.OkHttpClient; import okhttp3.Request; import okhttp3.RequestBody; +import okhttp3.Response; +import okhttp3.ResponseBody; /** * Buffers experimentation events and sends them to the Flagsmith events API, on a timer, when the - * buffer fills, and on {@link #close()}. Exposures are deduplicated within a flush window. + * buffer fills, and on {@link #close()}. A batch is retried on a 408, 429, 502, 503, 504 or a + * network error, and put back in the buffer for the next timed flush if it still fails; any other + * error drops it, and a 401 or 403 stops the processor. Exposures are deduplicated until they are + * delivered or dropped. */ public class EventProcessor { @@ -60,7 +64,10 @@ public class EventProcessor { private static final MediaType JSON_MEDIA_TYPE = MediaType.get("application/json; charset=utf-8"); static final int MAX_IN_FLIGHT_BATCHES = 2; - static final int MAX_BUFFERED_EVENTS = 1_000; + static final int MAX_ATTEMPTS = 3; + static final Set RETRYABLE_STATUSES = Set.of(408, 429, 502, 503, 504); + static final long BACKOFF_BASE_MILLIS = 1_000L; + static final long BACKOFF_CAP_MILLIS = 10_000L; static final long CLOSE_TIMEOUT_MILLIS = 25_000L; private static final long DROP_LOG_INTERVAL_NANOS = TimeUnit.SECONDS.toNanos(10); @@ -78,18 +85,23 @@ public class EventProcessor { private final ScheduledExecutorService scheduler; private final Set> inFlight = ConcurrentHashMap.newKeySet(); private CompletableFuture nextBatch = new CompletableFuture<>(); + private boolean held = false; // guarded by lock private int droppedSinceLastReport = 0; // guarded by lock private Long lastDropReportNanos = null; // guarded by lock + private final AtomicLong droppedEvents = new AtomicLong(); @Getter(AccessLevel.PACKAGE) private final RequestProcessor requestProcessor; @Getter(AccessLevel.PACKAGE) @Setter(AccessLevel.PACKAGE) private long closeTimeoutMillis; + @Setter(AccessLevel.PACKAGE) + private long backoffBaseMillis = BACKOFF_BASE_MILLIS; @Setter private FlagsmithSdk api; private FlagsmithLogger logger = new FlagsmithLogger(); private final AtomicBoolean closed = new AtomicBoolean(false); private final AtomicBoolean claimed = new AtomicBoolean(false); + private final AtomicBoolean stopped = new AtomicBoolean(false); private ScheduledFuture scheduledFlush; /** @@ -97,8 +109,8 @@ public class EventProcessor { * * @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 maxBufferItems number of buffered events that triggers an immediate flush, and + * the most the buffer holds; at least 1 * @param flushIntervalMillis interval between timed flushes; 0 disables the timer * @throws IllegalArgumentException when maxBufferItems is below 1 or flushIntervalMillis is * negative @@ -106,7 +118,7 @@ public class EventProcessor { public EventProcessor(OkHttpClient client, HttpUrl eventsUri, int maxBufferItems, int flushIntervalMillis) { this(eventsUri, maxBufferItems, flushIntervalMillis, - new RequestProcessor(withCallDeadline(client), new FlagsmithLogger(), buildRetry())); + new RequestProcessor(withCallDeadline(client), new FlagsmithLogger())); } EventProcessor(HttpUrl eventsUri, int maxBufferItems, int flushIntervalMillis, @@ -121,7 +133,7 @@ public EventProcessor(OkHttpClient client, HttpUrl eventsUri, int maxBufferItems this.maxBufferItems = maxBufferItems; this.flushIntervalMillis = flushIntervalMillis; this.requestProcessor = requestProcessor; - this.closeTimeoutMillis = worstCaseBatchMillis(requestProcessor.getClient(), buildRetry()); + this.closeTimeoutMillis = worstCaseBatchMillis(requestProcessor.getClient()); this.scheduler = Executors.newSingleThreadScheduledExecutor((runnable) -> { Thread thread = new Thread(runnable, "flagsmith-events"); thread.setDaemon(true); @@ -129,23 +141,14 @@ public EventProcessor(OkHttpClient client, HttpUrl eventsUri, int maxBufferItems }); } - private static Retry buildRetry() { - Retry retry = new Retry(2); - retry.setStatusForcelist( - IntStream.rangeClosed(500, 599).boxed().collect(Collectors.toSet())); - retry.setStatusForcelistOnly(Boolean.TRUE); - return retry; - } - - private static long worstCaseBatchMillis(OkHttpClient client, Retry retry) { + private static long worstCaseBatchMillis(OkHttpClient client) { long attemptMillis = attemptMillis(client); if (attemptMillis == 0) { return CLOSE_TIMEOUT_MILLIS; } - long total = retry.getTotal() * attemptMillis; - for (int retried = 1; retried < retry.getTotal(); retried++) { - total += (long) (Math.min(retry.getBackoffFactor() * 2 * retried, retry.getBackoffMax()) - * 1000); + long total = MAX_ATTEMPTS * attemptMillis; + for (int retry = 1; retry < MAX_ATTEMPTS; retry++) { + total += backoffCeilingMillis(BACKOFF_BASE_MILLIS, retry); } return total; } @@ -171,6 +174,14 @@ private static OkHttpClient withCallDeadline(OkHttpClient client) { return client.newBuilder().callTimeout(callTimeout, TimeUnit.MILLISECONDS).build(); } + static long backoffCeilingMillis(long baseMillis, int retry) { + return Math.min(BACKOFF_CAP_MILLIS, baseMillis << Math.min(retry - 1, 20)); + } + + static long backoffMillis(long baseMillis, int retry) { + return ThreadLocalRandom.current().nextLong(backoffCeilingMillis(baseMillis, retry) + 1); + } + /** * Reserve this processor for one client; called by {@code FlagsmithClient.Builder}. * @@ -192,6 +203,17 @@ public void setLogger(FlagsmithLogger logger) { requestProcessor.setLogger(logger); } + /** + * The number of events dropped so far: when the buffer overflowed, on a non-retryable error, + * when the events API rejected them, after a 401 or 403, and when a batch failed on close. It + * never decreases. + * + * @return the dropped event count + */ + public long getDroppedEventCount() { + return droppedEvents.get(); + } + /** * Buffer a custom event. * @@ -208,7 +230,7 @@ public void trackEvent(String event, String identifier, Object value, /** * Buffer a flag exposure event. Exposures equal in feature, identifier, value and experiment - * are only sent once per flush window. + * are only sent once until the events API accepts or drops them. * * @param featureName feature the identity was exposed to * @param identifier identity the exposure belongs to @@ -222,43 +244,59 @@ public void trackExposureEvent(String featureName, String identifier, Object val } /** - * Send everything buffered so far. + * Send everything buffered so far, including events kept after a failure. * - * @return a future completing once every event buffered so far has been sent or dropped + * @return a future completing once every event buffered so far has been sent, dropped or put + * back in the buffer; it never completes exceptionally */ public CompletableFuture flush() { - List> batch = null; - CompletableFuture tracked = null; + synchronized (lock) { + held = false; + } + return dispatch(); + } + + private CompletableFuture dispatch() { + Batch batch = null; CompletableFuture waiting = null; synchronized (lock) { - if (!buffer.isEmpty() && (inFlight.size() < MAX_IN_FLIGHT_BATCHES || closed.get())) { - batch = new ArrayList<>(buffer); - buffer.clear(); - dedupeKeys.clear(); - dedupeKeyByEvent.clear(); - tracked = nextBatch; - nextBatch = new CompletableFuture<>(); - inFlight.add(tracked); - } else if (!buffer.isEmpty()) { + if (!buffer.isEmpty() && (closed.get() || canSend())) { + batch = takeBatch(); + } else if (!buffer.isEmpty() && !held) { waiting = nextBatch; } } + CompletableFuture sent; if (batch != null) { - send(batch, tracked); + send(batch); + sent = CompletableFuture.allOf(awaitInFlight(), batch.tracked); + } else { + sent = awaitInFlight(); } - CompletableFuture sent = awaitInFlight(); return waiting == null ? sent : CompletableFuture.allOf(sent, waiting); } + private boolean canSend() { + return !held && inFlight.size() < MAX_IN_FLIGHT_BATCHES; + } + + private Batch takeBatch() { + Batch batch = new Batch(new ArrayList<>(buffer), nextBatch); + buffer.clear(); + nextBatch = new CompletableFuture<>(); + inFlight.add(batch.tracked); + return batch; + } + /** * 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. + * is already running, or once the processor is closed or stopped. */ public synchronized void start() { - if (flushIntervalMillis <= 0 || scheduledFlush != null) { + if (flushIntervalMillis <= 0 || scheduledFlush != null || stopped.get()) { return; } @@ -273,8 +311,9 @@ public synchronized void start() { /** * Stop the timer, flush what is left and release HTTP resources. Blocks until in-flight - * batches settle, for at most one batch's worst case under the client's timeouts, or - * {@link #CLOSE_TIMEOUT_MILLIS} if a timeout is off. + * batches settle, for at most one batch's worst case under the client's timeouts and the + * retries, or {@link #CLOSE_TIMEOUT_MILLIS} if a timeout is off. A batch that fails from here on + * is dropped, not kept. */ public void close() { synchronized (this) { @@ -303,6 +342,10 @@ private void bufferEvent(String event, String featureName, String identifier, Ob logClosed(event); return; } + if (stopped.get()) { + droppedEvents.incrementAndGet(); + return; + } try { final String stringValue = value == null ? null : String.valueOf(value); @@ -325,7 +368,7 @@ private void bufferEvent(String event, String featureName, String identifier, Ob eventPayload.put("metadata", toJson(eventMetadata)); eventPayload.put("timestamp", System.currentTimeMillis()); - boolean isFull = false; + Batch full = null; int droppedToReport = 0; synchronized (lock) { @@ -333,6 +376,10 @@ private void bufferEvent(String event, String featureName, String identifier, Ob logClosed(event); return; } + if (stopped.get()) { + droppedEvents.incrementAndGet(); + return; + } if (dedupe) { List key = dedupeKey(event, featureName, identifier, stringValue, experimentId); if (!dedupeKeys.add(key)) { @@ -340,20 +387,25 @@ private void bufferEvent(String event, String featureName, String identifier, Ob } dedupeKeyByEvent.put(eventPayload, key); } + if (buffer.size() >= maxBufferItems) { + if (canSend()) { + full = takeBatch(); + } else { + droppedToReport = drop(Collections.singletonList(buffer.remove(0))); + } + } buffer.add(eventPayload); - if (inFlight.size() < MAX_IN_FLIGHT_BATCHES) { - isFull = buffer.size() >= maxBufferItems; - } else if (buffer.size() > Math.max(maxBufferItems, MAX_BUFFERED_EVENTS)) { - dedupeKeys.remove(dedupeKeyByEvent.remove(buffer.remove(0))); - droppedToReport = recordDrop(1); + if (full == null && canSend() && buffer.size() >= maxBufferItems) { + full = takeBatch(); } } - logDrops(droppedToReport); - if (isFull) { - flush(); + reportDrops(droppedToReport, "the buffer is full"); + if (full != null) { + send(full); } } catch (JsonProcessingException | RuntimeException e) { + droppedEvents.incrementAndGet(); logger.error("Failed to buffer event " + event + ".", e); } } @@ -389,14 +441,15 @@ private static List dedupeKey(String event, String featureName, String i experimentId == null ? null : String.valueOf(experimentId)); } - private void send(List> batch, CompletableFuture tracked) { - final int batchSize = batch.size(); + private void send(Batch full) { + List> batch = full.events; + CompletableFuture tracked = full.tracked; boolean submitted = false; + String reason = "sending them failed"; try { if (api == null) { - logger.error("Dropping " + batchSize - + " events: the event processor has no API wrapper."); + reason = "the event processor has no API wrapper"; return; } @@ -413,45 +466,178 @@ private void send(List> batch, CompletableFuture track .header(ACCEPT_HEADER, "application/json") .build(); - requestProcessor - .submit(request, new TypeReference() {}, Boolean.FALSE, buildRetry()) - .whenComplete((response, error) -> { - try { - logRejections(response, batchSize); - } finally { - settle(tracked); - } - }); + requestProcessor.execute(() -> deliver(request, batch, tracked)); submitted = true; } catch (Exception e) { - logger.error("Dropping " + batchSize + " events: failed to send them.", e); + logger.error("Failed to send " + batch.size() + " events.", e); } finally { if (!submitted) { - settle(tracked); + settle(tracked, batch, Outcome.DROPPED, reason); + } + } + } + + private void deliver(Request request, List> batch, + CompletableFuture tracked) { + Outcome outcome = Outcome.DROPPED; + String reason = "sending them failed"; + + try { + for (int attempt = 1; ; attempt++) { + Integer status = null; + String body = null; + try (Response response = requestProcessor.getClient().newCall(request).execute()) { + status = response.code(); + ResponseBody responseBody = response.body(); + if (response.isSuccessful() && responseBody != null) { + body = responseBody.string(); + } + } catch (IOException e) { + reason = "sending them failed: " + e; + } + + if (status != null && status >= 200 && status < 300) { + countRejections(body, batch.size()); + outcome = Outcome.DELIVERED; + break; + } + if (status != null) { + reason = "the events API answered " + status; + } + if (status != null && (status == 401 || status == 403)) { + outcome = Outcome.UNAUTHORISED; + break; + } + if (status != null && !RETRYABLE_STATUSES.contains(status)) { + break; + } + if (attempt == MAX_ATTEMPTS || stopped.get()) { + outcome = Outcome.RETRY_LATER; + break; + } + Thread.sleep(backoffMillis(backoffBaseMillis, attempt)); } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + outcome = Outcome.RETRY_LATER; + } catch (RuntimeException e) { + logger.error("Failed to send " + batch.size() + " events.", e); + } finally { + settle(tracked, batch, outcome, reason); } } - private void logRejections(JsonNode response, int batchSize) { - JsonNode rejected = response == null ? null : response.get("rejected"); + private void countRejections(String body, int batchSize) { + JsonNode rejected; + try { + JsonNode response = body == null ? null : MapperFactory.getMapper().readTree(body); + rejected = response == null ? null : response.get("rejected"); + } catch (JsonProcessingException e) { + return; + } if (rejected != null && rejected.isArray() && rejected.size() > 0) { + droppedEvents.addAndGet(rejected.size()); logger.error("The events API rejected " + rejected.size() + " of " + batchSize - + " events. First rejection: " + rejected.get(0)); + + " events, which are not resent. First rejection: " + rejected.get(0)); } } - private void settle(CompletableFuture tracked) { + private void settle(CompletableFuture tracked, List> batch, + Outcome outcome, String reason) { + int droppedToReport = 0; + boolean kept = false; + boolean stopping = false; boolean waiting; - synchronized (lock) { - inFlight.remove(tracked); - waiting = !buffer.isEmpty(); + CompletableFuture released = null; + + try { + synchronized (lock) { + inFlight.remove(tracked); + switch (outcome) { + case DELIVERED: + release(batch); + break; + case RETRY_LATER: + if (closed.get() || stopped.get()) { + droppedToReport = drop(batch); + } else { + droppedToReport = requeue(batch); + kept = true; + held = true; + released = nextBatch; + nextBatch = new CompletableFuture<>(); + } + break; + case UNAUTHORISED: + stopping = stopped.compareAndSet(false, true); + droppedEvents.addAndGet(batch.size()); + if (stopping) { + droppedEvents.addAndGet(buffer.size()); + buffer.clear(); + dedupeKeys.clear(); + dedupeKeyByEvent.clear(); + released = nextBatch; + nextBatch = new CompletableFuture<>(); + } + break; + default: + droppedToReport = drop(batch); + } + waiting = !held && !buffer.isEmpty(); + } + + if (stopping) { + stopTimer(); + logger.error("The events API refused the environment key (" + reason + + "); events are dropped until the client is re-created with a valid key."); + } + if (kept) { + logger.error("Kept " + batch.size() + " events for the next flush: " + reason + "."); + } + reportDrops(droppedToReport, kept ? "the buffer is full" : reason); + } finally { + tracked.complete(null); + if (released != null) { + released.complete(null); + } } - tracked.complete(null); + if (waiting) { - flush(); + dispatch(); + } + } + + private int requeue(List> batch) { + buffer.addAll(0, batch); + int overflow = buffer.size() - maxBufferItems; + if (overflow <= 0) { + return 0; + } + List> oldest = buffer.subList(0, overflow); + int toReport = drop(new ArrayList<>(oldest)); + oldest.clear(); + return toReport; + } + + private void release(List> events) { + for (Map event : events) { + List key = dedupeKeyByEvent.remove(event); + if (key != null) { + dedupeKeys.remove(key); + } } } + private int drop(List> events) { + release(events); + droppedEvents.addAndGet(events.size()); + return recordDrop(events.size()); + } + + private synchronized void stopTimer() { + scheduler.shutdownNow(); + } + private int recordDrop(int count) { droppedSinceLastReport += count; long now = System.nanoTime(); @@ -464,11 +650,10 @@ private int recordDrop(int count) { return toReport; } - private void logDrops(int dropped) { - if (dropped > 0) { - logger.error("Dropped the " + dropped + " oldest events: " + MAX_IN_FLIGHT_BATCHES - + " batches are in flight to the events API and the buffer is full. Further drops are" - + " reported at most every " + TimeUnit.NANOSECONDS.toSeconds(DROP_LOG_INTERVAL_NANOS) + private void reportDrops(int dropped, String reason) { + if (dropped > 0 && !stopped.get()) { + logger.error("Dropped " + dropped + " events, latest because " + reason + ". Further drops" + + " are reported at most every " + TimeUnit.NANOSECONDS.toSeconds(DROP_LOG_INTERVAL_NANOS) + "s."); } } @@ -482,4 +667,17 @@ List> bufferedEvents() { return new ArrayList<>(buffer); } } + + private enum Outcome { DELIVERED, RETRY_LATER, DROPPED, UNAUTHORISED } + + private static final class Batch { + + private final List> events; + private final CompletableFuture tracked; + + private Batch(List> events, CompletableFuture tracked) { + this.events = events; + this.tracked = tracked; + } + } } diff --git a/src/main/java/com/flagsmith/threads/RequestProcessor.java b/src/main/java/com/flagsmith/threads/RequestProcessor.java index 36bc4bc..a244dd6 100644 --- a/src/main/java/com/flagsmith/threads/RequestProcessor.java +++ b/src/main/java/com/flagsmith/threads/RequestProcessor.java @@ -141,6 +141,10 @@ public CompletableFuture submit( return completableFuture; } + void execute(Runnable task) { + executor.execute(task); + } + public void close() { this.executor.shutdown(); } diff --git a/src/test/java/com/flagsmith/FlagsmithClientTest.java b/src/test/java/com/flagsmith/FlagsmithClientTest.java index 4e2df3a..265e827 100644 --- a/src/test/java/com/flagsmith/FlagsmithClientTest.java +++ b/src/test/java/com/flagsmith/FlagsmithClientTest.java @@ -1328,6 +1328,28 @@ public void testRebuildingClosesThePreviousEventProcessor() { assertEquals(2, batches.size()); } + @Test + public void testGetDroppedEventCount() { + FlagsmithClient disabled = FlagsmithClient.newBuilder().setApiKey("api-key").build(); + assertEquals(0, disabled.getDroppedEventCount()); + + MockInterceptor interceptor = new MockInterceptor(); + interceptor.addRule().post("http://events-uri/v1/events").anyTimes().respond(400); + FlagsmithClient client = FlagsmithClient.newBuilder() + .withConfiguration(FlagsmithConfig.newBuilder() + .addHttpInterceptor(interceptor) + .eventsUri("http://events-uri") + .withEnableEvents(Boolean.TRUE) + .build()) + .setApiKey("api-key") + .build(); + client.trackEvent("purchase", "user-1"); + client.trackEvent("purchase", "user-2"); + client.close(); + + assertEquals(2, client.getDroppedEventCount()); + } + @Test public void testEventsRequestLeavesOutCustomHeaders() { List requests = Collections.synchronizedList(new ArrayList<>()); diff --git a/src/test/java/com/flagsmith/threads/EventProcessorTest.java b/src/test/java/com/flagsmith/threads/EventProcessorTest.java index e161a43..0bf49bc 100644 --- a/src/test/java/com/flagsmith/threads/EventProcessorTest.java +++ b/src/test/java/com/flagsmith/threads/EventProcessorTest.java @@ -114,6 +114,7 @@ private EventProcessor newProcessor( flushIntervalMillis, new RequestProcessor(client, new FlagsmithLogger())); eventProcessor.setApi(api); + eventProcessor.setBackoffBaseMillis(1); return eventProcessor; } @@ -283,12 +284,12 @@ public void flush_completesOnlyAfterTheInFlightPostCompletes() { */ @ParameterizedTest @CsvSource({ - "503, 1, 2", "501, 1, 2", "500, 2, 2", - "400, 1, 1", - ", 1, 2", ", 2, 2"}) + "408, 1, 2, 0", "429, 1, 2, 0", "502, 1, 2, 0", "503, 2, 3, 0", "504, 1, 2, 0", + ", 1, 2, 0", ", 2, 3, 0", + "400, 1, 1, 1", "413, 1, 1, 1", "422, 1, 1, 1", "500, 1, 1, 1", "501, 1, 1, 1"}) @SneakyThrows - public void flush_retriesOnceOnAServerErrorOrConnectionFailure( - Integer status, int failures, int expectedAttempts) { + public void flush_retriesOnlyARetryableFailureAndDropsTheRest( + Integer status, int failures, int expectedAttempts, long expectedDropped) { EventProcessor processor = status == null ? newProcessor(1000, 0, new FailingInterceptor(failures)) : newProcessor(1000, 0); @@ -303,6 +304,233 @@ public void flush_retriesOnceOnAServerErrorOrConnectionFailure( assertEquals(expectedAttempts, recorder.count()); assertEquals(recorder.bodies().get(0), recorder.bodies().get(expectedAttempts - 1)); assertTrue(processor.bufferedEvents().isEmpty()); + assertEquals(expectedDropped, processor.getDroppedEventCount()); + } + + @ParameterizedTest + @CsvSource(value = {"503", "none"}, nullValues = "none") + @SneakyThrows + public void flush_keepsABatchThatFailsThreeTimesForTheNextFlush(Integer status) { + EventProcessor processor = status == null + ? newProcessor(2, 0, new FailingInterceptor(EventProcessor.MAX_ATTEMPTS)) + : newProcessor(2, 0); + FlagsmithLogger logger = mockLogger(processor); + if (status != null) { + interceptor.addRule().post(EVENTS_ENDPOINT).times(EventProcessor.MAX_ATTEMPTS) + .respond(status); + } + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackEvent("purchase", "kept", "1", null, null); + flushAndWait(processor); + + assertEquals(EventProcessor.MAX_ATTEMPTS, recorder.count()); + assertEquals(1, processor.bufferedEvents().size()); + verify(logger).error(startsWith("Kept 1 events for the next flush")); + + processor.trackEvent("purchase", "newer", "1", null, null); + assertEquals(2, processor.bufferedEvents().size()); + assertEquals(EventProcessor.MAX_ATTEMPTS, recorder.count()); + + flushAndWait(processor); + + assertEquals(EventProcessor.MAX_ATTEMPTS + 1, recorder.count()); + JsonNode events = MapperFactory.getMapper() + .readTree(recorder.bodies().get(EventProcessor.MAX_ATTEMPTS)).get("events"); + assertEquals("kept", events.get(0).get("identifier").asText()); + assertEquals("newer", events.get(1).get("identifier").asText()); + assertEquals(0, processor.getDroppedEventCount()); + } + + @Test + @SneakyThrows + public void flush_dropsTheOldestEventsWhenAKeptBatchOverflowsTheBuffer() { + CountDownLatch release = new CountDownLatch(1); + EventProcessor processor = newProcessor(3, 0, (chain) -> { + try { + release.await(WAIT_SECONDS, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + throw new IOException("connection refused"); + }); + mockLogger(processor); + + for (int i = 1; i <= 3; i++) { + processor.trackEvent("purchase", "old-" + i, "1", null, null); + } + processor.trackEvent("purchase", "new-1", "1", null, null); + processor.trackEvent("purchase", "new-2", "1", null, null); + release.countDown(); + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(WAIT_SECONDS); + while (processor.getDroppedEventCount() < 2 && System.nanoTime() < deadline) { + Thread.sleep(10); + } + + List> buffered = processor.bufferedEvents(); + assertEquals(3, buffered.size()); + assertEquals("old-3", buffered.get(0).get("identifier")); + assertEquals(2, processor.getDroppedEventCount()); + } + + @Test + @SneakyThrows + public void flush_keepsAnExposureDedupedUntilItsBatchIsDelivered() { + EventProcessor processor = newProcessor(1000, 0); + mockLogger(processor); + interceptor.addRule().post(EVENTS_ENDPOINT).times(EventProcessor.MAX_ATTEMPTS).respond(503); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + flushAndWait(processor); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + assertEquals(1, processor.bufferedEvents().size()); + + flushAndWait(processor); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + assertEquals(1, processor.bufferedEvents().size()); + } + + @Test + @SneakyThrows + public void flush_dropsAnExposureAndForgetsItsKeyOnANonRetryableStatus() { + EventProcessor processor = newProcessor(1000, 0); + mockLogger(processor); + interceptor.addRule().post(EVENTS_ENDPOINT).times(1).respond(400); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(ACCEPTED_BODY, MEDIATYPE_JSON); + + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + flushAndWait(processor); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + + assertEquals(1, processor.bufferedEvents().size()); + assertEquals(1, processor.getDroppedEventCount()); + } + + @ParameterizedTest + @CsvSource({"401", "403"}) + @SneakyThrows + public void flush_stopsTheProcessorWhenTheKeyIsRefused(int status) { + EventProcessor processor = newProcessor(1000, 60_000); + FlagsmithLogger logger = mockLogger(processor); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(status); + + processor.start(); + processor.trackEvent("purchase", "user-1", "1", null, null); + flushAndWait(processor); + + assertEquals(1, recorder.count()); + assertEquals(1, processor.getDroppedEventCount()); + verify(logger, times(1)).error(startsWith("The events API refused the environment key")); + assertTrue(processor.getScheduler().isShutdown()); + + processor.trackEvent("purchase", "user-2", "1", null, null); + processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); + flushAndWait(processor); + + assertTrue(processor.bufferedEvents().isEmpty()); + assertEquals(1, recorder.count()); + assertEquals(3, processor.getDroppedEventCount()); + verify(logger, times(1)).error(startsWith("The events API refused the environment key")); + } + + @Test + @SneakyThrows + public void flush_dropsTheBufferAndLogsOnceWhenTheKeyIsRefusedWithBatchesInFlight() { + CountDownLatch release = new CountDownLatch(1); + EventProcessor processor = newProcessor(1000, 0, (chain) -> { + try { + release.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(401) + .message("Unauthorized") + .body(ResponseBody.create("", MediaType.get("application/json"))) + .build(); + }); + FlagsmithLogger logger = mockLogger(processor); + + for (int i = 0; i < EventProcessor.MAX_IN_FLIGHT_BATCHES; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + processor.flush(); + } + processor.trackEvent("purchase", "waiting", "1", null, null); + CompletableFuture flushed = processor.flush(); + assertTrue(recorder.awaitCount(EventProcessor.MAX_IN_FLIGHT_BATCHES)); + + release.countDown(); + flushed.get(WAIT_SECONDS, TimeUnit.SECONDS); + + assertEquals(EventProcessor.MAX_IN_FLIGHT_BATCHES, recorder.count()); + assertEquals(EventProcessor.MAX_IN_FLIGHT_BATCHES + 1, processor.getDroppedEventCount()); + verify(logger, times(1)).error(startsWith("The events API refused the environment key")); + } + + @Test + @SneakyThrows + public void flush_completesWhenTheEventsItWaitsForAreKept() { + CountDownLatch release = new CountDownLatch(1); + EventProcessor processor = newProcessor(1000, 0, (chain) -> { + try { + release.await(WAIT_SECONDS, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IOException(e); + } + throw new IOException("connection refused"); + }); + mockLogger(processor); + + for (int i = 0; i < EventProcessor.MAX_IN_FLIGHT_BATCHES; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + processor.flush(); + } + processor.trackEvent("purchase", "waiting", "1", null, null); + CompletableFuture flushed = processor.flush(); + + release.countDown(); + flushed.get(WAIT_SECONDS, TimeUnit.SECONDS); + + assertEquals(EventProcessor.MAX_IN_FLIGHT_BATCHES + 1, processor.bufferedEvents().size()); + assertEquals(0, processor.getDroppedEventCount()); + } + + @Test + public void backoffMillis_isAFullJitterUpToTheCappedExponentialCeiling() { + long[] ceilings = {1_000, 2_000, 4_000, 8_000, 10_000, 10_000}; + for (int retry = 1; retry <= ceilings.length; retry++) { + assertEquals(ceilings[retry - 1], + EventProcessor.backoffCeilingMillis(EventProcessor.BACKOFF_BASE_MILLIS, retry)); + for (int i = 0; i < 100; i++) { + long backoff = EventProcessor.backoffMillis(EventProcessor.BACKOFF_BASE_MILLIS, retry); + assertTrue(backoff >= 0 && backoff <= ceilings[retry - 1], + "retry " + retry + " backed off " + backoff + "ms"); + } + } + assertEquals(EventProcessor.BACKOFF_CAP_MILLIS, + EventProcessor.backoffCeilingMillis(EventProcessor.BACKOFF_BASE_MILLIS, 1_000)); + } + + @Test + @SneakyThrows + public void close_dropsAndCountsABatchThatStillFails() { + EventProcessor processor = newProcessor(1000, 0); + eventProcessor = null; + FlagsmithLogger logger = mockLogger(processor); + interceptor.addRule().post(EVENTS_ENDPOINT).anyTimes().respond(503); + + processor.trackEvent("purchase", "user-1", "1", null, null); + processor.close(); + + assertEquals(EventProcessor.MAX_ATTEMPTS, recorder.count()); + assertTrue(processor.bufferedEvents().isEmpty()); + assertEquals(1, processor.getDroppedEventCount()); + verify(logger).error(startsWith("Dropped 1 events")); } @Test @@ -361,6 +589,7 @@ public void trackEvent_dropsTheBatchWithoutThrowingWhenItCannotBeSent( assertEquals(0, recorder.count()); assertTrue(processor.bufferedEvents().isEmpty()); + assertEquals(1, processor.getDroppedEventCount()); } @Test @@ -532,14 +761,14 @@ public void start_isANoOpAfterClose() { private static Stream clientTimeouts() { FlagsmithConfig longRead = FlagsmithConfig.newBuilder().readTimeout(30_000).build(); return Stream.of( - // Two attempts at 2s connect + 5s write + 5s read, and 200ms backoff before the second. - Arguments.of(FlagsmithConfig.newBuilder().build().getHttpClient(), 2 * 12_000 + 200, + // Three attempts at 2s connect + 5s write + 5s read, and backoffs of up to 1s and 2s. + Arguments.of(FlagsmithConfig.newBuilder().build().getHttpClient(), 3 * 12_000 + 3_000, 12_000), - Arguments.of(longRead.getHttpClient(), 2 * 37_000 + 200, 37_000), + Arguments.of(longRead.getHttpClient(), 3 * 37_000 + 3_000, 37_000), Arguments.of(new OkHttpClient.Builder().callTimeout(4, TimeUnit.SECONDS).build(), - 2 * 4_000 + 200, 4_000), + 3 * 4_000 + 3_000, 4_000), Arguments.of(new OkHttpClient.Builder().readTimeout(0, TimeUnit.SECONDS).build(), - 2 * EventProcessor.CLOSE_TIMEOUT_MILLIS + 200, EventProcessor.CLOSE_TIMEOUT_MILLIS)); + 3 * EventProcessor.CLOSE_TIMEOUT_MILLIS + 3_000, EventProcessor.CLOSE_TIMEOUT_MILLIS)); } @ParameterizedTest @@ -671,17 +900,18 @@ public void trackEvent_dropsTheOldestEventsOnceTheBufferIsFullBehindTheLimit() { for (int i = 0; i < inFlight; i++) { processor.trackEvent("purchase", "user-" + i, "1", null, null); } - for (int i = 0; i < EventProcessor.MAX_BUFFERED_EVENTS + 500; i++) { + for (int i = 0; i < 1000 + 500; i++) { processor.trackEvent("purchase", "overflow-" + i + "-", "1", null, null); } - assertEquals(EventProcessor.MAX_BUFFERED_EVENTS, processor.bufferedEvents().size()); - verify(logger).error(contains("Dropped the 1 oldest events")); + assertEquals(1000, processor.bufferedEvents().size()); + assertEquals(500, processor.getDroppedEventCount()); + verify(logger).error(startsWith("Dropped 1 events, latest because the buffer is full")); verify(logger, times(1)).error(startsWith("Dropped")); eventsApi.release(); - assertTrue(awaitDelivered(inFlight + EventProcessor.MAX_BUFFERED_EVENTS)); + assertTrue(awaitDelivered(inFlight + 1000)); for (String body : recorder.bodies()) { assertTrue(MapperFactory.getMapper().readTree(body).get("events").size() <= 1000); } @@ -701,39 +931,54 @@ public void trackExposureEvent_buffersAnExposureAgainOnceItsCopyWasDropped() { processor.trackEvent("purchase", "user-" + i, "1", null, null); } processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); - for (int i = 0; i < EventProcessor.MAX_BUFFERED_EVENTS; i++) { + for (int i = 0; i < 1000; i++) { processor.trackEvent("purchase", "overflow-" + i, "1", null, null); } processor.trackExposureEvent("checkout_cta", "user-1", "treatment", null, null); List> buffered = processor.bufferedEvents(); - assertEquals(EventProcessor.MAX_BUFFERED_EVENTS, buffered.size()); + assertEquals(1000, buffered.size()); assertEquals("$flag_exposure", buffered.get(buffered.size() - 1).get("event")); eventsApi.release(); } - private static Stream healthyApiLoads() { - return Stream.of( - Arguments.of(Integer.MAX_VALUE, EventProcessor.MAX_BUFFERED_EVENTS + 1), - Arguments.of(1, 2000)); - } - - @ParameterizedTest - @MethodSource("healthyApiLoads") + @Test @SneakyThrows - public void flush_neverDropsEventsOnAHealthyApi(int maxBufferItems, int events) { - EventProcessor processor = newProcessor(maxBufferItems, 0, AcceptingInterceptor.open()); + public void flush_neverDropsEventsOnAHealthyApiWithinTheBuffer() { + EventProcessor processor = newProcessor(Integer.MAX_VALUE, 0, AcceptingInterceptor.open()); FlagsmithLogger logger = mockLogger(processor); - for (int i = 0; i < events; i++) { + for (int i = 0; i < 1001; i++) { processor.trackEvent("purchase", "user-" + i, "1", null, null); } processor.flush(); - assertTrue(awaitDelivered(events)); + assertTrue(awaitDelivered(1001)); assertEquals(Collections.emptyList(), errorCalls(logger)); } + @Test + @SneakyThrows + public void flush_deliversOrCountsEveryEventWhenABurstOutrunsASmallBuffer() { + EventProcessor processor = newProcessor(1, 0, AcceptingInterceptor.open()); + mockLogger(processor); + + for (int i = 0; i < 2000; i++) { + processor.trackEvent("purchase", "user-" + i, "1", null, null); + } + flushAndWait(processor); + + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(WAIT_SECONDS); + while (deliveredEvents() + processor.getDroppedEventCount() < 2000 + && System.nanoTime() < deadline) { + Thread.sleep(10); + } + assertEquals(2000, deliveredEvents() + processor.getDroppedEventCount()); + for (String body : recorder.bodies()) { + assertEquals(1, MapperFactory.getMapper().readTree(body).get("events").size()); + } + } + /** * Every error-level call on a mocked logger. Reads invocations rather than verify(...), whose * varargs matching silently misses calls with a different argument count. @@ -789,6 +1034,7 @@ public void flush_logsOnlyEventsTheApiRejects() { verify(logger).error(contains("rejected 1 of 2 events")); verify(logger).error(contains("event too long")); + assertEquals(1, processor.getDroppedEventCount()); } /**