Repository navigation
feat(events): add flushAndWait with a bounded timeout (tier 2) - #403
abelonogov-ld wants to merge 43 commits into
Conversation
Tier 2 of the event durability spec: recoverable failure. flushAndWait(timeout) delivers what has been recorded and reports whether it got there inside the budget, so an application that knows it is about to go away — backgrounding, or an uncaught exception handler on its way out — gets an answer instead of a fire-and-forget flush. Delivery now reports its outcome so the bounded call can tell delivered from not. Spec: Event Durability, "Tier 2 — recoverable failure", §10. Co-authored-by: Cursor <cursoragent@cursor.com>
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # example/src/main/java/com/launchdarkly/example/MainActivity.java # example/src/main/res/layout/activity_main.xml
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush * andrey/event-durability-tier1-buffer: unserializable unit test
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush * andrey/event-durability-tier1-buffer: test(fdv2): stop requiring a changeset to be the first result
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com>
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/DirectEventProcessor.java
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java
…-durability-tier2-bounded-flush
…-durability-tier2-bounded-flush * andrey/event-durability-tier1-buffer: Document that close() discards events held while offline
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java
…-durability-tier2-bounded-flush
flushAndWait started from delivered = true and only updated it for the environments this client still belongs to. After close() that set is empty, so the call reported success for events it could not deliver. It now returns false when the client has been closed or replaced. Co-authored-by: Cursor <cursoragent@cursor.com>
…-durability-tier2-bounded-flush Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/DirectEventProcessorTest.java
…-durability-tier2-bounded-flush
| * timeout expired first or the events could not be delivered | ||
| * @since 5.17.0 | ||
| */ | ||
| default boolean blockingFlush(long timeout, TimeUnit unit) { |
There was a problem hiding this comment.
This isn't a valid default impl. In this case it just completely ignores the time param, so it violates the interface. You may need to version the interface in a breaking manner. Also I contend that this signature should not take a timeout and should return a future.
Had the previous developer done Future flush() + future.get() instead of blockingFlush, you would have had the flexibility to insert a timeout. You should set this up so that future devs have more flexibility.
| // 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); |
There was a problem hiding this comment.
Shouldn't these be running in parallel?
| } | ||
| // 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); |
There was a problem hiding this comment.
I am concerned it is possible for flush requests to queue behind one another and the last flush gets starved. Example:
- Many events are being recorded
- At the same time, flushes are being requested at some frequency that is faster than a flush can complete in
- The queue grows
- A shutdown happens and the flushAndWait cannot complete in time.
We should use a similar mechanism to a shedding queue that was used to implement client side identify. See the CSI-client-side-identify spec.
Review of the tier 2 PR raised four things, which are one change: - EventProcessor now offers Future<Boolean> flushAsync() in place of the blockingFlush(timeout, unit) added earlier in this branch. Only the public LDClient.flushAndWait takes a timeout; inside the SDK a flush hands back something to wait on, so waits can be composed. Nothing released carried the timeout overload, so nothing is broken by dropping it. - Its default implementation runs the old unbounded blockingFlush() on a pooled daemon thread and reports what it returned, instead of ignoring the timeout and claiming the events were delivered. - flushAndWait starts every environment's delivery before waiting on any of them, so environments deliver at once and share one deadline rather than spending it in series. - A flush request joins a delivery that is queued and has not started yet, because that delivery will take everything recorded up to the moment it begins, including the joining caller's events. Flushes arriving faster than a post completes used to queue one post each, leaving the flush that matters -- the one at shutdown, with a deadline -- waiting behind all of them. The delivery is still never cancelled when a caller stops waiting: by then its events have left the buffer, so interrupting the post would only make losing them certain. Co-authored-by: Cursor <cursoragent@cursor.com>
…-durability-tier2-bounded-flush
| Timber.e(t, "Could not deliver events from the crash handler"); | ||
| } finally { | ||
| if (next != null) { | ||
| next.uncaughtException(thread, throwable); |
There was a problem hiding this comment.
Do we want to run next before ours? Seems likely ours will be a longer handler in the worst case than most hanlders due to its wait nature, so we should run the longer ones later as a good citizen.
There was a problem hiding this comment.
It is part of test-app and yes during my testing I need to call next.uncaughtException(thread, throwable); after, otherwise network delivery doesn't workl
There was a problem hiding this comment.
Yes, it means tier-2 even for Android is weak
There was a problem hiding this comment.
I need to call next.uncaughtException(thread, throwable); after
Doesn't this mean there is some unknown here where some next.uncaughtException calls have unanticipatable side effects that can interfere with ours? Like, one of the nexts was the OS's handler?
There was a problem hiding this comment.
Next one is not OS handler, it is handler of another library.
yes, that's why I am introducing LDErrorHandler to allow customer to manage as it pleased
tanderson-ld
left a comment
There was a problem hiding this comment.
I did a deep pass on this tier: correctness and concurrency, plus the tier-2 sections of the event-durability spec. The hard parts hold up well. I walked the close()/flushAndWait interleavings (no hang; no CancellationException with the SDK's own wiring; the delivery abandoned on timeout can never see a closed sender because the release is queued behind it), interrupt restoration, the no-cancel-on-timeout choice (matches the spec's §10.1 word for word), the non-short-circuiting &= across environments, and there is no double-send path anywhere (drains are destructive and the sender's internal retry reuses one payload ID). Long.MAX_VALUE timeouts are also safe: the two's-complement wraparound in deadline - now cancels correctly.
The findings below are ordered by how much they change the design. 1 and 2 are one mechanism and are worth fixing together.
1. The boolean can report true for events that were drained and lost
deliverPayloadReportingOutcome returns true whenever payload == null, which conflates "nothing was ever recorded" with "a different delivery drained the caller's events and then failed". The sequence that matters for this API's headline use case:
- The periodic flush, the SDK's own flush-on-background in
ConnectivityManager, or anyflush()drains the buffer; the send fails (503, timeout). The run lives in locals and is never re-buffered, so the events are gone. - A crash handler calls
flushAndWait(2, SECONDS). Its delivery queues behind the failed one on the single-threaded scheduler, finds the buffer empty, and returnstrue.
The caller logs "events delivered" for events destroyed a moment earlier. Spec §10.1 is explicit here: a call made while a delivery is in flight must wait for it and report success only if that delivery and its own both got out — an empty buffer is not success.
Failing test (drop into DirectEventProcessorTest; fails today with true):
@Test
public void flushAndWaitDoesNotReportDeliveryForEventsAnEarlierFailedFlushDrained() throws Exception {
// Every request fails recoverably (503): the sender attempts, retries once, gives up.
try (HttpServer server = HttpServer.start(Handlers.status(503))) {
EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY);
try {
eventProcessor.recordCustomEvent(CONTEXT, "lost-event", LDValue.ofNull(), null);
// Stands in for the periodic flush or the SDK's flush-on-background: it drains
// the buffer and its delivery fails. The run is not restored to the buffer.
eventProcessor.blockingFlush();
server.getRecorder().requireRequest(10, TimeUnit.SECONDS);
// The event is irrecoverably gone, so "were my events delivered" is no.
assertFalse(eventProcessor.blockingFlush(10, TimeUnit.SECONDS));
} finally {
eventProcessor.close();
}
}
}A second, fully deterministic route to the same return true: if every event in the drained run fails to serialize, encode returns null after dropping them — the SDK logs "Dropping unserializable", posts nothing, and reports success (also fails today):
@Test
public void flushAndWaitDoesNotReportDeliveryWhenTheWholeRunWasDroppedAsUnserializable() throws Exception {
try (HttpServer server = startEventsServer()) {
EventProcessor eventProcessor = makeEventProcessor(server, DEFAULT_CAPACITY);
try {
// Same fixture as unserializableSummaryDoesNotPoisonLaterDeliveries: a null flag
// key enters the summarizer and Gson cannot write it, so the whole run is dropped.
eventProcessor.recordEvaluationEvent(CONTEXT, null, FLAG_VERSION, VARIATION,
FLAG_VALUE, null, DEFAULT_VALUE, false, null);
boolean reported = eventProcessor.blockingFlush(10, TimeUnit.SECONDS);
server.getRecorder().requireNoRequests(100, TimeUnit.MILLISECONDS);
assertFalse(reported);
} finally {
eventProcessor.close();
}
}
}(The partial variant also misreports: 4 of 5 events sent, 1 dropped unserializable, caller told true — and that drop isn't counted in the drop counter either.)
2. Queued flush deliveries pile up and can starve the delivery that matters
Every flush()/flushAndWait() submits a delivery task unconditionally. Tasks that find an empty buffer are near-free, but with events trickling in between calls, each queued task finds a non-empty buffer and pays a full HTTP round trip — on a failing network with the default sender that is ~20–40s per delivery (two attempts, 10s timeouts, 1s retry sleep). Worst case at exactly the wrong moment: the periodic flush is mid-failure, the backgrounding flush queues behind it, and a crash handler's flushAndWait(2s) queues third. Its delivery never gets scheduler time inside the budget; when the handler returns the process dies, and the queued delivery never runs at all — the events are lost while sitting untouched in the buffer.
The processor this branch replaces had this control: java-sdk-internal's DefaultEventProcessor bounds not-yet-picked-up flush payloads with an ArrayBlockingQueue<>(1) and, when the offer is refused, logs "Skipped flushing because all workers are busy" and restores the summaries to the outbox. The iOS SheddingQueue is the same family for identify (one in flight, one pending; flush wants merge rather than shed semantics, since flush requests don't supersede each other).
Findings 1 and 2 (and 3 below) fall out of one mechanism: track the in-flight delivery and at most one pending follow-up. Sketch that fits the existing single-thread scheduler with no new machinery: a generation counter bumped under recordLock on each record; the delivery task, on completion, checks under the lock whether a flush was requested past what it drained and resubmits itself once (the pending flag collapses any burst of requests); flushAndWait notes its target generation and waits up to its budget for "delivered generation >= mine, successfully". That yields the honest §10.1 answer, bounds the queue at one in-flight plus one pending, and batches better (one coalesced payload instead of N fragments).
3. Multi-environment: submit every environment's delivery before waiting on any
The loop waits env-by-env, so env 1 can burn the entire budget before env 2's delivery is even submitted — and iteration order is HashMap-arbitrary, so which environment gets the budget varies with the key set. Each environment has its own scheduler and sender, so the deliveries can run in parallel: submit all of them first, then wait on each against the shared deadline. Wall time becomes the max instead of the sum, and later environments get a real chance instead of a guaranteed false. (The non-short-circuiting &= is correct and worth a comment, since && would look like a harmless cleanup.)
4. Scope question: retrying a failed batch
The spec's tier-2 requirements also include R1 — a retryable delivery failure must leave the events intact for a later attempt. Today the drained run is discarded after the sender's single in-call retry, so one 503 during a flush interval silently destroys up to a full buffer. If R1 is deliberately a follow-up PR, that seems fine — but worth deciding now, because it changes the shape of fix 1: an honest false today also means "and the events are gone". The old processor's restoreTo-on-refusal move is precedent for where retention can live.
Smaller items
Inline comments on the specific lines for: the EventProcessor default method, a Long.MIN_VALUE overflow, the javadoc's missing reach statement, and two test gaps. Lower-severity items, listed here to keep threads down:
- A timeout returns silently;
close()logs "Gave up waiting…" for the structurally identical situation. A debug line naming the reason (timeout vs offline vs send failure) would make afalsediagnosable in the field. Future.get's fourth failure mode,CancellationException, is uncaught and would escape a public boolean API. Unreachable with the SDK's own wiring (I checked the shutdown policies), reachable if the injected scheduler is shut down externally; a one-line catch returningfalsecloses it.- A null
unitNPEs after the delivery has been submitted on the subsystem path (buffer already drained, exception instead of a boolean);identify(null)returns a failed future by contrast. - The timeout catch's comment ("the buffer has already been drained") is only true if the task started; while it is still queued behind another delivery nothing has been drained. The conclusion (don't cancel) stands either way — worth rewording so a future reader doesn't lean on the stronger claim.
- An error thrown by a delivery after its waiter timed out is captured in an abandoned
Futureand never logged (submit(Callable)skipsguardedby design, but on the timeout path nobody is there to receive it). Logging inside the callable keeps tier 1's nothing-escapes-unlogged property.
| * timeout expired first or the events could not be delivered | ||
| * @since 5.17.0 | ||
| */ | ||
| default boolean blockingFlush(long timeout, TimeUnit unit) { |
There was a problem hiding this comment.
LDConfig.Builder.events() documents custom EventProcessor implementations as supported, and for any of them this default turns the bounded public call unbounded — flushAndWait(100, MILLISECONDS) from a main-thread lifecycle callback blocks for as long as the custom flush takes — and then claims success unconditionally. Both halves are what the interface javadoc's "the timeout bounds the whole call" promises away.
The comment's premise (a pre-existing implementation has no way to honor a timeout) holds for the implementation but not for this default: it can bound the caller itself — run blockingFlush() on a worker and Future.get(timeout, unit), returning false on expiry. Failing that, flush(); return false; keeps the bound honest at the cost of a pessimistic answer. Either way, NullEventProcessor is worth an explicit return true override so Components.noEvents() keeps its correct answer deliberately rather than by inheriting whatever this default does.
|
|
||
| @Override | ||
| public boolean flushAndWait(long timeout, TimeUnit unit) { | ||
| long deadline = System.nanoTime() + unit.toNanos(timeout); |
There was a problem hiding this comment.
TimeUnit.toNanos saturates, so a timeout within ~1µs of Long.MIN_VALUE nanos makes deadline - System.nanoTime() underflow and wrap to ~Long.MAX_VALUE: an effectively unbounded wait on the caller's thread, and nondeterministic (only when nanoTime advanced between the two reads). Ordinary negatives are fine — they produce an immediate false.
One-line fix: clamp the timeout to >= 0 before computing the deadline. (Positive extremes are already safe: (t0 + MAX) - t1 wraps back to MAX - elapsed.)
| * | ||
| * @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 |
There was a problem hiding this comment.
Two doc gaps worth closing here, since this javadoc is the one place customers will read:
- The spec asks that the reach of a bounded flush be stated plainly wherever it is documented. This names the uncaught-exception handler as a caller without saying what the mechanism cannot reach —
SIGKILL, an ANR kill, a native crash, the system reclaiming a backgrounded process all run nothing. The right paragraph already exists in the test app'sFlushOnCrashHandlerheader; it belongs here. @return:falseon timeout does not mean the events were not sent — the delivery is deliberately left running and may still land. Callers planning compensating logic (persist-and-resend) need that stated, or every timeout that later succeeds becomes a duplicate. Also worth noting: the call blocks the calling thread (ANR consideration on main), and with the default sender a failing delivery takes ~21s, so budgets below that generally cannot observe atrueon an unhealthy network.
| } | ||
|
|
||
| @Test | ||
| public void flushWithTimeoutReportsFailureWhenTheTimeoutExpiresFirst() throws Exception { |
There was a problem hiding this comment.
This test can't distinguish the current implementation from one that cancels on timeout — adding pending.cancel(true) in the TimeoutException catch leaves all four new tests green, so the headline "left running rather than cancelled" decision is unpinned. Suggest releasing the semaphore inside the test body (not only in finally) and asserting the request still completes (server.getRecorder().requireRequest(...)): that pins the delivery surviving the caller's departure.
Same shape for the offline test: it passes against an implementation that returns false unconditionally. A follow-up setOffline(false) + flush asserting the event still arrives would also pin the "events survive an offline flush" claim it relies on.
| EventSender.Result result = eventSender.sendAnalyticsEvents(payload.getData(), | ||
| payload.getEventCount(), eventsUri); | ||
| handleResponse(result); | ||
| return result != null && result.isSuccess(); |
There was a problem hiding this comment.
This is the core new branch — false from a non-2xx — and it has no test. HttpServer.start(Handlers.status(503)) (recoverable, exercises the sender's retry) and Handlers.status(401) (the mustShutDown path, pattern already at beingToldToShutDownStopsRecordingAndDelivery) make both cases cheap. Worth covering, since everything else in this PR is about what this value means.
| * 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; |
There was a problem hiding this comment.
What information did you use to pick 2 seconds here? Seems like a long time.
There was a problem hiding this comment.
Network call time on cell phone connection. I was based that 5-15 sec would be also and chose the smallest 2 sec
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using default effort and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, have a team admin enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 3ab2ffc. Configure here.
Co-authored-by: Cursor <cursoragent@cursor.com>
…d delivery Co-authored-by: Cursor <cursoragent@cursor.com>
Bugbot couldn't run - usage limit reachedBugbot is counted against Cursor usage for this user or team, and this run hit a usage or spend limit. A user or team admin can review and increase usage limits in the Cursor dashboard. (requestId: c50db314-b25e-498e-aa0a-cb4c844f579d) |
…allel environments Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
…and use American spelling Co-authored-by: Cursor <cursoragent@cursor.com>
I added similar unit tests for tier-2 but they don't account the case of restoring connection. It is part of tier-3 |
|
Similar tests are |
…et is read as a bound Co-authored-by: Cursor <cursoragent@cursor.com>

Requirements
Related issues
Tier 2 of the Event Durability spec, "Tier 2 — recoverable failure" (§10). Stacked on #397 (tier 1), and targets that branch so this diff shows only tier 2. It is small and self-contained, so it can be reviewed in parallel with #397. It will be retargeted to
mainonce #397 merges.Describe the solution you've provided
A new public method,
LDClientInterface.flushAndWait(long timeout, TimeUnit unit). It sends what has been recorded and returns whether the events were delivered within the timeout.flush()is fire-and-forget, so an application that knows it is about to lose the chance to send (an uncaught exception handler, a move to the background) had no way to give its events that chance and find out whether it worked.flushAndWaitgives it both, with a bound the caller chooses.EventProcessorgainsblockingFlush(long, TimeUnit)as a default method, so custom implementations written before it still compile. The default falls back to the unboundedblockingFlush()and returns true.DirectEventProcessorimplements it. Delivery now reports its outcome: true if the events were accepted or there was nothing to send, false if offline, closed, or the send failed. If the timeout expires, the delivery is left running rather than cancelled, because the buffer has already been drained into the payload and interrupting it would only make the loss certain.LDClient.flushAndWaitshares one budget across them rather than giving each a fresh copy, so the timeout is the most the call can take.Describe alternatives you've considered
Make
flush()return aFuture. That changes the signature of an existing public method, and aFuturestill leaves the caller to pick a timeout and interpret the exceptions. A boolean with the timeout in the call is the shape a crash handler actually needs.Cancel the delivery on timeout. It frees the thread sooner but guarantees the events are lost, since they are no longer in the buffer. Letting it finish in the background gives them the best chance at no cost to the caller, who has already stopped waiting.
Additional context
Tests. Four new tests in
DirectEventProcessorTestcover delivered events, nothing to send, offline, and the timeout expiring first. The full unit suite and the instrumented suite (126 tests on an emulator) pass locally.Test app. Gains two controls that show what this tier can and cannot do:
FlushOnCrashHandlercallsflushAndWaitwith a two-second budget before handing the crash on, so the events are delivered.SIGKILL. Nothing can run, so the events are lost. That is the gap tier 3 (on-disk persistence) closes.Version. The new methods are marked
@since 5.17.0; adjust if the release lands under a different version.Note
Overview
Adds bounded, outcome-aware event flushing so apps can try to deliver analytics before the process exits or loses network access.
LDClient.flushAndWait(timeout, unit)(onLDClientInterface) starts delivery for every configured environment, waits on a single shared deadline, and returns whether events were delivered (or there were none). Closed clients, timeouts, offline mode, canceled custom futures, and negative-timeout edge cases are handled as documented failures without throwing.Event pipeline changes:
EventProcessorgainsflushAsync()(default wraps legacyblockingFlushviaLDFutures.fromBlockingCall).DirectEventProcessorcoalesces concurrent flush requests, reports success/failure through the future, tracks events lost by earlier deliveries or partial encoding (OutboundEventBuffer.Payload.isComplete()), and leaves in-flight HTTP work running when the caller stops waiting.close()reuses the same delivery path.LDCrashHandler(experimental) installs ahead of other uncaught-exception handlers and callsflushAndWaitbefore chaining to the previous handler.The test app adds immediate SIGKILL vs uncaught-crash controls to contrast unrecoverable loss with crash-handler rescue; tests cover multi-env parallelism, coalescing, HTTP failures, and serialization drops.
Reviewed by Cursor Bugbot for commit 7457a26. Bugbot is set up for automated code reviews on this repo. Configure here.