diff --git a/src/main/java/rx/internal/operators/OperatorSampleWithTime.java b/src/main/java/rx/internal/operators/OperatorSampleWithTime.java index ea94a7db21..7138d760d4 100644 --- a/src/main/java/rx/internal/operators/OperatorSampleWithTime.java +++ b/src/main/java/rx/internal/operators/OperatorSampleWithTime.java @@ -50,6 +50,7 @@ public Subscriber call(Subscriber child) { child.add(worker); SamplerSubscriber sampler = new SamplerSubscriber(s); + child.add(sampler); worker.schedulePeriodically(sampler, time, time, unit); return sampler; diff --git a/src/test/java/rx/internal/operators/OperatorSampleTest.java b/src/test/java/rx/internal/operators/OperatorSampleTest.java index 0a8c9da58d..815d002061 100644 --- a/src/test/java/rx/internal/operators/OperatorSampleTest.java +++ b/src/test/java/rx/internal/operators/OperatorSampleTest.java @@ -28,14 +28,12 @@ import org.junit.Test; import org.mockito.InOrder; -import rx.Observable; +import rx.*; import rx.Observable.OnSubscribe; -import rx.Observer; -import rx.Scheduler; -import rx.Subscriber; import rx.functions.Action0; import rx.schedulers.TestScheduler; import rx.subjects.PublishSubject; +import rx.subscriptions.Subscriptions; public class OperatorSampleTest { private TestScheduler scheduler; @@ -271,4 +269,19 @@ public void sampleWithSamplerThrows() { inOrder.verify(observer2, times(1)).onError(any(RuntimeException.class)); verify(observer, never()).onCompleted(); } + + @Test + public void testSampleUnsubscribe() { + final Subscription s = mock(Subscription.class); + Observable o = Observable.create( + new OnSubscribe() { + @Override + public void call(Subscriber subscriber) { + subscriber.add(s); + } + } + ); + o.throttleLast(1, TimeUnit.MILLISECONDS).subscribe().unsubscribe(); + verify(s).unsubscribe(); + } }