-
Notifications
You must be signed in to change notification settings - Fork 24
feat(events): add flushAndWait with a bounded timeout (tier 2) #403
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: andrey/event-durability-tier1-buffer
Are you sure you want to change the base?
Changes from all commits
0c98c92
d57313c
818d239
63ca732
d399f54
ef886dc
536a453
19a14b3
f233167
0bcef8b
644ac0f
68d8aa9
7bd8b07
7385975
f278820
c6386f7
f36aee4
23bba70
a408e75
113ee34
70babb2
ff5aeec
105c870
bdfd13a
6290c38
ee78359
c629046
2237489
df9e409
e113694
4ac8572
7fe25ed
6f1ae80
3ab2ffc
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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<Boolean> 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<Boolean> delivery = flushAsync(); | ||
| try { | ||
| delivery.get(); | ||
| } catch (InterruptedException e) { | ||
|
|
@@ -355,6 +359,61 @@ public void blockingFlush() { | |
| } | ||
| } | ||
|
|
||
| @Override | ||
| public Future<Boolean> flushAsync() { | ||
| if (isStopped()) { | ||
| return new LDSuccessFuture<>(false); | ||
| } | ||
| return queueDelivery(); | ||
| } | ||
|
|
||
| /** | ||
| * Queues a delivery, or hands back one that is already queued and has not started. | ||
| * <p> | ||
| * 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<Boolean> queueDelivery() { | ||
| synchronized (flushLock) { | ||
| if (pendingFlush != null) { | ||
| return pendingFlush; | ||
| } | ||
| LDAwaitFuture<Boolean> 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<Boolean> 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. | ||
| * <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 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. | ||
| * <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; | ||
|
|
@@ -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(); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This is the core new branch — |
||
| } catch (Exception e) { | ||
| logUnexpectedError(e); | ||
| return false; | ||
| } | ||
| } | ||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -779,6 +779,49 @@ private void flushInternal() { | |
| eventProcessor.flush(); | ||
| } | ||
|
|
||
| @Override | ||
| public boolean flushAndWait(long timeout, TimeUnit unit) { | ||
| long deadline = System.nanoTime() + unit.toNanos(timeout); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
One-line fix: clamp the timeout to >= 0 before computing the deadline. (Positive extremes are already safe: |
||
| 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; | ||
| } | ||
| // 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<Future<Boolean>> deliveries = new ArrayList<>(clients.size()); | ||
| for (LDClient client : clients.values()) { | ||
| deliveries.add(client.eventProcessor.flushAsync()); | ||
| } | ||
| boolean delivered = true; | ||
| for (Future<Boolean> 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; | ||
| } | ||
|
cursor[bot] marked this conversation as resolved.
|
||
|
|
||
| private boolean awaitDelivery(Future<Boolean> 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(); | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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. | ||
| * <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 | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Two doc gaps worth closing here, since this javadoc is the one place customers will read:
|
||
| * 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. | ||
| * | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Flush reports success after failed send
Medium Severity
flushAndWaitcan returntrueafter the caller's events were already taken by an in-progress or just-queued delivery that then failed. Coalescing only joins a delivery that has not started, so a laterflushAsyncwaits on a follow-up. That follow-up finds an empty buffer and treats it as success, even though the earlier send lost the events. A crash handler can then treat those events as delivered when they are gone.Additional Locations (2)
launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java#L507-L528launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDClient.java#L781-L805Reviewed by Cursor Bugbot for commit 3ab2ffc. Configure here.