From e53b0f0de994d0d3e5b1a71d804627df2968672b Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 27 Aug 2019 12:35:54 -0400 Subject: [PATCH] Fixing Kafka streams binder health indicator tests Resolves #731 --- .../KafkaStreamsBinderHealthIndicator.java | 2 +- ...afkaStreamsBinderHealthIndicatorTests.java | 40 ++++++++----------- 2 files changed, 17 insertions(+), 25 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java index 6d36ddd3b..ca3292a97 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java @@ -34,7 +34,7 @@ import org.springframework.boot.actuate.health.Status; * * @author Arnaud Jardiné */ -class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { +public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { private final KafkaStreamsRegistry kafkaStreamsRegistry; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index 60277c5d5..ce54f25fd 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -34,14 +34,15 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; +import org.springframework.boot.actuate.health.CompositeHealthContributor; import org.springframework.boot.actuate.health.Health; -import org.springframework.boot.actuate.health.HealthIndicator; import org.springframework.boot.actuate.health.Status; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.Output; import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsBinderHealthIndicator; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -114,23 +115,14 @@ public class KafkaStreamsBinderHealthIndicatorTests { } } - private static Status getStatusKStream(Map details) { - Health health = (Health) details.get("kstream"); - return health != null ? health.getStatus() : Status.DOWN; - } - - private static boolean waitFor(Map details) { - Health health = (Health) details.get("kstream"); - if (health.getStatus() == Status.UP) { - Map moreDetails = health.getDetails(); - Health kStreamHealth = (Health) moreDetails - .get("kafkaStreamsBinderHealthIndicator"); - String status = (String) kStreamHealth.getDetails().get("threadState"); - return status != null - && (status.equalsIgnoreCase(KafkaStreams.State.REBALANCING.name()) - || status.equalsIgnoreCase("PARTITIONS_REVOKED") - || status.equalsIgnoreCase("PARTITIONS_ASSIGNED") - || status.equalsIgnoreCase( + private static boolean waitFor(Status status, Map details) { + if (status == Status.UP) { + String threadState = (String) details.get("threadState"); + return threadState != null + && (threadState.equalsIgnoreCase(KafkaStreams.State.REBALANCING.name()) + || threadState.equalsIgnoreCase("PARTITIONS_REVOKED") + || threadState.equalsIgnoreCase("PARTITIONS_ASSIGNED") + || threadState.equalsIgnoreCase( KafkaStreams.State.PENDING_SHUTDOWN.name())); } return false; @@ -183,15 +175,15 @@ public class KafkaStreamsBinderHealthIndicatorTests { private static void checkHealth(ConfigurableApplicationContext context, Status expected) throws InterruptedException { - HealthIndicator healthIndicator = context.getBean("bindersHealthIndicator", - HealthIndicator.class); - Health health = healthIndicator.health(); - while (waitFor(health.getDetails())) { + CompositeHealthContributor healthIndicator = context + .getBean("bindersHealthContributor", CompositeHealthContributor.class); + KafkaStreamsBinderHealthIndicator kafkaStreamsBinderHealthIndicator = (KafkaStreamsBinderHealthIndicator) healthIndicator.getContributor("kstream"); + Health health = kafkaStreamsBinderHealthIndicator.health(); + while (waitFor(health.getStatus(), health.getDetails())) { TimeUnit.SECONDS.sleep(2); - health = healthIndicator.health(); + health = kafkaStreamsBinderHealthIndicator.health(); } assertThat(health.getStatus()).isEqualTo(expected); - assertThat(getStatusKStream(health.getDetails())).isEqualTo(expected); } private ConfigurableApplicationContext singleStream() {