/
githubmirror
/
RxJava
Обзор
Документация
Войти
/
githubmirror
/
RxJava
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
4.x
src/test/java/io/reactivex/rxjava4/schedulers/ExecutorSchedulerTest.java
482 строки
15 KB
David Karnok
4.x: Convert from JUnit 4 to JUnit 6 (#8202)
30 июн 2026, 12:01
Не верифицирован
30 июн 2026, 12:01
531388f
Код
Авторство
О чём код?
/* * Copyright (c) 2016-present, RxJava Contributors. * * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in * compliance with the License. You may obtain a copy of the License at * * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software distributed under the License is * distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See * the License for the specific language governing permissions and limitations under the License. */ package io.reactivex.rxjava4.schedulers; import static org.junit.jupiter.api.Assertions.*; import java.lang.management.*; import java.util.List; import java.util.concurrent.*; import java.util.concurrent.atomic.*; import org.junit.jupiter.api.Test; import io.reactivex.rxjava4.core.Scheduler; import io.reactivex.rxjava4.core.Scheduler.Worker; import io.reactivex.rxjava4.disposables.Disposable; import io.reactivex.rxjava4.internal.disposables.*; import io.reactivex.rxjava4.internal.functions.Functions; import io.reactivex.rxjava4.internal.schedulers.RxThreadFactory; import io.reactivex.rxjava4.plugins.RxJavaPlugins; import io.reactivex.rxjava4.testsupport.TestHelper; public class ExecutorSchedulerTest extends AbstractSchedulerConcurrencyTests { static final Executor executor = Executors.newFixedThreadPool(2, new RxThreadFactory("TestCustomPool")); @Override protected Scheduler getScheduler() { return Schedulers.from(executor); } @Test public final void handledErrorIsNotDeliveredToThreadHandler() throws InterruptedException { SchedulerTestHelper.handledErrorIsNotDeliveredToThreadHandler(getScheduler()); } public static void cancelledRetention(Scheduler.Worker w, boolean periodic) throws InterruptedException { System.out.println("Wait before GC"); Thread.sleep(1000); System.out.println("GC"); System.gc(); Thread.sleep(1000); MemoryMXBean memoryMXBean = ManagementFactory.getMemoryMXBean(); long initial = memoryMXBean.getHeapMemoryUsage().getUsed(); System.out.printf("Starting: %.3f MB%n", initial / 1024.0 / 1024.0); int n = 100 * 1000; if (periodic) { final CountDownLatch cdl = new CountDownLatch(n); final Runnable action = cdl::countDown; for (int i = 0; i < n; i++) { if (i % 50000 == 0) { System.out.println(" -> still scheduling: " + i); } w.schedulePeriodically(action, 0, 1, TimeUnit.DAYS); } System.out.println("Waiting for the first round to finish..."); cdl.await(); } else { for (int i = 0; i < n; i++) { if (i % 50000 == 0) { System.out.println(" -> still scheduling: " + i); } w.schedule(Functions.EMPTY_RUNNABLE, 1, TimeUnit.DAYS); } } long after = memoryMXBean.getHeapMemoryUsage().getUsed(); System.out.printf("Peak: %.3f MB%n", after / 1024.0 / 1024.0); w.dispose(); System.out.println("Wait before second GC"); System.out.println("JDK 6 purge is N log N because it removes and shifts one by one"); int t = (int)(n * Math.log(n) / 100) + 1000; int sleepStep = 100; while (t > 0) { System.out.printf(" >> Waiting for purge: %.2f s remaining%n", t / 1000d); System.gc(); long finish = memoryMXBean.getHeapMemoryUsage().getUsed(); System.out.printf("After: %.3f MB%n", finish / 1024.0 / 1024.0); if (finish <= initial * 5) { break; } Thread.sleep(sleepStep); t -= sleepStep; } System.out.println("Second GC"); System.gc(); t = 2000; long finish = memoryMXBean.getHeapMemoryUsage().getUsed(); while (t > 0) { System.out.printf("After: %.3f MB%n", finish / 1024.0 / 1024.0); if (finish <= initial * 5) { return; } Thread.sleep(sleepStep); t -= sleepStep; finish = memoryMXBean.getHeapMemoryUsage().getUsed(); } fail(String.format("Tasks retained: %.3f -> %.3f -> %.3f", initial / 1024 / 1024.0, after / 1024 / 1024.0, finish / 1024 / 1024d)); } @Test public void cancelledTaskRetention() throws InterruptedException { ExecutorService exec = Executors.newSingleThreadExecutor(); Scheduler s = Schedulers.from(exec); try { Scheduler.Worker w = s.createWorker(); try { cancelledRetention(w, false); } finally { w.dispose(); } w = s.createWorker(); try { cancelledRetention(w, true); } finally { w.dispose(); } } finally { exec.shutdownNow(); } } /** A simple executor which queues tasks and executes them one-by-one if executeOne() is called. */ static final class TestExecutor implements Executor { final ConcurrentLinkedQueue<Runnable> queue = new ConcurrentLinkedQueue<>(); @Override public void execute(Runnable command) { queue.offer(command); } public void executeOne() { Runnable r = queue.poll(); if (r != null) { r.run(); } } public void executeAll() { Runnable r; while ((r = queue.poll()) != null) { r.run(); } } } @Test public void cancelledTasksDontRun() { final AtomicInteger calls = new AtomicInteger(); Runnable task = calls::getAndIncrement; TestExecutor exec = new TestExecutor(); Scheduler custom = Schedulers.from(exec); Worker w = custom.createWorker(); try { Disposable d1 = w.schedule(task); Disposable d2 = w.schedule(task); Disposable d3 = w.schedule(task); d1.dispose(); d2.dispose(); d3.dispose(); exec.executeAll(); assertEquals(0, calls.get()); } finally { w.dispose(); } } @Test public void cancelledWorkerDoesntRunTasks() { final AtomicInteger calls = new AtomicInteger(); Runnable task = calls::getAndIncrement; TestExecutor exec = new TestExecutor(); Scheduler custom = Schedulers.from(exec); Worker w = custom.createWorker(); try { w.schedule(task); w.schedule(task); w.schedule(task); } finally { w.dispose(); } exec.executeAll(); assertEquals(0, calls.get()); } @Test public void plainExecutor() throws Exception { Scheduler s = Schedulers.from(Runnable::run); final CountDownLatch cdl = new CountDownLatch(5); Runnable r = cdl::countDown; s.scheduleDirect(r); s.scheduleDirect(r, 50, TimeUnit.MILLISECONDS); Disposable d = s.schedulePeriodicallyDirect(r, 10, 10, TimeUnit.MILLISECONDS); try { assertTrue(cdl.await(5, TimeUnit.SECONDS)); } finally { d.dispose(); } assertTrue(d.isDisposed()); } @Test public void rejectingExecutor() { ScheduledExecutorService exec = Executors.newSingleThreadScheduledExecutor(); exec.shutdown(); Scheduler s = Schedulers.from(exec); List<Throwable> errors = TestHelper.trackPluginErrors(); try { assertSame(EmptyDisposable.INSTANCE, s.scheduleDirect(Functions.EMPTY_RUNNABLE)); assertSame(EmptyDisposable.INSTANCE, s.scheduleDirect(Functions.EMPTY_RUNNABLE, 10, TimeUnit.MILLISECONDS)); assertSame(EmptyDisposable.INSTANCE, s.schedulePeriodicallyDirect(Functions.EMPTY_RUNNABLE, 10, 10, TimeUnit.MILLISECONDS)); TestHelper.assertUndeliverable(errors, 0, RejectedExecutionException.class); TestHelper.assertUndeliverable(errors, 1, RejectedExecutionException.class); TestHelper.assertUndeliverable(errors, 2, RejectedExecutionException.class); } finally { RxJavaPlugins.reset(); } } @Test public void rejectingExecutorWorker() { ScheduledExecutorService exec = Executors.newSingleThreadScheduledExecutor(); exec.shutdown(); List<Throwable> errors = TestHelper.trackPluginErrors(); try { Worker s = Schedulers.from(exec).createWorker(); assertSame(EmptyDisposable.INSTANCE, s.schedule(Functions.EMPTY_RUNNABLE)); s = Schedulers.from(exec).createWorker(); assertSame(EmptyDisposable.INSTANCE, s.schedule(Functions.EMPTY_RUNNABLE, 10, TimeUnit.MILLISECONDS)); s = Schedulers.from(exec).createWorker(); assertSame(EmptyDisposable.INSTANCE, s.schedulePeriodically(Functions.EMPTY_RUNNABLE, 10, 10, TimeUnit.MILLISECONDS)); TestHelper.assertUndeliverable(errors, 0, RejectedExecutionException.class); TestHelper.assertUndeliverable(errors, 1, RejectedExecutionException.class); TestHelper.assertUndeliverable(errors, 2, RejectedExecutionException.class); } finally { RxJavaPlugins.reset(); } } @Test public void reuseScheduledExecutor() throws Exception { ScheduledExecutorService exec = Executors.newSingleThreadScheduledExecutor(); try { Scheduler s = Schedulers.from(exec); final CountDownLatch cdl = new CountDownLatch(8); Runnable r = cdl::countDown; s.scheduleDirect(r); s.scheduleDirect(r, 10, TimeUnit.MILLISECONDS); Disposable d = s.schedulePeriodicallyDirect(r, 10, 10, TimeUnit.MILLISECONDS); try { assertTrue(cdl.await(5, TimeUnit.SECONDS)); } finally { d.dispose(); } } finally { exec.shutdown(); } } @Test public void reuseScheduledExecutorAsWorker() throws Exception { ScheduledExecutorService exec = Executors.newSingleThreadScheduledExecutor(); Worker s = Schedulers.from(exec).createWorker(); assertFalse(s.isDisposed()); try { final CountDownLatch cdl = new CountDownLatch(8); Runnable r = cdl::countDown; s.schedule(r); s.schedule(r, 10, TimeUnit.MILLISECONDS); Disposable d = s.schedulePeriodically(r, 10, 10, TimeUnit.MILLISECONDS); try { assertTrue(cdl.await(5, TimeUnit.SECONDS)); } finally { d.dispose(); } } finally { s.dispose(); exec.shutdown(); } assertTrue(s.isDisposed()); } @Test public void disposeRace() { ExecutorService exec = Executors.newSingleThreadExecutor(); final Scheduler s = Schedulers.from(exec); try { for (int i = 0; i < 500; i++) { final Worker w = s.createWorker(); final AtomicInteger c = new AtomicInteger(2); w.schedule(() -> { c.decrementAndGet(); while (c.get() != 0) { } }); c.decrementAndGet(); while (c.get() != 0) { } w.dispose(); } } finally { exec.shutdownNow(); } } @Test public void runnableDisposed() { final Scheduler s = Schedulers.from(Runnable::run); Disposable d = s.scheduleDirect(Functions.EMPTY_RUNNABLE); assertTrue(d.isDisposed()); } @Test public void runnableDisposedAsync() throws Exception { final Scheduler s = Schedulers.from(r -> new Thread(r).start()); Disposable d = s.scheduleDirect(Functions.EMPTY_RUNNABLE); while (!d.isDisposed()) { Thread.sleep(1); } } @Test public void runnableDisposedAsync2() throws Exception { final Scheduler s = Schedulers.from(executor); Disposable d = s.scheduleDirect(Functions.EMPTY_RUNNABLE); while (!d.isDisposed()) { Thread.sleep(1); } } @Test public void runnableDisposedAsyncCrash() throws Exception { final Scheduler s = Schedulers.from(r -> new Thread(r).start()); Disposable d = s.scheduleDirect(() -> { throw new IllegalStateException(); }); while (!d.isDisposed()) { Thread.sleep(1); } } @Test public void runnableDisposedAsyncTimed() throws Exception { final Scheduler s = Schedulers.from(r -> new Thread(r).start()); Disposable d = s.scheduleDirect(Functions.EMPTY_RUNNABLE, 1, TimeUnit.MILLISECONDS); while (!d.isDisposed()) { Thread.sleep(1); } } @Test public void runnableDisposedAsyncTimed2() throws Exception { ExecutorService executorScheduler = Executors.newScheduledThreadPool(1, new RxThreadFactory("TestCustomPoolTimed")); try { final Scheduler s = Schedulers.from(executorScheduler); Disposable d = s.scheduleDirect(Functions.EMPTY_RUNNABLE, 1, TimeUnit.MILLISECONDS); while (!d.isDisposed()) { Thread.sleep(1); } } finally { executorScheduler.shutdownNow(); } } @Test public void unwrapScheduleDirectTaskAfterDispose() { Scheduler scheduler = getScheduler(); final CountDownLatch cdl = new CountDownLatch(1); Runnable countDownRunnable = cdl::countDown; Disposable disposable = scheduler.scheduleDirect(countDownRunnable, 100, TimeUnit.MILLISECONDS); SchedulerRunnableIntrospection wrapper = (SchedulerRunnableIntrospection) disposable; assertSame(countDownRunnable, wrapper.getWrappedRunnable()); disposable.dispose(); assertSame(Functions.EMPTY_RUNNABLE, wrapper.getWrappedRunnable()); } @Test public void interruptibleRunnableRunDisposeRace() { ExecutorService exec = Executors.newSingleThreadExecutor(); try { Scheduler s = Schedulers.from(exec::execute, true); for (int i = 0; i < TestHelper.RACE_DEFAULT_LOOPS; i++) { SequentialDisposable sd = new SequentialDisposable(); TestHelper.race( () -> sd.update(s.scheduleDirect(() -> { })), sd::dispose ); } } finally { exec.shutdown(); } } @Test public void interruptibleRunnableRunDispose() { try { for (int i = 0; i < TestHelper.RACE_DEFAULT_LOOPS; i++) { AtomicReference<Runnable> runRef = new AtomicReference<>(); Scheduler s = Schedulers.from(runRef::set, true); Disposable d = s.scheduleDirect(() -> { }); TestHelper.race( () -> runRef.get().run(), d::dispose ); } } finally { Thread.interrupted(); } } }