/
githubmirror
/
RxJava
Обзор
Документация
Войти
/
githubmirror
/
RxJava
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
4.x
src/test/java/io/reactivex/rxjava4/internal/util/QueueDrainHelperTest.java
840 строк
21 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.internal.util; import static org.junit.jupiter.api.Assertions.*; import java.io.IOException; import java.util.*; import java.util.concurrent.Flow.*; import java.util.concurrent.atomic.AtomicLong; import org.junit.jupiter.api.Test; import io.reactivex.rxjava4.core.Observer; import io.reactivex.rxjava4.core.RxJavaTest; import io.reactivex.rxjava4.disposables.Disposable; import io.reactivex.rxjava4.exceptions.*; import io.reactivex.rxjava4.functions.BooleanSupplier; import io.reactivex.rxjava4.internal.subscriptions.BooleanSubscription; import io.reactivex.rxjava4.observers.TestObserver; import io.reactivex.rxjava4.operators.SpscArrayQueue; import io.reactivex.rxjava4.subscribers.TestSubscriber; import io.reactivex.rxjava4.testsupport.TestHelper; public class QueueDrainHelperTest extends RxJavaTest { @Test public void isCancelled() { assertTrue(QueueDrainHelper.isCancelled(() -> { throw new IOException(); })); } @Test public void requestMaxInt() { QueueDrainHelper.request(new Subscription() /* NFI */ { @Override public void request(long n) { assertEquals(Integer.MAX_VALUE, n); } @Override public void cancel() { } }, Integer.MAX_VALUE); } @Test public void requestMinInt() { QueueDrainHelper.request(new Subscription() /* NFI */ { @Override public void request(long n) { assertEquals(Long.MAX_VALUE, n); } @Override public void cancel() { } }, Integer.MIN_VALUE); } @Test public void requestAlmostMaxInt() { QueueDrainHelper.request(new Subscription() /* NFI */ { @Override public void request(long n) { assertEquals(Integer.MAX_VALUE - 1, n); } @Override public void cancel() { } }, Integer.MAX_VALUE - 1); } @Test public void postCompleteEmpty() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ArrayDeque<Integer> queue = new ArrayDeque<>(); AtomicLong state = new AtomicLong(); BooleanSupplier isCancelled = () -> false; ts.onSubscribe(new BooleanSubscription()); QueueDrainHelper.postComplete(ts, queue, state, isCancelled); ts.assertResult(); } @Test public void postCompleteWithRequest() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ArrayDeque<Integer> queue = new ArrayDeque<>(); AtomicLong state = new AtomicLong(); BooleanSupplier isCancelled = () -> false; ts.onSubscribe(new BooleanSubscription()); queue.offer(1); state.getAndIncrement(); QueueDrainHelper.postComplete(ts, queue, state, isCancelled); ts.assertResult(1); } @Test public void completeRequestRace() { for (int i = 0; i < TestHelper.RACE_DEFAULT_LOOPS; i++) { final TestSubscriber<Integer> ts = new TestSubscriber<>(); final ArrayDeque<Integer> queue = new ArrayDeque<>(); final AtomicLong state = new AtomicLong(); final BooleanSupplier isCancelled = () -> false; ts.onSubscribe(new BooleanSubscription()); queue.offer(1); Runnable r1 = () -> QueueDrainHelper.postCompleteRequest(1, ts, queue, state, isCancelled); Runnable r2 = () -> QueueDrainHelper.postComplete(ts, queue, state, isCancelled); TestHelper.race(r1, r2); ts.assertResult(1); } } @Test public void postCompleteCancelled() { final TestSubscriber<Integer> ts = new TestSubscriber<>(); ArrayDeque<Integer> queue = new ArrayDeque<>(); AtomicLong state = new AtomicLong(); BooleanSupplier isCancelled = ts::isCancelled; ts.onSubscribe(new BooleanSubscription()); queue.offer(1); state.getAndIncrement(); ts.cancel(); QueueDrainHelper.postComplete(ts, queue, state, isCancelled); ts.assertEmpty(); } @Test public void postCompleteCancelledAfterOne() { var ts = new TestSubscriber<Integer>() /* NFI */ { @Override public void onNext(Integer t) { super.onNext(t); cancel(); } }; ArrayDeque<Integer> queue = new ArrayDeque<>(); AtomicLong state = new AtomicLong(); BooleanSupplier isCancelled = ts::isCancelled; ts.onSubscribe(new BooleanSubscription()); queue.offer(1); state.getAndIncrement(); QueueDrainHelper.postComplete(ts, queue, state, isCancelled); ts.assertValue(1).assertNoErrors().assertNotComplete(); } @Test public void drainMaxLoopMissingBackpressure() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ts.onSubscribe(new BooleanSubscription()); var qd = new QueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public long requested() { return 0; } @Override public long produced(long n) { return 0; } @Override public int leave(int m) { return 0; } @Override public boolean accept(Subscriber<? super Integer> a, Integer v) { return false; } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); q.offer(1); QueueDrainHelper.drainMaxLoop(q, ts, false, null, qd); ts.assertFailure(MissingBackpressureException.class); } @Test public void drainMaxLoopMissingBackpressureWithResource() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ts.onSubscribe(new BooleanSubscription()); var qd = new QueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public long requested() { return 0; } @Override public long produced(long n) { return 0; } @Override public int leave(int m) { return 0; } @Override public boolean accept(Subscriber<? super Integer> a, Integer v) { return false; } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); q.offer(1); Disposable d = Disposable.empty(); QueueDrainHelper.drainMaxLoop(q, ts, false, d, qd); ts.assertFailure(MissingBackpressureException.class); assertTrue(d.isDisposed()); } @Test public void drainMaxLoopDontAccept() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ts.onSubscribe(new BooleanSubscription()); var qd = new QueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public long requested() { return 1; } @Override public long produced(long n) { return 0; } @Override public int leave(int m) { return 0; } @Override public boolean accept(Subscriber<? super Integer> a, Integer v) { return false; } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); q.offer(1); QueueDrainHelper.drainMaxLoop(q, ts, false, null, qd); ts.assertEmpty(); } @Test public void checkTerminatedDelayErrorEmpty() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ts.onSubscribe(new BooleanSubscription()); var qd = new QueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public long requested() { return 0; } @Override public long produced(long n) { return 0; } @Override public int leave(int m) { return 0; } @Override public boolean accept(Subscriber<? super Integer> a, Integer v) { return false; } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, true, ts, true, q, qd); ts.assertResult(); } @Test public void checkTerminatedDelayErrorNonEmpty() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ts.onSubscribe(new BooleanSubscription()); var qd = new QueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public long requested() { return 0; } @Override public long produced(long n) { return 0; } @Override public int leave(int m) { return 0; } @Override public boolean accept(Subscriber<? super Integer> a, Integer v) { return false; } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, false, ts, true, q, qd); ts.assertEmpty(); } @Test public void checkTerminatedDelayErrorEmptyError() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ts.onSubscribe(new BooleanSubscription()); var qd = new QueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return new TestException(); } @Override public boolean enter() { return true; } @Override public long requested() { return 0; } @Override public long produced(long n) { return 0; } @Override public int leave(int m) { return 0; } @Override public boolean accept(Subscriber<? super Integer> a, Integer v) { return false; } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, true, ts, true, q, qd); ts.assertFailure(TestException.class); } @Test public void checkTerminatedNonDelayErrorError() { TestSubscriber<Integer> ts = new TestSubscriber<>(); ts.onSubscribe(new BooleanSubscription()); var qd = new QueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return new TestException(); } @Override public boolean enter() { return true; } @Override public long requested() { return 0; } @Override public long produced(long n) { return 0; } @Override public int leave(int m) { return 0; } @Override public boolean accept(Subscriber<? super Integer> a, Integer v) { return false; } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, false, ts, false, q, qd); ts.assertFailure(TestException.class); } @Test public void observerCheckTerminatedDelayErrorEmpty() { TestObserver<Integer> to = new TestObserver<>(); to.onSubscribe(Disposable.empty()); var qd = new ObservableQueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public int leave(int m) { return 0; } @Override public void accept(Observer<? super Integer> a, Integer v) { } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, true, to, true, q, null, qd); to.assertResult(); } @Test public void observerCheckTerminatedDelayErrorEmptyResource() { TestObserver<Integer> to = new TestObserver<>(); to.onSubscribe(Disposable.empty()); var qd = new ObservableQueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public int leave(int m) { return 0; } @Override public void accept(Observer<? super Integer> a, Integer v) { } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); Disposable d = Disposable.empty(); QueueDrainHelper.checkTerminated(true, true, to, true, q, d, qd); to.assertResult(); assertTrue(d.isDisposed()); } @Test public void observerCheckTerminatedDelayErrorNonEmpty() { TestObserver<Integer> to = new TestObserver<>(); to.onSubscribe(Disposable.empty()); var qd = new ObservableQueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return null; } @Override public boolean enter() { return true; } @Override public int leave(int m) { return 0; } @Override public void accept(Observer<? super Integer> a, Integer v) { } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, false, to, true, q, null, qd); to.assertEmpty(); } @Test public void observerCheckTerminatedDelayErrorEmptyError() { TestObserver<Integer> to = new TestObserver<>(); to.onSubscribe(Disposable.empty()); var qd = new ObservableQueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return new TestException(); } @Override public boolean enter() { return true; } @Override public int leave(int m) { return 0; } @Override public void accept(Observer<? super Integer> a, Integer v) { } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, true, to, true, q, null, qd); to.assertFailure(TestException.class); } @Test public void observerCheckTerminatedNonDelayErrorError() { TestObserver<Integer> to = new TestObserver<>(); to.onSubscribe(Disposable.empty()); var qd = new ObservableQueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return new TestException(); } @Override public boolean enter() { return true; } @Override public int leave(int m) { return 0; } @Override public void accept(Observer<? super Integer> a, Integer v) { } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); QueueDrainHelper.checkTerminated(true, false, to, false, q, null, qd); to.assertFailure(TestException.class); } @Test public void observerCheckTerminatedNonDelayErrorErrorResource() { TestObserver<Integer> to = new TestObserver<>(); to.onSubscribe(Disposable.empty()); var qd = new ObservableQueueDrain<Integer, Integer>() /* NFI */ { @Override public boolean cancelled() { return false; } @Override public boolean done() { return false; } @Override public Throwable error() { return new TestException(); } @Override public boolean enter() { return true; } @Override public int leave(int m) { return 0; } @Override public void accept(Observer<? super Integer> a, Integer v) { } }; SpscArrayQueue<Integer> q = new SpscArrayQueue<>(32); Disposable d = Disposable.empty(); QueueDrainHelper.checkTerminated(true, false, to, false, q, d, qd); to.assertFailure(TestException.class); assertTrue(d.isDisposed()); } @Test public void postCompleteAlreadyComplete() { TestSubscriber<Integer> ts = new TestSubscriber<>(); Queue<Integer> q = new ArrayDeque<>(); q.offer(1); AtomicLong state = new AtomicLong(QueueDrainHelper.COMPLETED_MASK); QueueDrainHelper.postComplete(ts, q, state, () -> false); } }