Skip to content

Commit 1c90af0

Browse files
steveguryNiteshKant
authored andcommitted
Correctly wire subscription from Transport Connectors. (#97)
**Problem** The initialization of the Publisher chain is usually started from a `onSubscribe` call which was lacking in the transport implementations we care about. **Solution** Add a call to `onSubscribe` with an EmptySubscriber because there's no notion of back-pressure here, we only eagerly deliver one `ReactiveSocket`.
1 parent 524d6d5 commit 1c90af0

4 files changed

Lines changed: 37 additions & 34 deletions

File tree

reactivesocket-transport-tcp/src/main/java/io/reactivesocket/transport/tcp/client/ClientTcpDuplexConnection.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import io.reactivesocket.DuplexConnection;
2626
import io.reactivesocket.Frame;
2727
import io.reactivesocket.exceptions.TransportException;
28+
import io.reactivesocket.internal.rx.EmptySubscription;
2829
import io.reactivesocket.rx.Completable;
2930
import io.reactivesocket.rx.Observable;
3031
import io.reactivesocket.rx.Observer;
@@ -73,6 +74,7 @@ protected void initChannel(SocketChannel ch) throws Exception {
7374
connect.addListener(connectFuture -> {
7475
if (connectFuture.isSuccess()) {
7576
Channel ch = connect.channel();
77+
s.onSubscribe(EmptySubscription.INSTANCE);
7678
s.onNext(new ClientTcpDuplexConnection(ch, subjects));
7779
s.onComplete();
7880
} else {

reactivesocket-transport-tcp/src/main/java/io/reactivesocket/transport/tcp/client/TcpReactiveSocketConnector.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
import io.netty.channel.EventLoopGroup;
1919
import io.reactivesocket.*;
20+
import io.reactivesocket.internal.rx.EmptySubscription;
2021
import io.reactivesocket.rx.Completable;
2122
import org.reactivestreams.Publisher;
2223
import org.reactivestreams.Subscriber;
@@ -57,6 +58,7 @@ public void onNext(ClientTcpDuplexConnection connection) {
5758
reactiveSocket.start(new Completable() {
5859
@Override
5960
public void success() {
61+
s.onSubscribe(EmptySubscription.INSTANCE);
6062
s.onNext(reactiveSocket);
6163
s.onComplete();
6264
}

reactivesocket-transport-websocket/src/main/java/io/reactivesocket/transport/websocket/client/ClientWebSocketDuplexConnection.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import io.reactivesocket.DuplexConnection;
2929
import io.reactivesocket.Frame;
3030
import io.reactivesocket.exceptions.TransportException;
31+
import io.reactivesocket.internal.rx.EmptySubscription;
3132
import io.reactivesocket.rx.Completable;
3233
import io.reactivesocket.rx.Observable;
3334
import io.reactivesocket.rx.Observer;
@@ -91,6 +92,7 @@ protected void initChannel(SocketChannel ch) throws Exception {
9192
.getHandshakePromise()
9293
.addListener(handshakeFuture -> {
9394
if (handshakeFuture.isSuccess()) {
95+
s.onSubscribe(EmptySubscription.INSTANCE);
9496
s.onNext(new ClientWebSocketDuplexConnection(ch, subjects));
9597
s.onComplete();
9698
} else {

reactivesocket-transport-websocket/src/main/java/io/reactivesocket/transport/websocket/client/WebSocketReactiveSocketConnector.java

Lines changed: 31 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -17,14 +17,13 @@
1717

1818
import io.netty.channel.EventLoopGroup;
1919
import io.reactivesocket.*;
20+
import io.reactivesocket.internal.rx.EmptySubscription;
2021
import io.reactivesocket.rx.Completable;
2122
import org.reactivestreams.Publisher;
2223
import org.reactivestreams.Subscriber;
2324
import org.reactivestreams.Subscription;
2425
import org.slf4j.Logger;
2526
import org.slf4j.LoggerFactory;
26-
import rx.Observable;
27-
import rx.RxReactiveStreams;
2827

2928
import java.net.InetSocketAddress;
3029
import java.net.SocketAddress;
@@ -54,42 +53,40 @@ public Publisher<ReactiveSocket> connect(SocketAddress address) {
5453
Publisher<ClientWebSocketDuplexConnection> connection
5554
= ClientWebSocketDuplexConnection.create((InetSocketAddress)address, path, eventLoopGroup);
5655

57-
Observable<ReactiveSocket> result = Observable.create(s ->
58-
connection.subscribe(new Subscriber<ClientWebSocketDuplexConnection>() {
59-
@Override
60-
public void onSubscribe(Subscription s) {
61-
s.request(1);
62-
}
56+
return s -> connection.subscribe(new Subscriber<ClientWebSocketDuplexConnection>() {
57+
@Override
58+
public void onSubscribe(Subscription s) {
59+
s.request(1);
60+
}
6361

64-
@Override
65-
public void onNext(ClientWebSocketDuplexConnection connection) {
66-
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, connectionSetupPayload, errorStream);
67-
reactiveSocket.start(new Completable() {
68-
@Override
69-
public void success() {
70-
s.onNext(reactiveSocket);
71-
s.onCompleted();
72-
}
62+
@Override
63+
public void onNext(ClientWebSocketDuplexConnection connection) {
64+
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, connectionSetupPayload, errorStream);
65+
reactiveSocket.start(new Completable() {
66+
@Override
67+
public void success() {
68+
s.onSubscribe(EmptySubscription.INSTANCE);
69+
s.onNext(reactiveSocket);
70+
s.onComplete();
71+
}
7372

74-
@Override
75-
public void error(Throwable e) {
76-
s.onError(e);
77-
}
78-
});
79-
}
73+
@Override
74+
public void error(Throwable e) {
75+
s.onError(e);
76+
}
77+
});
78+
}
8079

81-
@Override
82-
public void onError(Throwable t) {
83-
s.onError(t);
84-
}
80+
@Override
81+
public void onError(Throwable t) {
82+
s.onError(t);
83+
}
8584

86-
@Override
87-
public void onComplete() {
88-
}
89-
})
90-
);
91-
92-
return RxReactiveStreams.toPublisher(result);
85+
@Override
86+
public void onComplete() {
87+
s.onComplete();
88+
}
89+
});
9390
} else {
9491
throw new IllegalArgumentException("unknown socket address type => " + address.getClass());
9592
}

0 commit comments

Comments
 (0)