From c28196dfa8d56cc2ba543575a13eaf48c9f2afce Mon Sep 17 00:00:00 2001 From: cuishuang Date: Sun, 6 Sep 2026 00:56:54 +0800 Subject: [PATCH] Fix timeout handling in SchedulerToExecutorService.awaitTermination Signed-off-by: cuishuang --- .../SchedulerToExecutorService.java | 26 +++++++++++++---- .../SchedulerToExecutorServiceTest.java | 28 +++++++++++++++++++ 2 files changed, 49 insertions(+), 5 deletions(-) diff --git a/src/main/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorService.java b/src/main/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorService.java index a6c86194e7..0205a1cc45 100644 --- a/src/main/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorService.java +++ b/src/main/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorService.java @@ -84,13 +84,29 @@ public boolean isTerminated() { @Override public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException { // FIXME no idea how to passively wait, not really applicable in Rx - long totalTime = unit.convert(timeout, TimeUnit.MILLISECONDS); + long timeoutNanos = unit.toNanos(timeout); + if (isTerminated()) { + return true; + } + if (timeoutNanos <= 0) { + return false; + } + + long start = System.nanoTime(); + for (;;) { + if (isTerminated()) { + return true; + } + + long elapsed = System.nanoTime() - start; + long remaining = timeoutNanos - elapsed; + if (remaining <= 0) { + return isTerminated(); + } - while (!isTerminated() && totalTime > 0) { - totalTime--; - Thread.sleep(1); + // There is no termination signal to await, so poll at a short interval. + TimeUnit.NANOSECONDS.sleep(Math.min(remaining, TimeUnit.MILLISECONDS.toNanos(1))); } - return totalTime > 0; } @Override diff --git a/src/test/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorServiceTest.java b/src/test/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorServiceTest.java index c37cafec91..f42d89e452 100644 --- a/src/test/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorServiceTest.java +++ b/src/test/java/io/reactivex/rxjava4/internal/schedulers/SchedulerToExecutorServiceTest.java @@ -216,4 +216,32 @@ public void invokeAllTimeoutDoesTimeout() throws Throwable { assertTrue(f.isCancelled(), "Task was not cancelled: " + f); } } + + @Test + public void awaitTerminationUsesRequestedTimeUnit() throws Exception { + Scheduler scheduler = Schedulers.computation(); + Scheduler.Worker worker = scheduler.createWorker(); + SchedulerToExecutorService executor = new SchedulerToExecutorService( + scheduler, new AtomicReference<>(worker)); + + try { + scheduler.scheduleDirect(worker::dispose, 50, TimeUnit.MILLISECONDS); + + assertTrue(executor.awaitTermination(1, TimeUnit.SECONDS)); + } finally { + worker.dispose(); + } + } + + @Test + public void awaitTerminationReturnsTrueIfAlreadyTerminated() throws Exception { + SchedulerToExecutorService executor = new SchedulerToExecutorService( + Schedulers.computation(), new AtomicReference<>(null)); + + assertFalse(executor.awaitTermination(0, TimeUnit.MILLISECONDS)); + + executor.shutdown(); + + assertTrue(executor.awaitTermination(0, TimeUnit.MILLISECONDS)); + } }