diff --git a/launchdarkly-android-client-sdk/src/androidTest/java/com/launchdarkly/sdk/android/LDClientEventTest.java b/launchdarkly-android-client-sdk/src/androidTest/java/com/launchdarkly/sdk/android/LDClientEventTest.java index 544c416b..86dd4230 100644 --- a/launchdarkly-android-client-sdk/src/androidTest/java/com/launchdarkly/sdk/android/LDClientEventTest.java +++ b/launchdarkly-android-client-sdk/src/androidTest/java/com/launchdarkly/sdk/android/LDClientEventTest.java @@ -17,7 +17,9 @@ import com.launchdarkly.sdk.android.DataModel.Flag; import com.launchdarkly.sdk.android.LDConfig.Builder.AutoEnvAttributes; import com.launchdarkly.sdk.android.integrations.DedupingHook; +import com.launchdarkly.sdk.android.integrations.LDCrashHandler; import com.launchdarkly.sdk.android.integrations.Hook; +import com.launchdarkly.sdk.android.subsystems.EventProcessor; import com.launchdarkly.sdk.android.subsystems.PersistentDataStore; import com.launchdarkly.sdk.internal.GsonHelpers; import com.launchdarkly.sdk.json.JsonSerialization; @@ -26,8 +28,14 @@ import org.junit.Test; import java.io.IOException; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.Future; +import java.util.concurrent.FutureTask; +import java.util.concurrent.TimeUnit; import okhttp3.HttpUrl; +import okhttp3.mockwebserver.Dispatcher; import okhttp3.mockwebserver.MockResponse; import okhttp3.mockwebserver.MockWebServer; import okhttp3.mockwebserver.RecordedRequest; @@ -95,6 +103,165 @@ public void testTrackData() throws IOException, InterruptedException { } } + @Test + public void flushAndWaitReportsDelivery() throws IOException, InterruptedException { + try (MockWebServer mockEventsServer = new MockWebServer()) { + mockEventsServer.start(); + mockEventsServer.enqueue(new MockResponse()); + + LDConfig ldConfig = baseConfigBuilder(mockEventsServer).build(); + try (LDClient client = LDClient.init(application, ldConfig, ldContext, 0)) { + client.track("test-event"); + + assertTrue(client.flushAndWait(5, TimeUnit.SECONDS)); + LDValue[] events = getEventsFromLastRequest(mockEventsServer, 2); + assertCustomEvent(events[1], ldContext, "test-event"); + } + } + } + + @Test + public void crashHandlerDeliversTheEventsBeforePassingTheCrashOn() throws IOException, InterruptedException { + Thread.UncaughtExceptionHandler originalDefault = Thread.getDefaultUncaughtExceptionHandler(); + try (MockWebServer mockEventsServer = new MockWebServer()) { + mockEventsServer.start(); + mockEventsServer.enqueue(new MockResponse()); + + LDConfig ldConfig = baseConfigBuilder(mockEventsServer).build(); + try (LDClient client = LDClient.init(application, ldConfig, ldContext, 0)) { + int[] requestsSeenByNextHandler = {-1}; + Thread.setDefaultUncaughtExceptionHandler((thread, throwable) -> + requestsSeenByNextHandler[0] = mockEventsServer.getRequestCount()); + LDCrashHandler.install(5, TimeUnit.SECONDS); + client.track("test-event"); + + Thread.getDefaultUncaughtExceptionHandler() + .uncaughtException(Thread.currentThread(), new RuntimeException("boom")); + + assertEquals(1, requestsSeenByNextHandler[0]); + LDValue[] events = getEventsFromLastRequest(mockEventsServer, 2); + assertCustomEvent(events[1], ldContext, "test-event"); + } + } finally { + Thread.setDefaultUncaughtExceptionHandler(originalDefault); + } + } + + @Test + public void flushAndWaitReportsFailureOnceClosed() throws IOException { + try (MockWebServer mockEventsServer = new MockWebServer()) { + mockEventsServer.start(); + + LDConfig ldConfig = baseConfigBuilder(mockEventsServer).build(); + LDClient client = LDClient.init(application, ldConfig, ldContext, 0); + client.close(); + + assertFalse(client.flushAndWait(5, TimeUnit.SECONDS)); + } + } + + @Test + public void flushAndWaitDeliversEveryEnvironmentAtOnce() throws Exception { + // Each post is held for longer than half the budget, so the call can only report true if the + // environments' deliveries ran side by side rather than one after the other. + try (MockWebServer mockEventsServer = new MockWebServer()) { + mockEventsServer.setDispatcher(new Dispatcher() { + @Override + public MockResponse dispatch(RecordedRequest request) { + return new MockResponse().setHeadersDelay(1_500, TimeUnit.MILLISECONDS); + } + }); + mockEventsServer.start(); + + Map secondaryKeys = new HashMap<>(); + secondaryKeys.put("second", "second-mobile-key"); + LDConfig ldConfig = baseConfigBuilder(mockEventsServer) + .secondaryMobileKeys(secondaryKeys) + .build(); + try (LDClient client = LDClient.init(application, ldConfig, ldContext, 0)) { + client.track("primary-event"); + LDClient.getForMobileKey("second").track("second-event"); + + assertTrue(client.flushAndWait(2_500, TimeUnit.MILLISECONDS)); + assertEquals(2, mockEventsServer.getRequestCount()); + } + } + } + + @Test + public void flushAndWaitWithTheMostNegativeTimeoutDoesNotWait() throws IOException { + // toNanos saturates at Long.MIN_VALUE, and unclamped that wraps round to a wait for as long + // as the delivery takes -- which would then be reported as delivered. + try (MockWebServer mockEventsServer = new MockWebServer()) { + mockEventsServer.start(); + mockEventsServer.enqueue(new MockResponse().setHeadersDelay(3, TimeUnit.SECONDS)); + + LDConfig ldConfig = baseConfigBuilder(mockEventsServer).build(); + try (LDClient client = LDClient.init(application, ldConfig, ldContext, 0)) { + client.track("test-event"); + + long started = System.nanoTime(); + assertFalse(client.flushAndWait(Long.MIN_VALUE, TimeUnit.NANOSECONDS)); + long waitedMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - started); + assertTrue("waited " + waitedMillis + "ms", waitedMillis < 1_000); + } + } + } + + @Test + public void flushAndWaitReportsFailureWhenTheDeliveryIsCanceled() throws IOException { + // A custom event processor may hand back a future that is canceled; Future.get then throws + // CancellationException, which is unchecked and must not escape a boolean answer. + try (MockWebServer mockEventsServer = new MockWebServer()) { + mockEventsServer.start(); + + LDConfig ldConfig = baseConfigBuilder(mockEventsServer) + .events(clientContext -> new CancelingEventProcessor()) + .build(); + try (LDClient client = LDClient.init(application, ldConfig, ldContext, 0)) { + assertFalse(client.flushAndWait(5, TimeUnit.SECONDS)); + } + } + } + + /** An event processor whose deliveries are always canceled before they can report. */ + private static final class CancelingEventProcessor implements EventProcessor { + @Override + public Future flushAsync() { + FutureTask delivery = new FutureTask<>(() -> true); + delivery.cancel(false); + return delivery; + } + + @Override + public void flush() {} + + @Override + public void blockingFlush() {} + + @Override + public void setInBackground(boolean inBackground) {} + + @Override + public void setOffline(boolean offline) {} + + @Override + public void close() {} + + @Override + public void recordEvaluationEvent(LDContext context, String flagKey, int flagVersion, + int variation, LDValue value, EvaluationReason reason, + LDValue defaultValue, boolean requireFullEvent, + Long debugEventsUntilDate) {} + + @Override + public void recordIdentifyEvent(LDContext context) {} + + @Override + public void recordCustomEvent(LDContext context, String eventKey, LDValue data, + Double metricValue) {} + } + @Test public void testTrackDataValueNull() throws IOException, InterruptedException { try (MockWebServer mockEventsServer = new MockWebServer()) { diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ComponentsImpl.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ComponentsImpl.java index 890e45dd..e996ab43 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ComponentsImpl.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ComponentsImpl.java @@ -25,6 +25,7 @@ import java.util.HashMap; import java.util.Map; +import java.util.concurrent.Future; /** * This class contains the package-private implementations of component factories and builders whose @@ -72,6 +73,12 @@ public void flush() {} @Override public void blockingFlush() {} + @Override + public Future flushAsync() { + // Nothing was recorded, so there is nothing undelivered to warn the caller about. + return new LDSuccessFuture<>(true); + } + @Override public void setInBackground(boolean inBackground) {} diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java index 147a9fb8..89d0da48 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java @@ -50,7 +50,7 @@ final class DirectEventProcessor implements EventProcessor { * trips, which two seconds covers up to roughly a 600ms RTT. A network slower than that is one the * post is likely to fail on anyway. *

- * Overshooting the budget is cheaper than it looks, because the delivery is not cancelled when the + * Overshooting the budget is cheaper than it looks, because the delivery is not canceled when the * budget expires; see {@link #close()}. */ static final long DEFAULT_CLOSE_BUDGET_MILLIS = 2_000; @@ -133,6 +133,37 @@ final class DirectEventProcessor implements EventProcessor { /** Set under {@link #submitLock} once close() has queued the release of the sender. */ private boolean shuttingDown = false; + /** + * Guards {@link #pendingFlush} and {@link #pendingFlushAnswersACaller}. Taken on a caller's + * thread and on the delivery thread, never while holding {@link #recordLock}, and nothing + * blocking happens under it. + */ + private final Object flushLock = new Object(); + + /** + * The delivery that is queued but has not started, which a flush request arriving now can wait + * on instead of queueing another. Null while nothing is queued, and cleared again as the queued + * delivery begins, which is the point past which it can no longer speak for what is recorded. + */ + private LDAwaitFuture pendingFlush; + + /** + * Whether a {@link #flushAsync()} caller, who is there to hear the outcome, is waiting on + * {@link #pendingFlush}, as opposed to only callers that discard it. + */ + private boolean pendingFlushAnswersACaller; + + /** + * Set when a delivery took events out of the buffer and did not get all of them to the service, + * and cleared once a {@link #flushAsync()} caller has been told so. + *

+ * Without it, a delivery that finds the buffer empty reports success even when the events it is + * being asked about were taken a moment earlier by another delivery that then lost them: the + * periodic flush, or a flush whose caller did not wait for the outcome. Only touched on the + * scheduler thread, which every delivery runs on. + */ + private boolean eventsLostSinceLastAnswer; + DirectEventProcessor( OutboundEventBuffer buffer, EventSender eventSender, @@ -334,7 +365,7 @@ public void flush() { if (isStopped()) { return; } - submit(this::deliverPayload); + queueDelivery(false); } @Override @@ -342,10 +373,7 @@ public void blockingFlush() { if (isStopped()) { return; } - Future delivery = submit(this::deliverPayload); - if (delivery == null) { - return; - } + Future delivery = queueDelivery(false); try { delivery.get(); } catch (InterruptedException e) { @@ -355,6 +383,78 @@ public void blockingFlush() { } } + @Override + public Future flushAsync() { + if (isStopped()) { + return new LDSuccessFuture<>(false); + } + return queueDelivery(true); + } + + /** + * Queues a delivery, or hands back one that is already queued and has not started. + *

+ * A delivery that has not started yet will send everything recorded before it starts. That + * includes the caller's events, so waiting on it is as good as queuing a new one. Without this, + * flushes arriving faster than a post completes each queue their own, and the one that matters -- + * the {@code flushAndWait} at shutdown -- waits behind all of them. + * + * @param answersACaller true if the caller will hear the outcome, so that the delivery reports + * any events lost since the last answer, and false if the caller discards it + */ + private Future queueDelivery(boolean answersACaller) { + synchronized (flushLock) { + if (pendingFlush != null) { + pendingFlushAnswersACaller |= answersACaller; + return pendingFlush; + } + LDAwaitFuture result = new LDAwaitFuture<>(); + if (submit(() -> runDelivery(result)) == null) { + // Shutting down, so there is no thread left to deliver on and nothing will be sent. + return new LDSuccessFuture<>(false); + } + pendingFlush = result; + pendingFlushAnswersACaller = answersACaller; + return result; + } + } + + /** + * Runs one delivery on behalf of every flush request that joined it, and tells them all how it + * went. + *

+ * Deliveries run one at a time, so any delivery that was in flight when a caller asked has + * finished before this one starts, and has already recorded whether it lost what it took. A + * caller who will hear the answer is told no if anything was lost since the last caller was + * told, as well as if this delivery fails: an empty buffer is not evidence that the events + * which used to be in it arrived. + */ + private void runDelivery(LDAwaitFuture result) { + boolean answersACaller = false; + synchronized (flushLock) { + // Requests arriving from here on need a delivery of their own: this one is about to take + // the buffer, and what it takes is all it can speak for. + if (pendingFlush == result) { + pendingFlush = null; + answersACaller = pendingFlushAnswersACaller; + pendingFlushAnswersACaller = false; + } + } + boolean delivered = false; + try { + delivered = deliverPayloadReportingOutcome(); + } catch (Throwable t) { + // Caught here rather than left to guarded(), because a caller is waiting on the future + // and completing it matters more than the stack reaching the executor. + logUnexpectedError(t); + } + if (answersACaller) { + delivered &= !eventsLostSinceLastAnswer; + eventsLostSinceLastAnswer = false; + } + result.set(delivered); + } + @Override public void close() throws IOException { if (!closed.compareAndSet(false, true)) { @@ -368,22 +468,23 @@ public void close() throws IOException { // once the processor is gone. While offline that chance is not taken, and whatever is held // is discarded. Offline is the application telling the SDK to stay off the network, and // shutting down does not revoke that. - Future delivery = submit(this::deliverPayload); - if (delivery != null) { - try { - delivery.get(closeBudgetMillis, TimeUnit.MILLISECONDS); - } catch (TimeoutException e) { - // Deliberately not cancelled. The run has already been drained into a payload, so - // interrupting now would make the loss certain, while leaving it to run costs - // nothing: the scheduler thread is a daemon, and returning from close() does not - // end an Android process. The budget bounds the caller, not the delivery. - logger.warn("Gave up waiting for the final event delivery after {}ms;" + - " it continues in the background", closeBudgetMillis); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } catch (ExecutionException e) { - logUnexpectedError(e.getCause() == null ? e : e.getCause()); - } + // + // Queued directly rather than through flushAsync(), which refuses once closed is set, but + // through the same coalescing: a delivery that has not started yet will take these events + // too, so there is no reason to queue a second one behind it. + try { + queueDelivery(false).get(closeBudgetMillis, TimeUnit.MILLISECONDS); + } catch (TimeoutException e) { + // Deliberately not canceled. The run has already been drained into a payload, so + // interrupting now would make the loss certain, while leaving it to run costs + // nothing: the scheduler thread is a daemon, and returning from close() does not + // end an Android process. The budget bounds the caller, not the delivery. + logger.warn("Gave up waiting for the final event delivery after {}ms;" + + " it continues in the background", closeBudgetMillis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } catch (ExecutionException e) { + logUnexpectedError(e.getCause() == null ? e : e.getCause()); } // Queued on both of the threads that post through the sender, so that it is released by // whichever of them finishes last. Closing it here instead would pull the HTTP client out @@ -425,17 +526,30 @@ private void releaseSenderWhenLast() { } /** - * Serializes and sends everything buffered. Runs on the scheduler thread, which is - * single-threaded, so only one payload is ever in flight and the run is taken exactly once per - * delivery. - *

- * The run and the counters are taken together under {@link #recordLock}, so an evaluation is - * never split across two payloads, and encoded outside it, so recording does not wait on the - * encoder. + * Serializes and sends everything buffered, for the periodic flush, which has nobody waiting to + * find out how it went. It is a fixed-delay series, so a run is only ever scheduled once the one + * before it has finished and these cannot pile up the way requested flushes could. */ private void deliverPayload() { + deliverPayloadReportingOutcome(); + } + + /** + * Delivers as {@link #deliverPayload()} does, and says whether it worked, for the callers of a + * requested flush, who are waiting to find out. + *

+ * Runs on the scheduler thread, which is single-threaded, so only one payload is ever in flight + * and the run is taken exactly once per delivery. The run and the counters are taken together + * under {@link #recordLock}, so an evaluation is never split across two payloads, and encoded + * outside it, so recording does not wait on the encoder. + * + * @return true if the events reached the service, or if there were none to send; false if they + * could not be sent, the service did not accept them, or some could not be serialized + */ + private boolean deliverPayloadReportingOutcome() { if (disabled || offline.get()) { - return; + // Nothing is taken, so nothing is lost: the events stay buffered for a later delivery. + return false; } List run; List summaries; @@ -445,24 +559,47 @@ private void deliverPayload() { summaries = buffer.takeSummaries(); summaryContextsExceeded.set(false); } + + boolean delivered = false; + try { + delivered = deliverTaken(run, summaries); + } finally { + // From here the events exist only in this delivery, so not delivering them loses them, + // including when something unexpected is thrown on the way. + if (!delivered) { + eventsLostSinceLastAnswer = true; + } + } + return delivered; + } + + private boolean deliverTaken(List run, List summaries) { OutboundEventBuffer.Payload payload; try { payload = buffer.encode(run, summaries); } catch (IOException e) { logUnexpectedError(e); - return; + return false; } if (payload == null) { - return; + return true; } + if (payload.getEventCount() == 0) { + return false; // everything taken was dropped as unserializable + } + if (diagnosticStore != null) { diagnosticStore.recordEventsInBatch(payload.getEventCount()); } + try { - handleResponse(eventSender.sendAnalyticsEvents(payload.getData(), - payload.getEventCount(), eventsUri)); + EventSender.Result result = eventSender.sendAnalyticsEvents(payload.getData(), + payload.getEventCount(), eventsUri); + handleResponse(result); + return result != null && result.isSuccess() && payload.isComplete(); } catch (Exception e) { logUnexpectedError(e); + return false; } } @@ -550,7 +687,7 @@ private void postDiagnostic(Runnable post) { * Unlike analytics events, diagnostics are not sent while offline or in the background. *

* {@link #updateScheduledTasks} cancels the periodic task when either becomes true, but that is - * not enough on its own. Cancelling does not stop a run already underway, and the init event is + * not enough on its own. Canceling does not stop a run already underway, and the init event is * submitted before it reaches the executor. Either can arrive here after the state changed. */ private boolean diagnosticsSuspended() { @@ -601,12 +738,12 @@ private boolean shouldDebugEvent(Long debugEventsUntilDate) { /** * Must be called holding {@code stateLock}. Once closed, this only ever cancels: close() sets the * flag and then calls this under the same lock, so a call that got here first has its tasks - * cancelled by close(), and any call after it finds the flag set. + * canceled by close(), and any call after it finds the flag set. */ private void updateScheduledTasks(boolean inBackground, boolean offline) { boolean stopped = closed.get(); - // Flushing stays scheduled whether or not we are offline or in the background; a run while - // offline returns without doing anything. Cancelling it for an outage would restart the + // Flushing stays scheduled even while we are offline or in the background; a run while + // offline returns without doing anything. Canceling it for an outage would restart the // interval on every reconnect, and a run of brief outages would then hold events back for // far longer than one interval. Left running, what an outage buffered goes out at the first // run after it ends. diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClient.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClient.java index f224668a..20537a6c 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClient.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClient.java @@ -37,6 +37,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CancellationException; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; @@ -779,6 +780,55 @@ private void flushInternal() { eventProcessor.flush(); } + @Override + public boolean flushAndWait(long timeout, TimeUnit unit) { + // Clamped because toNanos saturates: a timeout at Long.MIN_VALUE nanos would make the + // remaining time below underflow, and wrap round to a wait with no bound at all. + long deadline = System.nanoTime() + Math.max(0, unit.toNanos(timeout)); + Map clients = getInstancesIfTheyIncludeThisClient(); + if (clients.isEmpty()) { + // This client has been closed, or replaced by a later init; either way it can deliver + // nothing, and saying otherwise would tell the caller its events were safe. + return false; + } + // Every environment is started before any of them is waited on. Each has its own event + // processor and its own thread, so waiting on one before starting the next would spend the + // caller's budget on deliveries that could have been running all along. + List> deliveries = new ArrayList<>(clients.size()); + for (LDClient client : clients.values()) { + deliveries.add(client.eventProcessor.flushAsync()); + } + boolean delivered = true; + for (Future delivery : deliveries) { + // Each wait gets what is left of the one budget rather than a fresh copy of it, so that + // the timeout the caller asked for is the time this call can take. + delivered &= awaitDelivery(delivery, Math.max(0, deadline - System.nanoTime())); + } + return delivered; + } + + private boolean awaitDelivery(Future delivery, long remainingNanos) { + try { + return Boolean.TRUE.equals(delivery.get(remainingNanos, TimeUnit.NANOSECONDS)); + } catch (TimeoutException e) { + // Left running rather than canceled: the events have been taken out of the buffer by + // now, so interrupting the delivery would only make losing them certain. + return false; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return false; + } catch (CancellationException e) { + // Not something the SDK's own processor does, but a custom one can hand back a future + // that is canceled, and that must not escape a call whose answer is a boolean. + return false; + } catch (ExecutionException e) { + Throwable cause = e.getCause() == null ? e : e.getCause(); + logger.error("Exception caught when flushing events: {}", LogValues.exceptionSummary(cause)); + logger.debug("{}", LogValues.exceptionTrace(cause)); + return false; + } + } + @VisibleForTesting void blockingFlush() { eventProcessor.blockingFlush(); diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClientInterface.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClientInterface.java index dd62ec5c..f01df072 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClientInterface.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClientInterface.java @@ -11,6 +11,7 @@ import java.io.Closeable; import java.util.Map; import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; /** * The interface for the LaunchDarkly SDK client. @@ -146,6 +147,48 @@ public interface LDClientInterface extends Closeable { */ void flush(); + /** + * Sends all pending events to LaunchDarkly and waits for them to be delivered. + *

+ * Unlike {@link #flush()}, which returns before the events reach the network, this reports + * whether they arrived, which is what makes it usable at a point where the application is about + * to lose the ability to send them: an uncaught exception handler, a move to the background, or + * any other last chance. Events buffered in memory do not survive the process, so a caller that + * knows the process is ending can use this to give them one. + * {@link com.launchdarkly.sdk.android.integrations.LDCrashHandler} does this for + * uncaught exceptions. + *

+ * It can only help while the process is still running code. An uncaught exception runs its + * handler first, and a move to the background is announced, so both leave time for this call. + * A {@code SIGKILL}, an ANR kill, a native crash, and the system reclaiming a backgrounded process + * run nothing at all, and events still in memory at that moment are lost whatever the + * application does. + *

+ * The timeout bounds the whole call, including when the SDK is configured for more than one + * environment. Choose it with the caller in mind: a dying process is not a good place to wait on + * a network request that may never answer. The call blocks the thread it is made on, so on the + * main thread the timeout also counts towards an ANR. On a network that does not answer, a + * delivery takes about 21 seconds to give up with the default HTTP configuration: two attempts a + * second apart, each allowed the timeout set by + * {@link com.launchdarkly.sdk.android.integrations.HttpConfigurationBuilder#connectTimeoutMillis(int)}. + * A shorter timeout returns {@code false} before then, so there it bounds the wait rather than + * reporting how the delivery went. + * + * @param timeout how long to wait for delivery + * @param unit the time unit of {@code timeout} + * @return true if the events were delivered, or there were none to deliver; false if the timeout + * expired first, the SDK is offline, closed, or otherwise unable to deliver them, or events + * recorded since the last time this was answered were lost on the way, by this delivery or an + * earlier one. Events the service refuses, even with an error that may pass such as a 503, are + * retried once straight away and then lost: they are not kept for a later flush, so calling + * this again does not resend them. A {@code false} because the timeout expired does not mean + * the events were not sent: the delivery is left running when the caller stops waiting, and + * may still arrive if the process lives long enough. A caller that resends on {@code false} + * can therefore cause duplicates. + * @since 5.17.0 + */ + boolean flushAndWait(long timeout, TimeUnit unit); + /** * Returns a map of all feature flags for the current evaluation context. No events are sent to LaunchDarkly. * diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDFutures.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDFutures.java index c062d72e..8bed2024 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDFutures.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDFutures.java @@ -4,6 +4,7 @@ import java.util.ArrayList; import java.util.List; +import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -83,6 +84,28 @@ public static LDAwaitFuture fromFuture(Future future) { return result; } + /** + * Runs a blocking call on a pooled daemon thread and reports its result as a future. + *

+ * Use this where a caller has a deadline but the work it is waiting for has no way to take one. + * The call is left running if the caller stops waiting; nothing interrupts it. + * + * @param task the blocking call + * @param result type + * @return a future that completes with the call's result, or with whatever it threw + */ + public static Future fromBlockingCall(Callable task) { + LDAwaitFuture result = new LDAwaitFuture<>(); + getBridgeExecutor().execute(() -> { + try { + result.set(task.call()); + } catch (Throwable t) { + result.setException(t); + } + }); + return result; + } + /** * Returns a future that completes when the first of the given futures completes. * Equivalent to CompletableFuture.anyOf. Works with any {@link Future} (API-level safe). diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/OutboundEventBuffer.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/OutboundEventBuffer.java index 717dfa36..e1730568 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/OutboundEventBuffer.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/OutboundEventBuffer.java @@ -42,6 +42,7 @@ final class OutboundEventBuffer { private static final int INITIAL_OUTPUT_BUFFER_SIZE = 2000; private static final Event[] NO_EVENTS = new Event[0]; + private static final byte[] NO_DATA = new byte[0]; private static final List NO_SUMMARIES = Collections.emptyList(); private final EventOutputFormatter formatter; @@ -164,7 +165,8 @@ synchronized List takeSummaries() { * * @param run the full events to send, in the order they were recorded * @param summaries the counters taken alongside that run - * @return the payload to send, or null if there was nothing to send + * @return the payload to send, or null if there was nothing to send; a payload that had to drop + * something says so through {@link Payload#isComplete()}, and may then hold no events at all * @throws IOException if the events could not be serialized */ Payload encode(List run, List summaries) throws IOException { @@ -188,43 +190,55 @@ private Payload encodeAll(List run, List su if (outputEventCount == 0) { return null; } - return new Payload(buffer.toByteArray(), outputEventCount); + return new Payload(buffer.toByteArray(), outputEventCount, true); } private Payload encodeSkippingFailures(List run, List summaries) { List objects = new ArrayList<>(); int outputEventCount = 0; + boolean dropped = false; for (Event event : run) { EncodedPiece piece = tryEncode(new Event[] { event }, NO_SUMMARIES); if (piece == null) { logger.error("Dropping unserializable event of type {}", event.getClass().getSimpleName()); + dropped = true; continue; } - objects.add(piece.jsonObject); - outputEventCount += piece.eventCount; + if (piece != EncodedPiece.NOTHING) { + objects.add(piece.jsonObject); + outputEventCount += piece.eventCount; + } } for (EventSummarizer.EventSummary summary : summaries) { EncodedPiece piece = tryEncode(NO_EVENTS, Collections.singletonList(summary)); if (piece == null) { logger.error("Dropping unserializable summary event"); + dropped = true; continue; } - objects.add(piece.jsonObject); - outputEventCount += piece.eventCount; + if (piece != EncodedPiece.NOTHING) { + objects.add(piece.jsonObject); + outputEventCount += piece.eventCount; + } } if (objects.isEmpty()) { - return null; + return dropped ? new Payload(NO_DATA, 0, false) : null; } - return new Payload(joinObjects(objects), outputEventCount); + return new Payload(joinObjects(objects), outputEventCount, !dropped); } + /** + * @return the piece, {@link EncodedPiece#NOTHING} if the formatter had nothing to write for it, + * or null if it could not be serialized + */ private EncodedPiece tryEncode(Event[] events, List summaries) { try { ByteArrayOutputStream buffer = new ByteArrayOutputStream(INITIAL_OUTPUT_BUFFER_SIZE); int count = write(events, summaries, buffer); if (count == 0) { - return null; + // An empty summary, which the formatter skips. Nothing is lost by leaving it out. + return EncodedPiece.NOTHING; } byte[] jsonObject = objectFromArray(buffer.toByteArray()); if (jsonObject == null) { @@ -303,6 +317,8 @@ private static byte[] joinObjects(List objects) { } private static final class EncodedPiece { + static final EncodedPiece NOTHING = new EncodedPiece(NO_DATA, 0); + final byte[] jsonObject; final int eventCount; @@ -318,10 +334,12 @@ private static final class EncodedPiece { static final class Payload { private final byte[] data; private final int eventCount; + private final boolean complete; - Payload(byte[] data, int eventCount) { + Payload(byte[] data, int eventCount, boolean complete) { this.data = data; this.eventCount = eventCount; + this.complete = complete; } /** @@ -337,5 +355,13 @@ byte[] getData() { int getEventCount() { return eventCount; } + + /** + * @return false if something handed to the encoder could not be serialized and was dropped, + * so that even a successful post of this body leaves those events undelivered + */ + boolean isComplete() { + return complete; + } } } diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/integrations/LDCrashHandler.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/integrations/LDCrashHandler.java new file mode 100644 index 00000000..98e2bbb2 --- /dev/null +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/integrations/LDCrashHandler.java @@ -0,0 +1,126 @@ +package com.launchdarkly.sdk.android.integrations; + +import androidx.annotation.NonNull; +import androidx.annotation.Nullable; +import androidx.annotation.VisibleForTesting; + +import com.launchdarkly.sdk.android.LDClient; +import com.launchdarkly.sdk.android.LDClientInterface; +import com.launchdarkly.sdk.android.LaunchDarklyException; + +import java.util.Objects; +import java.util.concurrent.Callable; +import java.util.concurrent.TimeUnit; + +/** + * An uncaught exception handler that gives the SDK its last chance to act before the process ends: + * it delivers the analytics events still in memory. + *

+ * Events are kept in memory until they are sent, so a crash loses the ones recorded since the last + * flush, and those are often the ones that explain it. An uncaught exception runs its handlers while + * the process is still alive, and this one uses that time to call + * {@link LDClientInterface#flushAndWait(long, TimeUnit)} for every configured environment before + * passing the exception on to the handler that was installed before it. + *

+ * Install it before your crash reporter: + * + *


+ *     LDClient.init(application, config, context, 0);
+ *     LDCrashHandler.install(2, TimeUnit.SECONDS);
+ *     SentryAndroid.init(application, options -> { ... });
+ * 
+ *

+ * Crash reporters install their own handler in front of the one they find and call it once they + * have saved their report, so installing this one first lets the report be saved before the events + * are sent. Installed after the crash reporter, this handler runs first and holds the report back + * for as long as the delivery takes, and a process that dies in that time loses the report as well. + *

+ * The timeout is how long the crash is held open for the events. The thread that threw waits for + * it, so when that is the main thread the application stays frozen until the events are delivered + * or the timeout expires; a couple of seconds is a reasonable budget. That is less than a delivery + * takes to give up on a network that does not answer, so on such a network the crash is passed on + * when the timeout expires. A delivery still running then is not canceled, and may finish while the + * rest of the handlers run. + *

+ * This handler only helps with uncaught Java and Kotlin exceptions. A {@code SIGKILL}, an ANR kill, + * a native crash, and the system reclaiming a backgrounded process run no handlers, so the events + * still in memory at that moment are lost. + *

+ * This class is not stable, and not subject to any backwards compatibility guarantees or semantic versioning. + * It is experimental. + * + * @since 5.17.0 + */ +public final class LDCrashHandler implements Thread.UncaughtExceptionHandler { + private static final Object installLock = new Object(); + + private final Callable client; + private final long timeoutMillis; + @Nullable + private final Thread.UncaughtExceptionHandler next; + + @VisibleForTesting + LDCrashHandler( + @NonNull Callable client, + long timeoutMillis, + @Nullable Thread.UncaughtExceptionHandler next + ) { + this.client = client; + this.timeoutMillis = timeoutMillis; + this.next = next; + } + + /** + * Makes this handler the default uncaught exception handler, in front of whichever handler was + * the default before. + *

+ * It can be called before {@link LDClient#init}: a crash that happens before the client exists + * has no events to deliver, and is passed on straight away. Calling it again while this handler + * is still the default does nothing, so the timeout given first is the one that applies. + * + * @param timeout how long a crash waits for the events to be delivered + * @param unit the time unit of {@code timeout} + */ + public static void install(long timeout, @NonNull TimeUnit unit) { + Objects.requireNonNull(unit, "unit"); + long timeoutMillis = Math.max(0, unit.toMillis(timeout)); + synchronized (installLock) { + Thread.UncaughtExceptionHandler previous = Thread.getDefaultUncaughtExceptionHandler(); + if (previous instanceof LDCrashHandler) { + return; + } + Thread.setDefaultUncaughtExceptionHandler( + new LDCrashHandler(LDClient::get, timeoutMillis, previous)); + } + } + + @Override + public void uncaughtException(@NonNull Thread thread, @NonNull Throwable throwable) { + try { + flush(); + } finally { + if (next != null) { + next.uncaughtException(thread, throwable); + } else { + throwable.printStackTrace(); + } + } + } + + private void flush() { + try { + LDClientInterface ldClient = client.call(); + if (ldClient == null) { + return; + } + // Waiting on the crashing thread is safe because the SDK enforces the timeout instead of + // trusting the delivery to finish, which also covers a crash that is the reason the + // delivery cannot complete. + ldClient.flushAndWait(timeoutMillis, TimeUnit.MILLISECONDS); + } catch (LaunchDarklyException notInitialized) { + // A crash before LDClient.init has no events to deliver. + } catch (Throwable ignored) { + // Nothing that goes wrong here is worth losing the crash over. + } + } +} diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/subsystems/EventProcessor.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/subsystems/EventProcessor.java index e65a8e2b..4f8e8ce0 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/subsystems/EventProcessor.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/subsystems/EventProcessor.java @@ -4,8 +4,10 @@ import com.launchdarkly.sdk.EvaluationReason; import com.launchdarkly.sdk.LDContext; import com.launchdarkly.sdk.LDValue; +import com.launchdarkly.sdk.android.LDFutures; import java.io.Closeable; +import java.util.concurrent.Future; /** * Interface for an object that can send or store analytics events. @@ -99,4 +101,30 @@ void recordCustomEvent( * Specifies that any buffered events should be sent immediately, blocking until done. */ void blockingFlush(); + + /** + * Specifies that any buffered events should be sent immediately, and reports through the + * returned future whether they were delivered. + *

+ * This is the form the SDK itself uses, so that a caller with a deadline can wait for as long as + * it has and no longer, and so that several of these can be waited on together. Only the public + * API puts a timeout in a signature; see {@code LDClient.flushAndWait}. + *

+ * Nothing cancels the delivery when a caller stops waiting for it: by then the events have been + * taken out of the buffer, so interrupting the post would only make losing them certain. + * + * @return a future that completes with true if the events reached the service, or there were + * none to send; false if they could not be sent, including when a delivery that started + * earlier took them and then lost them + * @since 5.17.0 + */ + default Future flushAsync() { + // An implementation written before this method existed has only its unbounded blocking + // flush, so that runs on a thread of its own: the caller's deadline then bounds the wait + // rather than the flush, and the outcome it reports is still the flush's own. + return LDFutures.fromBlockingCall(() -> { + blockingFlush(); + return true; + }); + } } diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java index a4acfd39..2b7d039d 100644 --- a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java @@ -4,6 +4,7 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -13,6 +14,7 @@ import com.launchdarkly.sdk.internal.events.DiagnosticStore; import com.launchdarkly.sdk.internal.events.Event; import com.launchdarkly.sdk.internal.events.EventSender; +import com.launchdarkly.testhelpers.httptest.HandlerSwitcher; import com.launchdarkly.testhelpers.httptest.Handlers; import com.launchdarkly.testhelpers.httptest.HttpServer; import com.launchdarkly.testhelpers.httptest.RequestInfo; @@ -30,10 +32,13 @@ import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeoutException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -651,6 +656,22 @@ public void beingToldToShutDownStopsRecordingAndDelivery() throws Exception { } } + @Test + public void flushWithTimeoutReportsDeliveredEvents() throws Exception { + try (HttpServer server = startEventsServer()) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "an-event", LDValue.ofNull(), null); + + assertTrue(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + + assertEquals(1, countEventsOfKind(collectDelivered(server), "custom")); + } finally { + eventProcessor.close(); + } + } + } + @Test public void aFlushNeverSplitsAnEvaluationAcrossTwoPayloads() throws Exception { // The other half of the atomicity invariant. close() only ever delivers once, so it can show @@ -776,6 +797,21 @@ public Result sendAnalyticsEvents(byte[] data, int eventCount, URI eventsBaseUri } } + @Test + public void flushWithTimeoutReportsSuccessWhenThereIsNothingToSend() throws Exception { + try (HttpServer server = startEventsServer()) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + // Nothing was recorded, so the caller's events are not waiting anywhere. + assertTrue(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + + server.getRecorder().requireNoRequests(100, TimeUnit.MILLISECONDS); + } finally { + eventProcessor.close(); + } + } + } + @Test public void unexpectedRecordingErrorDoesNotBubbleToCallerAndLogs() throws Exception { ScheduledExecutorService scheduler = EventUtil.makeEventsTaskExecutor(); @@ -858,6 +894,290 @@ public Result sendAnalyticsEvents(byte[] data, int eventCount, URI eventsBaseUri } } + @Test + public void flushWithTimeoutReportsFailureWhileOffline() throws Exception { + try (HttpServer server = startEventsServer()) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.setOffline(true); + eventProcessor.recordCustomEvent(CONTEXT, "an-event", LDValue.ofNull(), null); + + // The events are still buffered rather than delivered, and no amount of waiting + // changes that, so the caller is told so instead of being told they are safe. + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + + server.getRecorder().requireNoRequests(100, TimeUnit.MILLISECONDS); + } finally { + eventProcessor.close(); + } + } + } + + @Test + public void flushesArrivingWhileADeliveryRunsShareOneFollowUpDelivery() throws Exception { + // Otherwise a flush called faster than a post completes queues a post per call, and the + // flush that matters -- the one at shutdown, with a deadline -- waits behind all of them. + CountDownLatch firstSendStarted = new CountDownLatch(1); + CountDownLatch releaseFirstSend = new CountDownLatch(1); + AtomicInteger sends = new AtomicInteger(0); + EventSender sender = new StubEventSender() { + @Override + public Result sendAnalyticsEvents(byte[] data, int eventCount, URI eventsBaseUri) { + if (sends.incrementAndGet() == 1) { + firstSendStarted.countDown(); + awaitQuietly(releaseFirstSend, 5, TimeUnit.SECONDS); + } + return new Result(true, false, null); + } + }; + + ScheduledExecutorService scheduler = EventUtil.makeEventsTaskExecutor(); + DirectEventProcessor eventProcessor = makeEventProcessor(sender, NO_PERIODIC_FLUSH_MILLIS, + scheduler); + try { + eventProcessor.setOffline(false); + eventProcessor.recordCustomEvent(CONTEXT, "first", LDValue.ofNull(), null); + Future first = eventProcessor.flushAsync(); + assertTrue("the first delivery never started", + firstSendStarted.await(2, TimeUnit.SECONDS)); + + // The delivery thread is inside that post, so none of these can start, and each of them + // has to be answered by the one delivery that is queued behind it. + eventProcessor.recordCustomEvent(CONTEXT, "second", LDValue.ofNull(), null); + Future queued = eventProcessor.flushAsync(); + for (int i = 0; i < 50; i++) { + assertSame(queued, eventProcessor.flushAsync()); + } + + releaseFirstSend.countDown(); + assertTrue(first.get(5, TimeUnit.SECONDS)); + assertTrue(queued.get(5, TimeUnit.SECONDS)); + + assertEquals("one post for the running delivery and one for the 51 that joined", + 2, sends.get()); + } finally { + releaseFirstSend.countDown(); + eventProcessor.close(); + scheduler.shutdownNow(); + } + } + + @Test + public void aFlushJoiningADeliveryStillCoversWhatTheCallerRecorded() throws Exception { + // Joining is only sound while the delivery it joins has not taken the buffer yet, so what + // the joining caller recorded has to come back in that delivery's payload. + Semaphore letFirstResponseFinish = new Semaphore(0); + try (HttpServer server = HttpServer.start(Handlers.sequential( + Handlers.all(Handlers.waitFor(letFirstResponseFinish), Handlers.status(202)), + Handlers.status(202)))) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "first", LDValue.ofNull(), null); + Future first = eventProcessor.flushAsync(); + server.getRecorder().requireRequest(5, TimeUnit.SECONDS); + + eventProcessor.recordCustomEvent(CONTEXT, "joined", LDValue.ofNull(), null); + Future queued = eventProcessor.flushAsync(); + + letFirstResponseFinish.release(1); + assertTrue(first.get(5, TimeUnit.SECONDS)); + assertTrue(queued.get(5, TimeUnit.SECONDS)); + + RequestInfo second = server.getRecorder().requireRequest(5, TimeUnit.SECONDS); + assertTrue("the joining caller's event was left behind", + second.getBody().contains("\"key\":\"joined\"")); + } finally { + letFirstResponseFinish.release(Integer.MAX_VALUE); + eventProcessor.close(); + } + } + } + + @Test + public void flushWithTimeoutReportsFailureWhenTheTimeoutExpiresFirst() throws Exception { + Semaphore letResponseFinish = new Semaphore(0); + try (HttpServer server = HttpServer.start(Handlers.all(Handlers.waitFor(letResponseFinish), + Handlers.status(202)))) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "an-event", LDValue.ofNull(), null); + Future delivery = eventProcessor.flushAsync(); + + try { + delivery.get(100, TimeUnit.MILLISECONDS); + fail("the delivery finished while its response was still being held"); + } catch (TimeoutException expected) { + // The caller gives up here, as flushAndWait does when its budget runs out. + } + + // Left running rather than canceled, the delivery still gets its events through + // once the service answers, and a canceled one could not report that. + letResponseFinish.release(1); + assertTrue(delivery.get(5, TimeUnit.SECONDS)); + server.getRecorder().requireRequest(1, TimeUnit.SECONDS); + } finally { + // Released before closing, so that the delivery still in flight can finish rather + // than hold up the shutdown that close() waits on. + letResponseFinish.drainPermits(); + letResponseFinish.release(Integer.MAX_VALUE); + eventProcessor.close(); + } + } + } + + @Test + public void flushReportsFailureWhenTheServiceRefusesTheEventsForNow() throws Exception { + // 503 is a failure that may pass, so the sender retries it once before giving up. + try (HttpServer server = HttpServer.start(Handlers.status(503))) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "an-event", LDValue.ofNull(), null); + + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + + server.getRecorder().requireRequest(1, TimeUnit.SECONDS); + server.getRecorder().requireRequest(1, TimeUnit.SECONDS); + } finally { + eventProcessor.close(); + } + } + } + + @Test + public void flushReportsFailureWhenTheServiceRefusesTheEventsForGood() throws Exception { + // 401 is not retried, and stops the processor for the life of the process, so the flush + // after it has nothing it can deliver on either. + try (HttpServer server = HttpServer.start(Handlers.status(401))) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "refused", LDValue.ofNull(), null); + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + server.getRecorder().requireRequest(1, TimeUnit.SECONDS); + + eventProcessor.recordCustomEvent(CONTEXT, "after", LDValue.ofNull(), null); + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + server.getRecorder().requireNoRequests(200, TimeUnit.MILLISECONDS); + } finally { + eventProcessor.close(); + } + } + } + + @Test + public void aFlushIsNotToldEventsArrivedThatAnEarlierDeliveryTookAndLost() throws Exception { + // By the time this flush runs the buffer is empty, which is also what it looks like when + // the events arrived, so only the earlier delivery's outcome can tell the two apart. + AtomicInteger sends = new AtomicInteger(0); + EventSender sender = new StubEventSender() { + @Override + public Result sendAnalyticsEvents(byte[] data, int eventCount, URI eventsBaseUri) { + return new Result(sends.incrementAndGet() > 1, false, null); + } + }; + ScheduledExecutorService scheduler = EventUtil.makeEventsTaskExecutor(); + DirectEventProcessor eventProcessor = makeEventProcessor(sender, NO_PERIODIC_FLUSH_MILLIS, + scheduler); + try { + eventProcessor.setOffline(false); + eventProcessor.recordCustomEvent(CONTEXT, "lost", LDValue.ofNull(), null); + eventProcessor.blockingFlush(); // takes the event, and its post fails unheard + + assertFalse(awaitFlush(eventProcessor, 5, TimeUnit.SECONDS)); + assertEquals("the second flush had nothing of its own to post", 1, sends.get()); + + // Once a caller has been told, the next is answered only for what came after. + eventProcessor.recordCustomEvent(CONTEXT, "delivered", LDValue.ofNull(), null); + assertTrue(awaitFlush(eventProcessor, 5, TimeUnit.SECONDS)); + } finally { + eventProcessor.close(); + scheduler.shutdownNow(); + } + } + + @Test + public void flushAndWaitDoesNotReportDeliveryForEventsAnEarlierFailedFlushDrained() throws Exception { + // Every request fails recoverably (503): the sender attempts, retries once, gives up. + try (HttpServer server = HttpServer.start(Handlers.status(503))) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "lost-event", LDValue.ofNull(), null); + + // Stands in for the periodic flush or the SDK's flush-on-background: it drains + // the buffer and its delivery fails. The run is not restored to the buffer. + eventProcessor.blockingFlush(); + server.getRecorder().requireRequest(10, TimeUnit.SECONDS); + + // The event is irrecoverably gone, so "were my events delivered" is no. + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + } finally { + eventProcessor.close(); + } + } + } + + @Test + public void aFlushAfterTheServiceRecoversIsAnsweredOnlyForWhatCameAfterTheLoss() throws Exception { + HandlerSwitcher service = new HandlerSwitcher(Handlers.status(503)); + try (HttpServer server = HttpServer.start(service)) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "lost-event", LDValue.ofNull(), null); + eventProcessor.blockingFlush(); + // The attempt and its one retry: the service only recovers once the run is lost. + server.getRecorder().requireRequest(10, TimeUnit.SECONDS); + server.getRecorder().requireRequest(10, TimeUnit.SECONDS); + service.setTarget(Handlers.status(202)); + + // The service is back, and the event it refused is still not delivered. + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + server.getRecorder().requireNoRequests(100, TimeUnit.MILLISECONDS); + + eventProcessor.recordCustomEvent(CONTEXT, "delivered-event", LDValue.ofNull(), null); + assertTrue(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + RequestInfo request = server.getRecorder().requireRequest(10, TimeUnit.SECONDS); + assertTrue(request.getBody().contains("\"key\":\"delivered-event\"")); + assertFalse("a lost run is not resent", request.getBody().contains("lost-event")); + } finally { + eventProcessor.close(); + } + } + } + + @Test + public void aFlushIsNotToldEventsArrivedThatCouldNotBeSerialized() throws Exception { + try (HttpServer server = startEventsServer()) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "poison", LDValue.ofNull(), Double.NaN); + + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + + server.getRecorder().requireNoRequests(100, TimeUnit.MILLISECONDS); + } finally { + eventProcessor.close(); + } + } + } + + @Test + public void aFlushIsNotToldEventsArrivedWhenSomeOfThemCouldNotBeSerialized() throws Exception { + try (HttpServer server = startEventsServer()) { + EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY); + try { + eventProcessor.recordCustomEvent(CONTEXT, "poison", LDValue.ofNull(), Double.NaN); + eventProcessor.recordCustomEvent(CONTEXT, "fine", LDValue.ofNull(), 1.0); + + // The post succeeds, and still not everything the caller recorded is in it. + assertFalse(awaitFlush(eventProcessor, 10, TimeUnit.SECONDS)); + + List events = collectDelivered(server); + assertEquals(LDValue.of("fine"), requireEventOfKind(events, "custom").get("key")); + } finally { + eventProcessor.close(); + } + } + } + @Test public void closeReleasesTheSenderOnlyAfterTheLastDeliveryFinishes() throws Exception { // Giving up on the wait must not turn into pulling the HTTP client out from under the @@ -1253,6 +1573,19 @@ private DiagnosticStore makeDiagnosticStore() { "android-client-sdk", "0.0.0", "Android", null, Collections.emptyMap(), null)); } + /** + * Flushes and waits for the outcome the way {@code LDClient.flushAndWait} does, which is the + * only place a timeout belongs. + */ + private static boolean awaitFlush(EventProcessor eventProcessor, long timeout, TimeUnit unit) + throws Exception { + try { + return Boolean.TRUE.equals(eventProcessor.flushAsync().get(timeout, unit)); + } catch (TimeoutException e) { + return false; + } + } + private static void awaitQuietly(CountDownLatch latch, long timeout, TimeUnit unit) { try { latch.await(timeout, unit); diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/EventProcessorFlushAsyncDefaultTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/EventProcessorFlushAsyncDefaultTest.java new file mode 100644 index 00000000..427fc39a --- /dev/null +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/EventProcessorFlushAsyncDefaultTest.java @@ -0,0 +1,90 @@ +package com.launchdarkly.sdk.android; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +import com.launchdarkly.sdk.EvaluationReason; +import com.launchdarkly.sdk.LDContext; +import com.launchdarkly.sdk.LDValue; +import com.launchdarkly.sdk.android.subsystems.EventProcessor; + +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.Timeout; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * Covers what {@link EventProcessor#flushAsync()} does for an implementation that predates it and + * has only its unbounded {@link EventProcessor#blockingFlush()}. The SDK has to be able to put a + * deadline on such a flush, and must not tell the caller its events are safe without knowing. + */ +public class EventProcessorFlushAsyncDefaultTest { + @Rule + public Timeout globalTimeout = Timeout.seconds(30); + + @Test + public void theCallersDeadlineBoundsTheWaitAndNotTheFlush() throws Exception { + LegacyEventProcessor eventProcessor = new LegacyEventProcessor(); + Future delivery = eventProcessor.flushAsync(); + + try { + delivery.get(100, TimeUnit.MILLISECONDS); + fail("the wait outlived the deadline"); + } catch (TimeoutException expected) { + // The flush is still going, which is why this is what the caller is told. + } + assertFalse(eventProcessor.flushReturned.get()); + + eventProcessor.letFlushFinish.countDown(); + assertTrue("the flush's own outcome was not reported", + delivery.get(5, TimeUnit.SECONDS)); + assertTrue(eventProcessor.flushReturned.get()); + } + + /** An implementation written before {@code flushAsync} existed. */ + private static final class LegacyEventProcessor implements EventProcessor { + final CountDownLatch letFlushFinish = new CountDownLatch(1); + final AtomicBoolean flushReturned = new AtomicBoolean(false); + + @Override + public void blockingFlush() { + try { + letFlushFinish.await(10, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + flushReturned.set(true); + } + + @Override + public void flush() {} + + @Override + public void setInBackground(boolean inBackground) {} + + @Override + public void setOffline(boolean offline) {} + + @Override + public void close() {} + + @Override + public void recordEvaluationEvent(LDContext context, String flagKey, int flagVersion, + int variation, LDValue value, EvaluationReason reason, + LDValue defaultValue, boolean requireFullEvent, + Long debugEventsUntilDate) {} + + @Override + public void recordIdentifyEvent(LDContext context) {} + + @Override + public void recordCustomEvent(LDContext context, String eventKey, LDValue data, + Double metricValue) {} + } +} diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/integrations/LDCrashHandlerTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/integrations/LDCrashHandlerTest.java new file mode 100644 index 00000000..97cb317e --- /dev/null +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/integrations/LDCrashHandlerTest.java @@ -0,0 +1,118 @@ +package com.launchdarkly.sdk.android.integrations; + +import static org.easymock.EasyMock.createMock; +import static org.easymock.EasyMock.eq; +import static org.easymock.EasyMock.expect; +import static org.easymock.EasyMock.replay; +import static org.easymock.EasyMock.verify; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertTrue; + +import com.launchdarkly.sdk.android.LDClientInterface; +import com.launchdarkly.sdk.android.LaunchDarklyException; + +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.concurrent.TimeUnit; + +public class LDCrashHandlerTest { + private final List calls = new ArrayList<>(); + private final Thread crashed = new Thread("crashed"); + private final Throwable crash = new RuntimeException("boom"); + private final Thread.UncaughtExceptionHandler next = (thread, throwable) -> { + assertSame(crashed, thread); + assertSame(crash, throwable); + calls.add("next"); + }; + + private Thread.UncaughtExceptionHandler originalDefault; + + @Before + public void saveDefaultHandler() { + originalDefault = Thread.getDefaultUncaughtExceptionHandler(); + } + + @After + public void restoreDefaultHandler() { + Thread.setDefaultUncaughtExceptionHandler(originalDefault); + } + + @Test + public void deliversTheEventsBeforePassingTheCrashOn() { + LDClientInterface client = createMock(LDClientInterface.class); + expect(client.flushAndWait(eq(2000L), eq(TimeUnit.MILLISECONDS))).andAnswer(() -> { + calls.add("flush"); + return true; + }); + replay(client); + + new LDCrashHandler(() -> client, 2000, next).uncaughtException(crashed, crash); + + verify(client); + assertEquals(Arrays.asList("flush", "next"), calls); + } + + @Test + public void passesTheCrashOnWhenTheEventsAreNotDelivered() { + LDClientInterface client = createMock(LDClientInterface.class); + expect(client.flushAndWait(eq(2000L), eq(TimeUnit.MILLISECONDS))).andReturn(false); + replay(client); + + new LDCrashHandler(() -> client, 2000, next).uncaughtException(crashed, crash); + + verify(client); + assertEquals(Arrays.asList("next"), calls); + } + + @Test + public void passesTheCrashOnWhenTheClientWasNeverInitialized() { + new LDCrashHandler(() -> { + throw new LaunchDarklyException("LDClient.get() was called before init()!"); + }, 2000, next).uncaughtException(crashed, crash); + + assertEquals(Arrays.asList("next"), calls); + } + + @Test + public void passesTheCrashOnWhenTheFlushThrows() { + LDClientInterface client = createMock(LDClientInterface.class); + expect(client.flushAndWait(eq(2000L), eq(TimeUnit.MILLISECONDS))) + .andThrow(new IllegalStateException("the crash broke the client too")); + replay(client); + + new LDCrashHandler(() -> client, 2000, next).uncaughtException(crashed, crash); + + assertEquals(Arrays.asList("next"), calls); + } + + @Test + public void installsInFrontOfTheCurrentDefaultHandler() { + Thread.setDefaultUncaughtExceptionHandler(next); + + LDCrashHandler.install(2, TimeUnit.SECONDS); + Thread.UncaughtExceptionHandler installed = Thread.getDefaultUncaughtExceptionHandler(); + installed.uncaughtException(crashed, crash); + + assertTrue(installed instanceof LDCrashHandler); + assertEquals(Arrays.asList("next"), calls); + } + + @Test + public void installingTwiceKeepsTheFirstHandler() { + Thread.setDefaultUncaughtExceptionHandler(next); + + LDCrashHandler.install(2, TimeUnit.SECONDS); + Thread.UncaughtExceptionHandler first = Thread.getDefaultUncaughtExceptionHandler(); + LDCrashHandler.install(5, TimeUnit.SECONDS); + + assertSame(first, Thread.getDefaultUncaughtExceptionHandler()); + first.uncaughtException(crashed, crash); + assertEquals(Arrays.asList("next"), calls); + } +} diff --git a/test-app/README.md b/test-app/README.md index e7d624ae..cd23fd36 100644 --- a/test-app/README.md +++ b/test-app/README.md @@ -11,9 +11,16 @@ launchdarkly.environment=production Set `launchdarkly.environment=staging` to use LaunchDarkly's staging endpoints. -## Tier 1 event-loss scenario +## Event-loss scenarios Create a boolean flag named `kill-flag`, or enter another flag key in the app. Tap **Eval+track+kill** to evaluate the flag, track a stand-in error event, request a flush, and terminate the process five seconds later. This exercises the interval between recording and delivery without Android lifecycle callbacks masking the result. + +The two immediate controls compare exits that application code can and cannot observe: + +- **Eval+Kill now** records the same pair and sends `SIGKILL` immediately. No handler or SDK code + can run before the process ends. +- **Eval+Crash now** throws an uncaught exception immediately after recording. The installed crash + handler calls `flushAndWait` with a two-second budget before delegating to Android's handler. diff --git a/test-app/src/main/java/com/launchdarkly/sdk/testapp/MainActivity.java b/test-app/src/main/java/com/launchdarkly/sdk/testapp/MainActivity.java index bbdd1e7a..03aaa310 100644 --- a/test-app/src/main/java/com/launchdarkly/sdk/testapp/MainActivity.java +++ b/test-app/src/main/java/com/launchdarkly/sdk/testapp/MainActivity.java @@ -23,9 +23,11 @@ import com.launchdarkly.sdk.android.LDFailure; import com.launchdarkly.sdk.android.LDStatusListener; import com.launchdarkly.sdk.android.integrations.DedupingHook; +import com.launchdarkly.sdk.android.integrations.LDCrashHandler; import java.util.Date; import java.util.Locale; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import timber.log.Timber; @@ -102,7 +104,12 @@ public void onCreate(Bundle savedInstanceState) { setupTrackButton(); setupIdentifyButton(); setupKillUnsentButton(); + setupKillNowButton(); + setupCrashNowButton(); setupOfflineSwitch(); + // Rescues the events for "Eval+Crash now" and cannot run for "Eval+Kill now", which is what + // makes the pair worth pressing. + LDCrashHandler.install(2, TimeUnit.SECONDS); setupListeners(); updateDedupeStatus(); @@ -212,6 +219,31 @@ private void setupTrackButton() { }); } + /** + * The flag the kill and crash buttons evaluate: whatever is typed in the feature key field, or a + * default, so the buttons work without anything being typed first. + */ + private String flagKeyToKillOver() { + String typedKey = ((EditText) findViewById(R.id.feature_flag_key)).getText().toString().trim(); + return typedKey.isEmpty() ? "kill-flag" : typedKey; + } + + /** + * Records the pair whose survival is in question: an evaluation, which is the exposure, and a + * track, standing in for the error an application reports just before it dies. + * + *

Returns false when there is no client, in which case nothing was recorded and ending the + * process would demonstrate nothing. + */ + private boolean recordExposureAndError(String flagKey) { + if (ldClient == null) { + return false; + } + ldClient.boolVariation(flagKey, false); + ldClient.track("$ld:telemetry:error"); + return true; + } + /** * Reproduces in-memory event loss: evaluate (exposure) and track (stand-in for an error), * wait 5s so both calls are queued, then kill the process before the 30s flush. @@ -221,17 +253,58 @@ private void setupTrackButton() { private void setupKillUnsentButton() { Button killUnsentButton = findViewById(R.id.kill_unsent_button); killUnsentButton.setOnClickListener(v -> { - final String typedKey = ((EditText) findViewById(R.id.feature_flag_key)).getText().toString().trim(); - final String flagKey = typedKey.isEmpty() ? "kill-flag" : typedKey; + final String flagKey = flagKeyToKillOver(); Timber.w("eval+track+kill flag=%s", flagKey); - doSafeClientAction(() -> { - ldClient.boolVariation(flagKey, false); - ldClient.track("$ld:telemetry:error"); - ldClient.flush(); - new Handler(Looper.getMainLooper()).postDelayed( - () -> android.os.Process.killProcess(android.os.Process.myPid()), - 5_000); - }); + if (!recordExposureAndError(flagKey)) { + return; + } + ldClient.flush(); + new Handler(Looper.getMainLooper()).postDelayed( + () -> android.os.Process.killProcess(android.os.Process.myPid()), + 5_000); + }); + } + + /** + * The same sequence with nothing at all between the track and the process dying: no flush to + * deliver the events, no delay for a timer to fire in, and SIGKILL to itself, which cannot be + * caught, so no part of the SDK gets to run on the way out. + * + *

Whether the exposure and the track are reported therefore says exactly one thing: whether + * recording them had already put them somewhere that outlives the process. They should arrive on + * the next launch of the app, not this one. + */ + private void setupKillNowButton() { + Button killNowButton = findViewById(R.id.kill_now_button); + killNowButton.setOnClickListener(v -> { + final String flagKey = flagKeyToKillOver(); + Timber.w("eval+track+kill now flag=%s", flagKey); + if (!recordExposureAndError(flagKey)) { + return; + } + android.os.Process.killProcess(android.os.Process.myPid()); + }); + } + + /** + * The same again, ending in an uncaught exception instead of a signal the process never sees. + * + *

This is the shape a customer report takes: app code fails immediately after reporting the + * failure. Unlike SIGKILL, an uncaught exception runs the default handler before the process + * goes, so this is the one variant an application can rescue on its own, which + * {@link LDCrashHandler} does by calling {@link LDClient#flushAndWait} from there. So these + * events should arrive and the ones from the button next to it should not. + */ + private void setupCrashNowButton() { + Button crashNowButton = findViewById(R.id.crash_now_button); + crashNowButton.setOnClickListener(v -> { + final String flagKey = flagKeyToKillOver(); + Timber.w("eval+track+crash now flag=%s", flagKey); + if (!recordExposureAndError(flagKey)) { + return; + } + throw new RuntimeException( + "Eval+Crash: deliberate uncaught exception immediately after track, to test event persistence"); }); } diff --git a/test-app/src/main/res/layout/activity_main.xml b/test-app/src/main/res/layout/activity_main.xml index dfa8fac0..3b377995 100644 --- a/test-app/src/main/res/layout/activity_main.xml +++ b/test-app/src/main/res/layout/activity_main.xml @@ -114,8 +114,12 @@ android:layout_alignParentRight="true" android:minLines="4" /> -