/** * Copyright 2015 Netflix, Inc. * * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in * compliance with the License. You may obtain a copy of the License at * * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software distributed under the License is * distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See * the License for the specific language governing permissions and limitations under the License. */ package rx; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; import static org.mockito.Matchers.eq; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyZeroInteractions; import static org.mockito.Mockito.when; import java.util.Arrays; import java.util.Collections; import java.util.LinkedHashMap; import java.util.List; import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import org.junit.Test; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; import rx.Single.OnSubscribe; import rx.exceptions.CompositeException; import rx.functions.Action0; import rx.functions.Action1; import rx.functions.Func1; import rx.functions.Func2; import rx.functions.Func3; import rx.functions.Func4; import rx.functions.Func5; import rx.functions.Func6; import rx.functions.Func7; import rx.functions.Func8; import rx.functions.Func9; import rx.functions.FuncN; import rx.schedulers.TestScheduler; import rx.singles.BlockingSingle; import rx.observers.TestSubscriber; import rx.schedulers.Schedulers; import rx.subscriptions.Subscriptions; public class SingleTest { @Test public void testHelloWorld() { TestSubscriber ts = new TestSubscriber(); Single.just("Hello World!").subscribe(ts); ts.assertReceivedOnNext(Arrays.asList("Hello World!")); } @Test public void testHelloWorld2() { final AtomicReference v = new AtomicReference(); Single.just("Hello World!").subscribe(new SingleSubscriber() { @Override public void onSuccess(String value) { v.set(value); } @Override public void onError(Throwable error) { } }); assertEquals("Hello World!", v.get()); } @Test public void testMap() { TestSubscriber ts = new TestSubscriber(); Single.just("A") .map(new Func1() { @Override public String call(String s) { return s + "B"; } }) .subscribe(ts); ts.assertReceivedOnNext(Arrays.asList("AB")); } @Test public void zip2Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single.zip(a, b, new Func2() { @Override public String call(Integer a, Integer b) { return "" + a + b; } }) .subscribe(ts); ts.assertValue("12"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zip3Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single c = Single.just(3); Single.zip(a, b, c, new Func3() { @Override public String call(Integer a, Integer b, Integer c) { return "" + a + b + c; } }) .subscribe(ts); ts.assertValue("123"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zip4Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single c = Single.just(3); Single d = Single.just(4); Single.zip(a, b, c, d, new Func4() { @Override public String call(Integer a, Integer b, Integer c, Integer d) { return "" + a + b + c + d; } }) .subscribe(ts); ts.assertValue("1234"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zip5Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single c = Single.just(3); Single d = Single.just(4); Single e = Single.just(5); Single.zip(a, b, c, d, e, new Func5() { @Override public String call(Integer a, Integer b, Integer c, Integer d, Integer e) { return "" + a + b + c + d + e; } }) .subscribe(ts); ts.assertValue("12345"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zip6Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single c = Single.just(3); Single d = Single.just(4); Single e = Single.just(5); Single f = Single.just(6); Single.zip(a, b, c, d, e, f, new Func6() { @Override public String call(Integer a, Integer b, Integer c, Integer d, Integer e, Integer f) { return "" + a + b + c + d + e + f; } }) .subscribe(ts); ts.assertValue("123456"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zip7Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single c = Single.just(3); Single d = Single.just(4); Single e = Single.just(5); Single f = Single.just(6); Single g = Single.just(7); Single.zip(a, b, c, d, e, f, g, new Func7() { @Override public String call(Integer a, Integer b, Integer c, Integer d, Integer e, Integer f, Integer g) { return "" + a + b + c + d + e + f + g; } }) .subscribe(ts); ts.assertValue("1234567"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zip8Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single c = Single.just(3); Single d = Single.just(4); Single e = Single.just(5); Single f = Single.just(6); Single g = Single.just(7); Single h = Single.just(8); Single.zip(a, b, c, d, e, f, g, h, new Func8() { @Override public String call(Integer a, Integer b, Integer c, Integer d, Integer e, Integer f, Integer g, Integer h) { return "" + a + b + c + d + e + f + g + h; } }) .subscribe(ts); ts.assertValue("12345678"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zip9Singles() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just(1); Single b = Single.just(2); Single c = Single.just(3); Single d = Single.just(4); Single e = Single.just(5); Single f = Single.just(6); Single g = Single.just(7); Single h = Single.just(8); Single i = Single.just(9); Single.zip(a, b, c, d, e, f, g, h, i, new Func9() { @Override public String call(Integer a, Integer b, Integer c, Integer d, Integer e, Integer f, Integer g, Integer h, Integer i) { return "" + a + b + c + d + e + f + g + h + i; } }) .subscribe(ts); ts.assertValue("123456789"); ts.assertCompleted(); ts.assertNoErrors(); } @Test public void zipIterableShouldZipListOfSingles() { TestSubscriber ts = new TestSubscriber(); Iterable> singles = Arrays.asList(Single.just(1), Single.just(2), Single.just(3)); Single .zip(singles, new FuncN() { @Override public String call(Object... args) { StringBuilder stringBuilder = new StringBuilder(); for (Object arg : args) { stringBuilder.append(arg); } return stringBuilder.toString(); } }).subscribe(ts); ts.assertValue("123"); ts.assertNoErrors(); ts.assertCompleted(); } @Test public void zipIterableShouldZipSetOfSingles() { TestSubscriber ts = new TestSubscriber(); Set> singlesSet = Collections.newSetFromMap(new LinkedHashMap, Boolean>(2)); Single s1 = Single.just("1"); Single s2 = Single.just("2"); Single s3 = Single.just("3"); singlesSet.add(s1); singlesSet.add(s2); singlesSet.add(s3); Single .zip(singlesSet, new FuncN() { @Override public String call(Object... args) { StringBuilder stringBuilder = new StringBuilder(); for (Object arg : args) { stringBuilder.append(arg); } return stringBuilder.toString(); } }).subscribe(ts); ts.assertValue("123"); ts.assertNoErrors(); ts.assertCompleted(); } @Test public void testZipWith() { TestSubscriber ts = new TestSubscriber(); Single.just("A").zipWith(Single.just("B"), new Func2() { @Override public String call(String a, String b) { return a + b; } }) .subscribe(ts); ts.assertReceivedOnNext(Arrays.asList("AB")); } @Test public void testMerge() { TestSubscriber ts = new TestSubscriber(); Single a = Single.just("A"); Single b = Single.just("B"); Single.merge(a, b).subscribe(ts); ts.assertReceivedOnNext(Arrays.asList("A", "B")); } @Test public void testMergeWith() { TestSubscriber ts = new TestSubscriber(); Single.just("A").mergeWith(Single.just("B")).subscribe(ts); ts.assertReceivedOnNext(Arrays.asList("A", "B")); } @Test public void testCreateSuccess() { TestSubscriber ts = new TestSubscriber(); Single.create(new OnSubscribe() { @Override public void call(SingleSubscriber s) { s.onSuccess("Hello"); } }).subscribe(ts); ts.assertReceivedOnNext(Arrays.asList("Hello")); } @Test public void testCreateError() { TestSubscriber ts = new TestSubscriber(); Single.create(new OnSubscribe() { @Override public void call(SingleSubscriber s) { s.onError(new RuntimeException("fail")); } }).subscribe(ts); assertEquals(1, ts.getOnErrorEvents().size()); } @Test public void testAsync() { TestSubscriber ts = new TestSubscriber(); Single.just("Hello") .subscribeOn(Schedulers.io()) .map(new Func1() { @Override public String call(String v) { System.out.println("SubscribeOn Thread: " + Thread.currentThread()); return v; } }) .observeOn(Schedulers.computation()) .map(new Func1() { @Override public String call(String v) { System.out.println("ObserveOn Thread: " + Thread.currentThread()); return v; } }) .subscribe(ts); ts.awaitTerminalEvent(); ts.assertReceivedOnNext(Arrays.asList("Hello")); } @Test public void testFlatMap() { TestSubscriber ts = new TestSubscriber(); Single.just("Hello").flatMap(new Func1>() { @Override public Single call(String s) { return Single.just(s + " World!").subscribeOn(Schedulers.computation()); } }).subscribe(ts); ts.awaitTerminalEvent(); ts.assertReceivedOnNext(Arrays.asList("Hello World!")); } @Test public void testTimeout() { TestSubscriber ts = new TestSubscriber(); Single s = Single.create(new OnSubscribe() { @Override public void call(SingleSubscriber s) { try { Thread.sleep(5000); } catch (InterruptedException e) { // ignore as we expect this for the test } s.onSuccess("success"); } }).subscribeOn(Schedulers.io()); s.timeout(100, TimeUnit.MILLISECONDS).subscribe(ts); ts.awaitTerminalEvent(); ts.assertError(TimeoutException.class); } @Test public void testTimeoutWithFallback() { TestSubscriber ts = new TestSubscriber(); Single s = Single.create(new OnSubscribe() { @Override public void call(SingleSubscriber s) { try { Thread.sleep(5000); } catch (InterruptedException e) { // ignore as we expect this for the test } s.onSuccess("success"); } }).subscribeOn(Schedulers.io()); s.timeout(100, TimeUnit.MILLISECONDS, Single.just("hello")).subscribe(ts); ts.awaitTerminalEvent(); ts.assertNoErrors(); ts.assertValue("hello"); } @Test public void testToBlocking() { Single s = Single.just("one"); BlockingSingle blocking = s.toBlocking(); assertNotNull(blocking); assertEquals("one", blocking.value()); } @Test public void testUnsubscribe() throws InterruptedException { TestSubscriber ts = new TestSubscriber(); final AtomicBoolean unsubscribed = new AtomicBoolean(); final AtomicBoolean interrupted = new AtomicBoolean(); final CountDownLatch latch = new CountDownLatch(2); Single s = Single.create(new OnSubscribe() { @Override public void call(final SingleSubscriber s) { final Thread t = new Thread(new Runnable() { @Override public void run() { try { Thread.sleep(5000); s.onSuccess("success"); } catch (InterruptedException e) { interrupted.set(true); latch.countDown(); } } }); s.add(Subscriptions.create(new Action0() { @Override public void call() { unsubscribed.set(true); t.interrupt(); latch.countDown(); } })); t.start(); } }); s.subscribe(ts); Thread.sleep(100); ts.unsubscribe(); if (latch.await(1000, TimeUnit.MILLISECONDS)) { assertTrue(unsubscribed.get()); assertTrue(interrupted.get()); } else { fail("timed out waiting for latch"); } } /** * Assert that unsubscribe propagates when passing in a SingleSubscriber and not a Subscriber */ @Test public void testUnsubscribe2() throws InterruptedException { SingleSubscriber ts = new SingleSubscriber() { @Override public void onSuccess(String value) { // not interested in value } @Override public void onError(Throwable error) { // not interested in value } }; final AtomicBoolean unsubscribed = new AtomicBoolean(); final AtomicBoolean interrupted = new AtomicBoolean(); final CountDownLatch latch = new CountDownLatch(2); Single s = Single.create(new OnSubscribe() { @Override public void call(final SingleSubscriber s) { final Thread t = new Thread(new Runnable() { @Override public void run() { try { Thread.sleep(5000); s.onSuccess("success"); } catch (InterruptedException e) { interrupted.set(true); latch.countDown(); } } }); s.add(Subscriptions.create(new Action0() { @Override public void call() { unsubscribed.set(true); t.interrupt(); latch.countDown(); } })); t.start(); } }); s.subscribe(ts); Thread.sleep(100); ts.unsubscribe(); if (latch.await(1000, TimeUnit.MILLISECONDS)) { assertTrue(unsubscribed.get()); assertTrue(interrupted.get()); } else { fail("timed out waiting for latch"); } } /** * Assert that unsubscribe propagates when passing in a SingleSubscriber and not a Subscriber */ @Test public void testUnsubscribeViaReturnedSubscription() throws InterruptedException { final AtomicBoolean unsubscribed = new AtomicBoolean(); final AtomicBoolean interrupted = new AtomicBoolean(); final CountDownLatch latch = new CountDownLatch(2); Single s = Single.create(new OnSubscribe() { @Override public void call(final SingleSubscriber s) { final Thread t = new Thread(new Runnable() { @Override public void run() { try { Thread.sleep(5000); s.onSuccess("success"); } catch (InterruptedException e) { interrupted.set(true); latch.countDown(); } } }); s.add(Subscriptions.create(new Action0() { @Override public void call() { unsubscribed.set(true); t.interrupt(); latch.countDown(); } })); t.start(); } }); Subscription subscription = s.subscribe(); Thread.sleep(100); subscription.unsubscribe(); if (latch.await(1000, TimeUnit.MILLISECONDS)) { assertTrue(unsubscribed.get()); assertTrue(interrupted.get()); } else { fail("timed out waiting for latch"); } } @Test public void testBackpressureAsObservable() { Single s = Single.create(new OnSubscribe() { @Override public void call(SingleSubscriber t) { t.onSuccess("hello"); } }); TestSubscriber ts = new TestSubscriber() { @Override public void onStart() { request(0); } }; s.subscribe(ts); ts.assertNoValues(); ts.requestMore(1); ts.assertValue("hello"); } @Test public void testToObservable() { Observable a = Single.just("a").toObservable(); TestSubscriber ts = TestSubscriber.create(); a.subscribe(ts); ts.assertValue("a"); ts.assertCompleted(); } @Test public void doOnErrorShouldNotCallActionIfNoErrorHasOccurred() { Action1 action = mock(Action1.class); TestSubscriber testSubscriber = new TestSubscriber(); Single .just("value") .doOnError(action) .subscribe(testSubscriber); testSubscriber.assertValue("value"); testSubscriber.assertNoErrors(); verifyZeroInteractions(action); } @Test public void doOnErrorShouldCallActionIfErrorHasOccurred() { Action1 action = mock(Action1.class); TestSubscriber testSubscriber = new TestSubscriber(); Throwable error = new IllegalStateException(); Single .error(error) .doOnError(action) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(error); verify(action).call(error); } @Test public void doOnErrorShouldThrowCompositeExceptionIfOnErrorActionThrows() { Action1 action = mock(Action1.class); Throwable error = new RuntimeException(); Throwable exceptionFromOnErrorAction = new IllegalStateException(); doThrow(exceptionFromOnErrorAction).when(action).call(error); TestSubscriber testSubscriber = new TestSubscriber(); Single .error(error) .doOnError(action) .subscribe(testSubscriber); testSubscriber.assertNoValues(); CompositeException compositeException = (CompositeException) testSubscriber.getOnErrorEvents().get(0); assertEquals(2, compositeException.getExceptions().size()); assertSame(error, compositeException.getExceptions().get(0)); assertSame(exceptionFromOnErrorAction, compositeException.getExceptions().get(1)); verify(action).call(error); } @Test public void shouldEmitValueFromCallable() throws Exception { Callable callable = mock(Callable.class); when(callable.call()).thenReturn("value"); TestSubscriber testSubscriber = new TestSubscriber(); Single .fromCallable(callable) .subscribe(testSubscriber); testSubscriber.assertValue("value"); testSubscriber.assertNoErrors(); verify(callable).call(); } @Test public void shouldPassErrorFromCallable() throws Exception { Callable callable = mock(Callable.class); Throwable error = new IllegalStateException(); when(callable.call()).thenThrow(error); TestSubscriber testSubscriber = new TestSubscriber(); Single .fromCallable(callable) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(error); verify(callable).call(); } @Test public void doOnSuccessShouldInvokeAction() { Action1 action = mock(Action1.class); TestSubscriber testSubscriber = new TestSubscriber(); Single .just("value") .doOnSuccess(action) .subscribe(testSubscriber); testSubscriber.assertValue("value"); testSubscriber.assertNoErrors(); verify(action).call(eq("value")); } @Test public void doOnSuccessShouldPassErrorFromActionToSubscriber() { Action1 action = mock(Action1.class); Throwable error = new IllegalStateException(); doThrow(error).when(action).call(eq("value")); TestSubscriber testSubscriber = new TestSubscriber(); Single .just("value") .doOnSuccess(action) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(error); verify(action).call(eq("value")); } @Test public void doOnSuccessShouldNotCallActionIfSingleThrowsError() { Action1 action = mock(Action1.class); Throwable error = new IllegalStateException(); TestSubscriber testSubscriber = new TestSubscriber(); Single .error(error) .doOnSuccess(action) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(error); verifyZeroInteractions(action); } @Test public void doOnSuccessShouldNotSwallowExceptionThrownByAction() { Action1 action = mock(Action1.class); Throwable exceptionFromAction = new IllegalStateException(); doThrow(exceptionFromAction).when(action).call(eq("value")); TestSubscriber testSubscriber = new TestSubscriber(); Single .just("value") .doOnSuccess(action) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(exceptionFromAction); verify(action).call(eq("value")); } @Test public void delayWithSchedulerShouldDelayCompletion() { TestScheduler scheduler = new TestScheduler(); Single single = Single.just(1).delay(100, TimeUnit.DAYS, scheduler); TestSubscriber subscriber = new TestSubscriber(); single.subscribe(subscriber); subscriber.assertNotCompleted(); scheduler.advanceTimeBy(99, TimeUnit.DAYS); subscriber.assertNotCompleted(); scheduler.advanceTimeBy(91, TimeUnit.DAYS); subscriber.assertCompleted(); subscriber.assertValue(1); } @Test public void delayWithSchedulerShouldShortCutWithFailure() { TestScheduler scheduler = new TestScheduler(); final RuntimeException expected = new RuntimeException(); Single single = Single.create(new OnSubscribe() { @Override public void call(SingleSubscriber singleSubscriber) { singleSubscriber.onSuccess(1); singleSubscriber.onError(expected); } }).delay(100, TimeUnit.DAYS, scheduler); TestSubscriber subscriber = new TestSubscriber(); single.subscribe(subscriber); subscriber.assertNotCompleted(); scheduler.advanceTimeBy(99, TimeUnit.DAYS); subscriber.assertNotCompleted(); scheduler.advanceTimeBy(91, TimeUnit.DAYS); subscriber.assertNoValues(); subscriber.assertError(expected); } @Test public void deferShouldNotCallFactoryFuncUntilSubscriberSubscribes() throws Exception { Callable> singleFactory = mock(Callable.class); Single.defer(singleFactory); verifyZeroInteractions(singleFactory); } @Test public void deferShouldSubscribeSubscriberToSingleFromFactoryFuncAndEmitValue() throws Exception { Callable> singleFactory = mock(Callable.class); Object value = new Object(); Single single = Single.just(value); when(singleFactory.call()).thenReturn(single); TestSubscriber testSubscriber = new TestSubscriber(); Single .defer(singleFactory) .subscribe(testSubscriber); testSubscriber.assertValue(value); testSubscriber.assertNoErrors(); verify(singleFactory).call(); } @Test public void deferShouldSubscribeSubscriberToSingleFromFactoryFuncAndEmitError() throws Exception { Callable> singleFactory = mock(Callable.class); Throwable error = new IllegalStateException(); Single single = Single.error(error); when(singleFactory.call()).thenReturn(single); TestSubscriber testSubscriber = new TestSubscriber(); Single .defer(singleFactory) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(error); verify(singleFactory).call(); } @Test public void deferShouldPassErrorFromSingleFactoryToTheSubscriber() throws Exception { Callable> singleFactory = mock(Callable.class); Throwable errorFromSingleFactory = new IllegalStateException(); when(singleFactory.call()).thenThrow(errorFromSingleFactory); TestSubscriber testSubscriber = new TestSubscriber(); Single .defer(singleFactory) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(errorFromSingleFactory); verify(singleFactory).call(); } @Test public void deferShouldCallSingleFactoryForEachSubscriber() throws Exception { Callable> singleFactory = mock(Callable.class); String[] values = {"1", "2", "3"}; final Single[] singles = new Single[]{Single.just(values[0]), Single.just(values[1]), Single.just(values[2])}; final AtomicInteger singleFactoryCallsCounter = new AtomicInteger(); when(singleFactory.call()).thenAnswer(new Answer>() { @Override public Single answer(InvocationOnMock invocation) throws Throwable { return singles[singleFactoryCallsCounter.getAndIncrement()]; } }); Single deferredSingle = Single.defer(singleFactory); for (int i = 0; i < singles.length; i ++) { TestSubscriber testSubscriber = new TestSubscriber(); deferredSingle.subscribe(testSubscriber); testSubscriber.assertValue(values[i]); testSubscriber.assertNoErrors(); } verify(singleFactory, times(3)).call(); } @Test public void deferShouldPassNullPointerExceptionToTheSubscriberIfSingleFactoryIsNull() { TestSubscriber testSubscriber = new TestSubscriber(); Single .defer(null) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(NullPointerException.class); } @Test public void deferShouldPassNullPointerExceptionToTheSubscriberIfSingleFactoryReturnsNull() throws Exception { Callable> singleFactory = mock(Callable.class); when(singleFactory.call()).thenReturn(null); TestSubscriber testSubscriber = new TestSubscriber(); Single .defer(singleFactory) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(NullPointerException.class); verify(singleFactory).call(); } @Test public void doOnUnsubscribeShouldInvokeActionAfterSuccess() { Action0 action = mock(Action0.class); Single single = Single .just("test") .doOnUnsubscribe(action); verifyZeroInteractions(action); TestSubscriber testSubscriber = new TestSubscriber(); single.subscribe(testSubscriber); testSubscriber.assertValue("test"); testSubscriber.assertCompleted(); verify(action).call(); } @Test public void doOnUnsubscribeShouldInvokeActionAfterError() { Action0 action = mock(Action0.class); Single single = Single .error(new RuntimeException("test")) .doOnUnsubscribe(action); verifyZeroInteractions(action); TestSubscriber testSubscriber = new TestSubscriber(); single.subscribe(testSubscriber); testSubscriber.assertError(RuntimeException.class); assertEquals("test", testSubscriber.getOnErrorEvents().get(0).getMessage()); verify(action).call(); } @Test public void doOnUnsubscribeShouldInvokeActionAfterExplicitUnsubscription() { Action0 action = mock(Action0.class); Single single = Single .create(new OnSubscribe() { @Override public void call(SingleSubscriber singleSubscriber) { // Broken Single that never ends itself (simulates long computation in one thread). } }) .doOnUnsubscribe(action); TestSubscriber testSubscriber = new TestSubscriber(); Subscription subscription = single.subscribe(testSubscriber); verifyZeroInteractions(action); subscription.unsubscribe(); verify(action).call(); testSubscriber.assertNoValues(); testSubscriber.assertNoTerminalEvent(); } @Test public void doAfterTerminateActionShouldBeInvokedAfterOnSuccess() { Action0 action = mock(Action0.class); TestSubscriber testSubscriber = new TestSubscriber(); Single .just("value") .doAfterTerminate(action) .subscribe(testSubscriber); testSubscriber.assertValue("value"); testSubscriber.assertNoErrors(); verify(action).call(); } @Test public void doAfterTerminateActionShouldBeInvokedAfterOnError() { Action0 action = mock(Action0.class); TestSubscriber testSubscriber = new TestSubscriber(); Throwable error = new IllegalStateException(); Single .error(error) .doAfterTerminate(action) .subscribe(testSubscriber); testSubscriber.assertNoValues(); testSubscriber.assertError(error); verify(action).call(); } @Test public void doAfterTerminateActionShouldNotBeInvokedUntilSubscriberSubscribes() { Action0 action = mock(Action0.class); Single .just("value") .doAfterTerminate(action); Single .error(new IllegalStateException()) .doAfterTerminate(action); verifyZeroInteractions(action); } @Test(expected = NullPointerException.class) public void iterableToArrayShouldThrowNullPointerExceptionIfIterableNull() { Single.iterableToArray(null); } @Test public void iterableToArrayShouldConvertList() { List> singlesList = Arrays.asList(Single.just("1"), Single.just("2")); Single[] singlesArray = Single.iterableToArray(singlesList); assertEquals(2, singlesArray.length); assertSame(singlesList.get(0), singlesArray[0]); assertSame(singlesList.get(1), singlesArray[1]); } @Test public void iterableToArrayShouldConvertSet() { // Just to trigger different path of the code that handles non-list iterables. Set> singlesSet = Collections.newSetFromMap(new LinkedHashMap, Boolean>(2)); Single s1 = Single.just("1"); Single s2 = Single.just("2"); singlesSet.add(s1); singlesSet.add(s2); Single[] singlesArray = Single.iterableToArray(singlesSet); assertEquals(2, singlesArray.length); assertSame(s1, singlesArray[0]); assertSame(s2, singlesArray[1]); } }