Skip to content

Commit ddba337

Browse files
committed
throwing a TransportException when a connection is closed
1 parent 5e295bc commit ddba337

8 files changed

Lines changed: 27 additions & 7 deletions

File tree

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-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) {

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,11 +80,13 @@ public ByteBuffer getMetadata() {
8080
.toObservable(
8181
reactiveSocket
8282
.requestResponse(keyPayload))
83+
.doOnError(t -> t.printStackTrace())
8384
.doOnNext(s -> {
8485
long diff = System.nanoTime() - start;
8586
histogram.recordValue(diff);
8687
});
8788
}, 16)
89+
.doOnError(t -> t.printStackTrace())
8890
.subscribe(new Subscriber<Payload>() {
8991
@Override
9092
public void onCompleted() {

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ public static void main(String... args) throws Exception {
4343
r.nextBytes(response);
4444

4545
ReactiveSocketServerHandler serverHandler =
46-
ReactiveSocketServerHandler.create(setupPayload -> new RequestHandler() {
46+
ReactiveSocketServerHandler.create((setupPayload, rs) -> new RequestHandler() {
4747
@Override
4848
public Publisher<Payload> handleRequestResponse(Payload payload) {
4949
return new Publisher<Payload>() {

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

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

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

reactivesocket-netty/src/test/java/io/reactivesocket/netty/websocket/Pong.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ public static void main(String... args) throws Exception {
4646
r.nextBytes(response);
4747

4848
ReactiveSocketServerHandler serverHandler =
49-
ReactiveSocketServerHandler.create(setupPayload -> new RequestHandler() {
49+
ReactiveSocketServerHandler.create((setupPayload, rs) -> new RequestHandler() {
5050
@Override
5151
public Publisher<Payload> handleRequestResponse(Payload payload) {
5252
return new Publisher<Payload>() {

0 commit comments

Comments
 (0)