minor fix and update to work on latest 2.1.0.BUILD-SNAPSHOT to use core check in CFUtils

This commit is contained in:
Stephane Maldini
2015-10-05 16:48:54 +01:00
parent 32214e0a49
commit ec1189b0b5

View File

@@ -24,6 +24,8 @@ import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;
import reactor.core.error.Exceptions;
import reactor.core.error.SpecificationExceptions;
import reactor.core.support.BackpressureUtils;
import reactor.rx.Stream;
import reactor.rx.action.Action;
import reactor.rx.subscription.ReactiveSubscription;
@@ -111,15 +113,19 @@ public class CompletableFutureUtils {
@Override
public void request(long elements) {
Action.checkRequest(elements);
try{
BackpressureUtils.checkRequest(elements);
}catch(SpecificationExceptions.Spec309_NullOrNegativeRequest iae){
subscriber.onError(iae);
return;
}
if (isComplete()) return;
try {
future.whenComplete((result, error) -> {
if (error != null) {
onError(error);
}
else {
} else {
subscriber.onNext(result);
onComplete();
}