|
15 | 15 | */
|
16 | 16 | package rx.operators;
|
17 | 17 |
|
18 |
| -import static org.mockito.Matchers.*; |
19 | 18 | import static org.mockito.Mockito.*;
|
20 | 19 |
|
21 | 20 | import java.util.concurrent.TimeUnit;
|
|
28 | 27 | import rx.Observer;
|
29 | 28 | import rx.Subscription;
|
30 | 29 | import rx.concurrency.TestScheduler;
|
| 30 | +import rx.subjects.PublishSubject; |
31 | 31 | import rx.subscriptions.Subscriptions;
|
32 | 32 | import rx.util.functions.Action0;
|
33 | 33 |
|
34 | 34 | public class OperationSampleTest {
|
35 | 35 | private TestScheduler scheduler;
|
36 | 36 | private Observer<Long> observer;
|
| 37 | + private Observer<Object> observer2; |
37 | 38 |
|
38 | 39 | @Before
|
39 | 40 | @SuppressWarnings("unchecked")
|
40 | 41 | // due to mocking
|
41 | 42 | public void before() {
|
42 | 43 | scheduler = new TestScheduler();
|
43 | 44 | observer = mock(Observer.class);
|
| 45 | + observer2 = mock(Observer.class); |
44 | 46 | }
|
45 | 47 |
|
46 | 48 | @Test
|
@@ -105,4 +107,159 @@ public void call() {
|
105 | 107 | verify(observer, times(1)).onCompleted();
|
106 | 108 | verify(observer, never()).onError(any(Throwable.class));
|
107 | 109 | }
|
| 110 | + @Test |
| 111 | + public void sampleWithSamplerNormal() { |
| 112 | + PublishSubject<Integer> source = PublishSubject.create(); |
| 113 | + PublishSubject<Integer> sampler = PublishSubject.create(); |
| 114 | + |
| 115 | + Observable<Integer> m = source.sample(sampler); |
| 116 | + m.subscribe(observer2); |
| 117 | + |
| 118 | + source.onNext(1); |
| 119 | + source.onNext(2); |
| 120 | + sampler.onNext(1); |
| 121 | + source.onNext(3); |
| 122 | + source.onNext(4); |
| 123 | + sampler.onNext(2); |
| 124 | + source.onCompleted(); |
| 125 | + sampler.onNext(3); |
| 126 | + |
| 127 | + |
| 128 | + InOrder inOrder = inOrder(observer2); |
| 129 | + inOrder.verify(observer2, never()).onNext(1); |
| 130 | + inOrder.verify(observer2, times(1)).onNext(2); |
| 131 | + inOrder.verify(observer2, never()).onNext(3); |
| 132 | + inOrder.verify(observer2, times(1)).onNext(4); |
| 133 | + inOrder.verify(observer2, times(1)).onCompleted(); |
| 134 | + verify(observer, never()).onError(any(Throwable.class)); |
| 135 | + } |
| 136 | + @Test |
| 137 | + public void sampleWithSamplerNoDuplicates() { |
| 138 | + PublishSubject<Integer> source = PublishSubject.create(); |
| 139 | + PublishSubject<Integer> sampler = PublishSubject.create(); |
| 140 | + |
| 141 | + Observable<Integer> m = source.sample(sampler); |
| 142 | + m.subscribe(observer2); |
| 143 | + |
| 144 | + source.onNext(1); |
| 145 | + source.onNext(2); |
| 146 | + sampler.onNext(1); |
| 147 | + sampler.onNext(1); |
| 148 | + |
| 149 | + source.onNext(3); |
| 150 | + source.onNext(4); |
| 151 | + sampler.onNext(2); |
| 152 | + sampler.onNext(2); |
| 153 | + |
| 154 | + source.onCompleted(); |
| 155 | + sampler.onNext(3); |
| 156 | + |
| 157 | + |
| 158 | + InOrder inOrder = inOrder(observer2); |
| 159 | + inOrder.verify(observer2, never()).onNext(1); |
| 160 | + inOrder.verify(observer2, times(1)).onNext(2); |
| 161 | + inOrder.verify(observer2, never()).onNext(3); |
| 162 | + inOrder.verify(observer2, times(1)).onNext(4); |
| 163 | + inOrder.verify(observer2, times(1)).onCompleted(); |
| 164 | + verify(observer, never()).onError(any(Throwable.class)); |
| 165 | + } |
| 166 | + @Test |
| 167 | + public void sampleWithSamplerTerminatingEarly() { |
| 168 | + PublishSubject<Integer> source = PublishSubject.create(); |
| 169 | + PublishSubject<Integer> sampler = PublishSubject.create(); |
| 170 | + |
| 171 | + Observable<Integer> m = source.sample(sampler); |
| 172 | + m.subscribe(observer2); |
| 173 | + |
| 174 | + source.onNext(1); |
| 175 | + source.onNext(2); |
| 176 | + sampler.onNext(1); |
| 177 | + sampler.onCompleted(); |
| 178 | + |
| 179 | + source.onNext(3); |
| 180 | + source.onNext(4); |
| 181 | + |
| 182 | + |
| 183 | + |
| 184 | + InOrder inOrder = inOrder(observer2); |
| 185 | + inOrder.verify(observer2, never()).onNext(1); |
| 186 | + inOrder.verify(observer2, times(1)).onNext(2); |
| 187 | + inOrder.verify(observer2, times(1)).onCompleted(); |
| 188 | + inOrder.verify(observer2, never()).onNext(any()); |
| 189 | + verify(observer, never()).onError(any(Throwable.class)); |
| 190 | + } |
| 191 | + @Test |
| 192 | + public void sampleWithSamplerEmitAndTerminate() { |
| 193 | + PublishSubject<Integer> source = PublishSubject.create(); |
| 194 | + PublishSubject<Integer> sampler = PublishSubject.create(); |
| 195 | + |
| 196 | + Observable<Integer> m = source.sample(sampler); |
| 197 | + m.subscribe(observer2); |
| 198 | + |
| 199 | + source.onNext(1); |
| 200 | + source.onNext(2); |
| 201 | + sampler.onNext(1); |
| 202 | + source.onNext(3); |
| 203 | + source.onCompleted(); |
| 204 | + sampler.onNext(2); |
| 205 | + sampler.onCompleted(); |
| 206 | + |
| 207 | + InOrder inOrder = inOrder(observer2); |
| 208 | + inOrder.verify(observer2, never()).onNext(1); |
| 209 | + inOrder.verify(observer2, times(1)).onNext(2); |
| 210 | + inOrder.verify(observer2, never()).onNext(3); |
| 211 | + inOrder.verify(observer2, times(1)).onCompleted(); |
| 212 | + inOrder.verify(observer2, never()).onNext(any()); |
| 213 | + verify(observer, never()).onError(any(Throwable.class)); |
| 214 | + } |
| 215 | + @Test |
| 216 | + public void sampleWithSamplerEmptySource() { |
| 217 | + PublishSubject<Integer> source = PublishSubject.create(); |
| 218 | + PublishSubject<Integer> sampler = PublishSubject.create(); |
| 219 | + |
| 220 | + Observable<Integer> m = source.sample(sampler); |
| 221 | + m.subscribe(observer2); |
| 222 | + |
| 223 | + source.onCompleted(); |
| 224 | + sampler.onNext(1); |
| 225 | + |
| 226 | + InOrder inOrder = inOrder(observer2); |
| 227 | + inOrder.verify(observer2, times(1)).onCompleted(); |
| 228 | + verify(observer2, never()).onNext(any()); |
| 229 | + verify(observer, never()).onError(any(Throwable.class)); |
| 230 | + } |
| 231 | + @Test |
| 232 | + public void sampleWithSamplerSourceThrows() { |
| 233 | + PublishSubject<Integer> source = PublishSubject.create(); |
| 234 | + PublishSubject<Integer> sampler = PublishSubject.create(); |
| 235 | + |
| 236 | + Observable<Integer> m = source.sample(sampler); |
| 237 | + m.subscribe(observer2); |
| 238 | + |
| 239 | + source.onNext(1); |
| 240 | + source.onError(new RuntimeException("Forced failure!")); |
| 241 | + sampler.onNext(1); |
| 242 | + |
| 243 | + InOrder inOrder = inOrder(observer2); |
| 244 | + inOrder.verify(observer2, times(1)).onError(any(Throwable.class)); |
| 245 | + verify(observer2, never()).onNext(any()); |
| 246 | + verify(observer, never()).onCompleted(); |
| 247 | + } |
| 248 | + @Test |
| 249 | + public void sampleWithSamplerThrows() { |
| 250 | + PublishSubject<Integer> source = PublishSubject.create(); |
| 251 | + PublishSubject<Integer> sampler = PublishSubject.create(); |
| 252 | + |
| 253 | + Observable<Integer> m = source.sample(sampler); |
| 254 | + m.subscribe(observer2); |
| 255 | + |
| 256 | + source.onNext(1); |
| 257 | + sampler.onNext(1); |
| 258 | + sampler.onError(new RuntimeException("Forced failure!")); |
| 259 | + |
| 260 | + InOrder inOrder = inOrder(observer2); |
| 261 | + inOrder.verify(observer2, times(1)).onNext(1); |
| 262 | + inOrder.verify(observer2, times(1)).onError(any(RuntimeException.class)); |
| 263 | + verify(observer, never()).onCompleted(); |
| 264 | + } |
108 | 265 | }
|
0 commit comments