Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
0c98c92
feat(events): add a bounded flush (tier 2)
abelonogov-ld Sep 4, 2026
d57313c
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 15, 2026
818d239
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 15, 2026
63ca732
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 15, 2026
d399f54
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 17, 2026
ef886dc
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 18, 2026
536a453
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 21, 2026
19a14b3
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 21, 2026
f233167
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 21, 2026
0bcef8b
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 21, 2026
644ac0f
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 21, 2026
68d8aa9
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
7bd8b07
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
7385975
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
f278820
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
c6386f7
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
f36aee4
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
23bba70
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
a408e75
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
113ee34
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
70babb2
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
ff5aeec
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
105c870
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
bdfd13a
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
6290c38
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 22, 2026
ee78359
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 23, 2026
c629046
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 23, 2026
2237489
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 23, 2026
df9e409
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 23, 2026
e113694
Report a failed flushAndWait from a closed client
abelonogov-ld Sep 23, 2026
4ac8572
Merge branch 'andrey/event-durability-tier1-buffer' into andrey/event…
abelonogov-ld Sep 23, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.junit.Test;

import java.io.IOException;
import java.util.concurrent.TimeUnit;

import okhttp3.HttpUrl;
import okhttp3.mockwebserver.MockResponse;
Expand Down Expand Up @@ -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()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Boolean> delivery = this::deliverPayloadReportingOutcome;
Future<Boolean> 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)) {
Expand Down Expand Up @@ -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.
* <p>
* 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.
* <p>
* 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<Event> run;
List<EventSummarizer.EventSummary> summaries;
Expand All @@ -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;
}
}

Expand Down Expand Up @@ -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 <T> Future<T> submit(Callable<T> 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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, LDClient> 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;
}
Comment thread
cursor[bot] marked this conversation as resolved.

private boolean flushAndWaitInternal(long timeout, TimeUnit unit) {
return eventProcessor.blockingFlush(timeout, unit);
}

@VisibleForTesting
void blockingFlush() {
eventProcessor.blockingFlush();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -146,6 +147,27 @@ public interface LDClientInterface extends Closeable {
*/
void flush();

/**
* Sends all pending events to LaunchDarkly and waits for them to be delivered.
* <p>
* 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.
* <p>
* 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.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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
Expand Down
9 changes: 8 additions & 1 deletion test-app/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Loading
Loading