package io.reactiverse.rxjava.pgclient;
import java.util.Map;
import rx.Observable;
import rx.Single;
import io.vertx.core.AsyncResult;
import io.vertx.core.Handler;
@io.vertx.lang.rx.RxGen(io.reactiverse.pgclient.PgStream.class)
public class PgStream<T> implements io.vertx.rxjava.core.streams.ReadStream<T> {
@Override
public String toString() {
return delegate.toString();
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
PgStream that = (PgStream) o;
return delegate.equals(that.delegate);
}
@Override
public int hashCode() {
return delegate.hashCode();
}
public static final io.vertx.lang.rx.TypeArg<PgStream> __TYPE_ARG = new io.vertx.lang.rx.TypeArg<>( obj -> new PgStream((io.reactiverse.pgclient.PgStream) obj),
PgStream::getDelegate
);
private final io.reactiverse.pgclient.PgStream<T> delegate;
public final io.vertx.lang.rx.TypeArg<T> __typeArg_0;
public PgStream(io.reactiverse.pgclient.PgStream delegate) {
this.delegate = delegate;
this.__typeArg_0 = io.vertx.lang.rx.TypeArg.unknown(); }
public PgStream(io.reactiverse.pgclient.PgStream delegate, io.vertx.lang.rx.TypeArg<T> typeArg_0) {
this.delegate = delegate;
this.__typeArg_0 = typeArg_0;
}
public io.reactiverse.pgclient.PgStream getDelegate() {
return delegate;
}
private rx.Observable<T> observable;
public synchronized rx.Observable<T> toObservable() {
if (observable == null) {
java.util.function.Function<T, T> conv = (java.util.function.Function<T, T>) __typeArg_0.wrap;
observable = io.vertx.rx.java.RxHelper.toObservable(delegate, conv);
}
return observable;
}
public io.vertx.rxjava.core.streams.ReadStream<T> fetch(long arg0) {
delegate.fetch(arg0);
return this;
}
public io.vertx.rxjava.core.streams.Pipe<T> pipe() {
io.vertx.rxjava.core.streams.Pipe<T> ret = io.vertx.rxjava.core.streams.Pipe.newInstance(delegate.pipe(), __typeArg_0);
return ret;
}
public void pipeTo(io.vertx.rxjava.core.streams.WriteStream<T> dst) {
delegate.pipeTo(dst.getDelegate());
}
public void pipeTo(io.vertx.rxjava.core.streams.WriteStream<T> dst, Handler<AsyncResult<Void>> handler) {
delegate.pipeTo(dst.getDelegate(), handler);
}
public Single<Void> rxPipeTo(io.vertx.rxjava.core.streams.WriteStream<T> dst) {
return Single.create(new io.vertx.rx.java.SingleOnSubscribeAdapter<>(fut -> {
pipeTo(dst, fut);
}));
}
public io.reactiverse.rxjava.pgclient.PgStream<T> exceptionHandler(Handler<Throwable> handler) {
delegate.exceptionHandler(handler);
return this;
}
public io.reactiverse.rxjava.pgclient.PgStream<T> handler(Handler<T> handler) {
delegate.handler(new Handler<T>() {
public void handle(T event) {
handler.handle((T)__typeArg_0.wrap(event));
}
});
return this;
}
public io.reactiverse.rxjava.pgclient.PgStream<T> pause() {
delegate.pause();
return this;
}
public io.reactiverse.rxjava.pgclient.PgStream<T> resume() {
delegate.resume();
return this;
}
public io.reactiverse.rxjava.pgclient.PgStream<T> endHandler(Handler<Void> endHandler) {
delegate.endHandler(endHandler);
return this;
}
public void close() {
delegate.close();
}
public void close(Handler<AsyncResult<Void>> completionHandler) {
delegate.close(completionHandler);
}
public Single<Void> rxClose() {
return Single.create(new io.vertx.rx.java.SingleOnSubscribeAdapter<>(fut -> {
close(fut);
}));
}
public static <T>PgStream<T> newInstance(io.reactiverse.pgclient.PgStream arg) {
return arg != null ? new PgStream<T>(arg) : null;
}
public static <T>PgStream<T> newInstance(io.reactiverse.pgclient.PgStream arg, io.vertx.lang.rx.TypeArg<T> __typeArg_T) {
return arg != null ? new PgStream<T>(arg, __typeArg_T) : null;
}
}