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..3590f9c3 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 @@ -26,6 +26,7 @@ import org.junit.Test; import java.io.IOException; +import java.util.concurrent.TimeUnit; import okhttp3.HttpUrl; import okhttp3.mockwebserver.MockResponse; @@ -95,6 +96,36 @@ 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 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 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..bb411a32 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 @@ -133,6 +133,19 @@ 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}. 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; + DirectEventProcessor( OutboundEventBuffer buffer, EventSender eventSender, @@ -331,21 +344,12 @@ public void setOffline(boolean offline) { @Override public void flush() { - if (isStopped()) { - return; - } - submit(this::deliverPayload); + flushAsync(); } @Override public void blockingFlush() { - if (isStopped()) { - return; - } - Future delivery = submit(this::deliverPayload); - if (delivery == null) { - return; - } + Future delivery = flushAsync(); try { delivery.get(); } catch (InterruptedException e) { @@ -355,6 +359,61 @@ public void blockingFlush() { } } + @Override + public Future flushAsync() { + if (isStopped()) { + return new LDSuccessFuture<>(false); + } + return queueDelivery(); + } + + /** + * Queues a delivery, or hands back one that is already queued and has not started. + *

+ * A delivery that has not started yet will take everything recorded up to the moment it does, + * which includes whatever the caller recorded before asking, so waiting on it answers the + * caller's question as well as a delivery of its own would. 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. + */ + private Future queueDelivery() { + synchronized (flushLock) { + if (pendingFlush != null) { + 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; + return result; + } + } + + /** + * Runs one delivery on behalf of every flush request that joined it, and tells them all how it + * went. + */ + private void runDelivery(LDAwaitFuture result) { + 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; + } + } + 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); + } + result.set(delivered); + } + @Override public void close() throws IOException { if (!closed.compareAndSet(false, true)) { @@ -368,22 +427,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().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 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 +485,29 @@ 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 or the service did not accept them + */ + private boolean deliverPayloadReportingOutcome() { if (disabled || offline.get()) { - return; + return false; } List run; List summaries; @@ -450,19 +522,22 @@ private void deliverPayload() { payload = buffer.encode(run, summaries); } catch (IOException e) { logUnexpectedError(e); - return; + return false; } if (payload == null) { - return; + return true; } 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(); } catch (Exception e) { logUnexpectedError(e); + return false; } } 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..2bb84155 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 @@ -779,6 +779,49 @@ private void flushInternal() { eventProcessor.flush(); } + @Override + public boolean flushAndWait(long timeout, TimeUnit unit) { + long deadline = System.nanoTime() + 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 cancelled: 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 (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..81bc5373 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,27 @@ 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. + *

+ * 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. + * + * @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, or the SDK is offline, closed, or otherwise unable to deliver them + * @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/subsystems/EventProcessor.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/subsystems/EventProcessor.java index e65a8e2b..31538b21 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,29 @@ 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 + * @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..989fb4d8 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; @@ -30,10 +31,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 +655,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 +796,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 +893,124 @@ 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); + + assertFalse(awaitFlush(eventProcessor, 100, TimeUnit.MILLISECONDS)); + } finally { + // Released before closing, so that the delivery still in flight can finish rather + // than hold up the shutdown that close() waits on. + letResponseFinish.release(Integer.MAX_VALUE); + 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 +1406,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/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/FlushOnCrashHandler.java b/test-app/src/main/java/com/launchdarkly/sdk/testapp/FlushOnCrashHandler.java new file mode 100644 index 00000000..3fa56c31 --- /dev/null +++ b/test-app/src/main/java/com/launchdarkly/sdk/testapp/FlushOnCrashHandler.java @@ -0,0 +1,67 @@ +package com.launchdarkly.sdk.testapp; + +import com.launchdarkly.sdk.android.LDClient; + +import java.util.concurrent.TimeUnit; + +import timber.log.Timber; + +/** + * Delivers buffered analytics events from the uncaught exception handler, which is the most an + * application can do about event loss while the SDK keeps its events only in memory. + *

+ * This is what makes the two instant buttons in {@link MainActivity} an experiment and its control. + * An uncaught exception runs this handler while the process is still alive and its other threads are + * still running, so the events recorded a moment earlier can still reach the network. + * {@code SIGKILL} runs nothing, and neither does an ANR, a native crash, or the system reclaiming a + * backgrounded process, so those lose the same events. The difference between the two buttons is the + * ground that on-disk persistence would cover and a crash handler cannot. + */ +final class FlushOnCrashHandler implements Thread.UncaughtExceptionHandler { + /** + * How long the crash is held open for the events. + *

+ * The SDK's HTTP timeouts are measured in seconds, and a request that hangs must not hold the + * process in a half-dead state for all of them: past this point the events are worth less than + * the delay, and the crash goes on to be reported. + */ + private static final long DELIVERY_BUDGET_MILLIS = 2_000; + + private final Thread.UncaughtExceptionHandler next; + + private FlushOnCrashHandler(Thread.UncaughtExceptionHandler next) { + this.next = next; + } + + /** + * Installs the handler in front of whatever was already there, which on a real application is + * the crash reporter, and on this one is the platform handler that prints the trace. + */ + static void install() { + Thread.UncaughtExceptionHandler previous = Thread.getDefaultUncaughtExceptionHandler(); + if (previous instanceof FlushOnCrashHandler) { + return; + } + Thread.setDefaultUncaughtExceptionHandler(new FlushOnCrashHandler(previous)); + } + + @Override + public void uncaughtException(Thread thread, Throwable throwable) { + try { + // Waiting here on the crashing thread is safe because the timeout is the SDK's to + // enforce: it stops waiting on the delivery rather than trusting it to finish. That also + // covers the case where this crash is the reason the delivery cannot complete, such as + // an exception thrown while the event buffer was locked. + boolean delivered = LDClient.get() + .flushAndWait(DELIVERY_BUDGET_MILLIS, TimeUnit.MILLISECONDS); + Timber.w("crash handler: events delivered = %b", delivered); + } catch (Throwable t) { + // Nothing that happens in here is worth losing the crash report over. + Timber.e(t, "Could not deliver events from the crash handler"); + } finally { + if (next != null) { + next.uncaughtException(thread, throwable); + } + } + } +} 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..c0a1fcb7 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 @@ -102,7 +102,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. + FlushOnCrashHandler.install(); setupListeners(); updateDedupeStatus(); @@ -212,6 +217,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 +251,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 FlushOnCrashHandler} 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" /> -