diff --git a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/cloud/netflix/turbine/amqp/TurbineAmqpConfiguration.java b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/cloud/netflix/turbine/amqp/TurbineAmqpConfiguration.java index 868303af..b7b5e81e 100644 --- a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/cloud/netflix/turbine/amqp/TurbineAmqpConfiguration.java +++ b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/cloud/netflix/turbine/amqp/TurbineAmqpConfiguration.java @@ -21,9 +21,12 @@ import io.reactivex.netty.RxNetty; import io.reactivex.netty.protocol.http.server.HttpServer; import io.reactivex.netty.protocol.text.sse.ServerSentEvent; +import java.util.Collections; import java.util.Map; +import java.util.concurrent.TimeUnit; import lombok.extern.apachecommons.CommonsLog; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.SmartLifecycle; @@ -71,6 +74,11 @@ public class TurbineAmqpConfiguration implements SmartLifecycle { .doOnUnsubscribe(() -> log.info("Unsubscribing aggregation.")) .doOnSubscribe(() -> log.info("Starting aggregation")).flatMap(o -> o) .publish().refCount(); + Observable> ping = Observable.timer(1, 10, TimeUnit.SECONDS) + .map(count -> { + return Collections.singletonMap("type", (Object) "Ping"); + }).publish().refCount(); + Observable> output = Observable.merge(publishedStreams, ping); this.turbinePort = this.turbine.getPort(); @@ -83,7 +91,7 @@ public class TurbineAmqpConfiguration implements SmartLifecycle { (request, response) -> { log.info("SSE Request Received"); response.getHeaders().setHeader("Content-Type", "text/event-stream"); - return publishedStreams.doOnUnsubscribe( + return output.doOnUnsubscribe( () -> log.info("Unsubscribing RxNetty server connection")) .flatMap( data -> response.writeAndFlush(new ServerSentEvent(