package io.reactivex.internal.operators.observable; import io.reactivex.ObservableSource; import io.reactivex.Observer; import io.reactivex.Scheduler; import io.reactivex.disposables.Disposable; import io.reactivex.internal.disposables.DisposableHelper; import io.reactivex.plugins.RxJavaPlugins; import java.util.concurrent.atomic.AtomicBoolean; /* loaded from: classes2.dex */ public final class ObservableUnsubscribeOn extends AbstractObservableWithUpstream { final Scheduler b; static final class UnsubscribeObserver extends AtomicBoolean implements Observer, Disposable { final Observer a; final Scheduler b; Disposable c; final class DisposeTask implements Runnable { DisposeTask() { } @Override // java.lang.Runnable public void run() { UnsubscribeObserver.this.c.dispose(); } } UnsubscribeObserver(Observer observer, Scheduler scheduler) { this.a = observer; this.b = scheduler; } @Override // io.reactivex.disposables.Disposable public void dispose() { if (compareAndSet(false, true)) { this.b.a(new DisposeTask()); } } @Override // io.reactivex.Observer public void onComplete() { if (get()) { return; } this.a.onComplete(); } @Override // io.reactivex.Observer public void onError(Throwable th) { if (get()) { RxJavaPlugins.b(th); } else { this.a.onError(th); } } @Override // io.reactivex.Observer public void onNext(T t) { if (get()) { return; } this.a.onNext(t); } @Override // io.reactivex.Observer public void onSubscribe(Disposable disposable) { if (DisposableHelper.validate(this.c, disposable)) { this.c = disposable; this.a.onSubscribe(this); } } } public ObservableUnsubscribeOn(ObservableSource observableSource, Scheduler scheduler) { super(observableSource); this.b = scheduler; } @Override // io.reactivex.Observable public void subscribeActual(Observer observer) { this.a.subscribe(new UnsubscribeObserver(observer, this.b)); } }