QueueDrainSubscriber.smali
.class public abstract Lio/reactivex/internal/subscribers/QueueDrainSubscriber;
.super Lio/reactivex/internal/subscribers/QueueDrainSubscriberPad4;
.source "QueueDrainSubscriber.java"
# interfaces
.implements Lio/reactivex/FlowableSubscriber;
.implements Lio/reactivex/internal/util/QueueDrain;
# annotations
.annotation system Ldalvik/annotation/Signature;
value = {
"<T:",
"Ljava/lang/Object;",
"U:",
"Ljava/lang/Object;",
"V:",
"Ljava/lang/Object;",
">",
"Lio/reactivex/internal/subscribers/QueueDrainSubscriberPad4;",
"Lio/reactivex/FlowableSubscriber<",
"TT;>;",
"Lio/reactivex/internal/util/QueueDrain<",
"TU;TV;>;"
}
.end annotation
# instance fields
.field protected final actual:Lorg/reactivestreams/Subscriber;
.annotation system Ldalvik/annotation/Signature;
value = {
"Lorg/reactivestreams/Subscriber<",
"-TV;>;"
}
.end annotation
.end field
.field protected volatile cancelled:Z
.field protected volatile done:Z
.field protected error:Ljava/lang/Throwable;
.field protected final queue:Lio/reactivex/internal/fuseable/SimplePlainQueue;
.annotation system Ldalvik/annotation/Signature;
value = {
"Lio/reactivex/internal/fuseable/SimplePlainQueue<",
"TU;>;"
}
.end annotation
.end field
# direct methods
.method static constructor <clinit>()V
.registers 1
return-void
.end method
.method public constructor <init>(Lorg/reactivestreams/Subscriber;Lio/reactivex/internal/fuseable/SimplePlainQueue;)V
.registers 3
.annotation system Ldalvik/annotation/Signature;
value = {
"(",
"Lorg/reactivestreams/Subscriber<",
"-TV;>;",
"Lio/reactivex/internal/fuseable/SimplePlainQueue<",
"TU;>;)V"
}
.end annotation
.line 44
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
.local p1, "actual":Lorg/reactivestreams/Subscriber;, "Lorg/reactivestreams/Subscriber<-TV;>;"
.local p2, "queue":Lio/reactivex/internal/fuseable/SimplePlainQueue;, "Lio/reactivex/internal/fuseable/SimplePlainQueue<TU;>;"
invoke-direct {p0}, Lio/reactivex/internal/subscribers/QueueDrainSubscriberPad4;-><init>()V
.line 45
iput-object p1, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->actual:Lorg/reactivestreams/Subscriber;
.line 46
iput-object p2, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->queue:Lio/reactivex/internal/fuseable/SimplePlainQueue;
.line 47
return-void
.end method
# virtual methods
.method public accept(Lorg/reactivestreams/Subscriber;Ljava/lang/Object;)Z
.registers 4
.annotation system Ldalvik/annotation/Signature;
value = {
"(",
"Lorg/reactivestreams/Subscriber<",
"-TV;>;TU;)Z"
}
.end annotation
.line 133
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
.local p1, "a":Lorg/reactivestreams/Subscriber;, "Lorg/reactivestreams/Subscriber<-TV;>;"
.local p2, "v":Ljava/lang/Object;, "TU;"
const/4 v0, 0x0
return v0
.end method
.method public final cancelled()Z
.registers 2
.line 51
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-boolean v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->cancelled:Z
return v0
.end method
.method public final done()Z
.registers 2
.line 56
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-boolean v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->done:Z
return v0
.end method
.method public final enter()Z
.registers 2
.line 61
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
invoke-virtual {v0}, Ljava/util/concurrent/atomic/AtomicInteger;->getAndIncrement()I
move-result v0
if-nez v0, :cond_a
const/4 v0, 0x1
goto :goto_b
:cond_a
const/4 v0, 0x0
:goto_b
return v0
.end method
.method public final error()Ljava/lang/Throwable;
.registers 2
.line 138
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->error:Ljava/lang/Throwable;
return-object v0
.end method
.method public final fastEnter()Z
.registers 4
.line 65
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
invoke-virtual {v0}, Ljava/util/concurrent/atomic/AtomicInteger;->get()I
move-result v0
const/4 v1, 0x1
const/4 v2, 0x0
if-nez v0, :cond_13
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
invoke-virtual {v0, v2, v1}, Ljava/util/concurrent/atomic/AtomicInteger;->compareAndSet(II)Z
move-result v0
if-eqz v0, :cond_13
goto :goto_14
:cond_13
const/4 v1, 0x0
:goto_14
return v1
.end method
.method protected final fastPathEmitMax(Ljava/lang/Object;ZLio/reactivex/disposables/Disposable;)V
.registers 11
.param p2, "delayError" # Z
.param p3, "dispose" # Lio/reactivex/disposables/Disposable;
.annotation system Ldalvik/annotation/Signature;
value = {
"(TU;Z",
"Lio/reactivex/disposables/Disposable;",
")V"
}
.end annotation
.line 69
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
.local p1, "value":Ljava/lang/Object;, "TU;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->actual:Lorg/reactivestreams/Subscriber;
.line 70
.local v0, "s":Lorg/reactivestreams/Subscriber;, "Lorg/reactivestreams/Subscriber<-TV;>;"
iget-object v1, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->queue:Lio/reactivex/internal/fuseable/SimplePlainQueue;
.line 72
.local v1, "q":Lio/reactivex/internal/fuseable/SimplePlainQueue;, "Lio/reactivex/internal/fuseable/SimplePlainQueue<TU;>;"
iget-object v2, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
invoke-virtual {v2}, Ljava/util/concurrent/atomic/AtomicInteger;->get()I
move-result v2
if-nez v2, :cond_4d
iget-object v2, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
const/4 v3, 0x0
const/4 v4, 0x1
invoke-virtual {v2, v3, v4}, Ljava/util/concurrent/atomic/AtomicInteger;->compareAndSet(II)Z
move-result v2
if-eqz v2, :cond_4d
.line 73
iget-object v2, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->requested:Ljava/util/concurrent/atomic/AtomicLong;
invoke-virtual {v2}, Ljava/util/concurrent/atomic/AtomicLong;->get()J
move-result-wide v2
.line 74
.local v2, "r":J
const-wide/16 v4, 0x0
cmp-long v6, v2, v4
if-eqz v6, :cond_3f
.line 75
invoke-virtual {p0, v0, p1}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->accept(Lorg/reactivestreams/Subscriber;Ljava/lang/Object;)Z
move-result v4
if-eqz v4, :cond_36
.line 76
const-wide v4, 0x7fffffffffffffffL
cmp-long v6, v2, v4
if-eqz v6, :cond_36
.line 77
const-wide/16 v4, 0x1
invoke-virtual {p0, v4, v5}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->produced(J)J
.line 80
:cond_36
const/4 v4, -0x1
invoke-virtual {p0, v4}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->leave(I)I
move-result v4
if-nez v4, :cond_3e
.line 81
return-void
.line 88
.end local v2 # "r":J
:cond_3e
goto :goto_57
.line 84
.restart local v2 # "r":J
:cond_3f
invoke-interface {p3}, Lio/reactivex/disposables/Disposable;->dispose()V
.line 85
new-instance v4, Lio/reactivex/exceptions/MissingBackpressureException;
const-string v5, "Could not emit buffer due to lack of requests"
invoke-direct {v4, v5}, Lio/reactivex/exceptions/MissingBackpressureException;-><init>(Ljava/lang/String;)V
invoke-interface {v0, v4}, Lorg/reactivestreams/Subscriber;->onError(Ljava/lang/Throwable;)V
.line 86
return-void
.line 89
.end local v2 # "r":J
:cond_4d
invoke-interface {v1, p1}, Lio/reactivex/internal/fuseable/SimplePlainQueue;->offer(Ljava/lang/Object;)Z
.line 90
invoke-virtual {p0}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->enter()Z
move-result v2
if-nez v2, :cond_57
.line 91
return-void
.line 94
:cond_57
:goto_57
invoke-static {v1, v0, p2, p3, p0}, Lio/reactivex/internal/util/QueueDrainHelper;->drainMaxLoop(Lio/reactivex/internal/fuseable/SimplePlainQueue;Lorg/reactivestreams/Subscriber;ZLio/reactivex/disposables/Disposable;Lio/reactivex/internal/util/QueueDrain;)V
.line 95
return-void
.end method
.method protected final fastPathOrderedEmitMax(Ljava/lang/Object;ZLio/reactivex/disposables/Disposable;)V
.registers 12
.param p2, "delayError" # Z
.param p3, "dispose" # Lio/reactivex/disposables/Disposable;
.annotation system Ldalvik/annotation/Signature;
value = {
"(TU;Z",
"Lio/reactivex/disposables/Disposable;",
")V"
}
.end annotation
.line 98
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
.local p1, "value":Ljava/lang/Object;, "TU;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->actual:Lorg/reactivestreams/Subscriber;
.line 99
.local v0, "s":Lorg/reactivestreams/Subscriber;, "Lorg/reactivestreams/Subscriber<-TV;>;"
iget-object v1, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->queue:Lio/reactivex/internal/fuseable/SimplePlainQueue;
.line 101
.local v1, "q":Lio/reactivex/internal/fuseable/SimplePlainQueue;, "Lio/reactivex/internal/fuseable/SimplePlainQueue<TU;>;"
iget-object v2, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
invoke-virtual {v2}, Ljava/util/concurrent/atomic/AtomicInteger;->get()I
move-result v2
if-nez v2, :cond_58
iget-object v2, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
const/4 v3, 0x0
const/4 v4, 0x1
invoke-virtual {v2, v3, v4}, Ljava/util/concurrent/atomic/AtomicInteger;->compareAndSet(II)Z
move-result v2
if-eqz v2, :cond_58
.line 102
iget-object v2, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->requested:Ljava/util/concurrent/atomic/AtomicLong;
invoke-virtual {v2}, Ljava/util/concurrent/atomic/AtomicLong;->get()J
move-result-wide v2
.line 103
.local v2, "r":J
const-wide/16 v5, 0x0
cmp-long v7, v2, v5
if-eqz v7, :cond_48
.line 104
invoke-interface {v1}, Lio/reactivex/internal/fuseable/SimplePlainQueue;->isEmpty()Z
move-result v4
if-eqz v4, :cond_44
.line 105
invoke-virtual {p0, v0, p1}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->accept(Lorg/reactivestreams/Subscriber;Ljava/lang/Object;)Z
move-result v4
if-eqz v4, :cond_3c
.line 106
const-wide v4, 0x7fffffffffffffffL
cmp-long v6, v2, v4
if-eqz v6, :cond_3c
.line 107
const-wide/16 v4, 0x1
invoke-virtual {p0, v4, v5}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->produced(J)J
.line 110
:cond_3c
const/4 v4, -0x1
invoke-virtual {p0, v4}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->leave(I)I
move-result v4
if-nez v4, :cond_47
.line 111
return-void
.line 114
:cond_44
invoke-interface {v1, p1}, Lio/reactivex/internal/fuseable/SimplePlainQueue;->offer(Ljava/lang/Object;)Z
.line 122
.end local v2 # "r":J
:cond_47
goto :goto_62
.line 117
.restart local v2 # "r":J
:cond_48
iput-boolean v4, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->cancelled:Z
.line 118
invoke-interface {p3}, Lio/reactivex/disposables/Disposable;->dispose()V
.line 119
new-instance v4, Lio/reactivex/exceptions/MissingBackpressureException;
const-string v5, "Could not emit buffer due to lack of requests"
invoke-direct {v4, v5}, Lio/reactivex/exceptions/MissingBackpressureException;-><init>(Ljava/lang/String;)V
invoke-interface {v0, v4}, Lorg/reactivestreams/Subscriber;->onError(Ljava/lang/Throwable;)V
.line 120
return-void
.line 123
.end local v2 # "r":J
:cond_58
invoke-interface {v1, p1}, Lio/reactivex/internal/fuseable/SimplePlainQueue;->offer(Ljava/lang/Object;)Z
.line 124
invoke-virtual {p0}, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->enter()Z
move-result v2
if-nez v2, :cond_62
.line 125
return-void
.line 128
:cond_62
:goto_62
invoke-static {v1, v0, p2, p3, p0}, Lio/reactivex/internal/util/QueueDrainHelper;->drainMaxLoop(Lio/reactivex/internal/fuseable/SimplePlainQueue;Lorg/reactivestreams/Subscriber;ZLio/reactivex/disposables/Disposable;Lio/reactivex/internal/util/QueueDrain;)V
.line 129
return-void
.end method
.method public final leave(I)I
.registers 3
.param p1, "m" # I
.line 143
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->wip:Ljava/util/concurrent/atomic/AtomicInteger;
invoke-virtual {v0, p1}, Ljava/util/concurrent/atomic/AtomicInteger;->addAndGet(I)I
move-result v0
return v0
.end method
.method public final produced(J)J
.registers 6
.param p1, "n" # J
.line 153
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->requested:Ljava/util/concurrent/atomic/AtomicLong;
neg-long v1, p1
invoke-virtual {v0, v1, v2}, Ljava/util/concurrent/atomic/AtomicLong;->addAndGet(J)J
move-result-wide v0
return-wide v0
.end method
.method public final requested()J
.registers 3
.line 148
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->requested:Ljava/util/concurrent/atomic/AtomicLong;
invoke-virtual {v0}, Ljava/util/concurrent/atomic/AtomicLong;->get()J
move-result-wide v0
return-wide v0
.end method
.method public final requested(J)V
.registers 4
.param p1, "n" # J
.line 157
.local p0, "this":Lio/reactivex/internal/subscribers/QueueDrainSubscriber;, "Lio/reactivex/internal/subscribers/QueueDrainSubscriber<TT;TU;TV;>;"
invoke-static {p1, p2}, Lio/reactivex/internal/subscriptions/SubscriptionHelper;->validate(J)Z
move-result v0
if-eqz v0, :cond_b
.line 158
iget-object v0, p0, Lio/reactivex/internal/subscribers/QueueDrainSubscriber;->requested:Ljava/util/concurrent/atomic/AtomicLong;
invoke-static {v0, p1, p2}, Lio/reactivex/internal/util/BackpressureHelper;->add(Ljava/util/concurrent/atomic/AtomicLong;J)J
.line 160
:cond_b
return-void
.end method