DIRTYFIX:subscribe
[demos/kafka/chat] / src / test / java / de / juplo / kafka / chat / backend / TestListener.java
index e413c52..f01e9b5 100644 (file)
@@ -23,7 +23,7 @@ public class TestListener
   static final ParameterizedTypeReference<ServerSentEvent<String>> SSE_TYPE = new ParameterizedTypeReference<>() {};
 
 
-  public Mono<Void> run()
+  public Flux<MessageTo> run()
   {
     return Flux
         .fromArray(chatRooms)
@@ -51,15 +51,10 @@ public class TestListener
                     "Received a message from chat-room {}: {}",
                     chatRoom.getName(),
                     message);
-              })
-              .take(10);
+              });
         })
-        .take(100)
         .takeUntil(message -> !running)
-        .doOnComplete(() -> log.info("TestListener is done"))
-        .parallel(chatRooms.length)
-        .runOn(Schedulers.parallel())
-        .then();
+        .doOnComplete(() -> log.info("TestListener is done"));
   }
 
   Flux<ServerSentEvent<String>> receiveMessages(ChatRoomInfoTo chatRoom)