|
18 | 18 | import static org.junit.Assert.assertEquals;
|
19 | 19 | import static org.junit.Assert.assertTrue;
|
20 | 20 |
|
| 21 | +import java.util.ArrayList; |
| 22 | +import java.util.Arrays; |
| 23 | +import java.util.List; |
21 | 24 | import java.util.concurrent.CountDownLatch;
|
22 | 25 | import java.util.concurrent.TimeUnit;
|
23 | 26 | import java.util.concurrent.atomic.AtomicInteger;
|
@@ -343,7 +346,6 @@ public void onError(Throwable e) {
|
343 | 346 |
|
344 | 347 | @Override
|
345 | 348 | public void onNext(Integer t) {
|
346 |
| - System.out.println(t); |
347 | 349 | request(1);
|
348 | 350 | }
|
349 | 351 |
|
@@ -375,7 +377,6 @@ public void onError(Throwable e) {
|
375 | 377 |
|
376 | 378 | @Override
|
377 | 379 | public void onNext(Integer t) {
|
378 |
| - System.out.println(t); |
379 | 380 | request(1);
|
380 | 381 | }
|
381 | 382 |
|
@@ -411,7 +412,6 @@ public void onError(Throwable e) {
|
411 | 412 |
|
412 | 413 | @Override
|
413 | 414 | public void onNext(Integer t) {
|
414 |
| - System.out.println(t); |
415 | 415 | child.onNext(t);
|
416 | 416 | request(1);
|
417 | 417 | }
|
@@ -454,4 +454,58 @@ public void onNext(Integer t) {
|
454 | 454 | assertTrue(latch.await(10, TimeUnit.SECONDS));
|
455 | 455 | assertTrue(exception.get() instanceof IllegalArgumentException);
|
456 | 456 | }
|
| 457 | + |
| 458 | + @Test |
| 459 | + public void testOnStartRequestsAreAdditive() { |
| 460 | + final List<Integer> list = new ArrayList<Integer>(); |
| 461 | + Observable.just(1,2,3,4,5).subscribe(new Subscriber<Integer>() { |
| 462 | + @Override |
| 463 | + public void onStart() { |
| 464 | + request(3); |
| 465 | + request(2); |
| 466 | + } |
| 467 | + |
| 468 | + @Override |
| 469 | + public void onCompleted() { |
| 470 | + |
| 471 | + } |
| 472 | + |
| 473 | + @Override |
| 474 | + public void onError(Throwable e) { |
| 475 | + |
| 476 | + } |
| 477 | + |
| 478 | + @Override |
| 479 | + public void onNext(Integer t) { |
| 480 | + list.add(t); |
| 481 | + }}); |
| 482 | + assertEquals(Arrays.asList(1,2,3,4,5), list); |
| 483 | + } |
| 484 | + |
| 485 | + @Test |
| 486 | + public void testOnStartRequestsAreAdditiveAndOverflowBecomesMaxValue() { |
| 487 | + final List<Integer> list = new ArrayList<Integer>(); |
| 488 | + Observable.just(1,2,3,4,5).subscribe(new Subscriber<Integer>() { |
| 489 | + @Override |
| 490 | + public void onStart() { |
| 491 | + request(2); |
| 492 | + request(Long.MAX_VALUE-1); |
| 493 | + } |
| 494 | + |
| 495 | + @Override |
| 496 | + public void onCompleted() { |
| 497 | + |
| 498 | + } |
| 499 | + |
| 500 | + @Override |
| 501 | + public void onError(Throwable e) { |
| 502 | + |
| 503 | + } |
| 504 | + |
| 505 | + @Override |
| 506 | + public void onNext(Integer t) { |
| 507 | + list.add(t); |
| 508 | + }}); |
| 509 | + assertEquals(Arrays.asList(1,2,3,4,5), list); |
| 510 | + } |
457 | 511 | }
|
0 commit comments