/**
* Copyright 2013 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.mockito.Matchers.*;
import static org.mockito.Mockito.*;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import org.junit.Before;
import org.junit.Test;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.MockitoAnnotations;
import rx.operators.OperationConcat;
import rx.operators.OperationFilter;
import rx.operators.OperationLast;
import rx.operators.OperationMap;
import rx.operators.OperationMaterialize;
import rx.operators.OperationMerge;
import rx.operators.OperationMergeDelayError;
import rx.operators.OperationOnErrorResumeNextViaFunction;
import rx.operators.OperationOnErrorResumeNextViaObservable;
import rx.operators.OperationOnErrorReturn;
import rx.operators.OperationScan;
import rx.operators.OperationSkip;
import rx.operators.OperationSynchronize;
import rx.operators.OperationTake;
import rx.operators.OperationToObservableFuture;
import rx.operators.OperationToObservableIterable;
import rx.operators.OperationToObservableList;
import rx.operators.OperationToObservableSortedList;
import rx.operators.OperationZip;
import rx.util.AtomicObservableSubscription;
import rx.util.AtomicObserver;
import rx.util.functions.Action0;
import rx.util.functions.Action1;
import rx.util.functions.Func1;
import rx.util.functions.Func2;
import rx.util.functions.Func3;
import rx.util.functions.Func4;
import rx.util.functions.FuncN;
import rx.util.functions.FunctionLanguageAdaptor;
import rx.util.functions.Functions;
/**
* The Observable interface that implements the Reactive Pattern.
*
* It provides overloaded methods for subscribing as well as delegate methods to the various operators.
*
* The documentation for this interface makes use of marble diagrams. The following legend explains
* these diagrams:
*
*
*
* For more information see the RxJava Wiki
*
* @param
*/
public class Observable {
private final Func1, Subscription> onSubscribe;
private final boolean isTrusted;
protected Observable(Func1, Subscription> onSubscribe) {
this(onSubscribe, false);
}
protected Observable() {
this(null, false);
}
private Observable(Func1, Subscription> onSubscribe, boolean isTrusted) {
this.onSubscribe = onSubscribe;
this.isTrusted = isTrusted;
}
/**
* an {@link Observer} must call an Observable's subscribe method in order to register itself
* to receive push-based notifications from the Observable. A typical implementation of the
* subscribe method does the following:
*
* It stores a reference to the Observer in a collection object, such as a List
* object.
*
* It returns a reference to the {@link Subscription} interface. This enables
* Observers to unsubscribe (that is, to stop receiving notifications) before the Observable has
* finished sending them and has called the Observer's {@link Observer#onCompleted()} method.
*
* At any given time, a particular instance of an Observable implementation is
* responsible for accepting all subscriptions and notifying all subscribers. Unless the
* documentation for a particular Observable implementation indicates otherwise,
* Observers should make no assumptions about the Observable implementation, such
* as the order of notifications that multiple Observers will receive.
*
* For more information see the RxJava Wiki
*
*
* @param Observer
* @return a {@link Subscription} reference that allows observers
* to stop receiving notifications before the provider has finished sending them
*/
public Subscription subscribe(Observer observer) {
if (onSubscribe == null) {
throw new IllegalStateException("onSubscribe function can not be null.");
// the subscribe function can also be overridden but generally that's not the appropriate approach so I won't mention that in the exception
}
if (isTrusted) {
return onSubscribe.call(observer);
} else {
AtomicObservableSubscription subscription = new AtomicObservableSubscription();
return subscription.wrap(onSubscribe.call(new AtomicObserver(subscription, observer)));
}
}
@SuppressWarnings({ "rawtypes", "unchecked" })
public Subscription subscribe(final Map callbacks) {
// lookup and memoize onNext
Object _onNext = callbacks.get("onNext");
if (_onNext == null) {
throw new RuntimeException("onNext must be implemented");
}
final FuncN onNext = Functions.from(_onNext);
return subscribe(new Observer() {
public void onCompleted() {
Object onComplete = callbacks.get("onCompleted");
if (onComplete != null) {
Functions.from(onComplete).call();
}
}
public void onError(Exception e) {
handleError(e);
Object onError = callbacks.get("onError");
if (onError != null) {
Functions.from(onError).call(e);
}
}
public void onNext(Object args) {
onNext.call(args);
}
});
}
@SuppressWarnings({ "rawtypes", "unchecked" })
public Subscription subscribe(final Object o) {
if (o instanceof Observer) {
// in case a dynamic language is not correctly handling the overloaded methods and we receive an Observer just forward to the correct method.
return subscribe((Observer) o);
}
// lookup and memoize onNext
if (o == null) {
throw new RuntimeException("onNext must be implemented");
}
final FuncN onNext = Functions.from(o);
return subscribe(new Observer() {
public void onCompleted() {
// do nothing
}
public void onError(Exception e) {
handleError(e);
// no callback defined
}
public void onNext(Object args) {
onNext.call(args);
}
});
}
public Subscription subscribe(final Action1 onNext) {
return subscribe(new Observer() {
public void onCompleted() {
// do nothing
}
public void onError(Exception e) {
handleError(e);
// no callback defined
}
public void onNext(T args) {
if (onNext == null) {
throw new RuntimeException("onNext must be implemented");
}
onNext.call(args);
}
});
}
@SuppressWarnings({ "rawtypes", "unchecked" })
public Subscription subscribe(final Object onNext, final Object onError) {
// lookup and memoize onNext
if (onNext == null) {
throw new RuntimeException("onNext must be implemented");
}
final FuncN onNextFunction = Functions.from(onNext);
return subscribe(new Observer() {
public void onCompleted() {
// do nothing
}
public void onError(Exception e) {
handleError(e);
if (onError != null) {
Functions.from(onError).call(e);
}
}
public void onNext(Object args) {
onNextFunction.call(args);
}
});
}
public Subscription subscribe(final Action1 onNext, final Action1 onError) {
return subscribe(new Observer() {
public void onCompleted() {
// do nothing
}
public void onError(Exception e) {
handleError(e);
if (onError != null) {
onError.call(e);
}
}
public void onNext(T args) {
if (onNext == null) {
throw new RuntimeException("onNext must be implemented");
}
onNext.call(args);
}
});
}
@SuppressWarnings({ "rawtypes", "unchecked" })
public Subscription subscribe(final Object onNext, final Object onError, final Object onComplete) {
// lookup and memoize onNext
if (onNext == null) {
throw new RuntimeException("onNext must be implemented");
}
final FuncN onNextFunction = Functions.from(onNext);
return subscribe(new Observer() {
public void onCompleted() {
if (onComplete != null) {
Functions.from(onComplete).call();
}
}
public void onError(Exception e) {
handleError(e);
if (onError != null) {
Functions.from(onError).call(e);
}
}
public void onNext(Object args) {
onNextFunction.call(args);
}
});
}
public Subscription subscribe(final Action1 onNext, final Action1 onError, final Action0 onComplete) {
return subscribe(new Observer() {
public void onCompleted() {
onComplete.call();
}
public void onError(Exception e) {
handleError(e);
if (onError != null) {
onError.call(e);
}
}
public void onNext(T args) {
if (onNext == null) {
throw new RuntimeException("onNext must be implemented");
}
onNext.call(args);
}
});
}
private void handleError(Exception e) {
// not implemented yet since open-sourcing
// intended for plugins to capture and log all errors
// even if Observers drop them on the floor
}
/**
* An Observable that never sends any information to an {@link Observer}.
*
* This Observable is useful primarily for testing purposes.
*
* @param
* the type of item emitted by the Observable
*/
private static class NeverObservable extends Observable {
public NeverObservable() {
super(new Func1, Subscription>() {
@Override
public Subscription call(Observer t1) {
return new NoOpObservableSubscription();
}
});
}
}
/**
* A {@link Subscription} that does nothing when its unsubscribe method is called.
*/
private static class NoOpObservableSubscription implements Subscription {
public void unsubscribe() {
}
}
/**
* an Observable that calls {@link Observer#onError(Exception)} when the Observer subscribes.
*
* @param
* the type of object returned by the Observable
*/
private static class ThrowObservable extends Observable {
public ThrowObservable(final Exception exception) {
super(new Func1, Subscription>() {
/**
* Accepts an {@link Observer} and calls its onError method.
*
* @param observer
* an {@link Observer} of this Observable
* @return a reference to the subscription
*/
@Override
public Subscription call(Observer observer) {
observer.onError(exception);
return new NoOpObservableSubscription();
}
});
}
}
/**
* Creates an Observable that will execute the given function when a {@link Observer} subscribes to it.
*
* Write the function you pass to create so that it behaves as an Observable - calling the passed-in
* onNext, onError, and onCompleted methods appropriately.
*
* A well-formed Observable must call either the {@link Observer}'s onCompleted method exactly once or its onError method exactly once.
*
* See Rx Design Guidelines (PDF) for detailed information.
*
* @param
* the type emitted by the Observable sequence
* @param func
* a function that accepts an Observer and calls its onNext, onError, and onCompleted methods
* as appropriate, and returns a {@link Subscription} to allow canceling the subscription (if applicable)
* @return an Observable that, when an {@link Observer} subscribes to it, will execute the given function
*/
public static Observable create(Func1, Subscription> func) {
return new Observable(func);
}
/*
* Private version that creates a 'trusted' Observable to allow performance optimizations.
*/
private static Observable _create(Func1, Subscription> func) {
return new Observable(func, true);
}
/**
* Creates an Observable that will execute the given function when a {@link Observer} subscribes to it.
*
* This method accept {@link Object} to allow different languages to pass in closures using {@link FunctionLanguageAdaptor}.
*
* Write the function you pass to create so that it behaves as an Observable - calling the passed-in
* onNext, onError, and onCompleted methods appropriately.
*
* A well-formed Observable must call either the {@link Observer}'s onCompleted method exactly once or its onError method exactly once.
*
* See Rx Design Guidelines (PDF) for detailed information.
*
* @param
* the type emitted by the Observable sequence
* @param func
* a function that accepts an Observer and calls its onNext, onError, and onCompleted methods
* as appropriate, and returns a {@link Subscription} to allow canceling the subscription (if applicable)
* @return an Observable that, when an {@link Observer} subscribes to it, will execute the given function
*/
public static Observable create(final Object callback) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(callback);
return create(new Func1, Subscription>() {
@Override
public Subscription call(Observer t1) {
return (Subscription) _f.call(t1);
}
});
}
/**
* Returns an Observable that returns no data to the {@link Observer} and immediately invokes its onCompleted method.
*
*
*
* @param
* the type of item emitted by the Observable
* @return an Observable that returns no data to the {@link Observer} and immediately invokes the {@link Observer}'s onCompleted method
*/
public static Observable empty() {
return toObservable(new ArrayList());
}
/**
* Returns an Observable that calls onError when an {@link Observer} subscribes to it.
*
*
* @param exception
* the error to throw
* @param
* the type of object returned by the Observable
* @return an Observable object that calls onError when an {@link Observer} subscribes
*/
public static Observable error(Exception exception) {
return new ThrowObservable(exception);
}
/**
* Filters an Observable by discarding any of its emissions that do not meet some test.
*
*
*
* @param that
* the Observable to filter
* @param predicate
* a function that evaluates the items emitted by the source Observable, returning true if they pass the filter
* @return an Observable that emits only those items in the original Observable that the filter evaluates as true
*/
public static Observable filter(Observable that, Func1 predicate) {
return _create(OperationFilter.filter(that, predicate));
}
/**
* Filters an Observable by discarding any of its emissions that do not meet some test.
*
*
*
* @param that
* the Observable to filter
* @param predicate
* a function that evaluates the items emitted by the source Observable, returning true if they pass the filter
* @return an Observable that emits only those items in the original Observable that the filter evaluates as true
*/
public static Observable filter(Observable that, final Object function) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(function);
return filter(that, new Func1() {
@Override
public Boolean call(T t1) {
return (Boolean) _f.call(t1);
}
});
}
/**
* Converts an {@link Iterable} sequence to an Observable sequence.
*
* @param iterable
* the source {@link Iterable} sequence
* @param
* the type of items in the {@link Iterable} sequence and the type emitted by the resulting Observable
* @return an Observable that emits each item in the source {@link Iterable} sequence
* @see {@link #toObservable(Iterable)}
*/
public static Observable from(Iterable iterable) {
return toObservable(iterable);
}
/**
* Converts an Array to an Observable sequence.
*
* @param iterable
* the source Array
* @param
* the type of items in the Array, and the type of items emitted by the resulting Observable
* @return an Observable that emits each item in the source Array
* @see {@link #toObservable(Object...)}
*/
public static Observable from(T... items) {
return toObservable(items);
}
/**
* Returns an Observable that notifies an {@link Observer} of a single value and then completes.
*
* To convert any object into an Observable that emits that object, pass that object into the just method.
*
* This is similar to the {@link toObservable} method, except that toObservable will convert
* an {@link Iterable} object into an Observable that emits each of the items in the {@link Iterable}, one
* at a time, while the just method would convert the {@link Iterable} into an Observable
* that emits the entire {@link Iterable} as a single item.
*
*
*
* @param value
* the value to pass to the Observer's onNext method
* @param
* the type of the value
* @return an Observable that notifies an {@link Observer} of a single value and then completes
*/
public static Observable just(T value) {
List list = new ArrayList();
list.add(value);
return toObservable(list);
}
/**
* Takes the last item emitted by a source Observable and returns an Observable that emits only
* that item as its sole emission.
*
*
*
* @param that
* the source Observable
* @return an Observable that emits a single item, which is identical to the last item emitted
* by the source Observable
*/
public static Observable last(final Observable that) {
return _create(OperationLast.last(that));
}
/**
* Applies a function of your choosing to every notification emitted by an Observable, and returns
* this transformation as a new Observable sequence.
*
*
*
* @param sequence
* the source Observable
* @param func
* a function to apply to each item in the sequence emitted by the source Observable
* @param
* the type of items emitted by the the source Observable
* @param
* the type of items returned by map function
* @return an Observable that is the result of applying the transformation function to each item
* in the sequence emitted by the source Observable
*/
public static Observable map(Observable sequence, Func1 func) {
return _create(OperationMap.map(sequence, func));
}
/**
* Applies a function of your choosing to every notification emitted by an Observable, and returns
* this transformation as a new Observable sequence.
*
*
*
* @param sequence
* the source Observable
* @param func
* a function to apply to each item in the sequence emitted by the source Observable
* @param
* the type of items emitted by the the source Observable
* @param
* the type of items returned by map function
* @return an Observable that is the result of applying the transformation function to each item
* in the sequence emitted by the source Observable
*/
public static Observable map(Observable sequence, final Object function) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(function);
return map(sequence, new Func1() {
@SuppressWarnings("unchecked")
@Override
public R call(T t1) {
return (R) _f.call(t1);
}
});
}
/**
* Creates a new Observable sequence by applying a function that you supply to each object in the
* original Observable sequence, where that function is itself an Observable that emits objects,
* and then merges the results of that function applied to every item emitted by the original
* Observable, emitting these merged results as its own sequence.
*
*
*
* @param sequence
* the source Observable
* @param func
* a function to apply to each item emitted by the source Observable, generating a
* Observable
* @param
* the type emitted by the source Observable
* @param
* the type emitted by the Observables emitted by func
* @return an Observable that emits a sequence that is the result of applying the transformation
* function to each item emitted by the source Observable and merging the results of
* the Observables obtained from this transformation
*/
public static Observable mapMany(Observable sequence, Func1> func) {
return _create(OperationMap.mapMany(sequence, func));
}
/**
* Creates a new Observable sequence by applying a function that you supply to each object in the
* original Observable sequence, where that function is itself an Observable that emits objects,
* and then merges the results of that function applied to every item emitted by the original
* Observable, emitting these merged results as its own sequence.
*
*
*
* @param sequence
* the source Observable
* @param func
* a function to apply to each item emitted by the source Observable, generating a
* Observable
* @param
* the type emitted by the source Observable
* @param
* the type emitted by the Observables emitted by func
* @return an Observable that emits a sequence that is the result of applying the transformation
* function to each item emitted by the source Observable and merging the results of
* the Observables obtained from this transformation
*/
public static Observable mapMany(Observable sequence, final Object function) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(function);
return mapMany(sequence, new Func1() {
@SuppressWarnings("unchecked")
@Override
public R call(T t1) {
return (R) _f.call(t1);
}
});
}
/**
* Materializes the implicit notifications of an observable sequence as explicit notification values.
*
*
*
* @param source
* An observable sequence of elements to project.
* @return An observable sequence whose elements are the result of materializing the notifications of the given sequence.
* @see http://msdn.microsoft.com/en-us/library/hh229453(v=VS.103).aspx
*/
public static Observable> materialize(final Observable sequence) {
return _create(OperationMaterialize.materialize(sequence));
}
/**
* Flattens the Observable sequences from a list of Observables into one Observable sequence
* without any transformation. You can combine the output of multiple Observables so that they
* act like a single Observable, by using the merge method.
*
*
*
* @param source
* a list of Observables that emit sequences of items
* @return an Observable that emits a sequence of elements that are the result of flattening the
* output from the source list of Observables
* @see MSDN: Observable.Merge Method
*/
public static Observable merge(List> source) {
return _create(OperationMerge.merge(source));
}
/**
* Flattens the Observable sequences emitted by a sequence of Observables that are emitted by a
* Observable into one Observable sequence without any transformation. You can combine the output
* of multiple Observables so that they act like a single Observable, by using the merge method.
*
*
*
* @param source
* an Observable that emits Observables
* @return an Observable that emits a sequence of elements that are the result of flattening the
* output from the Observables emitted by the source Observable
* @see MSDN: Observable.Merge Method
*/
public static Observable merge(Observable> source) {
return _create(OperationMerge.merge(source));
}
/**
* Flattens the Observable sequences from a series of Observables into one Observable sequence
* without any transformation. You can combine the output of multiple Observables so that they
* act like a single Observable, by using the merge method.
*
*
*
* @param source
* a series of Observables that emit sequences of items
* @return an Observable that emits a sequence of elements that are the result of flattening the
* output from the source Observables
* @see MSDN: Observable.Merge Method
*/
public static Observable merge(Observable... source) {
return _create(OperationMerge.merge(source));
}
/**
* Combines the objects emitted by two or more Observables, and emits the result as a single Observable,
* by using the concat method.
*
*
*
* @param source
* a series of Observables that emit sequences of items
* @return an Observable that emits a sequence of elements that are the result of combining the
* output from the source Observables
* @see MSDN: Observable.Concat Method
*/
public static Observable concat(Observable... source) {
return _create(OperationConcat.concat(source));
}
/**
* Same functionality as merge except that errors received to onError will be held until all sequences have finished (onComplete/onError) before sending the error.
*
* Only the first onError received will be sent.
*
* This enables receiving all successes from merged sequences without one onError from one sequence causing all onNext calls to be prevented.
*
*
*
* @param source
* a list of Observables that emit sequences of items
* @return an Observable that emits a sequence of elements that are the result of flattening the
* output from the source list of Observables
* @see MSDN: Observable.Merge Method
*/
public static Observable mergeDelayError(List> source) {
return _create(OperationMergeDelayError.mergeDelayError(source));
}
/**
* Same functionality as merge except that errors received to onError will be held until all sequences have finished (onComplete/onError) before sending the error.
*
* Only the first onError received will be sent.
*
* This enables receiving all successes from merged sequences without one onError from one sequence causing all onNext calls to be prevented.
*
*
*
* @param source
* an Observable that emits Observables
* @return an Observable that emits a sequence of elements that are the result of flattening the
* output from the Observables emitted by the source Observable
* @see MSDN: Observable.Merge Method
*/
public static Observable mergeDelayError(Observable> source) {
return _create(OperationMergeDelayError.mergeDelayError(source));
}
/**
* Same functionality as merge except that errors received to onError will be held until all sequences have finished (onComplete/onError) before sending the error.
*
* Only the first onError received will be sent.
*
* This enables receiving all successes from merged sequences without one onError from one sequence causing all onNext calls to be prevented.
*
*
*
* @param source
* a series of Observables that emit sequences of items
* @return an Observable that emits a sequence of elements that are the result of flattening the
* output from the source Observables
* @see MSDN: Observable.Merge Method
*/
public static Observable mergeDelayError(Observable... source) {
return _create(OperationMergeDelayError.mergeDelayError(source));
}
/**
* Returns an Observable that never sends any information to an {@link Observer}.
*
* This observable is useful primarily for testing purposes.
*
* @param
* the type of item (not) emitted by the Observable
* @return an Observable that never sends any information to an {@link Observer}
*/
public static Observable never() {
return new NeverObservable();
}
/**
* A {@link Subscription} that does nothing.
*
* //TODO should this be moved to a Subscriptions utility class?
*
* @return
*/
public static Subscription noOpSubscription() {
return new NoOpObservableSubscription();
}
/**
* A {@link Subscription} implemented via a Func
*
* //TODO should this be moved to a Subscriptions utility class?
*
* @return
*/
public static Subscription createSubscription(final Action0 unsubscribe) {
return new Subscription() {
@Override
public void unsubscribe() {
unsubscribe.call();
}
};
}
/**
* A {@link Subscription} implemented via an anonymous function (such as closures from other languages).
*
* //TODO should this be moved to a Subscriptions utility class?
*
* @return
*/
public static Subscription createSubscription(final Object unsubscribe) {
final FuncN> f = Functions.from(unsubscribe);
return new Subscription() {
@Override
public void unsubscribe() {
f.call();
}
};
}
/**
* Instruct an Observable to pass control to another Observable (the return value of a function)
* rather than calling onError if it encounters an error.
*
* By default, when an Observable encounters an error that prevents it from emitting the expected item to its Observer,
* the Observable calls its {@link Observer}'s onError function, and then quits without calling any more
* of its {@link Observer}'s closures. The onErrorResumeNext method changes this behavior. If you pass a
* function that emits an Observable (resumeFunction) to an Observable's onErrorResumeNext method,
* if the original Observable encounters an error, instead of calling its {@link Observer}'s onError function, it
* will instead relinquish control to this new Observable, which will call the {@link Observer}'s onNext method if
* it is able to do so. In such a case, because no Observable necessarily invokes onError, the Observer may
* never know that an error happened.
*
* You can use this to prevent errors from propagating or to supply fallback data should errors be encountered.
*
*
*
* @param that
* the source Observable
* @param resumeFunction
* a function that returns an Observable that will take over if the source Observable
* encounters an error
* @return the source Observable, with its behavior modified as described
*/
public static Observable onErrorResumeNext(final Observable that, final Func1> resumeFunction) {
return _create(OperationOnErrorResumeNextViaFunction.onErrorResumeNextViaFunction(that, resumeFunction));
}
/**
* Instruct an Observable to pass control to another Observable (the return value of a function)
* rather than calling onError if it encounters an error.
*
* By default, when an Observable encounters an error that prevents it from emitting the expected item to its Observer,
* the Observable calls its {@link Observer}'s onError function, and then quits without calling any more
* of its {@link Observer}'s closures. The onErrorResumeNext method changes this behavior. If you pass a
* function that emits an Observable (resumeFunction) to an Observable's onErrorResumeNext method,
* if the original Observable encounters an error, instead of calling its {@link Observer}'s onError function, it
* will instead relinquish control to this new Observable, which will call the {@link Observer}'s onNext method if
* it is able to do so. In such a case, because no Observable necessarily invokes onError, the Observer may
* never know that an error happened.
*
* You can use this to prevent errors from propagating or to supply fallback data should errors be encountered.
*
*
*
* @param that
* the source Observable
* @param resumeFunction
* a function that returns an Observable that will take over if the source Observable
* encounters an error
* @return the source Observable, with its behavior modified as described
*/
public static Observable onErrorResumeNext(final Observable that, final Object resumeFunction) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(resumeFunction);
return onErrorResumeNext(that, new Func1>() {
@SuppressWarnings("unchecked")
@Override
public Observable call(Exception e) {
return (Observable) _f.call(e);
}
});
}
/**
* Instruct an Observable to pass control to another Observable rather than calling onError if it encounters an error.
*
* By default, when an Observable encounters an error that prevents it from emitting the expected item to its Observer,
* the Observable calls its {@link Observer}'s onError function, and then quits without calling any more
* of its {@link Observer}'s closures. The onErrorResumeNext method changes this behavior. If you pass a
* function that emits an Observable (resumeFunction) to an Observable's onErrorResumeNext method,
* if the original Observable encounters an error, instead of calling its {@link Observer}'s onError function, it
* will instead relinquish control to this new Observable, which will call the {@link Observer}'s onNext method if
* it is able to do so. In such a case, because no Observable necessarily invokes onError, the Observer may
* never know that an error happened.
*
* You can use this to prevent errors from propagating or to supply fallback data should errors be encountered.
*
*
*
* @param that
* the source Observable
* @param resumeFunction
* a function that returns an Observable that will take over if the source Observable
* encounters an error
* @return the source Observable, with its behavior modified as described
*/
public static Observable onErrorResumeNext(final Observable that, final Observable resumeSequence) {
return _create(OperationOnErrorResumeNextViaObservable.onErrorResumeNextViaObservable(that, resumeSequence));
}
/**
* Instruct an Observable to emit a particular item to its Observer's onNext function
* rather than calling onError if it encounters an error.
*
* By default, when an Observable encounters an error that prevents it from emitting the expected item to its {@link Observer}, the Observable calls its {@link Observer}'s onError
* function, and then quits
* without calling any more of its {@link Observer}'s closures. The onErrorReturn method changes
* this behavior. If you pass a function (resumeFunction) to an Observable's onErrorReturn
* method, if the original Observable encounters an error, instead of calling its {@link Observer}'s
* onError function, it will instead pass the return value of resumeFunction to the {@link Observer}'s onNext method.
*
* You can use this to prevent errors from propagating or to supply fallback data should errors be encountered.
*
* @param that
* the source Observable
* @param resumeFunction
* a function that returns a value that will be passed into an {@link Observer}'s onNext function if the Observable encounters an error that would
* otherwise cause it to call onError
* @return the source Observable, with its behavior modified as described
*/
public static Observable onErrorReturn(final Observable that, Func1 resumeFunction) {
return _create(OperationOnErrorReturn.onErrorReturn(that, resumeFunction));
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the final result from the final call to your function as its sole
* output.
*
* This technique, which is called "reduce" here, is sometimes called "fold," "accumulate," "compress," or "inject" in other programming contexts. Groovy, for instance, has an inject
* method that does a similar operation on lists.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be used in the next accumulator call (if applicable)
*
* @return an Observable that emits a single element that is the result of accumulating the
* output from applying the accumulator to the sequence of items emitted by the source
* Observable
* @see MSDN: Observable.Aggregate
* @see Wikipedia: Fold (higher-order function)
*/
public static Observable reduce(Observable sequence, Func2 accumulator) {
return last(_create(OperationScan.scan(sequence, accumulator)));
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the final result from the final call to your function as its sole
* output.
*
* This technique, which is called "reduce" here, is sometimes called "fold," "accumulate," "compress," or "inject" in other programming contexts. Groovy, for instance, has an inject
* method that does a similar operation on lists.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be used in the next accumulator call (if applicable)
*
* @return an Observable that emits a single element that is the result of accumulating the
* output from applying the accumulator to the sequence of items emitted by the source
* Observable
* @see MSDN: Observable.Aggregate
* @see Wikipedia: Fold (higher-order function)
*/
public static Observable reduce(final Observable sequence, final Object accumulator) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(accumulator);
return reduce(sequence, new Func2() {
@SuppressWarnings("unchecked")
@Override
public T call(T t1, T t2) {
return (T) _f.call(t1, t2);
}
});
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the final result from the final call to your function as its sole
* output.
*
* This technique, which is called "reduce" here, is sometimes called "fold," "accumulate," "compress," or "inject" in other programming contexts. Groovy, for instance, has an inject
* method that does a similar operation on lists.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param initialValue
* a seed passed into the first execution of the accumulator function
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be used in the next accumulator call (if applicable)
*
* @return an Observable that emits a single element that is the result of accumulating the
* output from applying the accumulator to the sequence of items emitted by the source
* Observable
* @see MSDN: Observable.Aggregate
* @see Wikipedia: Fold (higher-order function)
*/
public static Observable reduce(Observable sequence, T initialValue, Func2 accumulator) {
return last(_create(OperationScan.scan(sequence, initialValue, accumulator)));
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the final result from the final call to your function as its sole
* output.
*
* This technique, which is called "reduce" here, is sometimes called "fold," "accumulate," "compress," or "inject" in other programming contexts. Groovy, for instance, has an inject
* method that does a similar operation on lists.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param initialValue
* a seed passed into the first execution of the accumulator function
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be used in the next accumulator call (if applicable)
* @return an Observable that emits a single element that is the result of accumulating the
* output from applying the accumulator to the sequence of items emitted by the source
* Observable
* @see MSDN: Observable.Aggregate
* @see Wikipedia: Fold (higher-order function)
*/
public static Observable reduce(final Observable sequence, final T initialValue, final Object accumulator) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(accumulator);
return reduce(sequence, initialValue, new Func2() {
@SuppressWarnings("unchecked")
@Override
public T call(T t1, T t2) {
return (T) _f.call(t1, t2);
}
});
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the result of each of these iterations as its own sequence.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be emitted and used in the next accumulator call (if applicable)
* @return an Observable that emits a sequence of items that are the result of accumulating the
* output from the sequence emitted by the source Observable
* @see MSDN: Observable.Scan
*/
public static Observable scan(Observable sequence, Func2 accumulator) {
return _create(OperationScan.scan(sequence, accumulator));
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the result of each of these iterations as its own sequence.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be emitted and used in the next accumulator call (if applicable)
* @return an Observable that emits a sequence of items that are the result of accumulating the
* output from the sequence emitted by the source Observable
* @see MSDN: Observable.Scan
*/
public static Observable scan(final Observable sequence, final Object accumulator) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(accumulator);
return scan(sequence, new Func2() {
@SuppressWarnings("unchecked")
@Override
public T call(T t1, T t2) {
return (T) _f.call(t1, t2);
}
});
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the result of each of these iterations as its own sequence.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param initialValue
* the initial (seed) accumulator value
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be emitted and used in the next accumulator call (if applicable)
* @return an Observable that emits a sequence of items that are the result of accumulating the
* output from the sequence emitted by the source Observable
* @see MSDN: Observable.Scan
*/
public static Observable scan(Observable sequence, T initialValue, Func2 accumulator) {
return _create(OperationScan.scan(sequence, initialValue, accumulator));
}
/**
* Returns an Observable that applies a function of your choosing to the first item emitted by a
* source Observable, then feeds the result of that function along with the second item emitted
* by an Observable into the same function, and so on until all items have been emitted by the
* source Observable, emitting the result of each of these iterations as its own sequence.
*
*
*
* @param
* the type item emitted by the source Observable
* @param sequence
* the source Observable
* @param initialValue
* the initial (seed) accumulator value
* @param accumulator
* an accumulator function to be invoked on each element from the sequence, whose
* result will be emitted and used in the next accumulator call (if applicable)
* @return an Observable that emits a sequence of items that are the result of accumulating the
* output from the sequence emitted by the source Observable
* @see MSDN: Observable.Scan
*/
public static Observable scan(final Observable sequence, final T initialValue, final Object accumulator) {
@SuppressWarnings("rawtypes")
final FuncN _f = Functions.from(accumulator);
return scan(sequence, initialValue, new Func2() {
@SuppressWarnings("unchecked")
@Override
public T call(T t1, T t2) {
return (T) _f.call(t1, t2);
}
});
}
/**
* Returns an Observable that skips the first num items emitted by the source
* Observable. You can ignore the first num items emitted by an Observable and attend
* only to those items that come after, by modifying the Observable with the skip method.
*
*
*
* @param items
* the source Observable
* @param num
* the number of items to skip
* @return an Observable that emits the same sequence of items emitted by the source Observable,
* except for the first num items
* @see MSDN: Observable.Skip Method
*/
public static Observable skip(final Observable items, int num) {
return _create(OperationSkip.skip(items, num));
}
/**
* Accepts an Observable and wraps it in another Observable that ensures that the resulting
* Observable is chronologically well-behaved.
*
* A well-behaved observable ensures onNext, onCompleted, or onError calls to its subscribers are not interleaved, onCompleted and
* onError are only called once respectively, and no
* onNext calls follow onCompleted and onError calls.
*
* @param observable
* the source Observable
* @param
* the type of item emitted by the source Observable
* @return an Observable that is a chronologically well-behaved version of the source Observable
*/
public static Observable synchronize(Observable observable) {
return _create(OperationSynchronize.synchronize(observable));
}
/**
* Returns an Observable that emits the first num items emitted by the source
* Observable.
*
* You can choose to pay attention only to the first num values emitted by an Observable by calling its take method. This method returns an Observable that will call a
* subscribing Observer's onNext function a
* maximum of num times before calling onCompleted.
*
*
*
* @param items
* the source Observable
* @param num
* the number of items from the start of the sequence emitted by the source
* Observable to emit
* @return an Observable that only emits the first num items emitted by the source
* Observable
*/
public static Observable take(final Observable items, final int num) {
return _create(OperationTake.take(items, num));
}
/**
* Returns an Observable that emits a single item, a list composed of all the items emitted by
* the source Observable.
*
* Normally, an Observable that returns multiple items will do so by calling its Observer's onNext function for each such item. You can change this behavior, instructing the
* Observable
* to
* compose a list of all of these multiple items and
* then to call the Observer's onNext function once, passing it the entire list, by calling the Observable object's toList method prior to calling its
* subscribe
* method.
*
*
*
* @param that
* the source Observable
* @return an Observable that emits a single item: a List containing all of the
* items emitted by the source Observable
*/
public static Observable> toList(final Observable that) {
return _create(OperationToObservableList.toObservableList(that));
}
/**
* Converts an Iterable sequence to an Observable sequence.
*
* Any object that supports the Iterable interface can be converted into an Observable that emits
* each iterable item in the object, by passing the object into the toObservable method.
*
*
*
* @param iterable
* the source Iterable sequence
* @param
* the type of items in the iterable sequence and the type emitted by the resulting
* Observable
* @return an Observable that emits each item in the source Iterable sequence
*/
public static Observable toObservable(Iterable iterable) {
return _create(OperationToObservableIterable.toObservableIterable(iterable));
}
/**
* Converts an Future to an Observable sequence.
*
* Any object that supports the {@link Future} interface can be converted into an Observable that emits
* the return value of the get() method in the object, by passing the object into the toObservable method.
* The subscribe method on this synchronously so the Subscription returned doesn't nothing.
*
* @param future
* the source {@link Future}
* @param
* the type of of object that the future's returns and the type emitted by the resulting
* Observable
* @return an Observable that emits the item from the source Future
*/
public static Observable toObservable(Future future) {
return _create(OperationToObservableFuture.toObservableFuture(future));
}
/**
* Converts an Future to an Observable sequence.
*
* Any object that supports the {@link Future} interface can be converted into an Observable that emits
* the return value of the get() method in the object, by passing the object into the toObservable method.
* The subscribe method on this synchronously so the Subscription returned doesn't nothing.
* If the future timesout the {@link TimeoutException} exception is passed to the onError.
*
* @param future
* the source {@link Future}
* @param time
* the maximum time to wait
* @param unit
* the time unit of the time argument
* @param
* the type of of object that the future's returns and the type emitted by the resulting
* Observable
* @return an Observable that emits the item from the source Future
*/
public static Observable toObservable(Future future, long time, TimeUnit unit) {
return _create(OperationToObservableFuture.toObservableFuture(future, time, unit));
}
/**
* Converts an Array sequence to an Observable sequence.
*
* An Array can be converted into an Observable that emits each item in the Array, by passing the
* Array into the toObservable method.
*
*
*
* @param iterable
* the source Array
* @param
* the type of items in the Array, and the type of items emitted by the resulting
* Observable
* @return an Observable that emits each item in the source Array
*/
public static Observable toObservable(T... items) {
return toObservable(Arrays.asList(items));
}
/**
* Sort T objects by their natural order (object must implement Comparable).
*
*
*
* @param sequence
* @throws ClassCastException
* if T objects do not implement Comparable
* @return
*/
public static Observable> toSortedList(Observable sequence) {
return _create(OperationToObservableSortedList.toSortedList(sequence));
}
/**
* Sort T objects using the defined sort function.
*