|
27 | 27 | import org.mockito.Mockito;
|
28 | 28 |
|
29 | 29 | import rx.Observable;
|
| 30 | +import rx.Observable.OnSubscribe; |
30 | 31 | import rx.Observer;
|
31 | 32 | import rx.Subscriber;
|
32 | 33 | import rx.Subscription;
|
@@ -105,6 +106,55 @@ public String call(String s) {
|
105 | 106 | verify(observer, times(1)).onNext("twoResume");
|
106 | 107 | verify(observer, times(1)).onNext("threeResume");
|
107 | 108 | }
|
| 109 | + |
| 110 | + @Test |
| 111 | + public void testResumeNextWithFailedOnSubscribe() { |
| 112 | + Subscription s = mock(Subscription.class); |
| 113 | + Observable<String> testObservable = Observable.create(new OnSubscribe<String>() { |
| 114 | + |
| 115 | + @Override |
| 116 | + public void call(Subscriber<? super String> t1) { |
| 117 | + throw new RuntimeException("force failure"); |
| 118 | + } |
| 119 | + |
| 120 | + }); |
| 121 | + Observable<String> resume = Observable.just("resume"); |
| 122 | + Observable<String> observable = testObservable.onErrorResumeNext(resume); |
| 123 | + |
| 124 | + @SuppressWarnings("unchecked") |
| 125 | + Observer<String> observer = mock(Observer.class); |
| 126 | + observable.subscribe(observer); |
| 127 | + |
| 128 | + verify(observer, Mockito.never()).onError(any(Throwable.class)); |
| 129 | + verify(observer, times(1)).onCompleted(); |
| 130 | + verify(observer, times(1)).onNext("resume"); |
| 131 | + } |
| 132 | + |
| 133 | + @Test |
| 134 | + public void testResumeNextWithFailedOnSubscribeAsync() { |
| 135 | + Subscription s = mock(Subscription.class); |
| 136 | + Observable<String> testObservable = Observable.create(new OnSubscribe<String>() { |
| 137 | + |
| 138 | + @Override |
| 139 | + public void call(Subscriber<? super String> t1) { |
| 140 | + throw new RuntimeException("force failure"); |
| 141 | + } |
| 142 | + |
| 143 | + }); |
| 144 | + Observable<String> resume = Observable.just("resume"); |
| 145 | + Observable<String> observable = testObservable.subscribeOn(Schedulers.io()).onErrorResumeNext(resume); |
| 146 | + |
| 147 | + @SuppressWarnings("unchecked") |
| 148 | + Observer<String> observer = mock(Observer.class); |
| 149 | + TestSubscriber<String> ts = new TestSubscriber<String>(observer); |
| 150 | + observable.subscribe(ts); |
| 151 | + |
| 152 | + ts.awaitTerminalEvent(); |
| 153 | + |
| 154 | + verify(observer, Mockito.never()).onError(any(Throwable.class)); |
| 155 | + verify(observer, times(1)).onCompleted(); |
| 156 | + verify(observer, times(1)).onNext("resume"); |
| 157 | + } |
108 | 158 |
|
109 | 159 | private static class TestObservable implements Observable.OnSubscribe<String> {
|
110 | 160 |
|
|
0 commit comments