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 java.util.concurrent.atomic.AtomicReference; /* loaded from: classes2.dex */ public final class ObservableSubscribeOn extends AbstractObservableWithUpstream { final Scheduler b; static final class SubscribeOnObserver extends AtomicReference implements Observer, Disposable { final Observer a; final AtomicReference b = new AtomicReference<>(); SubscribeOnObserver(Observer observer) { this.a = observer; } void a(Disposable disposable) { DisposableHelper.setOnce(this, disposable); } @Override // io.reactivex.disposables.Disposable public void dispose() { DisposableHelper.dispose(this.b); DisposableHelper.dispose(this); } @Override // io.reactivex.Observer public void onComplete() { this.a.onComplete(); } @Override // io.reactivex.Observer public void onError(Throwable th) { this.a.onError(th); } @Override // io.reactivex.Observer public void onNext(T t) { this.a.onNext(t); } @Override // io.reactivex.Observer public void onSubscribe(Disposable disposable) { DisposableHelper.setOnce(this.b, disposable); } } final class SubscribeTask implements Runnable { private final SubscribeOnObserver a; SubscribeTask(SubscribeOnObserver subscribeOnObserver) { this.a = subscribeOnObserver; } @Override // java.lang.Runnable public void run() { ObservableSubscribeOn.this.a.subscribe(this.a); } } public ObservableSubscribeOn(ObservableSource observableSource, Scheduler scheduler) { super(observableSource); this.b = scheduler; } @Override // io.reactivex.Observable public void subscribeActual(Observer observer) { SubscribeOnObserver subscribeOnObserver = new SubscribeOnObserver(observer); observer.onSubscribe(subscribeOnObserver); subscribeOnObserver.a(this.b.a(new SubscribeTask(subscribeOnObserver))); } }