The following issues were found
src/main/java/io/reactivex/rxjava3/disposables/CompositeDisposable.java
15 issues
Line: 29
*/
public final class CompositeDisposable implements Disposable, DisposableContainer {
OpenHashSet<Disposable> resources;
volatile boolean disposed;
/**
* Creates an empty {@code CompositeDisposable}.
Reported by PMD.
Line: 31
OpenHashSet<Disposable> resources;
volatile boolean disposed;
/**
* Creates an empty {@code CompositeDisposable}.
*/
public CompositeDisposable() {
Reported by PMD.
Line: 79
}
disposed = true;
set = resources;
resources = null;
}
dispose(set);
}
Reported by PMD.
Line: 104
synchronized (this) {
if (!disposed) {
OpenHashSet<Disposable> set = resources;
if (set == null) {
set = new OpenHashSet<>();
resources = set;
}
set.add(disposable);
return true;
Reported by PMD.
Line: 130
synchronized (this) {
if (!disposed) {
OpenHashSet<Disposable> set = resources;
if (set == null) {
set = new OpenHashSet<>(disposables.length + 1);
resources = set;
}
for (Disposable d : disposables) {
Objects.requireNonNull(d, "A Disposable in the disposables array is null");
Reported by PMD.
Line: 183
}
OpenHashSet<Disposable> set = resources;
if (set == null || !set.remove(disposable)) {
return false;
}
}
return true;
}
Reported by PMD.
Line: 204
}
set = resources;
resources = null;
}
dispose(set);
}
Reported by PMD.
Line: 223
return 0;
}
OpenHashSet<Disposable> set = resources;
return set != null ? set.size() : 0;
}
}
/**
* Dispose the contents of the {@link OpenHashSet} by suppressing non-fatal
Reported by PMD.
Line: 232
* {@link Throwable}s till the end.
* @param set the {@code OpenHashSet} to dispose elements of
*/
void dispose(@Nullable OpenHashSet<Disposable> set) {
if (set == null) {
return;
}
List<Throwable> errors = null;
Object[] array = set.keys();
Reported by PMD.
Line: 242
if (o instanceof Disposable) {
try {
((Disposable) o).dispose();
} catch (Throwable ex) {
Exceptions.throwIfFatal(ex);
if (errors == null) {
errors = new ArrayList<>();
}
errors.add(ex);
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableCollectSingle.java
14 issues
Line: 28
public final class ObservableCollectSingle<T, U> extends Single<U> implements FuseToObservable<U> {
final ObservableSource<T> source;
final Supplier<? extends U> initialSupplier;
final BiConsumer<? super U, ? super T> collector;
public ObservableCollectSingle(ObservableSource<T> source,
Reported by PMD.
Line: 30
final ObservableSource<T> source;
final Supplier<? extends U> initialSupplier;
final BiConsumer<? super U, ? super T> collector;
public ObservableCollectSingle(ObservableSource<T> source,
Supplier<? extends U> initialSupplier, BiConsumer<? super U, ? super T> collector) {
this.source = source;
Reported by PMD.
Line: 31
final ObservableSource<T> source;
final Supplier<? extends U> initialSupplier;
final BiConsumer<? super U, ? super T> collector;
public ObservableCollectSingle(ObservableSource<T> source,
Supplier<? extends U> initialSupplier, BiConsumer<? super U, ? super T> collector) {
this.source = source;
this.initialSupplier = initialSupplier;
Reported by PMD.
Line: 45
U u;
try {
u = Objects.requireNonNull(initialSupplier.get(), "The initialSupplier returned a null value");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
EmptyDisposable.error(e, t);
return;
}
Reported by PMD.
Line: 60
}
static final class CollectObserver<T, U> implements Observer<T>, Disposable {
final SingleObserver<? super U> downstream;
final BiConsumer<? super U, ? super T> collector;
final U u;
Disposable upstream;
Reported by PMD.
Line: 61
static final class CollectObserver<T, U> implements Observer<T>, Disposable {
final SingleObserver<? super U> downstream;
final BiConsumer<? super U, ? super T> collector;
final U u;
Disposable upstream;
boolean done;
Reported by PMD.
Line: 62
static final class CollectObserver<T, U> implements Observer<T>, Disposable {
final SingleObserver<? super U> downstream;
final BiConsumer<? super U, ? super T> collector;
final U u;
Disposable upstream;
boolean done;
Reported by PMD.
Line: 64
final BiConsumer<? super U, ? super T> collector;
final U u;
Disposable upstream;
boolean done;
CollectObserver(SingleObserver<? super U> actual, U u, BiConsumer<? super U, ? super T> collector) {
this.downstream = actual;
Reported by PMD.
Line: 66
Disposable upstream;
boolean done;
CollectObserver(SingleObserver<? super U> actual, U u, BiConsumer<? super U, ? super T> collector) {
this.downstream = actual;
this.collector = collector;
this.u = u;
Reported by PMD.
Line: 99
}
try {
collector.accept(u, t);
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
upstream.dispose();
onError(e);
}
}
Reported by PMD.
src/jmh/java/io/reactivex/rxjava3/core/ReducePerf.java
14 issues
Line: 32
@State(Scope.Thread)
public class ReducePerf implements BiFunction<Integer, Integer, Integer> {
@Param({ "1", "1000", "1000000" })
public int times;
Single<Integer> obsSingle;
Single<Integer> flowSingle;
Reported by PMD.
Line: 34
@Param({ "1", "1000", "1000000" })
public int times;
Single<Integer> obsSingle;
Single<Integer> flowSingle;
Maybe<Integer> obsMaybe;
Reported by PMD.
Line: 34
@Param({ "1", "1000", "1000000" })
public int times;
Single<Integer> obsSingle;
Single<Integer> flowSingle;
Maybe<Integer> obsMaybe;
Reported by PMD.
Line: 36
Single<Integer> obsSingle;
Single<Integer> flowSingle;
Maybe<Integer> obsMaybe;
Maybe<Integer> flowMaybe;
Reported by PMD.
Line: 36
Single<Integer> obsSingle;
Single<Integer> flowSingle;
Maybe<Integer> obsMaybe;
Maybe<Integer> flowMaybe;
Reported by PMD.
Line: 38
Single<Integer> flowSingle;
Maybe<Integer> obsMaybe;
Maybe<Integer> flowMaybe;
@Override
public Integer apply(Integer t1, Integer t2) {
Reported by PMD.
Line: 38
Single<Integer> flowSingle;
Maybe<Integer> obsMaybe;
Maybe<Integer> flowMaybe;
@Override
public Integer apply(Integer t1, Integer t2) {
Reported by PMD.
Line: 40
Maybe<Integer> obsMaybe;
Maybe<Integer> flowMaybe;
@Override
public Integer apply(Integer t1, Integer t2) {
return t1 + t2;
}
Reported by PMD.
Line: 40
Maybe<Integer> obsMaybe;
Maybe<Integer> flowMaybe;
@Override
public Integer apply(Integer t1, Integer t2) {
return t1 + t2;
}
Reported by PMD.
Line: 52
Integer[] array = new Integer[times];
Arrays.fill(array, 777);
obsSingle = Observable.fromArray(array).reduce(0, this);
obsMaybe = Observable.fromArray(array).reduce(this);
flowSingle = Flowable.fromArray(array).reduce(0, this);
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/completable/CompletableDelay.java
14 issues
Line: 25
public final class CompletableDelay extends Completable {
final CompletableSource source;
final long delay;
final TimeUnit unit;
Reported by PMD.
Line: 27
final CompletableSource source;
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
Reported by PMD.
Line: 29
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
Reported by PMD.
Line: 31
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
public CompletableDelay(CompletableSource source, long delay, TimeUnit unit, Scheduler scheduler, boolean delayError) {
this.source = source;
Reported by PMD.
Line: 33
final Scheduler scheduler;
final boolean delayError;
public CompletableDelay(CompletableSource source, long delay, TimeUnit unit, Scheduler scheduler, boolean delayError) {
this.source = source;
this.delay = delay;
this.unit = unit;
Reported by PMD.
Line: 53
private static final long serialVersionUID = 465972761105851022L;
final CompletableObserver downstream;
final long delay;
final TimeUnit unit;
Reported by PMD.
Line: 55
final CompletableObserver downstream;
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
Reported by PMD.
Line: 55
final CompletableObserver downstream;
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
Reported by PMD.
Line: 57
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
Reported by PMD.
Line: 59
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
Throwable error;
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/completable/CompletableMergeDelayErrorIterable.java
14 issues
Line: 26
import io.reactivex.rxjava3.internal.operators.completable.CompletableMergeArrayDelayError.*;
import io.reactivex.rxjava3.internal.util.AtomicThrowable;
public final class CompletableMergeDelayErrorIterable extends Completable {
final Iterable<? extends CompletableSource> sources;
public CompletableMergeDelayErrorIterable(Iterable<? extends CompletableSource> sources) {
this.sources = sources;
Reported by PMD.
Line: 26
import io.reactivex.rxjava3.internal.operators.completable.CompletableMergeArrayDelayError.*;
import io.reactivex.rxjava3.internal.util.AtomicThrowable;
public final class CompletableMergeDelayErrorIterable extends Completable {
final Iterable<? extends CompletableSource> sources;
public CompletableMergeDelayErrorIterable(Iterable<? extends CompletableSource> sources) {
this.sources = sources;
Reported by PMD.
Line: 28
public final class CompletableMergeDelayErrorIterable extends Completable {
final Iterable<? extends CompletableSource> sources;
public CompletableMergeDelayErrorIterable(Iterable<? extends CompletableSource> sources) {
this.sources = sources;
}
Reported by PMD.
Line: 35
}
@Override
public void subscribeActual(final CompletableObserver observer) {
final CompositeDisposable set = new CompositeDisposable();
observer.onSubscribe(set);
Iterator<? extends CompletableSource> iterator;
Reported by PMD.
Line: 35
}
@Override
public void subscribeActual(final CompletableObserver observer) {
final CompositeDisposable set = new CompositeDisposable();
observer.onSubscribe(set);
Iterator<? extends CompletableSource> iterator;
Reported by PMD.
Line: 35
}
@Override
public void subscribeActual(final CompletableObserver observer) {
final CompositeDisposable set = new CompositeDisposable();
observer.onSubscribe(set);
Iterator<? extends CompletableSource> iterator;
Reported by PMD.
Line: 35
}
@Override
public void subscribeActual(final CompletableObserver observer) {
final CompositeDisposable set = new CompositeDisposable();
observer.onSubscribe(set);
Iterator<? extends CompletableSource> iterator;
Reported by PMD.
Line: 44
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: 63
boolean b;
try {
b = iterator.hasNext();
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
errors.tryAddThrowableOrReport(e);
break;
}
Reported by PMD.
Line: 81
try {
c = Objects.requireNonNull(iterator.next(), "The iterator returned a null CompletableSource");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
errors.tryAddThrowableOrReport(e);
break;
}
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/flowable/FlowableDelay.java
14 issues
Line: 26
import io.reactivex.rxjava3.subscribers.SerializedSubscriber;
public final class FlowableDelay<T> extends AbstractFlowableWithUpstream<T, T> {
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
public FlowableDelay(Flowable<T> source, long delay, TimeUnit unit, Scheduler scheduler, boolean delayError) {
Reported by PMD.
Line: 27
public final class FlowableDelay<T> extends AbstractFlowableWithUpstream<T, T> {
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
public FlowableDelay(Flowable<T> source, long delay, TimeUnit unit, Scheduler scheduler, boolean delayError) {
super(source);
Reported by PMD.
Line: 28
public final class FlowableDelay<T> extends AbstractFlowableWithUpstream<T, T> {
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
public FlowableDelay(Flowable<T> source, long delay, TimeUnit unit, Scheduler scheduler, boolean delayError) {
super(source);
this.delay = delay;
Reported by PMD.
Line: 29
final long delay;
final TimeUnit unit;
final Scheduler scheduler;
final boolean delayError;
public FlowableDelay(Flowable<T> source, long delay, TimeUnit unit, Scheduler scheduler, boolean delayError) {
super(source);
this.delay = delay;
this.unit = unit;
Reported by PMD.
Line: 54
}
static final class DelaySubscriber<T> implements FlowableSubscriber<T>, Subscription {
final Subscriber<? super T> downstream;
final long delay;
final TimeUnit unit;
final Scheduler.Worker w;
final boolean delayError;
Reported by PMD.
Line: 55
static final class DelaySubscriber<T> implements FlowableSubscriber<T>, Subscription {
final Subscriber<? super T> downstream;
final long delay;
final TimeUnit unit;
final Scheduler.Worker w;
final boolean delayError;
Subscription upstream;
Reported by PMD.
Line: 56
static final class DelaySubscriber<T> implements FlowableSubscriber<T>, Subscription {
final Subscriber<? super T> downstream;
final long delay;
final TimeUnit unit;
final Scheduler.Worker w;
final boolean delayError;
Subscription upstream;
Reported by PMD.
Line: 57
final Subscriber<? super T> downstream;
final long delay;
final TimeUnit unit;
final Scheduler.Worker w;
final boolean delayError;
Subscription upstream;
DelaySubscriber(Subscriber<? super T> actual, long delay, TimeUnit unit, Worker w, boolean delayError) {
Reported by PMD.
Line: 58
final long delay;
final TimeUnit unit;
final Scheduler.Worker w;
final boolean delayError;
Subscription upstream;
DelaySubscriber(Subscriber<? super T> actual, long delay, TimeUnit unit, Worker w, boolean delayError) {
super();
Reported by PMD.
Line: 60
final Scheduler.Worker w;
final boolean delayError;
Subscription upstream;
DelaySubscriber(Subscriber<? super T> actual, long delay, TimeUnit unit, Worker w, boolean delayError) {
super();
this.downstream = actual;
this.delay = delay;
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/flowable/FlowableDematerialize.java
14 issues
Line: 28
public final class FlowableDematerialize<T, R> extends AbstractFlowableWithUpstream<T, R> {
final Function<? super T, ? extends Notification<R>> selector;
public FlowableDematerialize(Flowable<T> source, Function<? super T, ? extends Notification<R>> selector) {
super(source);
this.selector = selector;
}
Reported by PMD.
Line: 42
static final class DematerializeSubscriber<T, R> implements FlowableSubscriber<T>, Subscription {
final Subscriber<? super R> downstream;
final Function<? super T, ? extends Notification<R>> selector;
boolean done;
Reported by PMD.
Line: 44
final Subscriber<? super R> downstream;
final Function<? super T, ? extends Notification<R>> selector;
boolean done;
Subscription upstream;
Reported by PMD.
Line: 46
final Function<? super T, ? extends Notification<R>> selector;
boolean done;
Subscription upstream;
DematerializeSubscriber(Subscriber<? super R> downstream, Function<? super T, ? extends Notification<R>> selector) {
this.downstream = downstream;
Reported by PMD.
Line: 48
boolean done;
Subscription upstream;
DematerializeSubscriber(Subscriber<? super R> downstream, Function<? super T, ? extends Notification<R>> selector) {
this.downstream = downstream;
this.selector = selector;
}
Reported by PMD.
Line: 68
if (done) {
if (item instanceof Notification) {
Notification<?> notification = (Notification<?>)item;
if (notification.isOnError()) {
RxJavaPlugins.onError(notification.getError());
}
}
return;
}
Reported by PMD.
Line: 68
if (done) {
if (item instanceof Notification) {
Notification<?> notification = (Notification<?>)item;
if (notification.isOnError()) {
RxJavaPlugins.onError(notification.getError());
}
}
return;
}
Reported by PMD.
Line: 69
if (item instanceof Notification) {
Notification<?> notification = (Notification<?>)item;
if (notification.isOnError()) {
RxJavaPlugins.onError(notification.getError());
}
}
return;
}
Reported by PMD.
Line: 79
try {
notification = Objects.requireNonNull(selector.apply(item), "The selector returned a null Notification");
} catch (Throwable ex) {
Exceptions.throwIfFatal(ex);
upstream.cancel();
onError(ex);
return;
}
Reported by PMD.
Line: 85
onError(ex);
return;
}
if (notification.isOnError()) {
upstream.cancel();
onError(notification.getError());
} else if (notification.isOnComplete()) {
upstream.cancel();
onComplete();
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/flowable/FlowableDistinctUntilChanged.java
14 issues
Line: 26
public final class FlowableDistinctUntilChanged<T, K> extends AbstractFlowableWithUpstream<T, T> {
final Function<? super T, K> keySelector;
final BiPredicate<? super K, ? super K> comparer;
public FlowableDistinctUntilChanged(Flowable<T> source, Function<? super T, K> keySelector, BiPredicate<? super K, ? super K> comparer) {
super(source);
Reported by PMD.
Line: 28
final Function<? super T, K> keySelector;
final BiPredicate<? super K, ? super K> comparer;
public FlowableDistinctUntilChanged(Flowable<T> source, Function<? super T, K> keySelector, BiPredicate<? super K, ? super K> comparer) {
super(source);
this.keySelector = keySelector;
this.comparer = comparer;
Reported by PMD.
Line: 49
static final class DistinctUntilChangedSubscriber<T, K> extends BasicFuseableSubscriber<T, T>
implements ConditionalSubscriber<T> {
final Function<? super T, K> keySelector;
final BiPredicate<? super K, ? super K> comparer;
K last;
Reported by PMD.
Line: 51
final Function<? super T, K> keySelector;
final BiPredicate<? super K, ? super K> comparer;
K last;
boolean hasValue;
Reported by PMD.
Line: 53
final BiPredicate<? super K, ? super K> comparer;
K last;
boolean hasValue;
DistinctUntilChangedSubscriber(Subscriber<? super T> actual,
Function<? super T, K> keySelector,
Reported by PMD.
Line: 55
K last;
boolean hasValue;
DistinctUntilChangedSubscriber(Subscriber<? super T> actual,
Function<? super T, K> keySelector,
BiPredicate<? super K, ? super K> comparer) {
super(actual);
Reported by PMD.
Line: 96
hasValue = true;
last = key;
}
} catch (Throwable ex) {
fail(ex);
return true;
}
downstream.onNext(t);
Reported by PMD.
Line: 140
static final class DistinctUntilChangedConditionalSubscriber<T, K> extends BasicFuseableConditionalSubscriber<T, T> {
final Function<? super T, K> keySelector;
final BiPredicate<? super K, ? super K> comparer;
K last;
Reported by PMD.
Line: 142
final Function<? super T, K> keySelector;
final BiPredicate<? super K, ? super K> comparer;
K last;
boolean hasValue;
Reported by PMD.
Line: 144
final BiPredicate<? super K, ? super K> comparer;
K last;
boolean hasValue;
DistinctUntilChangedConditionalSubscriber(ConditionalSubscriber<? super T> actual,
Function<? super T, K> keySelector,
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/flowable/FlowableIntervalRange.java
14 issues
Line: 31
import io.reactivex.rxjava3.internal.util.BackpressureHelper;
public final class FlowableIntervalRange extends Flowable<Long> {
final Scheduler scheduler;
final long start;
final long end;
final long initialDelay;
final long period;
final TimeUnit unit;
Reported by PMD.
Line: 32
public final class FlowableIntervalRange extends Flowable<Long> {
final Scheduler scheduler;
final long start;
final long end;
final long initialDelay;
final long period;
final TimeUnit unit;
Reported by PMD.
Line: 33
public final class FlowableIntervalRange extends Flowable<Long> {
final Scheduler scheduler;
final long start;
final long end;
final long initialDelay;
final long period;
final TimeUnit unit;
public FlowableIntervalRange(long start, long end, long initialDelay, long period, TimeUnit unit, Scheduler scheduler) {
Reported by PMD.
Line: 34
final Scheduler scheduler;
final long start;
final long end;
final long initialDelay;
final long period;
final TimeUnit unit;
public FlowableIntervalRange(long start, long end, long initialDelay, long period, TimeUnit unit, Scheduler scheduler) {
this.initialDelay = initialDelay;
Reported by PMD.
Line: 35
final long start;
final long end;
final long initialDelay;
final long period;
final TimeUnit unit;
public FlowableIntervalRange(long start, long end, long initialDelay, long period, TimeUnit unit, Scheduler scheduler) {
this.initialDelay = initialDelay;
this.period = period;
Reported by PMD.
Line: 36
final long end;
final long initialDelay;
final long period;
final TimeUnit unit;
public FlowableIntervalRange(long start, long end, long initialDelay, long period, TimeUnit unit, Scheduler scheduler) {
this.initialDelay = initialDelay;
this.period = period;
this.unit = unit;
Reported by PMD.
Line: 57
if (sch instanceof TrampolineScheduler) {
Worker worker = sch.createWorker();
is.setResource(worker);
worker.schedulePeriodically(is, initialDelay, period, unit);
} else {
Disposable d = sch.schedulePeriodicallyDirect(is, initialDelay, period, unit);
is.setResource(d);
}
}
Reported by PMD.
Line: 69
private static final long serialVersionUID = -2809475196591179431L;
final Subscriber<? super Long> downstream;
final long end;
long count;
final AtomicReference<Disposable> resource = new AtomicReference<>();
Reported by PMD.
Line: 70
private static final long serialVersionUID = -2809475196591179431L;
final Subscriber<? super Long> downstream;
final long end;
long count;
final AtomicReference<Disposable> resource = new AtomicReference<>();
Reported by PMD.
Line: 72
final Subscriber<? super Long> downstream;
final long end;
long count;
final AtomicReference<Disposable> resource = new AtomicReference<>();
IntervalRangeSubscriber(Subscriber<? super Long> actual, long start, long end) {
this.downstream = actual;
Reported by PMD.
src/main/java/io/reactivex/rxjava3/internal/operators/flowable/FlowableMapNotification.java
14 issues
Line: 27
public final class FlowableMapNotification<T, R> extends AbstractFlowableWithUpstream<T, R> {
final Function<? super T, ? extends R> onNextMapper;
final Function<? super Throwable, ? extends R> onErrorMapper;
final Supplier<? extends R> onCompleteSupplier;
public FlowableMapNotification(
Flowable<T> source,
Reported by PMD.
Line: 28
public final class FlowableMapNotification<T, R> extends AbstractFlowableWithUpstream<T, R> {
final Function<? super T, ? extends R> onNextMapper;
final Function<? super Throwable, ? extends R> onErrorMapper;
final Supplier<? extends R> onCompleteSupplier;
public FlowableMapNotification(
Flowable<T> source,
Function<? super T, ? extends R> onNextMapper,
Reported by PMD.
Line: 29
final Function<? super T, ? extends R> onNextMapper;
final Function<? super Throwable, ? extends R> onErrorMapper;
final Supplier<? extends R> onCompleteSupplier;
public FlowableMapNotification(
Flowable<T> source,
Function<? super T, ? extends R> onNextMapper,
Function<? super Throwable, ? extends R> onErrorMapper,
Reported by PMD.
Line: 51
extends SinglePostCompleteSubscriber<T, R> {
private static final long serialVersionUID = 2757120512858778108L;
final Function<? super T, ? extends R> onNextMapper;
final Function<? super Throwable, ? extends R> onErrorMapper;
final Supplier<? extends R> onCompleteSupplier;
MapNotificationSubscriber(Subscriber<? super R> actual,
Function<? super T, ? extends R> onNextMapper,
Reported by PMD.
Line: 52
private static final long serialVersionUID = 2757120512858778108L;
final Function<? super T, ? extends R> onNextMapper;
final Function<? super Throwable, ? extends R> onErrorMapper;
final Supplier<? extends R> onCompleteSupplier;
MapNotificationSubscriber(Subscriber<? super R> actual,
Function<? super T, ? extends R> onNextMapper,
Function<? super Throwable, ? extends R> onErrorMapper,
Reported by PMD.
Line: 53
private static final long serialVersionUID = 2757120512858778108L;
final Function<? super T, ? extends R> onNextMapper;
final Function<? super Throwable, ? extends R> onErrorMapper;
final Supplier<? extends R> onCompleteSupplier;
MapNotificationSubscriber(Subscriber<? super R> actual,
Function<? super T, ? extends R> onNextMapper,
Function<? super Throwable, ? extends R> onErrorMapper,
Supplier<? extends R> onCompleteSupplier) {
Reported by PMD.
Line: 71
try {
p = Objects.requireNonNull(onNextMapper.apply(t), "The onNext publisher returned is null");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
downstream.onError(e);
return;
}
Reported by PMD.
Line: 87
try {
p = Objects.requireNonNull(onErrorMapper.apply(t), "The onError publisher returned is null");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
downstream.onError(new CompositeException(t, e));
return;
}
Reported by PMD.
Line: 102
try {
p = Objects.requireNonNull(onCompleteSupplier.get(), "The onComplete publisher returned is null");
} catch (Throwable e) {
Exceptions.throwIfFatal(e);
downstream.onError(e);
return;
}
Reported by PMD.
Line: 19
import org.reactivestreams.Subscriber;
import io.reactivex.rxjava3.core.Flowable;
import io.reactivex.rxjava3.exceptions.*;
import io.reactivex.rxjava3.functions.*;
import io.reactivex.rxjava3.internal.subscribers.SinglePostCompleteSubscriber;
import java.util.Objects;
Reported by PMD.