Use Declarable to initialize the bus queues and exchanges

Also tidies up conditionals a bit and adds a new @ConfigurationProperties.

Fixes gh-14
This commit is contained in:
Dave Syer
2015-03-16 14:39:14 +00:00
parent f5360dd419
commit 42d189ada2
9 changed files with 86 additions and 43 deletions

View File

@@ -20,14 +20,17 @@
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-integration</artifactId>
<artifactId>spring-boot-starter-web</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.integration</groupId>
@@ -36,6 +39,7 @@
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-amqp</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.integration</groupId>
@@ -48,14 +52,17 @@
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-spring-service-connector</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-localconfig-connector</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-cloudfoundry-connector</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>

View File

@@ -39,8 +39,6 @@ import org.springframework.messaging.SubscribableChannel;
@ConditionalOnBusEnabled
public class BusAutoConfiguration {
public static final String SPRING_CLOUD_BUS_ENABLED = "spring.cloud.bus.enabled";
@Autowired
private ConfigurableApplicationContext context;

View File

@@ -1,17 +1,19 @@
package org.springframework.cloud.bus;
import org.springframework.context.annotation.Conditional;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
/**
* @author Spencer Gibb
*/
@Conditional(OnBusEnabled.class)
@ConditionalOnProperty(value = ConditionalOnBusEnabled.SPRING_CLOUD_BUS_ENABLED, matchIfMissing = true)
@Retention(RetentionPolicy.RUNTIME)
@Target({ ElementType.TYPE, ElementType.METHOD })
public @interface ConditionalOnBusEnabled {
public static String SPRING_CLOUD_BUS_ENABLED = "spring.cloud.bus.amqp.enabled";
}

View File

@@ -1,26 +0,0 @@
package org.springframework.cloud.bus;
import org.springframework.boot.autoconfigure.condition.ConditionOutcome;
import org.springframework.boot.autoconfigure.condition.SpringBootCondition;
import org.springframework.boot.bind.RelaxedPropertyResolver;
import org.springframework.context.annotation.ConditionContext;
import org.springframework.core.type.AnnotatedTypeMetadata;
import static org.springframework.cloud.bus.BusAutoConfiguration.SPRING_CLOUD_BUS_ENABLED;
/**
* Match if spring.cloud.bus.enabled is missing or not false
* @author Spencer Gibb
*/
class OnBusEnabled extends SpringBootCondition {
@Override
public ConditionOutcome getMatchOutcome(ConditionContext context, AnnotatedTypeMetadata metadata) {
RelaxedPropertyResolver resolver = new RelaxedPropertyResolver(context.getEnvironment());
String enabled = resolver.getProperty(SPRING_CLOUD_BUS_ENABLED);
if (!"false".equalsIgnoreCase(enabled)) {
return ConditionOutcome.match();
}
return ConditionOutcome.noMatch(SPRING_CLOUD_BUS_ENABLED + " is " + enabled);
}
}

View File

@@ -7,12 +7,14 @@ import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.FanoutExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.cloud.bus.ConditionalOnBusEnabled;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -43,9 +45,11 @@ import com.fasterxml.jackson.databind.ObjectMapper;
@ConditionalOnBusEnabled
@ConditionalOnClass({ AmqpTemplate.class, RabbitTemplate.class })
@ConditionalOnProperty(value = "spring.cloud.bus.amqp.enabled", matchIfMissing = true)
@EnableConfigurationProperties(AmqpBusProperties.class)
public class AmqpBusAutoConfiguration {
public static final String SPRING_CLOUD_BUS = "spring.cloud.bus";
@Autowired
private AmqpBusProperties bus;
@Autowired(required = false)
@BusConnectionFactory
@@ -56,6 +60,8 @@ public class AmqpBusAutoConfiguration {
private RabbitTemplate amqpTemplate;
private RabbitAdmin amqpAdmin;
@Autowired(required = false)
private ObjectMapper objectMapper;
@@ -65,6 +71,11 @@ public class AmqpBusAutoConfiguration {
Jackson2JsonMessageConverter converter = messageConverter();
amqpTemplate.setMessageConverter(converter);
this.amqpTemplate = amqpTemplate;
this.amqpAdmin = new RabbitAdmin(connectionFactory());
cloudBusExchange().setAdminsThatShouldDeclare(this.amqpAdmin);
localCloudBusQueueBinding().setAdminsThatShouldDeclare(this.amqpAdmin);
localCloudBusQueue().setAdminsThatShouldDeclare(this.amqpAdmin);
this.amqpAdmin.afterPropertiesSet();
}
return amqpTemplate;
}
@@ -72,7 +83,7 @@ public class AmqpBusAutoConfiguration {
@Bean
protected FanoutExchange cloudBusExchange() {
// TODO: change to TopicExchange?
FanoutExchange exchange = new FanoutExchange(SPRING_CLOUD_BUS);
FanoutExchange exchange = new FanoutExchange(bus.getExchange());
return exchange;
}
@@ -92,7 +103,7 @@ public class AmqpBusAutoConfiguration {
return IntegrationFlows
.from(cloudBusOutboundChannel)
.handle(Amqp.outboundAdapter(amqpTemplate()).exchangeName(
SPRING_CLOUD_BUS)).get();
bus.getExchange())).get();
}
@Bean

View File

@@ -0,0 +1,37 @@
/*
* Copyright 2014-2015 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.cloud.bus.amqp;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
/**
* @author Dave Syer
*
*/
@ConfigurationProperties("spring.cloud.bus.amqp")
@Data
public class AmqpBusProperties {
public static final String SPRING_CLOUD_BUS = "spring.cloud.bus";
private boolean enabled;
private String exchange = SPRING_CLOUD_BUS;
}

View File

@@ -1,5 +1,8 @@
package org.springframework.cloud.bus;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import org.junit.After;
import org.junit.Rule;
import org.junit.Test;
@@ -9,10 +12,6 @@ import org.springframework.context.annotation.AnnotationConfigApplicationContext
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.springframework.cloud.bus.BusAutoConfiguration.SPRING_CLOUD_BUS_ENABLED;
/**
* @author Spencer Gibb
*/
@@ -32,7 +31,7 @@ public class ConditionalOnBusEnabledTests {
@Test
public void busEnabledTrue() {
load(MyBusEnabledConfig.class, SPRING_CLOUD_BUS_ENABLED+":true");
load(MyBusEnabledConfig.class, ConditionalOnBusEnabled.SPRING_CLOUD_BUS_ENABLED+":true");
assertTrue("missing bean from @ConditionalOnBusEnabled config",
this.context.containsBean("foo"));
}
@@ -46,7 +45,7 @@ public class ConditionalOnBusEnabledTests {
@Test
public void busDisabled() {
load(MyBusEnabledConfig.class, SPRING_CLOUD_BUS_ENABLED+":false");
load(MyBusEnabledConfig.class, ConditionalOnBusEnabled.SPRING_CLOUD_BUS_ENABLED+":false");
assertFalse("bean exists from disabled @ConditionalOnBusEnabled config",
this.context.containsBean("foo"));
}

View File

@@ -1,10 +1,11 @@
package org.springframework.cloud.bus.amqp;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.springframework.cloud.bus.BusAutoConfiguration.SPRING_CLOUD_BUS_ENABLED;
import org.junit.Test;
import org.springframework.amqp.core.Declarable;
import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.boot.autoconfigure.PropertyPlaceholderAutoConfiguration;
@@ -12,6 +13,7 @@ import org.springframework.boot.autoconfigure.amqp.RabbitAutoConfiguration;
import org.springframework.boot.autoconfigure.amqp.RabbitProperties;
import org.springframework.boot.test.EnvironmentTestUtils;
import org.springframework.cloud.bus.BusAutoConfiguration;
import org.springframework.cloud.bus.ConditionalOnBusEnabled;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -38,6 +40,18 @@ public class AmqpBusAutoConfigurationTests {
context.close();
}
@Test
public void qualifiedConnectionFactoryAdmin() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(
getConfigClasses(QualifiedConnectionFactory.class));
assertTrue(context.containsBean("amqpAdmin"));
Object admin = context.getBean("amqpAdmin");
Declarable declarable = context.getBean("cloudBusExchange", Declarable.class);
assertEquals(1, declarable.getDeclaringAdmins().size());
assertFalse(declarable.getDeclaringAdmins().contains(admin));
context.close();
}
@Test
public void unqualifiedConnectionFactory() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(
@@ -57,7 +71,7 @@ public class AmqpBusAutoConfigurationTests {
@Test
public void notStartedIfBusDisabled() {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext();
EnvironmentTestUtils.addEnvironment(context, SPRING_CLOUD_BUS_ENABLED+":false");
EnvironmentTestUtils.addEnvironment(context, ConditionalOnBusEnabled.SPRING_CLOUD_BUS_ENABLED+":false");
context.register(getConfigClasses());
context.refresh();

View File

@@ -0,0 +1 @@
spring.main.web-environment: false