Replace of ReplayProcessor in RSocket tests
See gh-25085
This commit is contained in:
@@ -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<PayloadLeakInfo> payloads = ReplayProcessor.create();
|
||||
private Sinks.StandaloneFluxSink<PayloadLeakInfo> payloads = Sinks.replayAll();
|
||||
|
||||
PayloadSavingDecorator(RSocket delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
ReplayProcessor<PayloadLeakInfo> getPayloads() {
|
||||
Sinks.StandaloneFluxSink<PayloadLeakInfo> 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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String> fireForgetPayloads = ReplayProcessor.create();
|
||||
final Sinks.StandaloneFluxSink<String> fireForgetPayloads = Sinks.replayAll();
|
||||
|
||||
final ReplayProcessor<String> metadataPushPayloads = ReplayProcessor.create();
|
||||
final Sinks.StandaloneFluxSink<String> 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
|
||||
|
||||
@@ -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<String> fireForgetPayloads = ReplayProcessor.create();
|
||||
final Sinks.StandaloneFluxSink<String> fireForgetPayloads = Sinks.replayAll();
|
||||
|
||||
@MessageMapping("receive")
|
||||
void receive(String payload) {
|
||||
this.fireForgetPayloads.onNext(payload);
|
||||
this.fireForgetPayloads.next(payload);
|
||||
}
|
||||
|
||||
@MessageMapping("echo")
|
||||
|
||||
@@ -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<String>()
|
||||
val fireForgetPayloads = Sinks.replayAll<String>()
|
||||
|
||||
@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")
|
||||
|
||||
Reference in New Issue
Block a user