package io.reactivex.internal.operators.observable; import io.reactivex.ObservableSource; import io.reactivex.Observer; import io.reactivex.disposables.Disposable; import io.reactivex.internal.disposables.DisposableHelper; import java.util.ArrayDeque; /* loaded from: classes2.dex */ public final class ObservableTakeLast extends AbstractObservableWithUpstream { final int b; static final class TakeLastObserver extends ArrayDeque implements Observer, Disposable { final Observer a; final int b; Disposable c; volatile boolean d; TakeLastObserver(Observer observer, int i) { this.a = observer; this.b = i; } @Override // io.reactivex.disposables.Disposable public void dispose() { if (this.d) { return; } this.d = true; this.c.dispose(); } @Override // io.reactivex.Observer public void onComplete() { Observer observer = this.a; while (!this.d) { T poll = poll(); if (poll == null) { if (this.d) { return; } observer.onComplete(); return; } observer.onNext(poll); } } @Override // io.reactivex.Observer public void onError(Throwable th) { this.a.onError(th); } @Override // io.reactivex.Observer public void onNext(T t) { if (this.b == size()) { poll(); } offer(t); } @Override // io.reactivex.Observer public void onSubscribe(Disposable disposable) { if (DisposableHelper.validate(this.c, disposable)) { this.c = disposable; this.a.onSubscribe(this); } } } public ObservableTakeLast(ObservableSource observableSource, int i) { super(observableSource); this.b = i; } @Override // io.reactivex.Observable public void subscribeActual(Observer observer) { this.a.subscribe(new TakeLastObserver(observer, this.b)); } }