Added logging and making kafka tests less brittle
This commit is contained in:
@@ -22,6 +22,9 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Map.Entry;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
|
||||
import org.springframework.boot.autoconfigure.AutoConfigureBefore;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
|
||||
@@ -60,11 +63,16 @@ import org.springframework.util.StringUtils;
|
||||
@AutoConfigureBefore(ContractVerifierKafkaConfiguration.class)
|
||||
public class StubRunnerKafkaConfiguration {
|
||||
|
||||
private static final Log log = LogFactory.getLog(StubRunnerKafkaConfiguration.class);
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
@ConditionalOnProperty(name = "stubrunner.kafka.initializer.enabled",
|
||||
havingValue = "true", matchIfMissing = true)
|
||||
KafkaStubMessagesInitializer stubRunnerKafkaStubMessagesInitializer() {
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug("Registering a noop kafka messages initializer");
|
||||
}
|
||||
return (broker, kafkaProperties) -> new HashMap<>();
|
||||
}
|
||||
|
||||
@@ -100,6 +108,9 @@ public class StubRunnerKafkaConfiguration {
|
||||
matchingContracts, beanFactory);
|
||||
StubRunnerKafkaRouter listener = (StubRunnerKafkaRouter) beanFactory
|
||||
.initializeBean(router, flowName);
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug("Initialized kafka router with name [" + flowName + "]");
|
||||
}
|
||||
beanFactory.registerSingleton(flowName, listener);
|
||||
registerContainers(beanFactory, matchingContracts, flowName, listener);
|
||||
}
|
||||
@@ -127,6 +138,9 @@ public class StubRunnerKafkaConfiguration {
|
||||
Object initializedContainer = beanFactory.initializeBean(container,
|
||||
containerName);
|
||||
beanFactory.registerSingleton(containerName, initializedContainer);
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug("Initialized kafka message container with name [" + containerName + "] listening to destination [" + destination + "]");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -16,6 +16,9 @@
|
||||
|
||||
package org.springframework.cloud.contract.verifier.messaging.kafka;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.boot.autoconfigure.AutoConfigureBefore;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
|
||||
@@ -38,13 +41,15 @@ import org.springframework.messaging.Message;
|
||||
*/
|
||||
@Configuration
|
||||
@ConditionalOnClass({ KafkaTemplate.class, EmbeddedKafkaBroker.class })
|
||||
@ConditionalOnProperty(name = "stubrunner.kafka.enabled", havingValue = "true",
|
||||
matchIfMissing = true)
|
||||
@ConditionalOnProperty(name = "stubrunner.kafka.enabled", havingValue = "true", matchIfMissing = true)
|
||||
@AutoConfigureBefore({ ContractVerifierIntegrationConfiguration.class,
|
||||
NoOpContractVerifierAutoConfiguration.class })
|
||||
@ConditionalOnBean(EmbeddedKafkaBroker.class)
|
||||
public class ContractVerifierKafkaConfiguration {
|
||||
|
||||
private static final Log log = LogFactory
|
||||
.getLog(ContractVerifierKafkaConfiguration.class);
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
MessageVerifier<Message<?>> contractVerifierKafkaMessageExchange(
|
||||
@@ -56,6 +61,9 @@ public class ContractVerifierKafkaConfiguration {
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
KafkaStubMessagesInitializer contractVerifierKafkaStubMessagesInitializer() {
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug("Registering contract verifier stub messages initializer");
|
||||
}
|
||||
return new ContractVerifierKafkaStubMessagesInitializer();
|
||||
}
|
||||
|
||||
|
||||
@@ -21,8 +21,10 @@ import java.util.concurrent.TimeUnit
|
||||
|
||||
import groovy.json.JsonOutput
|
||||
import groovy.json.JsonSlurper
|
||||
import groovy.util.logging.Commons
|
||||
import spock.lang.IgnoreIf
|
||||
import spock.lang.Specification
|
||||
import spock.lang.Stepwise
|
||||
import spock.util.concurrent.PollingConditions
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
@@ -54,6 +56,8 @@ import org.springframework.test.context.ContextConfiguration
|
||||
@AutoConfigureStubRunner
|
||||
@IgnoreIf({ os.windows })
|
||||
@EmbeddedKafka(topics = ["input", "output", "delete"])
|
||||
@Commons
|
||||
@Stepwise
|
||||
class KafkaStubRunnerSpec extends Specification {
|
||||
|
||||
@Autowired
|
||||
@@ -79,15 +83,19 @@ class KafkaStubRunnerSpec extends Specification {
|
||||
|
||||
def 'should download the stub and register a route for it'() {
|
||||
when:
|
||||
log.info("Sending the message")
|
||||
// tag::client_send[]
|
||||
Message message = MessageBuilder.createMessage(new BookReturned('foo'), new MessageHeaders([sample: "header",]))
|
||||
kafkaTemplate.setDefaultTopic('input')
|
||||
kafkaTemplate.send(message)
|
||||
// end::client_send[]
|
||||
log.info("Message sent")
|
||||
then:
|
||||
log.info("Receiving the message")
|
||||
// tag::client_receive[]
|
||||
Message receivedMessage = receiveFromOutput()
|
||||
// end::client_receive[]
|
||||
log.info("Message received [" + receivedMessage + "]")
|
||||
and:
|
||||
await.eventually {
|
||||
// tag::client_receive_message[]
|
||||
|
||||
Reference in New Issue
Block a user