Skip to content

Commit ce207ba

Browse files
author
Ryland Degnan
authored
Merge pull request #493 from rdegnan/rsocket-factory
Remove nested flatMap in RSocketFactory
2 parents e45327c + 6d98a5e commit ce207ba

1 file changed

Lines changed: 29 additions & 31 deletions

File tree

rsocket-core/src/main/java/io/rsocket/RSocketFactory.java

Lines changed: 29 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -228,28 +228,21 @@ public Mono<RSocket> start() {
228228
ackTimeout,
229229
missedAcks);
230230

231-
Mono<RSocket> wrappedRSocketClient =
232-
Mono.just(rSocketClient).map(plugins::applyClient);
233-
234-
DuplexConnection finalConnection = connection;
235-
return wrappedRSocketClient.flatMap(
236-
wrappedClientRSocket -> {
237-
RSocket unwrappedServerSocket = acceptor.get().apply(wrappedClientRSocket);
238-
239-
Mono<RSocket> wrappedRSocketServer =
240-
Mono.just(unwrappedServerSocket).map(plugins::applyServer);
241-
242-
return wrappedRSocketServer
243-
.doOnNext(
244-
rSocket ->
245-
new RSocketServer(
246-
multiplexer.asServerConnection(),
247-
rSocket,
248-
frameDecoder,
249-
errorConsumer))
250-
.then(finalConnection.sendOne(setupFrame))
251-
.then(wrappedRSocketClient);
252-
});
231+
RSocket wrappedRSocketClient = plugins.applyClient(rSocketClient);
232+
233+
RSocket unwrappedServerSocket = acceptor.get().apply(wrappedRSocketClient);
234+
235+
RSocket wrappedRSocketServer = plugins.applyServer(unwrappedServerSocket);
236+
237+
RSocketServer rSocketServer = new RSocketServer(
238+
multiplexer.asServerConnection(),
239+
wrappedRSocketServer,
240+
frameDecoder,
241+
errorConsumer);
242+
243+
return connection
244+
.sendOne(setupFrame)
245+
.thenReturn(wrappedRSocketClient);
253246
});
254247
}
255248
}
@@ -332,7 +325,7 @@ public Mono<T> start() {
332325
});
333326
}
334327

335-
private Mono<? extends Void> processSetupFrame(
328+
private Mono<Void> processSetupFrame(
336329
ClientServerInputMultiplexer multiplexer, Frame setupFrame) {
337330
int version = Frame.Setup.version(setupFrame);
338331
if (version != SetupFrameFlyweight.CURRENT_VERSION) {
@@ -355,15 +348,20 @@ private Mono<? extends Void> processSetupFrame(
355348
errorConsumer,
356349
StreamIdSupplier.serverSupplier());
357350

358-
Mono<RSocket> wrappedRSocketClient = Mono.just(rSocketClient).map(plugins::applyClient);
351+
RSocket wrappedRSocketClient = plugins.applyClient(rSocketClient);
359352

360-
return wrappedRSocketClient
361-
.flatMap(
362-
sender -> acceptor.get().accept(setupPayload, sender).map(plugins::applyServer))
363-
.map(
364-
handler ->
365-
new RSocketServer(
366-
multiplexer.asClientConnection(), handler, frameDecoder, errorConsumer))
353+
return acceptor
354+
.get()
355+
.accept(setupPayload, wrappedRSocketClient)
356+
.doOnNext(unwrappedServerSocket -> {
357+
RSocket wrappedRSocketServer = plugins.applyServer(unwrappedServerSocket);
358+
359+
RSocketServer rSocketServer = new RSocketServer(
360+
multiplexer.asClientConnection(),
361+
wrappedRSocketServer,
362+
frameDecoder,
363+
errorConsumer);
364+
})
367365
.then();
368366
}
369367
}

0 commit comments

Comments
 (0)