The following issues were found
src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableOnErrorNext.java
11 issues
Line: 24
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
public final class ObservableOnErrorNext<T> extends AbstractObservableWithUpstream<T, T> {
final Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier;
public ObservableOnErrorNext(ObservableSource<T> source,
Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier) {
super(source);
this.nextSupplier = nextSupplier;
Reported by PMD.
Line: 40
}
static final class OnErrorNextObserver<T> implements Observer<T> {
final Observer<? super T> downstream;
final Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier;
final SequentialDisposable arbiter;
boolean once;
Reported by PMD.
Line: 41
static final class OnErrorNextObserver<T> implements Observer<T> {
final Observer<? super T> downstream;
final Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier;
final SequentialDisposable arbiter;
boolean once;
boolean done;
Reported by PMD.
Line: 42
static final class OnErrorNextObserver<T> implements Observer<T> {
final Observer<? super T> downstream;
final Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier;
final SequentialDisposable arbiter;
boolean once;
boolean done;
Reported by PMD.
Line: 44
final Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier;
final SequentialDisposable arbiter;
boolean once;
boolean done;
OnErrorNextObserver(Observer<? super T> actual, Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier) {
this.downstream = actual;
Reported by PMD.
Line: 46
boolean once;
boolean done;
OnErrorNextObserver(Observer<? super T> actual, Function<? super Throwable, ? extends ObservableSource<? extends T>> nextSupplier) {
this.downstream = actual;
this.nextSupplier = nextSupplier;
this.arbiter = new SequentialDisposable();
Reported by PMD.
Line: 83
try {
p = nextSupplier.apply(t);
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
downstream.onError(new CompositeException(t, e));
return;
}
Reported by PMD.
Line: 96
return;
}
p.subscribe(this);
}
@Override
public void onComplete() {
if (done) {
Reported by PMD.
Line: 16
package io.reactivex.rxjava3.internal.operators.observable;
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.disposables.Disposable;
import io.reactivex.rxjava3.exceptions.*;
import io.reactivex.rxjava3.functions.Function;
import io.reactivex.rxjava3.internal.disposables.SequentialDisposable;
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
Reported by PMD.
Line: 18
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.disposables.Disposable;
import io.reactivex.rxjava3.exceptions.*;
import io.reactivex.rxjava3.functions.Function;
import io.reactivex.rxjava3.internal.disposables.SequentialDisposable;
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
public final class ObservableOnErrorNext<T> extends AbstractObservableWithUpstream<T, T> {
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableToListSingle.java
11 issues
Line: 31
public final class ObservableToListSingle<T, U extends Collection<? super T>>
extends Single<U> implements FuseToObservable<U> {
final ObservableSource<T> source;
final Supplier<U> collectionSupplier;
@SuppressWarnings({ "unchecked", "rawtypes" })
public ObservableToListSingle(ObservableSource<T> source, final int defaultCapacityHint) {
Reported by PMD.
Line: 33
final ObservableSource<T> source;
final Supplier<U> collectionSupplier;
@SuppressWarnings({ "unchecked", "rawtypes" })
public ObservableToListSingle(ObservableSource<T> source, final int defaultCapacityHint) {
this.source = source;
this.collectionSupplier = (Supplier)Functions.createArrayList(defaultCapacityHint);
Reported by PMD.
Line: 51
U coll;
try {
coll = ExceptionHelper.nullCheck(collectionSupplier.get(), "The collectionSupplier returned a null Collection.");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
EmptyDisposable.error(e, t);
return;
}
source.subscribe(new ToListObserver<>(t, coll));
Reported by PMD.
Line: 65
}
static final class ToListObserver<T, U extends Collection<? super T>> implements Observer<T>, Disposable {
final SingleObserver<? super U> downstream;
U collection;
Disposable upstream;
Reported by PMD.
Line: 67
static final class ToListObserver<T, U extends Collection<? super T>> implements Observer<T>, Disposable {
final SingleObserver<? super U> downstream;
U collection;
Disposable upstream;
ToListObserver(SingleObserver<? super U> actual, U collection) {
this.downstream = actual;
Reported by PMD.
Line: 69
U collection;
Disposable upstream;
ToListObserver(SingleObserver<? super U> actual, U collection) {
this.downstream = actual;
this.collection = collection;
}
Reported by PMD.
Line: 101
@Override
public void onError(Throwable t) {
collection = null;
downstream.onError(t);
}
@Override
public void onComplete() {
Reported by PMD.
Line: 108
@Override
public void onComplete() {
U c = collection;
collection = null;
downstream.onSuccess(c);
}
}
}
Reported by PMD.
Line: 18
import java.util.Collection;
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.disposables.Disposable;
import io.reactivex.rxjava3.exceptions.Exceptions;
import io.reactivex.rxjava3.functions.Supplier;
import io.reactivex.rxjava3.internal.disposables.*;
import io.reactivex.rxjava3.internal.functions.Functions;
Reported by PMD.
Line: 22
import io.reactivex.rxjava3.disposables.Disposable;
import io.reactivex.rxjava3.exceptions.Exceptions;
import io.reactivex.rxjava3.functions.Supplier;
import io.reactivex.rxjava3.internal.disposables.*;
import io.reactivex.rxjava3.internal.functions.Functions;
import io.reactivex.rxjava3.internal.fuseable.FuseToObservable;
import io.reactivex.rxjava3.internal.util.ExceptionHelper;
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableThrottleFirstTimed.java
11 issues
Line: 26
import io.reactivex.rxjava3.observers.SerializedObserver;
public final class ObservableThrottleFirstTimed<T> extends AbstractObservableWithUpstream<T, T> {
final long timeout;
final TimeUnit unit;
final Scheduler scheduler;
public ObservableThrottleFirstTimed(ObservableSource<T> source,
long timeout, TimeUnit unit, Scheduler scheduler) {
Reported by PMD.
Line: 27
public final class ObservableThrottleFirstTimed<T> extends AbstractObservableWithUpstream<T, T> {
final long timeout;
final TimeUnit unit;
final Scheduler scheduler;
public ObservableThrottleFirstTimed(ObservableSource<T> source,
long timeout, TimeUnit unit, Scheduler scheduler) {
super(source);
Reported by PMD.
Line: 28
public final class ObservableThrottleFirstTimed<T> extends AbstractObservableWithUpstream<T, T> {
final long timeout;
final TimeUnit unit;
final Scheduler scheduler;
public ObservableThrottleFirstTimed(ObservableSource<T> source,
long timeout, TimeUnit unit, Scheduler scheduler) {
super(source);
this.timeout = timeout;
Reported by PMD.
Line: 50
implements Observer<T>, Disposable, Runnable {
private static final long serialVersionUID = 786994795061867455L;
final Observer<? super T> downstream;
final long timeout;
final TimeUnit unit;
final Scheduler.Worker worker;
Disposable upstream;
Reported by PMD.
Line: 51
private static final long serialVersionUID = 786994795061867455L;
final Observer<? super T> downstream;
final long timeout;
final TimeUnit unit;
final Scheduler.Worker worker;
Disposable upstream;
Reported by PMD.
Line: 52
final Observer<? super T> downstream;
final long timeout;
final TimeUnit unit;
final Scheduler.Worker worker;
Disposable upstream;
volatile boolean gate;
Reported by PMD.
Line: 53
final Observer<? super T> downstream;
final long timeout;
final TimeUnit unit;
final Scheduler.Worker worker;
Disposable upstream;
volatile boolean gate;
Reported by PMD.
Line: 55
final TimeUnit unit;
final Scheduler.Worker worker;
Disposable upstream;
volatile boolean gate;
DebounceTimedObserver(Observer<? super T> actual, long timeout, TimeUnit unit, Worker worker) {
this.downstream = actual;
Reported by PMD.
Line: 57
Disposable upstream;
volatile boolean gate;
DebounceTimedObserver(Observer<? super T> actual, long timeout, TimeUnit unit, Worker worker) {
this.downstream = actual;
this.timeout = timeout;
this.unit = unit;
Reported by PMD.
Line: 83
Disposable d = get();
if (d != null) {
d.dispose();
}
DisposableHelper.replace(this, worker.schedule(this, timeout, unit));
}
}
Reported by PMD.
src/jmh/java/io/reactivex/rxjava3/core/TakeUntilPerf.java
11 issues
Line: 33
@AuxCounters
public class TakeUntilPerf implements Consumer<Integer> {
public volatile int items;
static final int count = 10000;
Flowable<Integer> flowable;
Reported by PMD.
Line: 37
static final int count = 10000;
Flowable<Integer> flowable;
Observable<Integer> observable;
@Override
public void accept(Integer t) {
Reported by PMD.
Line: 37
static final int count = 10000;
Flowable<Integer> flowable;
Observable<Integer> observable;
@Override
public void accept(Integer t) {
Reported by PMD.
Line: 39
Flowable<Integer> flowable;
Observable<Integer> observable;
@Override
public void accept(Integer t) {
items++;
}
Reported by PMD.
Line: 39
Flowable<Integer> flowable;
Observable<Integer> observable;
@Override
public void accept(Integer t) {
items++;
}
Reported by PMD.
Line: 53
@Override
public Object call() {
int c = count;
while (items < c) { }
return 1;
}
}).subscribeOn(Schedulers.single()));
observable = Observable.range(1, 1000 * 1000).takeUntil(Observable.fromCallable(new Callable<Object>() {
Reported by PMD.
Line: 62
@Override
public Object call() {
int c = count;
while (items < c) { }
return 1;
}
}).subscribeOn(Schedulers.single()));
}
Reported by PMD.
Line: 79
}
});
while (cdl.getCount() != 0) { }
}
@Benchmark
public void observable() {
final CountDownLatch cdl = new CountDownLatch(1);
Reported by PMD.
Line: 93
}
});
while (cdl.getCount() != 0) { }
}
}
Reported by PMD.
Line: 18
import java.util.concurrent.*;
import org.openjdk.jmh.annotations.*;
import io.reactivex.rxjava3.functions.*;
import io.reactivex.rxjava3.internal.functions.Functions;
import io.reactivex.rxjava3.schedulers.Schedulers;
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableReduceMaybe.java
11 issues
Line: 33
*/
public final class ObservableReduceMaybe<T> extends Maybe<T> {
final ObservableSource<T> source;
final BiFunction<T, T, T> reducer;
public ObservableReduceMaybe(ObservableSource<T> source, BiFunction<T, T, T> reducer) {
this.source = source;
Reported by PMD.
Line: 35
final ObservableSource<T> source;
final BiFunction<T, T, T> reducer;
public ObservableReduceMaybe(ObservableSource<T> source, BiFunction<T, T, T> reducer) {
this.source = source;
this.reducer = reducer;
}
Reported by PMD.
Line: 49
static final class ReduceObserver<T> implements Observer<T>, Disposable {
final MaybeObserver<? super T> downstream;
final BiFunction<T, T, T> reducer;
boolean done;
Reported by PMD.
Line: 51
final MaybeObserver<? super T> downstream;
final BiFunction<T, T, T> reducer;
boolean done;
T value;
Reported by PMD.
Line: 53
final BiFunction<T, T, T> reducer;
boolean done;
T value;
Disposable upstream;
Reported by PMD.
Line: 55
boolean done;
T value;
Disposable upstream;
ReduceObserver(MaybeObserver<? super T> observer, BiFunction<T, T, T> reducer) {
this.downstream = observer;
Reported by PMD.
Line: 57
T value;
Disposable upstream;
ReduceObserver(MaybeObserver<? super T> observer, BiFunction<T, T, T> reducer) {
this.downstream = observer;
this.reducer = reducer;
}
Reported by PMD.
Line: 83
} else {
try {
this.value = Objects.requireNonNull(reducer.apply(v, value), "The reducer returned a null value");
} catch (Throwable ex) {
Exceptions.throwIfFatal(ex);
upstream.dispose();
onError(ex);
}
}
Reported by PMD.
Line: 99
return;
}
done = true;
value = null;
downstream.onError(e);
}
@Override
public void onComplete() {
Reported by PMD.
Line: 110
}
done = true;
T v = value;
value = null;
if (v != null) {
downstream.onSuccess(v);
} else {
downstream.onComplete();
}
Reported by PMD.
src/jmh/java/io/reactivex/rxjava3/xmapz/FlowableFlatMapMaybeEmptyPerf.java
11 issues
Line: 34
@State(Scope.Thread)
public class FlowableFlatMapMaybeEmptyPerf {
@Param({ "1", "10", "100", "1000", "10000", "100000", "1000000" })
public int count;
Flowable<Integer> flowableConvert;
Flowable<Integer> flowableDedicated;
Reported by PMD.
Line: 36
@Param({ "1", "10", "100", "1000", "10000", "100000", "1000000" })
public int count;
Flowable<Integer> flowableConvert;
Flowable<Integer> flowableDedicated;
Flowable<Integer> flowablePlain;
Reported by PMD.
Line: 36
@Param({ "1", "10", "100", "1000", "10000", "100000", "1000000" })
public int count;
Flowable<Integer> flowableConvert;
Flowable<Integer> flowableDedicated;
Flowable<Integer> flowablePlain;
Reported by PMD.
Line: 38
Flowable<Integer> flowableConvert;
Flowable<Integer> flowableDedicated;
Flowable<Integer> flowablePlain;
@Setup
public void setup() {
Reported by PMD.
Line: 38
Flowable<Integer> flowableConvert;
Flowable<Integer> flowableDedicated;
Flowable<Integer> flowablePlain;
@Setup
public void setup() {
Reported by PMD.
Line: 40
Flowable<Integer> flowableDedicated;
Flowable<Integer> flowablePlain;
@Setup
public void setup() {
Integer[] sourceArray = new Integer[count];
Arrays.fill(sourceArray, 777);
Reported by PMD.
Line: 40
Flowable<Integer> flowableDedicated;
Flowable<Integer> flowablePlain;
@Setup
public void setup() {
Integer[] sourceArray = new Integer[count];
Arrays.fill(sourceArray, 777);
Reported by PMD.
Line: 59
flowableConvert = source.flatMap(new Function<Integer, Publisher<? extends Integer>>() {
@Override
public Publisher<? extends Integer> apply(Integer v) {
return Maybe.<Integer>empty().toFlowable();
}
});
flowableDedicated = source.flatMapMaybe(new Function<Integer, Maybe<? extends Integer>>() {
@Override
Reported by PMD.
Line: 59
flowableConvert = source.flatMap(new Function<Integer, Publisher<? extends Integer>>() {
@Override
public Publisher<? extends Integer> apply(Integer v) {
return Maybe.<Integer>empty().toFlowable();
}
});
flowableDedicated = source.flatMapMaybe(new Function<Integer, Maybe<? extends Integer>>() {
@Override
Reported by PMD.
Line: 19
import java.util.Arrays;
import java.util.concurrent.TimeUnit;
import org.openjdk.jmh.annotations.*;
import org.openjdk.jmh.infra.Blackhole;
import org.reactivestreams.Publisher;
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.functions.Function;
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/completable/CompletableConcatArray.java
11 issues
Line: 23
import io.reactivex.rxjava3.internal.disposables.SequentialDisposable;
public final class CompletableConcatArray extends Completable {
final CompletableSource[] sources;
public CompletableConcatArray(CompletableSource[] sources) {
this.sources = sources;
}
Reported by PMD.
Line: 25
public final class CompletableConcatArray extends Completable {
final CompletableSource[] sources;
public CompletableConcatArray(CompletableSource[] sources) {
this.sources = sources;
}
@Override
public void subscribeActual(CompletableObserver observer) {
Reported by PMD.
Line: 40
private static final long serialVersionUID = -7965400327305809232L;
final CompletableObserver downstream;
final CompletableSource[] sources;
int index;
final SequentialDisposable sd;
Reported by PMD.
Line: 41
private static final long serialVersionUID = -7965400327305809232L;
final CompletableObserver downstream;
final CompletableSource[] sources;
int index;
final SequentialDisposable sd;
Reported by PMD.
Line: 43
final CompletableObserver downstream;
final CompletableSource[] sources;
int index;
final SequentialDisposable sd;
ConcatInnerObserver(CompletableObserver actual, CompletableSource[] sources) {
this.downstream = actual;
Reported by PMD.
Line: 45
int index;
final SequentialDisposable sd;
ConcatInnerObserver(CompletableObserver actual, CompletableSource[] sources) {
this.downstream = actual;
this.sources = sources;
this.sd = new SequentialDisposable();
Reported by PMD.
Line: 47
final SequentialDisposable sd;
ConcatInnerObserver(CompletableObserver actual, CompletableSource[] sources) {
this.downstream = actual;
this.sources = sources;
this.sd = new SequentialDisposable();
}
Reported by PMD.
Line: 18
import java.util.concurrent.atomic.AtomicInteger;
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.disposables.Disposable;
import io.reactivex.rxjava3.internal.disposables.SequentialDisposable;
public final class CompletableConcatArray extends Completable {
final CompletableSource[] sources;
Reported by PMD.
Line: 25
public final class CompletableConcatArray extends Completable {
final CompletableSource[] sources;
public CompletableConcatArray(CompletableSource[] sources) {
this.sources = sources;
}
@Override
public void subscribeActual(CompletableObserver observer) {
Reported by PMD.
Line: 47
final SequentialDisposable sd;
ConcatInnerObserver(CompletableObserver actual, CompletableSource[] sources) {
this.downstream = actual;
this.sources = sources;
this.sd = new SequentialDisposable();
}
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/single/SingleEquals.java
11 issues
Line: 25
public final class SingleEquals<T> extends Single<Boolean> {
final SingleSource<? extends T> first;
final SingleSource<? extends T> second;
public SingleEquals(SingleSource<? extends T> first, SingleSource<? extends T> second) {
this.first = first;
this.second = second;
Reported by PMD.
Line: 26
public final class SingleEquals<T> extends Single<Boolean> {
final SingleSource<? extends T> first;
final SingleSource<? extends T> second;
public SingleEquals(SingleSource<? extends T> first, SingleSource<? extends T> second) {
this.first = first;
this.second = second;
}
Reported by PMD.
Line: 47
}
static class InnerObserver<T> implements SingleObserver<T> {
final int index;
final CompositeDisposable set;
final Object[] values;
final SingleObserver<? super Boolean> downstream;
final AtomicInteger count;
Reported by PMD.
Line: 48
static class InnerObserver<T> implements SingleObserver<T> {
final int index;
final CompositeDisposable set;
final Object[] values;
final SingleObserver<? super Boolean> downstream;
final AtomicInteger count;
InnerObserver(int index, CompositeDisposable set, Object[] values, SingleObserver<? super Boolean> observer, AtomicInteger count) {
Reported by PMD.
Line: 49
static class InnerObserver<T> implements SingleObserver<T> {
final int index;
final CompositeDisposable set;
final Object[] values;
final SingleObserver<? super Boolean> downstream;
final AtomicInteger count;
InnerObserver(int index, CompositeDisposable set, Object[] values, SingleObserver<? super Boolean> observer, AtomicInteger count) {
this.index = index;
Reported by PMD.
Line: 50
final int index;
final CompositeDisposable set;
final Object[] values;
final SingleObserver<? super Boolean> downstream;
final AtomicInteger count;
InnerObserver(int index, CompositeDisposable set, Object[] values, SingleObserver<? super Boolean> observer, AtomicInteger count) {
this.index = index;
this.set = set;
Reported by PMD.
Line: 51
final CompositeDisposable set;
final Object[] values;
final SingleObserver<? super Boolean> downstream;
final AtomicInteger count;
InnerObserver(int index, CompositeDisposable set, Object[] values, SingleObserver<? super Boolean> observer, AtomicInteger count) {
this.index = index;
this.set = set;
this.values = values;
Reported by PMD.
Line: 53
final SingleObserver<? super Boolean> downstream;
final AtomicInteger count;
InnerObserver(int index, CompositeDisposable set, Object[] values, SingleObserver<? super Boolean> observer, AtomicInteger count) {
this.index = index;
this.set = set;
this.values = values;
this.downstream = observer;
this.count = count;
Reported by PMD.
Line: 70
public void onSuccess(T value) {
values[index] = value;
if (count.incrementAndGet() == 2) {
downstream.onSuccess(Objects.equals(values[0], values[1]));
}
}
@Override
Reported by PMD.
Line: 19
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.disposables.*;
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
public final class SingleEquals<T> extends Single<Boolean> {
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableReduceSeedSingle.java
11 issues
Line: 34
*/
public final class ObservableReduceSeedSingle<T, R> extends Single<R> {
final ObservableSource<T> source;
final R seed;
final BiFunction<R, ? super T, R> reducer;
Reported by PMD.
Line: 36
final ObservableSource<T> source;
final R seed;
final BiFunction<R, ? super T, R> reducer;
public ObservableReduceSeedSingle(ObservableSource<T> source, R seed, BiFunction<R, ? super T, R> reducer) {
this.source = source;
Reported by PMD.
Line: 38
final R seed;
final BiFunction<R, ? super T, R> reducer;
public ObservableReduceSeedSingle(ObservableSource<T> source, R seed, BiFunction<R, ? super T, R> reducer) {
this.source = source;
this.seed = seed;
this.reducer = reducer;
Reported by PMD.
Line: 53
static final class ReduceSeedObserver<T, R> implements Observer<T>, Disposable {
final SingleObserver<? super R> downstream;
final BiFunction<R, ? super T, R> reducer;
R value;
Reported by PMD.
Line: 55
final SingleObserver<? super R> downstream;
final BiFunction<R, ? super T, R> reducer;
R value;
Disposable upstream;
Reported by PMD.
Line: 57
final BiFunction<R, ? super T, R> reducer;
R value;
Disposable upstream;
ReduceSeedObserver(SingleObserver<? super R> actual, BiFunction<R, ? super T, R> reducer, R value) {
this.downstream = actual;
Reported by PMD.
Line: 59
R value;
Disposable upstream;
ReduceSeedObserver(SingleObserver<? super R> actual, BiFunction<R, ? super T, R> reducer, R value) {
this.downstream = actual;
this.value = value;
this.reducer = reducer;
Reported by PMD.
Line: 82
if (v != null) {
try {
this.value = Objects.requireNonNull(reducer.apply(v, value), "The reducer returned a null value");
} catch (Throwable ex) {
Exceptions.throwIfFatal(ex);
upstream.dispose();
onError(ex);
}
}
Reported by PMD.
Line: 94
public void onError(Throwable e) {
R v = value;
if (v != null) {
value = null;
downstream.onError(e);
} else {
RxJavaPlugins.onError(e);
}
}
Reported by PMD.
Line: 105
public void onComplete() {
R v = value;
if (v != null) {
value = null;
downstream.onSuccess(v);
}
}
@Override
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/completable/CompletableMergeIterable.java
11 issues
Line: 26
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
public final class CompletableMergeIterable extends Completable {
final Iterable<? extends CompletableSource> sources;
public CompletableMergeIterable(Iterable<? extends CompletableSource> sources) {
this.sources = sources;
}
Reported by PMD.
Line: 45
try {
iterator = Objects.requireNonNull(sources.iterator(), "The source iterator returned is null");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
observer.onError(e);
return;
}
Reported by PMD.
Line: 59
boolean b;
try {
b = iterator.hasNext();
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
set.dispose();
shared.onError(e);
return;
}
Reported by PMD.
Line: 78
try {
c = Objects.requireNonNull(iterator.next(), "The iterator returned a null CompletableSource");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
set.dispose();
shared.onError(e);
return;
}
Reported by PMD.
Line: 91
wip.getAndIncrement();
c.subscribe(shared);
}
shared.onComplete();
}
Reported by PMD.
Line: 101
private static final long serialVersionUID = -7730517613164279224L;
final CompositeDisposable set;
final CompletableObserver downstream;
final AtomicInteger wip;
Reported by PMD.
Line: 103
final CompositeDisposable set;
final CompletableObserver downstream;
final AtomicInteger wip;
MergeCompletableObserver(CompletableObserver actual, CompositeDisposable set, AtomicInteger wip) {
this.downstream = actual;
Reported by PMD.
Line: 105
final CompletableObserver downstream;
final AtomicInteger wip;
MergeCompletableObserver(CompletableObserver actual, CompositeDisposable set, AtomicInteger wip) {
this.downstream = actual;
this.set = set;
this.wip = wip;
Reported by PMD.
Line: 20
import java.util.Objects;
import java.util.concurrent.atomic.*;
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.disposables.*;
import io.reactivex.rxjava3.exceptions.Exceptions;
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
public final class CompletableMergeIterable extends Completable {
Reported by PMD.
Line: 21
import java.util.concurrent.atomic.*;
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.disposables.*;
import io.reactivex.rxjava3.exceptions.Exceptions;
import io.reactivex.rxjava3.plugins.RxJavaPlugins;
public final class CompletableMergeIterable extends Completable {
final Iterable<? extends CompletableSource> sources;
Reported by PMD.