Remove EmitterProcessor#connect (dropped upstream)
This commit is contained in:
@@ -146,7 +146,6 @@ public class ReactiveTypeHandlerTests {
|
|||||||
|
|
||||||
EmitterProcessor<String> emitter = EmitterProcessor.create();
|
EmitterProcessor<String> emitter = EmitterProcessor.create();
|
||||||
testDeferredResultSubscriber(emitter, Flux.class, () -> {
|
testDeferredResultSubscriber(emitter, Flux.class, () -> {
|
||||||
emitter.connect();
|
|
||||||
emitter.onNext("foo");
|
emitter.onNext("foo");
|
||||||
emitter.onNext("bar");
|
emitter.onNext("bar");
|
||||||
emitter.onNext("baz");
|
emitter.onNext("baz");
|
||||||
@@ -233,7 +232,6 @@ public class ReactiveTypeHandlerTests {
|
|||||||
EmitterHandler emitterHandler = new EmitterHandler();
|
EmitterHandler emitterHandler = new EmitterHandler();
|
||||||
sseEmitter.initialize(emitterHandler);
|
sseEmitter.initialize(emitterHandler);
|
||||||
|
|
||||||
processor.connect();
|
|
||||||
processor.onNext("foo");
|
processor.onNext("foo");
|
||||||
processor.onNext("bar");
|
processor.onNext("bar");
|
||||||
processor.onNext("baz");
|
processor.onNext("baz");
|
||||||
@@ -253,7 +251,6 @@ public class ReactiveTypeHandlerTests {
|
|||||||
EmitterHandler emitterHandler = new EmitterHandler();
|
EmitterHandler emitterHandler = new EmitterHandler();
|
||||||
sseEmitter.initialize(emitterHandler);
|
sseEmitter.initialize(emitterHandler);
|
||||||
|
|
||||||
processor.connect();
|
|
||||||
processor.onNext(ServerSentEvent.builder("foo").id("1").build());
|
processor.onNext(ServerSentEvent.builder("foo").id("1").build());
|
||||||
processor.onNext(ServerSentEvent.builder("bar").id("2").build());
|
processor.onNext(ServerSentEvent.builder("bar").id("2").build());
|
||||||
processor.onNext(ServerSentEvent.builder("baz").id("3").build());
|
processor.onNext(ServerSentEvent.builder("baz").id("3").build());
|
||||||
@@ -277,7 +274,6 @@ public class ReactiveTypeHandlerTests {
|
|||||||
ServletServerHttpResponse message = new ServletServerHttpResponse(this.servletResponse);
|
ServletServerHttpResponse message = new ServletServerHttpResponse(this.servletResponse);
|
||||||
emitter.extendResponse(message);
|
emitter.extendResponse(message);
|
||||||
|
|
||||||
processor.connect();
|
|
||||||
processor.onNext("[\"foo\",\"bar\"]");
|
processor.onNext("[\"foo\",\"bar\"]");
|
||||||
processor.onNext("[\"bar\",\"baz\"]");
|
processor.onNext("[\"bar\",\"baz\"]");
|
||||||
processor.onComplete();
|
processor.onComplete();
|
||||||
@@ -295,7 +291,6 @@ public class ReactiveTypeHandlerTests {
|
|||||||
EmitterHandler emitterHandler = new EmitterHandler();
|
EmitterHandler emitterHandler = new EmitterHandler();
|
||||||
emitter.initialize(emitterHandler);
|
emitter.initialize(emitterHandler);
|
||||||
|
|
||||||
processor.connect();
|
|
||||||
processor.onNext("The quick");
|
processor.onNext("The quick");
|
||||||
processor.onNext(" brown fox jumps over ");
|
processor.onNext(" brown fox jumps over ");
|
||||||
processor.onNext("the lazy dog");
|
processor.onNext("the lazy dog");
|
||||||
|
|||||||
Reference in New Issue
Block a user