Skip to content

Commit df1d295

Browse files
committed
Verify precondition in TcpReactiveSocketFactory
1 parent 40e35c2 commit df1d295

4 files changed

Lines changed: 81 additions & 88 deletions

File tree

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

Lines changed: 0 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -56,18 +56,6 @@ private ClientTcpDuplexConnection(Channel channel, Bootstrap bootstrap, CopyOnWr
5656
this.bootstrap = bootstrap;
5757
}
5858

59-
public static Publisher<ClientTcpDuplexConnection> create(SocketAddress socketAddress, EventLoopGroup eventLoopGroup) {
60-
if (socketAddress instanceof InetSocketAddress) {
61-
try {
62-
return create(socketAddress, eventLoopGroup);
63-
} catch (Exception e) {
64-
throw new IllegalArgumentException(e.getMessage(), e);
65-
}
66-
} else {
67-
throw new IllegalArgumentException("unknown socket address type => " + socketAddress.getClass());
68-
}
69-
}
70-
7159
public static Publisher<ClientTcpDuplexConnection> create(InetSocketAddress address, EventLoopGroup eventLoopGroup) {
7260
return s -> {
7361
CopyOnWriteArrayList<Observer<Frame>> subjects = new CopyOnWriteArrayList<>();

reactivesocket-netty/src/main/java/io/reactivesocket/netty/tcp/client/TcpReactiveSocketFactory.java

Lines changed: 38 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import rx.Observable;
3030
import rx.RxReactiveStreams;
3131

32+
import java.net.InetSocketAddress;
3233
import java.net.SocketAddress;
3334
import java.util.function.Consumer;
3435

@@ -50,44 +51,48 @@ public TcpReactiveSocketFactory(EventLoopGroup eventLoopGroup, ConnectionSetupPa
5051

5152
@Override
5253
public Publisher<ReactiveSocket> call(SocketAddress address) {
53-
Publisher<ClientTcpDuplexConnection> connection
54-
= ClientTcpDuplexConnection.create(address, eventLoopGroup);
54+
if (address instanceof InetSocketAddress) {
55+
Publisher<ClientTcpDuplexConnection> connection
56+
= ClientTcpDuplexConnection.create((InetSocketAddress)address, eventLoopGroup);
5557

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

63-
@Override
64-
public void onNext(ClientTcpDuplexConnection connection) {
65-
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, connectionSetupPayload, errorStream);
66-
reactiveSocket.start(new Completable() {
67-
@Override
68-
public void success() {
69-
s.onNext(reactiveSocket);
70-
s.onCompleted();
71-
}
65+
@Override
66+
public void onNext(ClientTcpDuplexConnection connection) {
67+
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, connectionSetupPayload, errorStream);
68+
reactiveSocket.start(new Completable() {
69+
@Override
70+
public void success() {
71+
s.onNext(reactiveSocket);
72+
s.onCompleted();
73+
}
7274

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

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

85-
@Override
86-
public void onComplete() {
87-
}
88-
})
89-
);
87+
@Override
88+
public void onComplete() {
89+
}
90+
})
91+
);
9092

91-
return RxReactiveStreams.toPublisher(result);
93+
return RxReactiveStreams.toPublisher(result);
94+
} else {
95+
throw new IllegalArgumentException("unknown socket address type => " + address.getClass());
96+
}
9297
}
9398
}

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

Lines changed: 5 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -56,16 +56,11 @@ private ClientWebSocketDuplexConnection(Channel channel, Bootstrap bootstrap, C
5656
this.bootstrap = bootstrap;
5757
}
5858

59-
public static Publisher<ClientWebSocketDuplexConnection> create(SocketAddress socketAddress, String path, EventLoopGroup eventLoopGroup) {
60-
if (socketAddress instanceof InetSocketAddress) {
61-
InetSocketAddress address = (InetSocketAddress)socketAddress;
62-
try {
63-
return create(new URI("ws", null, address.getHostName(), address.getPort(), path, null, null), eventLoopGroup);
64-
} catch (URISyntaxException e) {
65-
throw new IllegalArgumentException(e.getMessage(), e);
66-
}
67-
} else {
68-
throw new IllegalArgumentException("unknown socket address type => " + socketAddress.getClass());
59+
public static Publisher<ClientWebSocketDuplexConnection> create(InetSocketAddress address, String path, EventLoopGroup eventLoopGroup) {
60+
try {
61+
return create(new URI("ws", null, address.getHostName(), address.getPort(), path, null, null), eventLoopGroup);
62+
} catch (URISyntaxException e) {
63+
throw new IllegalArgumentException(e.getMessage(), e);
6964
}
7065
}
7166

reactivesocket-netty/src/main/java/io/reactivesocket/netty/websocket/client/WebSocketReactiveSocketFactory.java

Lines changed: 38 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import rx.Observable;
3030
import rx.RxReactiveStreams;
3131

32+
import java.net.InetSocketAddress;
3233
import java.net.SocketAddress;
3334
import java.util.function.Consumer;
3435

@@ -52,44 +53,48 @@ public WebSocketReactiveSocketFactory(String path, EventLoopGroup eventLoopGroup
5253

5354
@Override
5455
public Publisher<ReactiveSocket> call(SocketAddress address) {
55-
Publisher<ClientWebSocketDuplexConnection> connection
56-
= ClientWebSocketDuplexConnection.create(address, path, eventLoopGroup);
56+
if (address instanceof InetSocketAddress) {
57+
Publisher<ClientWebSocketDuplexConnection> connection
58+
= ClientWebSocketDuplexConnection.create((InetSocketAddress)address, path, eventLoopGroup);
5759

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

65-
@Override
66-
public void onNext(ClientWebSocketDuplexConnection connection) {
67-
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, connectionSetupPayload, errorStream);
68-
reactiveSocket.start(new Completable() {
69-
@Override
70-
public void success() {
71-
s.onNext(reactiveSocket);
72-
s.onCompleted();
73-
}
67+
@Override
68+
public void onNext(ClientWebSocketDuplexConnection connection) {
69+
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, connectionSetupPayload, errorStream);
70+
reactiveSocket.start(new Completable() {
71+
@Override
72+
public void success() {
73+
s.onNext(reactiveSocket);
74+
s.onCompleted();
75+
}
7476

75-
@Override
76-
public void error(Throwable e) {
77-
s.onError(e);
78-
}
79-
});
80-
}
77+
@Override
78+
public void error(Throwable e) {
79+
s.onError(e);
80+
}
81+
});
82+
}
8183

82-
@Override
83-
public void onError(Throwable t) {
84-
s.onError(t);
85-
}
84+
@Override
85+
public void onError(Throwable t) {
86+
s.onError(t);
87+
}
8688

87-
@Override
88-
public void onComplete() {
89-
}
90-
})
91-
);
89+
@Override
90+
public void onComplete() {
91+
}
92+
})
93+
);
9294

93-
return RxReactiveStreams.toPublisher(result);
95+
return RxReactiveStreams.toPublisher(result);
96+
} else {
97+
throw new IllegalArgumentException("unknown socket address type => " + address.getClass());
98+
}
9499
}
95100
}

0 commit comments

Comments
 (0)