Code cleanup
This commit is contained in:
parent
24b4387b08
commit
171f07ccec
@ -37,13 +37,11 @@ public class MyRSocketClient implements RSocketChannelManager {
|
|||||||
private final Empty<Void> disposeRequest = Sinks.empty();
|
private final Empty<Void> disposeRequest = Sinks.empty();
|
||||||
|
|
||||||
public MyRSocketClient(HostAndPort baseHost) {
|
public MyRSocketClient(HostAndPort baseHost) {
|
||||||
RetryBackoffSpec retryStrategy = Retry.backoff(Long.MAX_VALUE, Duration.ofSeconds(1)).maxBackoff(Duration.ofSeconds(16)).jitter(1.0);
|
|
||||||
var transport = TcpClientTransport.create(baseHost.getHost(), baseHost.getPort());
|
var transport = TcpClientTransport.create(baseHost.getHost(), baseHost.getPort());
|
||||||
|
|
||||||
this.nextClient = RSocketConnector.create()
|
this.nextClient = RSocketConnector.create()
|
||||||
.setupPayload(DefaultPayload.create("client", "setup-info"))
|
.setupPayload(DefaultPayload.create("client", "setup-info"))
|
||||||
.payloadDecoder(PayloadDecoder.ZERO_COPY)
|
.payloadDecoder(PayloadDecoder.ZERO_COPY)
|
||||||
//.reconnect(retryStrategy)
|
|
||||||
.connect(transport)
|
.connect(transport)
|
||||||
.doOnNext(lastClient::set)
|
.doOnNext(lastClient::set)
|
||||||
.cacheInvalidateIf(RSocket::isDisposed);
|
.cacheInvalidateIf(RSocket::isDisposed);
|
||||||
|
Loading…
Reference in New Issue
Block a user