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/DirectEventProcessor.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java index 147a9fb8..02e0c775 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 @@ -19,6 +19,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; @@ -355,6 +356,32 @@ public void blockingFlush() { } } + @Override + public boolean blockingFlush(long timeout, TimeUnit unit) { + if (isStopped()) { + return false; + } + // Typed rather than inlined, so that it is unambiguously submitted as work with a result. + Callable delivery = this::deliverPayloadReportingOutcome; + Future pending = submit(delivery); + if (pending == null) { + return false; + } + try { + return Boolean.TRUE.equals(pending.get(timeout, unit)); + } catch (TimeoutException e) { + // Left running rather than cancelled: the buffer has already been drained into the + // payload, so interrupting the delivery now would only make the loss certain. + return false; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return false; + } catch (ExecutionException e) { + logUnexpectedError(e.getCause() == null ? e : e.getCause()); + return false; + } + } + @Override public void close() throws IOException { if (!closed.compareAndSet(false, true)) { @@ -425,17 +452,28 @@ 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 a caller that is not waiting to find out how it + * went. */ private void deliverPayload() { + deliverPayloadReportingOutcome(); + } + + /** + * Delivers as {@link #deliverPayload()} does, and says whether it worked, for a caller that is + * 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 +488,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; } } @@ -685,6 +726,25 @@ private Future submit(Runnable task) { } } + /** + * As {@link #submit(Runnable)}, for a delivery whose outcome the caller waits for. Not wrapped in + * {@link #guarded}, because here the caller is there to receive what escapes. + * + * @return the submitted task, or null if the processor is shutting down or already has + */ + private Future submit(Callable task) { + synchronized (submitLock) { + if (shuttingDown) { + return null; + } + try { + return scheduler.submit(task); + } catch (RuntimeException e) { // the executor was shut down under us + return null; + } + } + } + /** * Keeps an unexpected failure from killing a repeating task or bubbling out of the executor. * Anything that escapes a run suppresses the rest of a {@code scheduleWithFixedDelay} series, 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..8eb539a0 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,29 @@ 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; + } + boolean delivered = true; + for (LDClient client : clients.values()) { + // Each environment 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. + long remaining = Math.max(0, deadline - System.nanoTime()); + delivered &= client.flushAndWaitInternal(remaining, TimeUnit.NANOSECONDS); + } + return delivered; + } + + private boolean flushAndWaitInternal(long timeout, TimeUnit unit) { + return eventProcessor.blockingFlush(timeout, unit); + } + @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/subsystems/EventProcessor.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/subsystems/EventProcessor.java index e65a8e2b..18f68efc 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 @@ -6,6 +6,7 @@ import com.launchdarkly.sdk.LDValue; import java.io.Closeable; +import java.util.concurrent.TimeUnit; /** * Interface for an object that can send or store analytics events. @@ -99,4 +100,21 @@ 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, blocking until they have been + * delivered or until the timeout expires, whichever comes first. + * + * @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 events could not be delivered + * @since 5.17.0 + */ + default boolean blockingFlush(long timeout, TimeUnit unit) { + // An implementation written before this method existed has no way to honor a timeout, so it + // gets its unbounded flush and reports success, having no way to tell otherwise. + 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..2ee21372 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 @@ -34,6 +34,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -651,6 +652,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(eventProcessor.blockingFlush(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 +793,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(eventProcessor.blockingFlush(10, TimeUnit.SECONDS)); + + server.getRecorder().requireNoRequests(100, TimeUnit.MILLISECONDS); + } finally { + eventProcessor.close(); + } + } + } + @Test public void unexpectedRecordingErrorDoesNotBubbleToCallerAndLogs() throws Exception { ScheduledExecutorService scheduler = EventUtil.makeEventsTaskExecutor(); @@ -858,6 +890,44 @@ 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(eventProcessor.blockingFlush(10, TimeUnit.SECONDS)); + + server.getRecorder().requireNoRequests(100, TimeUnit.MILLISECONDS); + } finally { + 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(eventProcessor.blockingFlush(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 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" /> -