Skip to content

Commit

Permalink
Add sizeOf(ObservableCollection) factory methods.
Browse files Browse the repository at this point in the history
  • Loading branch information
TomasMikula committed Jun 17, 2014
1 parent 2e7115d commit 1515b01
Show file tree
Hide file tree
Showing 2 changed files with 80 additions and 0 deletions.
44 changes: 44 additions & 0 deletions reactfx/src/main/java/org/reactfx/EventStreams.java
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
package org.reactfx;

import java.time.Duration;
import java.util.Collection;
import java.util.concurrent.Executor;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.Consumer;
import java.util.function.Function;
import java.util.function.Supplier;

import javafx.beans.InvalidationListener;
import javafx.beans.Observable;
Expand Down Expand Up @@ -142,6 +144,48 @@ protected Subscription subscribeToInputs() {
};
}

public static <C extends Collection<?> & Observable> EventStream<Integer> sizeOf(C collection) {
return create(() -> collection.size(), collection);
}

public static EventStream<Integer> sizeOf(ObservableMap<?, ?> map) {
return create(() -> map.size(), map);
}

private static <T> EventStream<T> create(Supplier<? extends T> computeValue, Observable... dependencies) {
return new LazilyBoundStream<T>() {
private T previousValue;

@Override
protected Subscription subscribeToInputs() {
InvalidationListener listener = obs -> {
T value = computeValue.get();
if(value != previousValue) {
previousValue = value;
emit(value);
}
};
for(Observable dep: dependencies) {
dep.addListener(listener);
}

return () -> {
for(Observable dep: dependencies) {
dep.removeListener(listener);
}
};
}

@Override
protected void newSubscriber(Consumer<? super T> subscriber) {
if(!isBound()) { // this is the first subscriber
previousValue = computeValue.get();
}
subscriber.accept(previousValue);
}
};
}

public static <T extends Event> EventStream<T> eventsOf(Node node, EventType<T> eventType) {
return new LazilyBoundStream<T>() {
@Override
Expand Down
36 changes: 36 additions & 0 deletions reactfx/src/test/java/org/reactfx/SizeOfTest.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
package org.reactfx;

import static org.junit.Assert.*;

import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;

import javafx.collections.FXCollections;
import javafx.collections.ObservableList;

import org.junit.Test;

public class SizeOfTest {

@Test
public void test() {
ObservableList<Integer> list = FXCollections.observableArrayList();
EventStream<Integer> size = EventStreams.sizeOf(list);
List<Integer> sizes = new ArrayList<>();
Subscription sub = size.subscribe(sizes::add);
list.add(1);
list.addAll(2, 3, 4);
assertEquals(Arrays.asList(0, 1, 4), sizes);

sub.unsubscribe();
sizes.clear();
list.addAll(5, 6);
assertEquals(Arrays.asList(), sizes);

size.subscribe(sizes::add);
list.addAll(7, 8);
assertEquals(Arrays.asList(6, 8), sizes);
}

}

0 comments on commit 1515b01

Please sign in to comment.