/
githubmirror
/
RxJava
Обзор
Документация
Войти
/
githubmirror
/
RxJava
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
4.x
src/test/java/io/reactivex/rxjava4/schedulers/SchedulerWorkerTest.java
145 строк
4 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.assertTrue; import java.util.*; import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.Test; import io.reactivex.rxjava4.annotations.NonNull; import io.reactivex.rxjava4.core.*; import io.reactivex.rxjava4.disposables.Disposable; public class SchedulerWorkerTest extends RxJavaTest { static final class CustomDriftScheduler extends Scheduler { public volatile long drift; @NonNull @Override public Worker createWorker() { final Worker w = Schedulers.computation().createWorker(); return new Worker() /* NFI */ { @Override public void dispose() { w.dispose(); } @Override public boolean isDisposed() { return w.isDisposed(); } @NonNull @Override public Disposable schedule(@NonNull Runnable action) { return w.schedule(action); } @NonNull @Override public Disposable schedule(@NonNull Runnable action, long delayTime, @NonNull TimeUnit unit) { return w.schedule(action, delayTime, unit); } @Override public long now(TimeUnit unit) { return super.now(unit) + unit.convert(drift, TimeUnit.NANOSECONDS); } }; } @Override public long now(@NonNull TimeUnit unit) { return super.now(unit) + unit.convert(drift, TimeUnit.NANOSECONDS); } } @Test public void currentTimeDriftBackwards() throws Exception { CustomDriftScheduler s = new CustomDriftScheduler(); Scheduler.Worker w = s.createWorker(); try { final List<Long> times = new ArrayList<>(); Disposable d = w.schedulePeriodically(() -> times.add(System.currentTimeMillis()), 100, 100, TimeUnit.MILLISECONDS); Thread.sleep(150); s.drift = -TimeUnit.SECONDS.toNanos(1) - Scheduler.clockDriftTolerance(); Thread.sleep(400); d.dispose(); Thread.sleep(150); System.out.println("Runs: " + times.size()); for (int i = 0; i < times.size() - 1 ; i++) { long diff = times.get(i + 1) - times.get(i); System.out.println("Diff #" + i + ": " + diff); assertTrue(diff < 150 && diff > 50, "" + i + ":" + diff); } assertTrue(times.size() > 2, "Too few invocations: " + times.size()); } finally { w.dispose(); } } @Test public void currentTimeDriftForwards() throws Exception { CustomDriftScheduler s = new CustomDriftScheduler(); Scheduler.Worker w = s.createWorker(); try { final List<Long> times = new ArrayList<>(); Disposable d = w.schedulePeriodically(() -> times.add(System.currentTimeMillis()), 100, 100, TimeUnit.MILLISECONDS); Thread.sleep(150); s.drift = TimeUnit.SECONDS.toNanos(1) + Scheduler.clockDriftTolerance(); Thread.sleep(400); d.dispose(); Thread.sleep(150); System.out.println("Runs: " + times.size()); assertTrue(times.size() > 2); for (int i = 0; i < times.size() - 1 ; i++) { long diff = times.get(i + 1) - times.get(i); System.out.println("Diff #" + i + ": " + diff); assertTrue(diff < 250 && diff > 50, "Diff out of range: " + diff); } } finally { w.dispose(); } } }