Skip to content

Commit f04b796

Browse files
committed
Merge pull request #5 from robertroeser/master
throwing a TransportException when a connection is closed
2 parents 5e295bc + b30241f commit f04b796

14 files changed

Lines changed: 37 additions & 15 deletions

File tree

reactivesocket-aeron/src/examples/java/io/reactivesocket/aeron/example/fireandforget/Forget.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
import io.reactivesocket.ConnectionSetupHandler;
1919
import io.reactivesocket.ConnectionSetupPayload;
2020
import io.reactivesocket.Payload;
21+
import io.reactivesocket.ReactiveSocket;
2122
import io.reactivesocket.RequestHandler;
2223
import io.reactivesocket.aeron.server.ReactiveSocketAeronServer;
2324
import io.reactivesocket.exceptions.SetupException;
@@ -33,7 +34,7 @@ public static void main(String... args) {
3334

3435
ReactiveSocketAeronServer server = ReactiveSocketAeronServer.create(host, 39790, new ConnectionSetupHandler() {
3536
@Override
36-
public RequestHandler apply(ConnectionSetupPayload setupPayload) throws SetupException {
37+
public RequestHandler apply(ConnectionSetupPayload setupPayload, ReactiveSocket rs) throws SetupException {
3738
return new RequestHandler() {
3839
@Override
3940
public Publisher<Payload> handleRequestResponse(Payload payload) {

reactivesocket-aeron/src/examples/java/io/reactivesocket/aeron/example/requestreply/Pong.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
import io.reactivesocket.ConnectionSetupHandler;
1919
import io.reactivesocket.ConnectionSetupPayload;
2020
import io.reactivesocket.Payload;
21+
import io.reactivesocket.ReactiveSocket;
2122
import io.reactivesocket.RequestHandler;
2223
import io.reactivesocket.aeron.server.ReactiveSocketAeronServer;
2324
import io.reactivesocket.exceptions.SetupException;
@@ -44,7 +45,7 @@ public static void main(String... args) {
4445

4546
ReactiveSocketAeronServer server = ReactiveSocketAeronServer.create(host, 39790, new ConnectionSetupHandler() {
4647
@Override
47-
public RequestHandler apply(ConnectionSetupPayload setupPayload) throws SetupException {
48+
public RequestHandler apply(ConnectionSetupPayload setupPayload, ReactiveSocket rs) throws SetupException {
4849
return new RequestHandler() {
4950
@Override
5051
public Publisher<Payload> handleRequestResponse(Payload payload) {

reactivesocket-aeron/src/main/java/io/reactivesocket/aeron/client/AeronClientDuplexConnection.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@
1919
import io.reactivesocket.DuplexConnection;
2020
import io.reactivesocket.Frame;
2121
import io.reactivesocket.aeron.internal.Loggable;
22+
import io.reactivesocket.aeron.internal.NotConnectedException;
23+
import io.reactivesocket.exceptions.TransportException;
2224
import io.reactivesocket.rx.Completable;
2325
import io.reactivesocket.rx.Disposable;
2426
import io.reactivesocket.rx.Observable;
@@ -101,7 +103,11 @@ public void onNext(Frame frame) {
101103

102104
@Override
103105
public void onError(Throwable t) {
104-
callback.error(t);
106+
if (t instanceof NotConnectedException) {
107+
callback.error(new TransportException(t));
108+
} else {
109+
callback.error(t);
110+
}
105111
}
106112

107113
@Override

reactivesocket-aeron/src/test/java/io/reactivesocket/aeron/client/ReactiveSocketAeronTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -112,7 +112,7 @@ public void requestResponseN(int count) throws Exception {
112112
AtomicLong counter = new AtomicLong();
113113
ReactiveSocketAeronServer.create(new ConnectionSetupHandler() {
114114
@Override
115-
public RequestHandler apply(ConnectionSetupPayload setupPayload) throws SetupException {
115+
public RequestHandler apply(ConnectionSetupPayload setupPayload, ReactiveSocket rs) throws SetupException {
116116
return new RequestHandler.Builder()
117117
.withRequestResponse(new Function<Payload, Publisher<Payload>>() {
118118
Frame frame = Frame.from(ByteBuffer.allocate(1));
@@ -194,7 +194,7 @@ public void onNext(Payload payload) {
194194
}
195195

196196
public void requestStreamN(int count) throws Exception {
197-
ReactiveSocketAeronServer.create(setupPayload ->
197+
ReactiveSocketAeronServer.create((setupPayload, rs) ->
198198
new RequestHandler.Builder()
199199
.withRequestStream(payload -> {
200200
ByteBuffer data = payload.getData();
@@ -264,7 +264,7 @@ public void testReconnection() throws Exception {
264264

265265
ReactiveSocketAeronServer server = ReactiveSocketAeronServer.create(new ConnectionSetupHandler() {
266266
@Override
267-
public RequestHandler apply(ConnectionSetupPayload setupPayload) throws SetupException {
267+
public RequestHandler apply(ConnectionSetupPayload setupPayload, ReactiveSocket rs) throws SetupException {
268268
return new RequestHandler() {
269269
Frame frame = Frame.from(ByteBuffer.allocate(1));
270270

reactivesocket-jsr-356/src/test/java/io/reactivesocket/javax/websocket/ClientServerEndpoint.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@
2525

2626
public class ClientServerEndpoint extends ReactiveSocketWebSocketServer {
2727
public ClientServerEndpoint() {
28-
super(setupPayload -> new RequestHandler() {
28+
super((setupPayload, rs) -> new RequestHandler() {
2929
@Override
3030
public Publisher<Payload> handleRequestResponse(Payload payload) {
3131
return s -> {

reactivesocket-jsr-356/src/test/java/io/reactivesocket/javax/websocket/PongEndpoint.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ public class PongEndpoint extends ReactiveSocketWebSocketServer {
3434
}
3535

3636
public PongEndpoint() {
37-
super(setupPayload -> new RequestHandler() {
37+
super((setupPayload, rs) -> new RequestHandler() {
3838
@Override
3939
public Publisher<Payload> handleRequestResponse(Payload payload) {
4040
return new Publisher<Payload>() {

reactivesocket-local/src/test/java/io/reactivesocket/local/ClientServerTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,7 @@ public class ClientServerTest {
4141
public static void setup() throws Exception {
4242
server = LocalServerReactiveSocketFactory.INSTANCE.callAndWait(new LocalServerReactiveSocketFactory.Config("test", new ConnectionSetupHandler() {
4343
@Override
44-
public RequestHandler apply(ConnectionSetupPayload setupPayload) throws SetupException {
44+
public RequestHandler apply(ConnectionSetupPayload setupPayload, ReactiveSocket rs) throws SetupException {
4545
return new RequestHandler() {
4646
@Override
4747
public Publisher<Payload> handleRequestResponse(Payload payload) {

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

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import io.netty.handler.codec.LengthFieldBasedFrameDecoder;
2929
import io.reactivesocket.DuplexConnection;
3030
import io.reactivesocket.Frame;
31+
import io.reactivesocket.exceptions.TransportException;
3132
import io.reactivesocket.rx.Completable;
3233
import io.reactivesocket.rx.Observable;
3334
import io.reactivesocket.rx.Observer;
@@ -39,6 +40,7 @@
3940
import java.io.IOException;
4041
import java.net.InetSocketAddress;
4142
import java.net.SocketAddress;
43+
import java.nio.channels.ClosedChannelException;
4244
import java.util.concurrent.CopyOnWriteArrayList;
4345

4446
public class ClientTcpDuplexConnection implements DuplexConnection {
@@ -121,7 +123,11 @@ public void onNext(Frame frame) {
121123
channelFuture.addListener(future -> {
122124
Throwable cause = future.cause();
123125
if (cause != null) {
124-
callback.error(cause);
126+
if (cause instanceof ClosedChannelException) {
127+
onError(new TransportException(cause));
128+
} else {
129+
onError(cause);
130+
}
125131
}
126132
});
127133
} catch (Throwable t) {

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

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import io.netty.handler.codec.http.websocketx.*;
2828
import io.reactivesocket.DuplexConnection;
2929
import io.reactivesocket.Frame;
30+
import io.reactivesocket.exceptions.TransportException;
3031
import io.reactivesocket.rx.Completable;
3132
import io.reactivesocket.rx.Observable;
3233
import io.reactivesocket.rx.Observer;
@@ -39,6 +40,7 @@
3940
import java.net.SocketAddress;
4041
import java.net.URI;
4142
import java.net.URISyntaxException;
43+
import java.nio.channels.ClosedChannelException;
4244
import java.util.concurrent.CopyOnWriteArrayList;
4345

4446
public class ClientWebSocketDuplexConnection implements DuplexConnection {
@@ -136,7 +138,11 @@ public void onNext(Frame frame) {
136138
channelFuture.addListener(future -> {
137139
Throwable cause = future.cause();
138140
if (cause != null) {
139-
callback.error(cause);
141+
if (cause instanceof ClosedChannelException) {
142+
onError(new TransportException(cause));
143+
} else {
144+
onError(cause);
145+
}
140146
}
141147
});
142148
} catch (Throwable t) {

reactivesocket-netty/src/test/java/io/reactivesocket/netty/tcp/ClientServerTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,7 @@ public class ClientServerTest {
5252
static EventLoopGroup bossGroup = new NioEventLoopGroup(1);
5353
static EventLoopGroup workerGroup = new NioEventLoopGroup(4);
5454

55-
static ReactiveSocketServerHandler serverHandler = ReactiveSocketServerHandler.create(setupPayload ->
55+
static ReactiveSocketServerHandler serverHandler = ReactiveSocketServerHandler.create((setupPayload, rs) ->
5656
new RequestHandler() {
5757
@Override
5858
public Publisher<Payload> handleRequestResponse(Payload payload) {

0 commit comments

Comments
 (0)