Skip to content

Commit 89e08fb

Browse files
steveguryNiteshKant
authored andcommitted
ReactiveSocket: Remove startAndWait method (#143)
* ReactiveSocket: Remove `startAndWait` method ***Problem*** ReactiveSocket shouldn't have a public API encouraging usage of blocking code. ***Solution*** Remove the method, replace it by `start` where it was easy and by `Unsafe.startAndWait` elsewhere.
1 parent 75a126d commit 89e08fb

13 files changed

Lines changed: 37 additions & 55 deletions

File tree

reactivesocket-core/src/main/java/io/reactivesocket/ReactiveSocket.java

Lines changed: 0 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -67,36 +67,6 @@ public interface ReactiveSocket {
6767
*/
6868
void start(Completable c);
6969

70-
/**
71-
* Start and block the current thread until startup is finished.
72-
*
73-
* @throws RuntimeException
74-
* of InterruptedException
75-
*/
76-
default void startAndWait() {
77-
CountDownLatch latch = new CountDownLatch(1);
78-
AtomicReference<Throwable> err = new AtomicReference<>();
79-
start(new Completable() {
80-
@Override
81-
public void success() {
82-
latch.countDown();
83-
}
84-
85-
@Override
86-
public void error(Throwable e) {
87-
latch.countDown();
88-
}
89-
});
90-
try {
91-
latch.await();
92-
} catch (InterruptedException e) {
93-
throw new RuntimeException(e);
94-
}
95-
if (err.get() != null) {
96-
throw new RuntimeException(err.get());
97-
}
98-
}
99-
10070
/**
10171
* Invoked when Requester is ready. Non-null exception if error. Null if success.
10272
*

reactivesocket-core/src/main/java/io/reactivesocket/util/Unsafe.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,6 @@ public void error(Throwable e) {
2626
};
2727
rsc.start(completable);
2828
latch.await();
29-
// awaitAvailability(rsc);
3029

3130
return rsc;
3231
}

reactivesocket-core/src/test/java/io/reactivesocket/TestTransportRequestN.java

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

1818
import io.reactivesocket.internal.Publishers;
1919
import io.reactivesocket.lease.FairLeaseGovernor;
20+
import io.reactivesocket.util.Unsafe;
2021
import io.reactivex.subscribers.TestSubscriber;
2122
import org.junit.After;
2223
import org.junit.Ignore;
@@ -225,8 +226,8 @@ public Publisher<Void> handleMetadataPush(Payload payload) {
225226
err -> err.printStackTrace());
226227

227228
// start both the server and client and monitor for errors
228-
socketServer.startAndWait();
229-
socketClient.startAndWait();
229+
Unsafe.startAndWait(socketServer);
230+
Unsafe.startAndWait(socketClient);
230231
}
231232

232233
@After

reactivesocket-stats-servo/src/main/java/io/reactivesocket/loadbalancer/servo/AvailabilityMetricReactiveSocket.java

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -78,11 +78,6 @@ public void start(Completable c) {
7878
child.start(c);
7979
}
8080

81-
@Override
82-
public void startAndWait() {
83-
child.startAndWait();
84-
}
85-
8681
@Override
8782
public void onRequestReady(Consumer<Throwable> c) {
8883
child.onRequestReady(c);

reactivesocket-transport-aeron/src/examples/java/io/reactivesocket/aeron/example/fireandforget/Fire.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import io.reactivesocket.aeron.client.AeronClientDuplexConnection;
2424
import io.reactivesocket.aeron.client.AeronClientDuplexConnectionFactory;
2525
import io.reactivesocket.aeron.client.FrameHolder;
26+
import io.reactivesocket.util.Unsafe;
2627
import org.HdrHistogram.Recorder;
2728
import org.reactivestreams.Publisher;
2829
import org.reactivestreams.Subscription;
@@ -62,8 +63,9 @@ public static void main(String... args) throws Exception {
6263
AeronClientDuplexConnection connection = RxReactiveStreams.toObservable(udpConnection).toBlocking().single();
6364
System.out.println("Created duplex connection");
6465

65-
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS));
66-
reactiveSocket.startAndWait();
66+
ConnectionSetupPayload setupPayload = ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS);
67+
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, setupPayload);
68+
Unsafe.startAndWait(reactiveSocket);
6769

6870
CountDownLatch latch = new CountDownLatch(Integer.MAX_VALUE);
6971

reactivesocket-transport-aeron/src/examples/java/io/reactivesocket/aeron/example/requestreply/Ping.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import io.reactivesocket.aeron.client.AeronClientDuplexConnection;
2323
import io.reactivesocket.aeron.client.AeronClientDuplexConnectionFactory;
2424
import io.reactivesocket.aeron.client.FrameHolder;
25+
import io.reactivesocket.util.Unsafe;
2526
import org.HdrHistogram.Recorder;
2627
import org.reactivestreams.Publisher;
2728
import rx.Observable;
@@ -64,8 +65,9 @@ public static void main(String... args) throws Exception {
6465
AeronClientDuplexConnection connection = RxReactiveStreams.toObservable(udpConnection).toBlocking().single();
6566
System.out.println("Created duplex connection");
6667

67-
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS));
68-
reactiveSocket.startAndWait();
68+
ConnectionSetupPayload setupPayload = ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS);
69+
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, setupPayload);
70+
Unsafe.startAndWait(reactiveSocket);
6971

7072
CountDownLatch latch = new CountDownLatch(Integer.MAX_VALUE);
7173

reactivesocket-transport-aeron/src/main/java/io/reactivesocket/aeron/server/ReactiveSocketAeronServer.java

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import io.reactivesocket.aeron.internal.Loggable;
3030
import io.reactivesocket.aeron.internal.MessageType;
3131
import io.reactivesocket.rx.Observer;
32+
import io.reactivesocket.util.Unsafe;
3233
import org.agrona.BitUtil;
3334
import org.agrona.DirectBuffer;
3435
import org.agrona.concurrent.UnsafeBuffer;
@@ -174,7 +175,11 @@ public void accept(Throwable throwable) {
174175

175176
sockets.put(sessionId, socket);
176177

177-
socket.startAndWait();
178+
try {
179+
Unsafe.startAndWait(socket);
180+
} catch (InterruptedException e) {
181+
e.printStackTrace();
182+
}
178183
} else {
179184
debug("Unsupported stream id {}", streamId);
180185
}

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

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import io.reactivesocket.aeron.server.ReactiveSocketAeronServer;
2727
import io.reactivesocket.exceptions.SetupException;
2828
import io.reactivesocket.test.TestUtil;
29+
import io.reactivesocket.util.Unsafe;
2930
import org.junit.Assert;
3031
import org.junit.BeforeClass;
3132
import org.junit.Ignore;
@@ -151,8 +152,9 @@ public Publisher<Payload> apply(Payload payload) {
151152
AeronClientDuplexConnection connection = RxReactiveStreams.toObservable(udpConnection).toBlocking().single();
152153
System.out.println("Created duplex connection");
153154

154-
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS));
155-
reactiveSocket.startAndWait();
155+
ConnectionSetupPayload setupPayload = ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS);
156+
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, setupPayload);
157+
Unsafe.startAndWait(reactiveSocket);
156158

157159
CountDownLatch latch = new CountDownLatch(count);
158160

@@ -225,8 +227,9 @@ public void requestStreamN(int count) throws Exception {
225227
AeronClientDuplexConnection connection = RxReactiveStreams.toObservable(udpConnection).toBlocking().single();
226228
System.out.println("Created duplex connection");
227229

228-
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS));
229-
reactiveSocket.startAndWait();
230+
ConnectionSetupPayload setupPayload = ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS);
231+
ReactiveSocket reactiveSocket = DefaultReactiveSocket.fromClientConnection(connection, setupPayload);
232+
Unsafe.startAndWait(reactiveSocket);
230233

231234
CountDownLatch latch = new CountDownLatch(count);
232235
Payload payload = TestUtil.utf8EncodedPayload("client_request", "client_metadata");
@@ -323,8 +326,9 @@ public Publisher<Void> handleMetadataPush(Payload payload) {
323326
AeronClientDuplexConnection connection = RxReactiveStreams.toObservable(udpConnection).toBlocking().single();
324327
System.out.println("Created duplex connection => " + j);
325328

326-
ReactiveSocket client = DefaultReactiveSocket.fromClientConnection(connection, ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS));
327-
client.startAndWait();
329+
ConnectionSetupPayload setupPayload = ConnectionSetupPayload.create("UTF-8", "UTF-8", ConnectionSetupPayload.NO_FLAGS);
330+
ReactiveSocket client = DefaultReactiveSocket.fromClientConnection(connection, setupPayload);
331+
Unsafe.startAndWait(client);
328332

329333
Observable
330334
.range(1, 10)

reactivesocket-transport-local/src/main/java/io/reactivesocket/local/LocalClientReactiveSocketConnector.java

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

1818
import io.reactivesocket.*;
1919
import io.reactivesocket.internal.rx.EmptySubscription;
20+
import io.reactivesocket.util.Unsafe;
2021
import org.reactivestreams.Publisher;
2122

2223
public class LocalClientReactiveSocketConnector implements ReactiveSocketConnector<LocalClientReactiveSocketConnector.Config> {
@@ -35,7 +36,7 @@ public Publisher<ReactiveSocket> connect(Config config) {
3536
ReactiveSocket reactiveSocket = DefaultReactiveSocket
3637
.fromClientConnection(clientConnection, ConnectionSetupPayload.create(config.getMetadataMimeType(), config.getDataMimeType()));
3738

38-
reactiveSocket.startAndWait();
39+
Unsafe.startAndWait(reactiveSocket);
3940

4041
s.onNext(reactiveSocket);
4142
s.onComplete();

reactivesocket-transport-local/src/main/java/io/reactivesocket/local/LocalServerReactiveSocketConnector.java

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

1818
import io.reactivesocket.*;
1919
import io.reactivesocket.internal.rx.EmptySubscription;
20+
import io.reactivesocket.util.Unsafe;
2021
import org.reactivestreams.Publisher;
2122

2223
public class LocalServerReactiveSocketConnector implements ReactiveSocketConnector<LocalServerReactiveSocketConnector.Config> {
@@ -35,7 +36,7 @@ public Publisher<ReactiveSocket> connect(Config config) {
3536
ReactiveSocket reactiveSocket = DefaultReactiveSocket
3637
.fromServerConnection(clientConnection, config.getConnectionSetupHandler());
3738

38-
reactiveSocket.startAndWait();
39+
Unsafe.startAndWait(reactiveSocket);
3940
s.onNext(reactiveSocket);
4041
s.onComplete();
4142
} catch (Throwable t) {

0 commit comments

Comments
 (0)