|
29 | 29 | import rx.Observer;
|
30 | 30 | import rx.exceptions.TestException;
|
31 | 31 | import rx.functions.*;
|
| 32 | +import rx.internal.util.RxRingBuffer; |
32 | 33 | import rx.observers.TestSubscriber;
|
33 | 34 | import rx.schedulers.Schedulers;
|
34 | 35 |
|
@@ -544,4 +545,109 @@ public Observable<Integer> call(Integer t) {
|
544 | 545 | ts.assertValueCount(n * 2);
|
545 | 546 | }
|
546 | 547 | }
|
| 548 | + |
| 549 | + @Test |
| 550 | + public void justEmptyMixture() { |
| 551 | + TestSubscriber<Integer> ts = TestSubscriber.create(); |
| 552 | + |
| 553 | + Observable.range(0, 4 * RxRingBuffer.SIZE) |
| 554 | + .flatMap(new Func1<Integer, Observable<Integer>>() { |
| 555 | + @Override |
| 556 | + public Observable<Integer> call(Integer v) { |
| 557 | + return (v & 1) == 0 ? Observable.<Integer>empty() : Observable.just(v); |
| 558 | + } |
| 559 | + }) |
| 560 | + .subscribe(ts); |
| 561 | + |
| 562 | + ts.assertValueCount(2 * RxRingBuffer.SIZE); |
| 563 | + ts.assertNoErrors(); |
| 564 | + ts.assertCompleted(); |
| 565 | + |
| 566 | + int j = 1; |
| 567 | + for (Integer v : ts.getOnNextEvents()) { |
| 568 | + Assert.assertEquals(j, v.intValue()); |
| 569 | + |
| 570 | + j += 2; |
| 571 | + } |
| 572 | + } |
| 573 | + |
| 574 | + @Test |
| 575 | + public void rangeEmptyMixture() { |
| 576 | + TestSubscriber<Integer> ts = TestSubscriber.create(); |
| 577 | + |
| 578 | + Observable.range(0, 4 * RxRingBuffer.SIZE) |
| 579 | + .flatMap(new Func1<Integer, Observable<Integer>>() { |
| 580 | + @Override |
| 581 | + public Observable<Integer> call(Integer v) { |
| 582 | + return (v & 1) == 0 ? Observable.<Integer>empty() : Observable.range(v, 2); |
| 583 | + } |
| 584 | + }) |
| 585 | + .subscribe(ts); |
| 586 | + |
| 587 | + ts.assertValueCount(4 * RxRingBuffer.SIZE); |
| 588 | + ts.assertNoErrors(); |
| 589 | + ts.assertCompleted(); |
| 590 | + |
| 591 | + int j = 1; |
| 592 | + List<Integer> list = ts.getOnNextEvents(); |
| 593 | + for (int i = 0; i < list.size(); i += 2) { |
| 594 | + Assert.assertEquals(j, list.get(i).intValue()); |
| 595 | + Assert.assertEquals(j + 1, list.get(i + 1).intValue()); |
| 596 | + |
| 597 | + j += 2; |
| 598 | + } |
| 599 | + } |
| 600 | + |
| 601 | + @Test |
| 602 | + public void justEmptyMixtureMaxConcurrent() { |
| 603 | + TestSubscriber<Integer> ts = TestSubscriber.create(); |
| 604 | + |
| 605 | + Observable.range(0, 4 * RxRingBuffer.SIZE) |
| 606 | + .flatMap(new Func1<Integer, Observable<Integer>>() { |
| 607 | + @Override |
| 608 | + public Observable<Integer> call(Integer v) { |
| 609 | + return (v & 1) == 0 ? Observable.<Integer>empty() : Observable.just(v); |
| 610 | + } |
| 611 | + }, 16) |
| 612 | + .subscribe(ts); |
| 613 | + |
| 614 | + ts.assertValueCount(2 * RxRingBuffer.SIZE); |
| 615 | + ts.assertNoErrors(); |
| 616 | + ts.assertCompleted(); |
| 617 | + |
| 618 | + int j = 1; |
| 619 | + for (Integer v : ts.getOnNextEvents()) { |
| 620 | + Assert.assertEquals(j, v.intValue()); |
| 621 | + |
| 622 | + j += 2; |
| 623 | + } |
| 624 | + } |
| 625 | + |
| 626 | + @Test |
| 627 | + public void rangeEmptyMixtureMaxConcurrent() { |
| 628 | + TestSubscriber<Integer> ts = TestSubscriber.create(); |
| 629 | + |
| 630 | + Observable.range(0, 4 * RxRingBuffer.SIZE) |
| 631 | + .flatMap(new Func1<Integer, Observable<Integer>>() { |
| 632 | + @Override |
| 633 | + public Observable<Integer> call(Integer v) { |
| 634 | + return (v & 1) == 0 ? Observable.<Integer>empty() : Observable.range(v, 2); |
| 635 | + } |
| 636 | + }, 16) |
| 637 | + .subscribe(ts); |
| 638 | + |
| 639 | + ts.assertValueCount(4 * RxRingBuffer.SIZE); |
| 640 | + ts.assertNoErrors(); |
| 641 | + ts.assertCompleted(); |
| 642 | + |
| 643 | + int j = 1; |
| 644 | + List<Integer> list = ts.getOnNextEvents(); |
| 645 | + for (int i = 0; i < list.size(); i += 2) { |
| 646 | + Assert.assertEquals(j, list.get(i).intValue()); |
| 647 | + Assert.assertEquals(j + 1, list.get(i + 1).intValue()); |
| 648 | + |
| 649 | + j += 2; |
| 650 | + } |
| 651 | + } |
| 652 | + |
547 | 653 | }
|
0 commit comments