ObservableFlatMap$MergeObserver.smali
.class final Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;
.super Ljava/util/concurrent/atomic/AtomicInteger;
.source "ObservableFlatMap.java"
# interfaces
.implements Lio/reactivex/disposables/b;
.implements Lio/reactivex/r;
# annotations
.annotation system Ldalvik/annotation/Signature;
value = {
"<T:",
"Ljava/lang/Object;",
"U:",
"Ljava/lang/Object;",
">",
"Ljava/util/concurrent/atomic/AtomicInteger;",
"Lio/reactivex/disposables/b;",
"Lio/reactivex/r",
"<TT;>;"
}
.end annotation
# static fields
.field static final CANCELLED:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.annotation system Ldalvik/annotation/Signature;
value = {
"[",
"Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver",
"<**>;"
}
.end annotation
.end field
.field static final EMPTY:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.annotation system Ldalvik/annotation/Signature;
value = {
"[",
"Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver",
"<**>;"
}
.end annotation
.end field
.field private static final serialVersionUID:J = -0x1d634c9cafb5cc5aL
# instance fields
.field final actual:Lio/reactivex/r;
.annotation system Ldalvik/annotation/Signature;
value = {
"Lio/reactivex/r",
"<-TU;>;"
}
.end annotation
.end field
.field final bufferSize:I
.field volatile cancelled:Z
.field final delayErrors:Z
.field volatile done:Z
.field final errors:Lio/reactivex/internal/util/AtomicThrowable;
.field lastId:J
.field lastIndex:I
.field final mapper:Lio/reactivex/b/h;
.annotation system Ldalvik/annotation/Signature;
value = {
"Lio/reactivex/b/h",
"<-TT;+",
"Lio/reactivex/p",
"<+TU;>;>;"
}
.end annotation
.end field
.field final maxConcurrency:I
.field final observers:Ljava/util/concurrent/atomic/AtomicReference;
.annotation system Ldalvik/annotation/Signature;
value = {
"Ljava/util/concurrent/atomic/AtomicReference",
"<[",
"Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver",
"<**>;>;"
}
.end annotation
.end field
.field volatile queue:Lio/reactivex/internal/a/f;
.annotation system Ldalvik/annotation/Signature;
value = {
"Lio/reactivex/internal/a/f",
"<TU;>;"
}
.end annotation
.end field
.field s:Lio/reactivex/disposables/b;
.field sources:Ljava/util/Queue;
.annotation system Ldalvik/annotation/Signature;
value = {
"Ljava/util/Queue",
"<",
"Lio/reactivex/p",
"<+TU;>;>;"
}
.end annotation
.end field
.field uniqueId:J
.field wip:I
# direct methods
.method static constructor <clinit>()V
.registers 2
.prologue
const/4 v1, 0x0
.line 78
new-array v0, v1, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
sput-object v0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->EMPTY:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 80
new-array v0, v1, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
sput-object v0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->CANCELLED:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
return-void
.end method
.method constructor <init>(Lio/reactivex/r;Lio/reactivex/b/h;ZII)V
.registers 8
.annotation system Ldalvik/annotation/Signature;
value = {
"(",
"Lio/reactivex/r",
"<-TU;>;",
"Lio/reactivex/b/h",
"<-TT;+",
"Lio/reactivex/p",
"<+TU;>;>;ZII)V"
}
.end annotation
.prologue
.line 93
invoke-direct {p0}, Ljava/util/concurrent/atomic/AtomicInteger;-><init>()V
.line 72
new-instance v0, Lio/reactivex/internal/util/AtomicThrowable;
invoke-direct {v0}, Lio/reactivex/internal/util/AtomicThrowable;-><init>()V
iput-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
.line 94
iput-object p1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->actual:Lio/reactivex/r;
.line 95
iput-object p2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->mapper:Lio/reactivex/b/h;
.line 96
iput-boolean p3, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->delayErrors:Z
.line 97
iput p4, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->maxConcurrency:I
.line 98
iput p5, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->bufferSize:I
.line 99
const v0, 0x7fffffff
if-eq p4, v0, :cond_20
.line 100
new-instance v0, Ljava/util/ArrayDeque;
invoke-direct {v0, p4}, Ljava/util/ArrayDeque;-><init>(I)V
iput-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->sources:Ljava/util/Queue;
.line 102
:cond_20
new-instance v0, Ljava/util/concurrent/atomic/AtomicReference;
sget-object v1, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->EMPTY:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
invoke-direct {v0, v1}, Ljava/util/concurrent/atomic/AtomicReference;-><init>(Ljava/lang/Object;)V
iput-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
.line 103
return-void
.end method
# virtual methods
.method final addInner(Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;)Z
.registers 6
.annotation system Ldalvik/annotation/Signature;
value = {
"(",
"Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver",
"<TT;TU;>;)Z"
}
.end annotation
.prologue
const/4 v1, 0x0
.line 171
:cond_1
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
invoke-virtual {v0}, Ljava/util/concurrent/atomic/AtomicReference;->get()Ljava/lang/Object;
move-result-object v0
check-cast v0, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 172
sget-object v2, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->CANCELLED:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
if-ne v0, v2, :cond_12
.line 173
invoke-virtual {p1}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->dispose()V
move v0, v1
.line 181
:goto_11
return v0
.line 176
:cond_12
array-length v2, v0
.line 177
add-int/lit8 v3, v2, 0x1
new-array v3, v3, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 178
invoke-static {v0, v1, v3, v1, v2}, Ljava/lang/System;->arraycopy(Ljava/lang/Object;ILjava/lang/Object;II)V
.line 179
aput-object p1, v3, v2
.line 180
iget-object v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
invoke-virtual {v2, v0, v3}, Ljava/util/concurrent/atomic/AtomicReference;->compareAndSet(Ljava/lang/Object;Ljava/lang/Object;)Z
move-result v0
if-eqz v0, :cond_1
.line 181
const/4 v0, 0x1
goto :goto_11
.end method
.method final checkTerminate()Z
.registers 4
.prologue
const/4 v1, 0x1
.line 487
iget-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->cancelled:Z
if-eqz v0, :cond_7
move v0, v1
.line 499
:goto_6
return v0
.line 490
:cond_7
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
invoke-virtual {v0}, Lio/reactivex/internal/util/AtomicThrowable;->get()Ljava/lang/Object;
move-result-object v0
check-cast v0, Ljava/lang/Throwable;
.line 491
iget-boolean v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->delayErrors:Z
if-nez v2, :cond_29
if-eqz v0, :cond_29
.line 492
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->disposeAll()Z
.line 493
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
invoke-virtual {v0}, Lio/reactivex/internal/util/AtomicThrowable;->terminate()Ljava/lang/Throwable;
move-result-object v0
.line 494
sget-object v2, Lio/reactivex/internal/util/ExceptionHelper;->bTd:Ljava/lang/Throwable;
if-eq v0, v2, :cond_27
.line 495
iget-object v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->actual:Lio/reactivex/r;
invoke-interface {v2, v0}, Lio/reactivex/r;->onError(Ljava/lang/Throwable;)V
:cond_27
move v0, v1
.line 497
goto :goto_6
.line 499
:cond_29
const/4 v0, 0x0
goto :goto_6
.end method
.method public final dispose()V
.registers 3
.prologue
.line 305
iget-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->cancelled:Z
if-nez v0, :cond_1c
.line 306
const/4 v0, 0x1
iput-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->cancelled:Z
.line 307
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->disposeAll()Z
move-result v0
if-eqz v0, :cond_1c
.line 308
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
invoke-virtual {v0}, Lio/reactivex/internal/util/AtomicThrowable;->terminate()Ljava/lang/Throwable;
move-result-object v0
.line 309
if-eqz v0, :cond_1c
sget-object v1, Lio/reactivex/internal/util/ExceptionHelper;->bTd:Ljava/lang/Throwable;
if-eq v0, v1, :cond_1c
.line 310
invoke-static {v0}, Lio/reactivex/d/a;->onError(Ljava/lang/Throwable;)V
.line 314
:cond_1c
return-void
.end method
.method final disposeAll()Z
.registers 5
.prologue
const/4 v1, 0x0
.line 503
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->s:Lio/reactivex/disposables/b;
invoke-interface {v0}, Lio/reactivex/disposables/b;->dispose()V
.line 504
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
invoke-virtual {v0}, Ljava/util/concurrent/atomic/AtomicReference;->get()Ljava/lang/Object;
move-result-object v0
check-cast v0, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 505
sget-object v2, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->CANCELLED:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
if-eq v0, v2, :cond_2d
.line 506
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
sget-object v2, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->CANCELLED:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
invoke-virtual {v0, v2}, Ljava/util/concurrent/atomic/AtomicReference;->getAndSet(Ljava/lang/Object;)Ljava/lang/Object;
move-result-object v0
check-cast v0, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 507
sget-object v2, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->CANCELLED:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
if-eq v0, v2, :cond_2d
.line 508
array-length v2, v0
:goto_21
if-ge v1, v2, :cond_2b
aget-object v3, v0, v1
.line 509
invoke-virtual {v3}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->dispose()V
.line 508
add-int/lit8 v1, v1, 0x1
goto :goto_21
.line 511
:cond_2b
const/4 v0, 0x1
.line 514
:goto_2c
return v0
:cond_2d
move v0, v1
goto :goto_2c
.end method
.method final drain()V
.registers 2
.prologue
.line 322
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->getAndIncrement()I
move-result v0
if-nez v0, :cond_9
.line 323
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->drainLoop()V
.line 325
:cond_9
return-void
.end method
.method final drainLoop()V
.registers 15
.prologue
const/4 v2, 0x1
const/4 v4, 0x0
.line 328
iget-object v7, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->actual:Lio/reactivex/r;
move v1, v2
.line 331
:cond_5
:goto_5
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->checkTerminate()Z
move-result v0
if-eqz v0, :cond_c
.line 484
:cond_b
:goto_b
return-void
.line 334
:cond_c
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->queue:Lio/reactivex/internal/a/f;
.line 336
if-eqz v0, :cond_22
.line 340
:cond_10
:goto_10
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->checkTerminate()Z
move-result v3
if-nez v3, :cond_b
.line 344
invoke-interface {v0}, Lio/reactivex/internal/a/f;->poll()Ljava/lang/Object;
move-result-object v3
.line 346
if-eqz v3, :cond_20
.line 350
invoke-interface {v7, v3}, Lio/reactivex/r;->onNext(Ljava/lang/Object;)V
goto :goto_10
.line 352
:cond_20
if-nez v3, :cond_10
.line 358
:cond_22
iget-boolean v3, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->done:Z
.line 359
iget-object v5, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->queue:Lio/reactivex/internal/a/f;
.line 360
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
invoke-virtual {v0}, Ljava/util/concurrent/atomic/AtomicReference;->get()Ljava/lang/Object;
move-result-object v0
check-cast v0, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 361
array-length v8, v0
.line 363
if-eqz v3, :cond_4f
if-eqz v5, :cond_39
invoke-interface {v5}, Lio/reactivex/internal/a/f;->isEmpty()Z
move-result v3
if-eqz v3, :cond_4f
:cond_39
if-nez v8, :cond_4f
.line 364
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
invoke-virtual {v0}, Lio/reactivex/internal/util/AtomicThrowable;->terminate()Ljava/lang/Throwable;
move-result-object v0
.line 365
sget-object v1, Lio/reactivex/internal/util/ExceptionHelper;->bTd:Ljava/lang/Throwable;
if-eq v0, v1, :cond_b
.line 366
if-nez v0, :cond_4b
.line 367
invoke-interface {v7}, Lio/reactivex/r;->onComplete()V
goto :goto_b
.line 369
:cond_4b
invoke-interface {v7, v0}, Lio/reactivex/r;->onError(Ljava/lang/Throwable;)V
goto :goto_b
.line 376
:cond_4f
if-eqz v8, :cond_118
.line 377
iget-wide v10, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->lastId:J
.line 378
iget v3, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->lastIndex:I
.line 380
if-le v8, v3, :cond_5f
aget-object v5, v0, v3
iget-wide v12, v5, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->id:J
cmp-long v5, v12, v10
if-eqz v5, :cond_7d
.line 381
:cond_5f
if-gt v8, v3, :cond_62
move v3, v4
:cond_62
move v5, v4
.line 385
:goto_63
if-ge v5, v8, :cond_75
.line 386
aget-object v6, v0, v3
iget-wide v12, v6, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->id:J
cmp-long v6, v12, v10
if-eqz v6, :cond_75
.line 389
add-int/lit8 v3, v3, 0x1
.line 390
if-ne v3, v8, :cond_72
move v3, v4
.line 385
:cond_72
add-int/lit8 v5, v5, 0x1
goto :goto_63
.line 395
:cond_75
iput v3, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->lastIndex:I
.line 396
aget-object v5, v0, v3
iget-wide v10, v5, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->id:J
iput-wide v10, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->lastId:J
:cond_7d
move v5, v4
move v6, v4
.line 401
:goto_7f
if-ge v5, v8, :cond_df
.line 402
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->checkTerminate()Z
move-result v9
if-nez v9, :cond_b
.line 406
aget-object v9, v0, v3
.line 409
:cond_89
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->checkTerminate()Z
move-result v10
if-nez v10, :cond_b
.line 412
iget-object v10, v9, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->queue:Lio/reactivex/internal/a/g;
.line 413
if-eqz v10, :cond_c1
.line 419
:cond_93
:try_start_93
invoke-interface {v10}, Lio/reactivex/internal/a/g;->poll()Ljava/lang/Object;
:try_end_96
.catch Ljava/lang/Throwable; {:try_start_93 .. :try_end_96} :catch_a4
move-result-object v11
.line 432
if-eqz v11, :cond_bf
.line 436
invoke-interface {v7, v11}, Lio/reactivex/r;->onNext(Ljava/lang/Object;)V
.line 438
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->checkTerminate()Z
move-result v11
if-eqz v11, :cond_93
goto/16 :goto_b
.line 420
:catch_a4
move-exception v6
.line 421
invoke-static {v6}, Lio/reactivex/exceptions/d;->throwIfFatal(Ljava/lang/Throwable;)V
.line 422
invoke-virtual {v9}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->dispose()V
.line 423
iget-object v10, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
invoke-virtual {v10, v6}, Lio/reactivex/internal/util/AtomicThrowable;->addThrowable(Ljava/lang/Throwable;)Z
.line 424
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->checkTerminate()Z
move-result v6
if-nez v6, :cond_b
.line 427
invoke-virtual {p0, v9}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->removeInner(Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;)V
.line 429
add-int/lit8 v5, v5, 0x1
move v6, v2
.line 401
:cond_bc
:goto_bc
add-int/lit8 v5, v5, 0x1
goto :goto_7f
.line 442
:cond_bf
if-nez v11, :cond_89
.line 446
:cond_c1
iget-boolean v10, v9, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->done:Z
.line 447
iget-object v11, v9, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->queue:Lio/reactivex/internal/a/g;
.line 448
if-eqz v10, :cond_d9
if-eqz v11, :cond_cf
invoke-interface {v11}, Lio/reactivex/internal/a/g;->isEmpty()Z
move-result v10
if-eqz v10, :cond_d9
.line 449
:cond_cf
invoke-virtual {p0, v9}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->removeInner(Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;)V
.line 450
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->checkTerminate()Z
move-result v6
if-nez v6, :cond_b
move v6, v2
.line 456
:cond_d9
add-int/lit8 v3, v3, 0x1
.line 457
if-ne v3, v8, :cond_bc
move v3, v4
.line 458
goto :goto_bc
.line 461
:cond_df
iput v3, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->lastIndex:I
.line 462
aget-object v0, v0, v3
iget-wide v8, v0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->id:J
iput-wide v8, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->lastId:J
move v0, v6
.line 465
:goto_e8
if-eqz v0, :cond_10e
.line 466
iget v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->maxConcurrency:I
const v3, 0x7fffffff
if-eq v0, v3, :cond_5
.line 468
monitor-enter p0
.line 469
:try_start_f2
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->sources:Ljava/util/Queue;
invoke-interface {v0}, Ljava/util/Queue;->poll()Ljava/lang/Object;
move-result-object v0
check-cast v0, Lio/reactivex/p;
.line 470
if-nez v0, :cond_108
.line 471
iget v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->wip:I
add-int/lit8 v0, v0, -0x1
iput v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->wip:I
.line 472
monitor-exit p0
goto/16 :goto_5
.line 474
:catchall_105
move-exception v0
monitor-exit p0
:try_end_107
.catchall {:try_start_f2 .. :try_end_107} :catchall_105
throw v0
:cond_108
:try_start_108
monitor-exit p0
:try_end_109
.catchall {:try_start_108 .. :try_end_109} :catchall_105
.line 475
invoke-virtual {p0, v0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->subscribeInner(Lio/reactivex/p;)V
goto/16 :goto_5
.line 479
:cond_10e
neg-int v0, v1
invoke-virtual {p0, v0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->addAndGet(I)I
move-result v0
.line 480
if-eqz v0, :cond_b
move v1, v0
goto/16 :goto_5
:cond_118
move v0, v4
goto :goto_e8
.end method
.method public final isDisposed()Z
.registers 2
.prologue
.line 318
iget-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->cancelled:Z
return v0
.end method
.method public final onComplete()V
.registers 2
.prologue
.line 296
iget-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->done:Z
if-eqz v0, :cond_5
.line 301
:goto_4
return-void
.line 299
:cond_5
const/4 v0, 0x1
iput-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->done:Z
.line 300
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->drain()V
goto :goto_4
.end method
.method public final onError(Ljava/lang/Throwable;)V
.registers 3
.prologue
.line 282
iget-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->done:Z
if-eqz v0, :cond_8
.line 283
invoke-static {p1}, Lio/reactivex/d/a;->onError(Ljava/lang/Throwable;)V
.line 292
:goto_7
return-void
.line 286
:cond_8
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
invoke-virtual {v0, p1}, Lio/reactivex/internal/util/AtomicThrowable;->addThrowable(Ljava/lang/Throwable;)Z
move-result v0
if-eqz v0, :cond_17
.line 287
const/4 v0, 0x1
iput-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->done:Z
.line 288
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->drain()V
goto :goto_7
.line 290
:cond_17
invoke-static {p1}, Lio/reactivex/d/a;->onError(Ljava/lang/Throwable;)V
goto :goto_7
.end method
.method public final onNext(Ljava/lang/Object;)V
.registers 5
.annotation system Ldalvik/annotation/Signature;
value = {
"(TT;)V"
}
.end annotation
.prologue
.line 116
iget-boolean v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->done:Z
if-eqz v0, :cond_5
.line 140
:goto_4
return-void
.line 121
:cond_5
:try_start_5
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->mapper:Lio/reactivex/b/h;
invoke-interface {v0, p1}, Lio/reactivex/b/h;->apply(Ljava/lang/Object;)Ljava/lang/Object;
move-result-object v0
const-string v1, "The mapper returned a null ObservableSource"
invoke-static {v0, v1}, Lio/reactivex/internal/functions/aj;->requireNonNull(Ljava/lang/Object;Ljava/lang/String;)Ljava/lang/Object;
move-result-object v0
check-cast v0, Lio/reactivex/p;
:try_end_13
.catch Ljava/lang/Throwable; {:try_start_5 .. :try_end_13} :catch_2b
.line 129
iget v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->maxConcurrency:I
const v2, 0x7fffffff
if-eq v1, v2, :cond_3f
.line 130
monitor-enter p0
.line 131
:try_start_1b
iget v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->wip:I
iget v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->maxConcurrency:I
if-ne v1, v2, :cond_38
.line 132
iget-object v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->sources:Ljava/util/Queue;
invoke-interface {v1, v0}, Ljava/util/Queue;->offer(Ljava/lang/Object;)Z
.line 133
monitor-exit p0
goto :goto_4
.line 136
:catchall_28
move-exception v0
monitor-exit p0
:try_end_2a
.catchall {:try_start_1b .. :try_end_2a} :catchall_28
throw v0
.line 122
:catch_2b
move-exception v0
.line 123
invoke-static {v0}, Lio/reactivex/exceptions/d;->throwIfFatal(Ljava/lang/Throwable;)V
.line 124
iget-object v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->s:Lio/reactivex/disposables/b;
invoke-interface {v1}, Lio/reactivex/disposables/b;->dispose()V
.line 125
invoke-virtual {p0, v0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->onError(Ljava/lang/Throwable;)V
goto :goto_4
.line 135
:cond_38
:try_start_38
iget v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->wip:I
add-int/lit8 v1, v1, 0x1
iput v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->wip:I
.line 136
monitor-exit p0
:try_end_3f
.catchall {:try_start_38 .. :try_end_3f} :catchall_28
.line 139
:cond_3f
invoke-virtual {p0, v0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->subscribeInner(Lio/reactivex/p;)V
goto :goto_4
.end method
.method public final onSubscribe(Lio/reactivex/disposables/b;)V
.registers 3
.prologue
.line 107
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->s:Lio/reactivex/disposables/b;
invoke-static {v0, p1}, Lio/reactivex/internal/disposables/DisposableHelper;->validate(Lio/reactivex/disposables/b;Lio/reactivex/disposables/b;)Z
move-result v0
if-eqz v0, :cond_f
.line 108
iput-object p1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->s:Lio/reactivex/disposables/b;
.line 109
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->actual:Lio/reactivex/r;
invoke-interface {v0, p0}, Lio/reactivex/r;->onSubscribe(Lio/reactivex/disposables/b;)V
.line 111
:cond_f
return-void
.end method
.method final removeInner(Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;)V
.registers 8
.annotation system Ldalvik/annotation/Signature;
value = {
"(",
"Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver",
"<TT;TU;>;)V"
}
.end annotation
.prologue
const/4 v3, 0x0
.line 188
:cond_1
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
invoke-virtual {v0}, Ljava/util/concurrent/atomic/AtomicReference;->get()Ljava/lang/Object;
move-result-object v0
check-cast v0, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 189
array-length v4, v0
.line 190
if-nez v4, :cond_d
.line 212
:cond_c
:goto_c
return-void
.line 193
:cond_d
const/4 v2, -0x1
move v1, v3
.line 194
:goto_f
if-ge v1, v4, :cond_16
.line 195
aget-object v5, v0, v1
if-ne v5, p1, :cond_26
move v2, v1
.line 200
:cond_16
if-ltz v2, :cond_c
.line 204
const/4 v1, 0x1
if-ne v4, v1, :cond_29
.line 205
sget-object v1, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->EMPTY:[Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 211
:goto_1d
iget-object v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->observers:Ljava/util/concurrent/atomic/AtomicReference;
invoke-virtual {v2, v0, v1}, Ljava/util/concurrent/atomic/AtomicReference;->compareAndSet(Ljava/lang/Object;Ljava/lang/Object;)Z
move-result v0
if-eqz v0, :cond_1
goto :goto_c
.line 194
:cond_26
add-int/lit8 v1, v1, 0x1
goto :goto_f
.line 207
:cond_29
add-int/lit8 v1, v4, -0x1
new-array v1, v1, [Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
.line 208
invoke-static {v0, v3, v1, v3, v2}, Ljava/lang/System;->arraycopy(Ljava/lang/Object;ILjava/lang/Object;II)V
.line 209
add-int/lit8 v5, v2, 0x1
sub-int/2addr v4, v2
add-int/lit8 v4, v4, -0x1
invoke-static {v0, v5, v1, v2, v4}, Ljava/lang/System;->arraycopy(Ljava/lang/Object;ILjava/lang/Object;II)V
goto :goto_1d
.end method
.method final subscribeInner(Lio/reactivex/p;)V
.registers 8
.annotation system Ldalvik/annotation/Signature;
value = {
"(",
"Lio/reactivex/p",
"<+TU;>;)V"
}
.end annotation
.prologue
.line 145
move-object v0, p1
:goto_1
instance-of v1, v0, Ljava/util/concurrent/Callable;
if-eqz v1, :cond_29
.line 146
check-cast v0, Ljava/util/concurrent/Callable;
invoke-virtual {p0, v0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->tryEmitScalar(Ljava/util/concurrent/Callable;)V
.line 148
iget v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->maxConcurrency:I
const v1, 0x7fffffff
if-eq v0, v1, :cond_23
.line 149
monitor-enter p0
.line 150
:try_start_12
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->sources:Ljava/util/Queue;
invoke-interface {v0}, Ljava/util/Queue;->poll()Ljava/lang/Object;
move-result-object v0
check-cast v0, Lio/reactivex/p;
.line 151
if-nez v0, :cond_24
.line 152
iget v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->wip:I
add-int/lit8 v0, v0, -0x1
iput v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->wip:I
.line 153
monitor-exit p0
.line 167
:cond_23
:goto_23
return-void
.line 155
:cond_24
monitor-exit p0
goto :goto_1
:catchall_26
move-exception v0
monitor-exit p0
:try_end_28
.catchall {:try_start_12 .. :try_end_28} :catchall_26
throw v0
.line 160
:cond_29
new-instance v1, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;
iget-wide v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->uniqueId:J
const-wide/16 v4, 0x1
add-long/2addr v4, v2
iput-wide v4, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->uniqueId:J
invoke-direct {v1, p0, v2, v3}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;-><init>(Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;J)V
.line 161
invoke-virtual {p0, v1}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->addInner(Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;)Z
move-result v2
if-eqz v2, :cond_23
.line 162
invoke-interface {v0, v1}, Lio/reactivex/p;->subscribe(Lio/reactivex/r;)V
goto :goto_23
.end method
.method final tryEmit(Ljava/lang/Object;Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;)V
.registers 5
.annotation system Ldalvik/annotation/Signature;
value = {
"(TU;",
"Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver",
"<TT;TU;>;)V"
}
.end annotation
.prologue
.line 261
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->get()I
move-result v0
if-nez v0, :cond_1a
const/4 v0, 0x0
const/4 v1, 0x1
invoke-virtual {p0, v0, v1}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->compareAndSet(II)Z
move-result v0
if-eqz v0, :cond_1a
.line 262
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->actual:Lio/reactivex/r;
invoke-interface {v0, p1}, Lio/reactivex/r;->onNext(Ljava/lang/Object;)V
.line 263
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->decrementAndGet()I
move-result v0
if-nez v0, :cond_30
.line 278
:cond_19
:goto_19
return-void
.line 267
:cond_1a
iget-object v0, p2, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->queue:Lio/reactivex/internal/a/g;
.line 268
if-nez v0, :cond_27
.line 269
new-instance v0, Lio/reactivex/internal/queue/a;
iget v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->bufferSize:I
invoke-direct {v0, v1}, Lio/reactivex/internal/queue/a;-><init>(I)V
.line 270
iput-object v0, p2, Lio/reactivex/internal/operators/observable/ObservableFlatMap$InnerObserver;->queue:Lio/reactivex/internal/a/g;
.line 272
:cond_27
invoke-interface {v0, p1}, Lio/reactivex/internal/a/g;->offer(Ljava/lang/Object;)Z
.line 273
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->getAndIncrement()I
move-result v0
if-nez v0, :cond_19
.line 277
:cond_30
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->drainLoop()V
goto :goto_19
.end method
.method final tryEmitScalar(Ljava/util/concurrent/Callable;)V
.registers 5
.annotation system Ldalvik/annotation/Signature;
value = {
"(",
"Ljava/util/concurrent/Callable",
"<+TU;>;)V"
}
.end annotation
.prologue
.line 220
:try_start_0
invoke-interface {p1}, Ljava/util/concurrent/Callable;->call()Ljava/lang/Object;
:try_end_3
.catch Ljava/lang/Throwable; {:try_start_0 .. :try_end_3} :catch_7
move-result-object v1
.line 228
if-nez v1, :cond_14
.line 258
:cond_6
:goto_6
return-void
.line 221
:catch_7
move-exception v0
.line 222
invoke-static {v0}, Lio/reactivex/exceptions/d;->throwIfFatal(Ljava/lang/Throwable;)V
.line 223
iget-object v1, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->errors:Lio/reactivex/internal/util/AtomicThrowable;
invoke-virtual {v1, v0}, Lio/reactivex/internal/util/AtomicThrowable;->addThrowable(Ljava/lang/Throwable;)Z
.line 224
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->drain()V
goto :goto_6
.line 233
:cond_14
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->get()I
move-result v0
if-nez v0, :cond_31
const/4 v0, 0x0
const/4 v2, 0x1
invoke-virtual {p0, v0, v2}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->compareAndSet(II)Z
move-result v0
if-eqz v0, :cond_31
.line 234
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->actual:Lio/reactivex/r;
invoke-interface {v0, v1}, Lio/reactivex/r;->onNext(Ljava/lang/Object;)V
.line 235
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->decrementAndGet()I
move-result v0
if-eqz v0, :cond_6
.line 257
:cond_2d
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->drainLoop()V
goto :goto_6
.line 239
:cond_31
iget-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->queue:Lio/reactivex/internal/a/f;
.line 240
if-nez v0, :cond_45
.line 241
iget v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->maxConcurrency:I
const v2, 0x7fffffff
if-ne v0, v2, :cond_56
.line 242
new-instance v0, Lio/reactivex/internal/queue/a;
iget v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->bufferSize:I
invoke-direct {v0, v2}, Lio/reactivex/internal/queue/a;-><init>(I)V
.line 246
:goto_43
iput-object v0, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->queue:Lio/reactivex/internal/a/f;
.line 249
:cond_45
invoke-interface {v0, v1}, Lio/reactivex/internal/a/f;->offer(Ljava/lang/Object;)Z
move-result v0
if-nez v0, :cond_5e
.line 250
new-instance v0, Ljava/lang/IllegalStateException;
const-string v1, "Scalar queue full?!"
invoke-direct {v0, v1}, Ljava/lang/IllegalStateException;-><init>(Ljava/lang/String;)V
invoke-virtual {p0, v0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->onError(Ljava/lang/Throwable;)V
goto :goto_6
.line 244
:cond_56
new-instance v0, Lio/reactivex/internal/queue/SpscArrayQueue;
iget v2, p0, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->maxConcurrency:I
invoke-direct {v0, v2}, Lio/reactivex/internal/queue/SpscArrayQueue;-><init>(I)V
goto :goto_43
.line 253
:cond_5e
invoke-virtual {p0}, Lio/reactivex/internal/operators/observable/ObservableFlatMap$MergeObserver;->getAndIncrement()I
move-result v0
if-eqz v0, :cond_2d
goto :goto_6
.end method