Skip to content

Commit 71a0ee9

Browse files
committed
Merge remote-tracking branch 'origin/master' into ci/harden-build-publish
2 parents 2c48dbc + 6227664 commit 71a0ee9

4 files changed

Lines changed: 362 additions & 79 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,11 @@
1+
# Unreleased
2+
- [Change] Batch upload tasks are handed to the network executor with `execute()` rather than `submit()`, so `shutdownNow()` returns them and the batches inside can be reported. A `Callback` or `Log` implementation that throws is now caught and logged rather than being absorbed by the discarded `Future`; it no longer reaches the thread's uncaught-exception handler either. Note that a caller supplying a `ForkJoinPool` or `ScheduledThreadPoolExecutor` through `Analytics.Builder.networkExecutor` still gets no shutdown callbacks: the former returns nothing from `shutdownNow()`, the latter wraps tasks regardless.
3+
- [Note] `maxRateLimitDuration` is measured against the system clock, so a large clock adjustment during a rate-limit episode can shorten or extend it.
4+
- [Fix] Batches abandoned at shutdown now report failure through `Callback`. Note this is a failure notification for a batch that may in fact have been delivered: a batch interrupted after its request went out is reported as failed because the client cannot know the outcome. Delivery has always been at-least-once; this makes the uncertainty visible rather than silent, so a `Callback` that counts failures will see some that were not lost. A batch interrupted while waiting out a `Retry-After` or a backoff returned without notifying anyone, and batches still queued when the executor stopped were counted in a log line and otherwise discarded. Both were unreachable while the network executor was never stopped; both became reachable with the fix below.
5+
- [Fix] `shutdown()` now stops the network executor rather than leaving it running. It was asked to stop and then given up to 75 seconds to finish on its own; if it had not, `shutdown()` returned anyway. Because `ExecutorService.shutdown()` does not interrupt running tasks, a thread waiting out a `Retry-After` kept running, and as these threads are not daemon threads, the JVM would not exit. `shutdown()` reported success while this happened. It now interrupts the executor once the timeout elapses, including when shutdown is itself interrupted.
6+
- [Change] `maxRateLimitDuration` now defaults to 30 minutes, was 12 hours. Retries prompted by a `Retry-After` header do not consume the retry count, so this duration is the only limit on how long they continue; at 12 hours a server that kept sending the header could hold a batch for that long, and because uploads run on a single thread, hold every other batch behind it. Pass a longer value to `Analytics.Builder.maxRateLimitDuration` to restore the previous limit.
7+
- [Fix] A `Retry-After` that will not fit in what is left of `maxRateLimitDuration` ends the episode rather than being shortened to fit. The budget was tested and then the full `Retry-After` slept on top of it, so an episode could run past its limit by up to one wait; shortening it instead would resume inside the window the server named, which it has already declined to serve, and the budget would be spent by then regardless.
8+
19
# Version 3.5.5 (June 30, 2026)
210
- [New](https://github.com/segmentio/analytics-java/pull/531) Unified HTTP response handling and retry behavior
311
- Retryable statuses (429, 408, 410, 460, 5xx except 501/505/511) check Retry-After header first, fall back to exponential backoff

‎analytics/src/main/java/com/segment/analytics/Analytics.java‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -455,7 +455,14 @@ public Analytics build() {
455455
maxTotalBackoffDurationMs = 43200 * 1000L; // 12 hours
456456
}
457457
if (maxRateLimitDurationMs == 0) {
458-
maxRateLimitDurationMs = 43200 * 1000L; // 12 hours
458+
// Retry-After retries are deliberately uncounted, so this duration is the
459+
// only thing bounding them. The network executor is single-threaded, so it
460+
// also bounds how long one stuck batch holds up every other one.
461+
//
462+
// Deliberately several times the Retry-After cap. When the two are equal a
463+
// single maximal Retry-After consumes the whole budget, and because the
464+
// elapsed check runs before the wait the batch is dropped after one attempt.
465+
maxRateLimitDurationMs = 1800 * 1000L; // 30 minutes
459466
}
460467

461468
HttpLoggingInterceptor interceptor =

‎analytics/src/main/java/com/segment/analytics/internal/AnalyticsClient.java‎

Lines changed: 138 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -53,8 +53,14 @@ public class AnalyticsClient {
5353
private static final String instanceId = UUID.randomUUID().toString();
5454
private static final int WAIT_FOR_THREAD_COMPLETE_S = 5;
5555
private static final int TERMINATION_TIMEOUT_S = 1;
56-
private static final int NETWORK_TERMINATION_TIMEOUT_S =
57-
75; // base Retry-After cap is 60s + headroom
56+
// Deliberately shorter than a maximal Retry-After wait. Shutdown does not wait out
57+
// the retry schedule; it interrupts it, and this is only the grace period before it
58+
// does.
59+
private static final int NETWORK_TERMINATION_TIMEOUT_S = 75;
60+
// A guard against an absurd header, not a second budget: waiting less than the server
61+
// asked for does not make the next attempt more likely to succeed, it just sends more
62+
// requests at something already rate-limiting us. How long we keep trying is
63+
// maxRateLimitDuration's job.
5864
private static final long MAX_RATE_LIMITED_SECONDS = 300L;
5965

6066
static {
@@ -259,24 +265,32 @@ public void flush() {
259265
}
260266
}
261267

262-
synchronized void setRateLimitState(long retryAfterSeconds) {
268+
/** Returns the clock reading it used, so a caller can measure from the same instant. */
269+
synchronized long setRateLimitState(long retryAfterSeconds) {
263270
long now = System.currentTimeMillis();
264271
if (rateLimitStartTime == 0) {
265272
rateLimitStartTime = now;
266273
}
267274
rateLimitWaitUntil = now + (retryAfterSeconds * 1000);
268275
rateLimited = true;
276+
return now;
269277
}
270278

271279
/**
272-
* Sets rate-limit state and atomically checks whether maxRateLimitDuration has been exceeded.
273-
* Returns true if the duration has been exceeded and the batch should be dropped.
280+
* Sets rate-limit state and returns how much of {@code maxRateLimitDuration} is left, in
281+
* milliseconds. Zero or less means the budget is spent.
282+
*
283+
* <p>Returning the remaining time rather than a boolean lets one clock reading serve both the
284+
* budget test and the wait that follows it. Testing and then sleeping a full Retry-After on top
285+
* would otherwise overshoot the budget by up to that much.
274286
*/
275-
synchronized boolean setRateLimitStateAndCheckDuration(
287+
synchronized long setRateLimitStateAndRemaining(
276288
long retryAfterSeconds, long maxRateLimitDurationMs) {
277-
setRateLimitState(retryAfterSeconds);
278-
return rateLimitStartTime > 0
279-
&& System.currentTimeMillis() - rateLimitStartTime > maxRateLimitDurationMs;
289+
long now = setRateLimitState(retryAfterSeconds);
290+
if (rateLimitStartTime <= 0) {
291+
return maxRateLimitDurationMs;
292+
}
293+
return maxRateLimitDurationMs - (now - rateLimitStartTime);
280294
}
281295

282296
synchronized void clearRateLimitState() {
@@ -337,8 +351,25 @@ private void waitForLooperCompletion() {
337351
}
338352
}
339353

354+
/**
355+
* Reports batches that were queued and never attempted.
356+
*
357+
* <p>Only reaches tasks the executor hands back as they were submitted. A {@code ForkJoinPool}
358+
* returns an empty list from {@code shutdownNow()} whatever is queued, and a {@code
359+
* ScheduledThreadPoolExecutor} wraps even {@code execute()}, so a caller supplying either through
360+
* {@code Analytics.Builder#networkExecutor} gets no callbacks here and a task count that reads
361+
* zero.
362+
*/
363+
private void notifyDroppedBatches(List<Runnable> dropped) {
364+
for (Runnable task : dropped) {
365+
if (task instanceof BatchUploadTask) {
366+
((BatchUploadTask) task)
367+
.notifyDropped(new IOException("Dropped at shutdown without being attempted"));
368+
}
369+
}
370+
}
371+
340372
public void shutdownAndWait(ExecutorService executor, String name) {
341-
boolean isLooperExecutor = name != null && name.equalsIgnoreCase("looper");
342373
boolean isNetworkExecutor = name != null && name.equalsIgnoreCase("network");
343374
int timeoutSeconds = isNetworkExecutor ? NETWORK_TERMINATION_TIMEOUT_S : TERMINATION_TIMEOUT_S;
344375
try {
@@ -348,49 +379,64 @@ public void shutdownAndWait(ExecutorService executor, String name) {
348379
log.print(VERBOSE, "%s executor terminated normally.", name);
349380
return;
350381
}
351-
if (isLooperExecutor) { // Handle looper - network should finish on its own
352-
// not terminated within timeout -> force shutdown
353-
log.print(
354-
VERBOSE,
355-
"%s did not terminate in %d seconds; requesting shutdownNow().",
356-
name,
357-
TERMINATION_TIMEOUT_S);
358-
List<Runnable> dropped = executor.shutdownNow(); // interrupts running tasks
359-
log.print(
360-
VERBOSE,
361-
"%s shutdownNow returned %d queued tasks that never started.",
362-
name,
363-
dropped.size());
364-
365-
// optional short wait to give interrupted tasks a chance to exit
366-
boolean terminatedAfterForce =
367-
executor.awaitTermination(TERMINATION_TIMEOUT_S, TimeUnit.SECONDS);
368-
log.print(
369-
VERBOSE,
370-
"%s executor %s after shutdownNow().",
371-
name,
372-
terminatedAfterForce ? "terminated" : "still running (did not terminate)");
373382

374-
if (!terminatedAfterForce) {
375-
// final warning — investigate tasks that ignore interrupts
376-
log.print(
377-
ERROR,
378-
"%s executor still did not terminate; tasks may be ignoring interrupts.",
379-
name);
380-
}
383+
// Both executors are force-stopped, the network one included. shutdown() does
384+
// not interrupt a running task, its task can be a whole rate-limit budget deep
385+
// in a sleep, and these threads are non-daemon — so leaving it to finish on its
386+
// own lets shutdown() return while a thread holds the JVM open.
387+
//
388+
// The interrupt only reaches a thread parked in a sleep. OkHttp's reads are
389+
// governed by SO_TIMEOUT, so a thread inside the HTTP call is bounded by the
390+
// client's own timeouts instead, and by nothing at all if a caller supplies an
391+
// OkHttpClient without them.
392+
log.print(
393+
VERBOSE,
394+
"%s did not terminate in %d seconds; requesting shutdownNow().",
395+
name,
396+
timeoutSeconds);
397+
List<Runnable> dropped = executor.shutdownNow(); // interrupts running tasks
398+
log.print(
399+
VERBOSE,
400+
"%s shutdownNow returned %d queued tasks that never started.",
401+
name,
402+
dropped.size());
403+
404+
// Submitted and never run, so their callbacks are still owed. Counting them in
405+
// a log line is not the same as telling the caller the messages did not go.
406+
notifyDroppedBatches(dropped);
407+
408+
// optional short wait to give interrupted tasks a chance to exit
409+
boolean terminatedAfterForce =
410+
executor.awaitTermination(TERMINATION_TIMEOUT_S, TimeUnit.SECONDS);
411+
log.print(
412+
VERBOSE,
413+
"%s executor %s after shutdownNow().",
414+
name,
415+
terminatedAfterForce ? "terminated" : "still running (did not terminate)");
416+
417+
if (!terminatedAfterForce) {
418+
// final warning — investigate tasks that ignore interrupts
419+
log.print(
420+
ERROR, "%s executor still did not terminate; tasks may be ignoring interrupts.", name);
381421
}
382422
} catch (InterruptedException e) {
383423
// Preserve interrupt status and attempt forceful shutdown
384424
log.print(ERROR, e, "Interrupted while stopping %s executor.", name);
425+
// Same reasoning as above: this applied to the looper only, leaving the network
426+
// executor running after an interrupted shutdown.
427+
List<Runnable> dropped = executor.shutdownNow();
428+
log.print(
429+
VERBOSE,
430+
"%s shutdownNow invoked after interrupt; %d tasks returned.",
431+
name,
432+
dropped.size());
433+
// These are owed a callback just as much as the ones dropped above; an
434+
// interrupted shutdown is still a shutdown. Reported before the interrupt flag
435+
// goes back on: callbacks run on this thread, and with the flag already set any
436+
// interruptible call inside one -- a queue put, an await, a Future.get -- throws
437+
// InterruptedException the moment it starts.
438+
notifyDroppedBatches(dropped);
385439
Thread.currentThread().interrupt();
386-
if (isLooperExecutor) {
387-
List<Runnable> dropped = executor.shutdownNow();
388-
log.print(
389-
VERBOSE,
390-
"%s shutdownNow invoked after interrupt; %d tasks returned.",
391-
name,
392-
dropped.size());
393-
}
394440
}
395441
}
396442

@@ -468,7 +514,11 @@ public void run() {
468514
batch.batch().size(),
469515
batch.sequence());
470516
try {
471-
networkExecutor.submit(
517+
// execute, not submit: submit wraps the task in a FutureTask, and the
518+
// work queue then holds that wrapper, so shutdownNow() hands back
519+
// FutureTasks and the batches inside them cannot be identified or
520+
// reported. The Future was discarded anyway.
521+
networkExecutor.execute(
472522
BatchUploadTask.create(AnalyticsClient.this, batch, maximumRetries));
473523
} catch (RejectedExecutionException e) {
474524
log.print(
@@ -541,6 +591,11 @@ static BatchUploadTask create(AnalyticsClient client, Batch batch, int maxRetrie
541591
this.maxRetries = maxRetries;
542592
}
543593

594+
/** Reports a batch that was discarded from the queue without ever being attempted. */
595+
void notifyDropped(Exception exception) {
596+
notifyCallbacksWithException(batch, exception);
597+
}
598+
544599
private void notifyCallbacksWithException(Batch batch, Exception exception) {
545600
for (Message message : batch.batch()) {
546601
for (Callback callback : client.callbacks) {
@@ -682,6 +737,22 @@ private static Long parseRetryAfterSeconds(String headerValue) {
682737

683738
@Override
684739
public void run() {
740+
// Handed to the executor with execute() rather than submit(), so that
741+
// shutdownNow() gives back the task itself and the batch inside it can be
742+
// reported. That also means nothing catches what escapes here: under submit()
743+
// a FutureTask absorbed it into a result nobody read, and the worker survived.
744+
// Callback and Log are supplied by the caller and are invoked below outside any
745+
// try, so one that throws would now kill and replace the pool's worker and reach
746+
// the application's uncaught-exception handler -- from a library that could not
747+
// previously raise one. Keep that property.
748+
try {
749+
runUploadLoop();
750+
} catch (Throwable t) {
751+
client.log.print(ERROR, t, "Batch %s upload task failed unexpectedly.", batch.sequence());
752+
}
753+
}
754+
755+
private void runUploadLoop() {
685756
int totalAttempts = 0; // counts every HTTP attempt (for header and error message)
686757
int backoffAttempts = 0; // counts attempts that consume backoff-based retries
687758
int maxBackoffAttempts = maxRetries + 1; // preserve existing semantics
@@ -697,11 +768,11 @@ public void run() {
697768
}
698769

699770
if (result.strategy == RetryStrategy.RATE_LIMITED) {
700-
// Atomically set rate-limit state and check whether maxRateLimitDuration is exceeded.
701-
boolean durationExceeded =
702-
client.setRateLimitStateAndCheckDuration(
771+
// Atomically set rate-limit state and take what is left of maxRateLimitDuration.
772+
long remainingMs =
773+
client.setRateLimitStateAndRemaining(
703774
result.retryAfterSeconds, client.maxRateLimitDurationMs);
704-
if (durationExceeded) {
775+
if (remainingMs <= 0) {
705776
client.clearRateLimitState();
706777
break;
707778
}
@@ -711,15 +782,28 @@ public void run() {
711782
break;
712783
}
713784

785+
long retryAfterMs = TimeUnit.SECONDS.toMillis(result.retryAfterSeconds);
786+
if (retryAfterMs > remainingMs) {
787+
// A wait that will not fit ends the episode. Shortening it would resume
788+
// inside the window the server named -- one it has already said it will
789+
// not serve -- and the budget is spent by then, so that attempt would be
790+
// the last either way.
791+
client.clearRateLimitState();
792+
break;
793+
}
794+
714795
try {
715-
TimeUnit.SECONDS.sleep(result.retryAfterSeconds);
796+
TimeUnit.MILLISECONDS.sleep(retryAfterMs);
716797
} catch (InterruptedException e) {
717798
client.log.print(
718799
DEBUG,
719800
"Thread interrupted while waiting for Retry-After for batch %s.",
720801
batch.sequence());
721802
client.clearRateLimitState();
722803
Thread.currentThread().interrupt();
804+
// Every exit from this loop reports the batch. Returning without this
805+
// loses it silently, with no callback at all.
806+
notifyCallbacksWithException(batch, new IOException("Interrupted during shutdown", e));
723807
return;
724808
}
725809
// Retry-After does not count against maxRetries.
@@ -743,6 +827,7 @@ public void run() {
743827
client.log.print(
744828
DEBUG, "Thread interrupted while backing off for batch %s.", batch.sequence());
745829
Thread.currentThread().interrupt();
830+
notifyCallbacksWithException(batch, new IOException("Interrupted during shutdown", e));
746831
return;
747832
}
748833
}

0 commit comments

Comments
 (0)