diff --git a/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketBufferLeakTests.java b/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketBufferLeakTests.java index 1906ba00b6..a7d6f7a8fb 100644 --- a/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketBufferLeakTests.java +++ b/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketBufferLeakTests.java @@ -42,7 +42,7 @@ import org.junit.jupiter.api.TestInstance.Lifecycle; import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import reactor.core.publisher.ReplayProcessor; +import reactor.core.publisher.Sinks; import reactor.test.StepVerifier; import org.springframework.context.annotation.AnnotationConfigApplicationContext; @@ -246,8 +246,8 @@ class RSocketBufferLeakTests { void checkForLeaks() { this.rsockets.stream().map(PayloadSavingDecorator::getPayloads) .forEach(payloadInfoProcessor -> { - payloadInfoProcessor.onComplete(); - payloadInfoProcessor + payloadInfoProcessor.complete(); + payloadInfoProcessor.asFlux() .doOnNext(this::checkForLeak) .blockLast(); }); @@ -291,18 +291,18 @@ class RSocketBufferLeakTests { private final RSocket delegate; - private ReplayProcessor payloads = ReplayProcessor.create(); + private Sinks.StandaloneFluxSink payloads = Sinks.replayAll(); PayloadSavingDecorator(RSocket delegate) { this.delegate = delegate; } - ReplayProcessor getPayloads() { + Sinks.StandaloneFluxSink getPayloads() { return this.payloads; } void reset() { - this.payloads = ReplayProcessor.create(); + this.payloads = Sinks.replayAll(); } @Override @@ -328,7 +328,7 @@ class RSocketBufferLeakTests { } private io.rsocket.Payload addPayload(io.rsocket.Payload payload) { - this.payloads.onNext(new PayloadLeakInfo(payload)); + this.payloads.next(new PayloadLeakInfo(payload)); return payload; } diff --git a/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketClientToServerIntegrationTests.java b/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketClientToServerIntegrationTests.java index a4426c6f11..ceaaf839f5 100644 --- a/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketClientToServerIntegrationTests.java +++ b/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketClientToServerIntegrationTests.java @@ -34,7 +34,7 @@ import org.junit.jupiter.api.Test; import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import reactor.core.publisher.ReplayProcessor; +import reactor.core.publisher.Sinks; import reactor.test.StepVerifier; import org.springframework.context.annotation.AnnotationConfigApplicationContext; @@ -108,7 +108,7 @@ public class RSocketClientToServerIntegrationTests { .concatMap(i -> requester.route("receive").data("Hello " + i).send()) .blockLast(); - StepVerifier.create(context.getBean(ServerController.class).fireForgetPayloads) + StepVerifier.create(context.getBean(ServerController.class).fireForgetPayloads.asFlux()) .expectNext("Hello 1") .expectNext("Hello 2") .expectNext("Hello 3") @@ -171,7 +171,7 @@ public class RSocketClientToServerIntegrationTests { .concatMap(s -> requester.route("foo-updates").metadata(s, FOO_MIME_TYPE).sendMetadata()) .blockLast(); - StepVerifier.create(context.getBean(ServerController.class).metadataPushPayloads) + StepVerifier.create(context.getBean(ServerController.class).metadataPushPayloads.asFlux()) .expectNext("bar") .expectNext("baz") .thenAwait(Duration.ofMillis(50)) @@ -225,14 +225,14 @@ public class RSocketClientToServerIntegrationTests { @Controller static class ServerController { - final ReplayProcessor fireForgetPayloads = ReplayProcessor.create(); + final Sinks.StandaloneFluxSink fireForgetPayloads = Sinks.replayAll(); - final ReplayProcessor metadataPushPayloads = ReplayProcessor.create(); + final Sinks.StandaloneFluxSink metadataPushPayloads = Sinks.replayAll(); @MessageMapping("receive") void receive(String payload) { - this.fireForgetPayloads.onNext(payload); + this.fireForgetPayloads.next(payload); } @MessageMapping("echo") @@ -274,7 +274,7 @@ public class RSocketClientToServerIntegrationTests { @ConnectMapping("foo-updates") public void handleMetadata(@Header("foo") String foo) { - this.metadataPushPayloads.onNext(foo); + this.metadataPushPayloads.next(foo); } @MessageExceptionHandler diff --git a/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketServerToClientIntegrationTests.java b/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketServerToClientIntegrationTests.java index c74058573a..d4e45439ae 100644 --- a/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketServerToClientIntegrationTests.java +++ b/spring-messaging/src/test/java/org/springframework/messaging/rsocket/RSocketServerToClientIntegrationTests.java @@ -29,7 +29,7 @@ import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.MonoProcessor; -import reactor.core.publisher.ReplayProcessor; +import reactor.core.publisher.Sinks; import reactor.core.scheduler.Schedulers; import reactor.test.StepVerifier; @@ -212,11 +212,11 @@ public class RSocketServerToClientIntegrationTests { private static class ClientHandler { - final ReplayProcessor fireForgetPayloads = ReplayProcessor.create(); + final Sinks.StandaloneFluxSink fireForgetPayloads = Sinks.replayAll(); @MessageMapping("receive") void receive(String payload) { - this.fireForgetPayloads.onNext(payload); + this.fireForgetPayloads.next(payload); } @MessageMapping("echo") diff --git a/spring-messaging/src/test/kotlin/org/springframework/messaging/rsocket/RSocketClientToServerCoroutinesIntegrationTests.kt b/spring-messaging/src/test/kotlin/org/springframework/messaging/rsocket/RSocketClientToServerCoroutinesIntegrationTests.kt index 23dff780b6..cdcea8e46c 100644 --- a/spring-messaging/src/test/kotlin/org/springframework/messaging/rsocket/RSocketClientToServerCoroutinesIntegrationTests.kt +++ b/spring-messaging/src/test/kotlin/org/springframework/messaging/rsocket/RSocketClientToServerCoroutinesIntegrationTests.kt @@ -39,7 +39,7 @@ import org.springframework.messaging.handler.annotation.MessageMapping import org.springframework.messaging.rsocket.annotation.support.RSocketMessageHandler import org.springframework.stereotype.Controller import reactor.core.publisher.Flux -import reactor.core.publisher.ReplayProcessor +import reactor.core.publisher.Sinks import reactor.test.StepVerifier import java.time.Duration @@ -56,7 +56,7 @@ class RSocketClientToServerCoroutinesIntegrationTests { Flux.range(1, 3) .concatMap { requester.route("receive").data("Hello $it").send() } .blockLast() - StepVerifier.create(context.getBean(ServerController::class.java).fireForgetPayloads) + StepVerifier.create(context.getBean(ServerController::class.java).fireForgetPayloads.asFlux()) .expectNext("Hello 1") .expectNext("Hello 2") .expectNext("Hello 3") @@ -70,7 +70,7 @@ class RSocketClientToServerCoroutinesIntegrationTests { Flux.range(1, 3) .concatMap { i: Int -> requester.route("receive-async").data("Hello $i").send() } .blockLast() - StepVerifier.create(context.getBean(ServerController::class.java).fireForgetPayloads) + StepVerifier.create(context.getBean(ServerController::class.java).fireForgetPayloads.asFlux()) .expectNext("Hello 1") .expectNext("Hello 2") .expectNext("Hello 3") @@ -145,17 +145,17 @@ class RSocketClientToServerCoroutinesIntegrationTests { @Controller class ServerController { - val fireForgetPayloads = ReplayProcessor.create() + val fireForgetPayloads = Sinks.replayAll() @MessageMapping("receive") fun receive(payload: String) { - fireForgetPayloads.onNext(payload) + fireForgetPayloads.next(payload) } @MessageMapping("receive-async") suspend fun receiveAsync(payload: String) { delay(10) - fireForgetPayloads.onNext(payload) + fireForgetPayloads.next(payload) } @MessageMapping("echo-async")