From de66728dd70c78320fb870fdba71234cdf6842a9 Mon Sep 17 00:00:00 2001 From: jmhofer Date: Wed, 28 Aug 2013 20:37:47 +0200 Subject: [PATCH 1/3] re-added combineLatest methods that got lost due to too optimistic super/extends generics --- rxjava-core/src/main/java/rx/Observable.java | 38 +++++++++++++++++--- 1 file changed, 33 insertions(+), 5 deletions(-) diff --git a/rxjava-core/src/main/java/rx/Observable.java b/rxjava-core/src/main/java/rx/Observable.java index f00bf47f7a..2b9ba46fcb 100644 --- a/rxjava-core/src/main/java/rx/Observable.java +++ b/rxjava-core/src/main/java/rx/Observable.java @@ -15,10 +15,6 @@ */ package rx; -import static org.junit.Assert.*; -import static org.mockito.Matchers.*; -import static org.mockito.Mockito.*; - import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; @@ -27,7 +23,6 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; - import rx.concurrency.Schedulers; import rx.observables.BlockingObservable; import rx.observables.ConnectableObservable; @@ -35,6 +30,7 @@ import rx.operators.OperationAll; import rx.operators.OperationBuffer; import rx.operators.OperationCache; +import rx.operators.OperationCombineLatest; import rx.operators.OperationConcat; import rx.operators.OperationDefer; import rx.operators.OperationDematerialize; @@ -1085,6 +1081,38 @@ public static Observable zip(Observable w0, Observabl return create(OperationZip.zip(w0, w1, w2, w3, reduceFunction)); } + /** + * Combines the given observables, emitting an event containing an aggregation of the latest values of each of the source observables + * each time an event is received from one of the source observables, where the aggregation is defined by the given function. + *

+ * + * + * @param w0 + * The first source observable. + * @param w1 + * The second source observable. + * @param combineFunction + * The aggregation function used to combine the source observable values. + * @return An Observable that combines the source Observables with the given combine function + */ + public static Observable combineLatest(Observable w0, Observable w1, Func2 combineFunction) { + return create(OperationCombineLatest.combineLatest(w0, w1, combineFunction)); + } + + /** + * @see #combineLatest(Observable, Observable, Func2) + */ + public static Observable combineLatest(Observable w0, Observable w1, Observable w2, Func3 combineFunction) { + return create(OperationCombineLatest.combineLatest(w0, w1, w2, combineFunction)); + } + + /** + * @see #combineLatest(Observable, Observable, Func2) + */ + public static Observable combineLatest(Observable w0, Observable w1, Observable w2, Observable w3, Func4 combineFunction) { + return create(OperationCombineLatest.combineLatest(w0, w1, w2, w3, combineFunction)); + } + /** * Creates an Observable which produces buffers of collected values. * From bd10365756b5d74f9faa9f50512735d3f4d7905f Mon Sep 17 00:00:00 2001 From: jmhofer Date: Wed, 28 Aug 2013 21:20:06 +0200 Subject: [PATCH 2/3] Reactivated the rxjava-core tests --- build.gradle | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/build.gradle b/build.gradle index c7378d469a..440cac0acc 100644 --- a/build.gradle +++ b/build.gradle @@ -63,3 +63,7 @@ subprojects { } } +project(':rxjava-core') { + sourceSets.test.java.srcDir 'src/test/java' +} + From e3e148a195927582102b0fc076d49a27b16c3de2 Mon Sep 17 00:00:00 2001 From: jmhofer Date: Wed, 28 Aug 2013 21:46:58 +0200 Subject: [PATCH 3/3] repaired rxjava-swing to work with new scheduler and observable api --- .../src/main/java/rx/concurrency/SwingScheduler.java | 8 ++------ .../src/main/java/rx/observables/SwingObservable.java | 3 +-- .../src/main/java/rx/swing/sources/KeyEventSource.java | 2 +- 3 files changed, 4 insertions(+), 9 deletions(-) diff --git a/rxjava-contrib/rxjava-swing/src/main/java/rx/concurrency/SwingScheduler.java b/rxjava-contrib/rxjava-swing/src/main/java/rx/concurrency/SwingScheduler.java index ced759f6f5..2ce298c590 100644 --- a/rxjava-contrib/rxjava-swing/src/main/java/rx/concurrency/SwingScheduler.java +++ b/rxjava-contrib/rxjava-swing/src/main/java/rx/concurrency/SwingScheduler.java @@ -38,7 +38,6 @@ import rx.subscriptions.CompositeSubscription; import rx.subscriptions.Subscriptions; import rx.util.functions.Action0; -import rx.util.functions.Func0; import rx.util.functions.Func2; /** @@ -187,14 +186,12 @@ public void testPeriodicScheduling() throws Exception { final CountDownLatch latch = new CountDownLatch(4); final Action0 innerAction = mock(Action0.class); - final Action0 unsubscribe = mock(Action0.class); - final Func0 action = new Func0() { + final Action0 action = new Action0() { @Override - public Subscription call() { + public void call() { try { innerAction.call(); assertTrue(SwingUtilities.isEventDispatchThread()); - return Subscriptions.create(unsubscribe); } finally { latch.countDown(); } @@ -210,7 +207,6 @@ public Subscription call() { sub.unsubscribe(); waitForEmptyEventQueue(); verify(innerAction, times(4)).call(); - verify(unsubscribe, times(4)).call(); } @Test diff --git a/rxjava-contrib/rxjava-swing/src/main/java/rx/observables/SwingObservable.java b/rxjava-contrib/rxjava-swing/src/main/java/rx/observables/SwingObservable.java index 174c529d2d..b280638191 100644 --- a/rxjava-contrib/rxjava-swing/src/main/java/rx/observables/SwingObservable.java +++ b/rxjava-contrib/rxjava-swing/src/main/java/rx/observables/SwingObservable.java @@ -26,7 +26,6 @@ import javax.swing.AbstractButton; import rx.Observable; -import static rx.Observable.filter; import rx.swing.sources.AbstractButtonSource; import rx.swing.sources.ComponentEventSource; import rx.swing.sources.KeyEventSource; @@ -68,7 +67,7 @@ public static Observable fromKeyEvents(Component component) { * @return Observable of key events. */ public static Observable fromKeyEvents(Component component, final Set keyCodes) { - return filter(fromKeyEvents(component), new Func1() { + return fromKeyEvents(component).filter(new Func1() { @Override public Boolean call(KeyEvent event) { return keyCodes.contains(event.getKeyCode()); diff --git a/rxjava-contrib/rxjava-swing/src/main/java/rx/swing/sources/KeyEventSource.java b/rxjava-contrib/rxjava-swing/src/main/java/rx/swing/sources/KeyEventSource.java index 3716b599f9..291e0202aa 100644 --- a/rxjava-contrib/rxjava-swing/src/main/java/rx/swing/sources/KeyEventSource.java +++ b/rxjava-contrib/rxjava-swing/src/main/java/rx/swing/sources/KeyEventSource.java @@ -85,7 +85,7 @@ public void call() { * @see SwingObservable.fromKeyEvents(Component, Set) */ public static Observable> currentlyPressedKeysOf(Component component) { - return Observable.>scan(fromKeyEventsOf(component), new HashSet(), new Func2, KeyEvent, Set>() { + return fromKeyEventsOf(component).>scan(new HashSet(), new Func2, KeyEvent, Set>() { @Override public Set call(Set pressedKeys, KeyEvent event) { Set afterEvent = new HashSet(pressedKeys);