From abab014f393abb5e2533ecdd88580f6353289acd Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 12 Jun 2019 21:47:07 +0200 Subject: [PATCH] GH-1707 Removed all references to aggregator builder - removed deprecations connected to InterceptableChannel - minor polishing to get rid of IDE warnings Resolves #1707 --- .../aggregate/AggregateApplicationTests.java | 65 --- .../aggregate/processor/TestProcessor.java | 44 -- .../config/aggregate/source/TestSource.java | 54 -- .../aggregate/bean/AggregateWithBeanTest.java | 110 ---- .../aggregate/main/AggregateWithMainTest.java | 105 ---- .../aggregate/AggregateApplication.java | 40 -- .../AggregateApplicationBuilder.java | 531 ------------------ .../aggregate/AggregateApplicationUtils.java | 91 --- .../SharedBindingTargetRegistry.java | 58 -- .../binder/AbstractMessageChannelBinder.java | 5 +- .../stream/binding/BindableProxyFactory.java | 45 +- .../config/SmartPayloadArgumentResolver.java | 4 +- .../stream/aggregation/AggregationTest.java | 458 --------------- .../BinderAwareChannelResolverTests.java | 8 +- .../stream/binding/BindingLifecycleTests.java | 1 + .../CustomPartitionedProducerTest.java | 8 +- .../stream/function/RoutingFunctionTests.java | 1 + 17 files changed, 23 insertions(+), 1605 deletions(-) delete mode 100644 spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/AggregateApplicationTests.java delete mode 100644 spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/processor/TestProcessor.java delete mode 100644 spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/source/TestSource.java delete mode 100644 spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/bean/AggregateWithBeanTest.java delete mode 100644 spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/main/AggregateWithMainTest.java delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationBuilder.java delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationUtils.java delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/SharedBindingTargetRegistry.java delete mode 100644 spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/AggregationTest.java diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/AggregateApplicationTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/AggregateApplicationTests.java deleted file mode 100644 index 66a368528..000000000 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/AggregateApplicationTests.java +++ /dev/null @@ -1,65 +0,0 @@ -/* - * Copyright 2017-2018 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 - * - * https://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.stream.config.aggregate; - -import java.util.concurrent.TimeUnit; - -import org.junit.Test; - -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.cloud.stream.aggregate.AggregateApplicationBuilder; -import org.springframework.cloud.stream.binder.BinderFactory; -import org.springframework.cloud.stream.config.aggregate.processor.TestProcessor; -import org.springframework.cloud.stream.config.aggregate.source.TestSource; -import org.springframework.cloud.stream.test.binder.TestSupportBinder; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Configuration; -import org.springframework.messaging.Message; -import org.springframework.messaging.MessageChannel; - -import static org.assertj.core.api.Assertions.assertThat; - -/** - * @author Ilayaperumal Gopinathan - * @author Oleg Zhurakousky - */ -public class AggregateApplicationTests { - - @Test - @SuppressWarnings("unchecked") - public void testAggregateApplication() throws Exception { - ConfigurableApplicationContext context = new AggregateApplicationBuilder( - AggregateApplicationTestConfig.class).web(false).from(TestSource.class) - .to(TestProcessor.class).run(); - TestSupportBinder testSupportBinder = (TestSupportBinder) context - .getBean(BinderFactory.class).getBinder(null, MessageChannel.class); - MessageChannel processorOutput = testSupportBinder.getChannelForName("output"); - Message received = (Message) (testSupportBinder.messageCollector() - .forChannel(processorOutput).poll(5, TimeUnit.SECONDS)); - assertThat(received).isNotNull(); - assertThat(received.getPayload().endsWith("processed")).isTrue(); - - context.close(); - } - - @Configuration - @EnableAutoConfiguration - static class AggregateApplicationTestConfig { - - } - -} diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/processor/TestProcessor.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/processor/TestProcessor.java deleted file mode 100644 index bd1867c7d..000000000 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/processor/TestProcessor.java +++ /dev/null @@ -1,44 +0,0 @@ -/* - * Copyright 2017-2019 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 - * - * https://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.stream.config.aggregate.processor; - -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.messaging.Processor; -import org.springframework.context.annotation.Configuration; -import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.handler.annotation.SendTo; -import org.springframework.messaging.support.MessageBuilder; - -/** - * @author Ilayaperumal Gopinathan - */ -@EnableBinding(Processor.class) -@EnableAutoConfiguration -@Configuration -public class TestProcessor { - - @StreamListener(Processor.INPUT) - @SendTo(Processor.OUTPUT) - public Message process(String message) { - return MessageBuilder.withPayload(message + " processed") - .setHeader(MessageHeaders.CONTENT_TYPE, "text/plain").build(); - } - -} diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/source/TestSource.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/source/TestSource.java deleted file mode 100644 index 3f5a36934..000000000 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/aggregate/source/TestSource.java +++ /dev/null @@ -1,54 +0,0 @@ -/* - * Copyright 2017-2019 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 - * - * https://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.stream.config.aggregate.source; - -import java.text.SimpleDateFormat; -import java.util.Date; - -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.messaging.Source; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; -import org.springframework.integration.annotation.InboundChannelAdapter; -import org.springframework.integration.core.MessageSource; -import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.support.MessageBuilder; - -/** - * @author Ilayaperumal Gopinathan - */ -@EnableBinding(Source.class) -@EnableAutoConfiguration -@Configuration -public class TestSource { - - @Bean - @InboundChannelAdapter(Source.OUTPUT) - public MessageSource timerMessageSource() { - return new MessageSource() { - @Override - public Message receive() { - return MessageBuilder - .withPayload(new SimpleDateFormat("DDMMMYYYY").format(new Date())) - .setHeader(MessageHeaders.CONTENT_TYPE, "text/plain").build(); - } - }; - } - -} diff --git a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/bean/AggregateWithBeanTest.java b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/bean/AggregateWithBeanTest.java deleted file mode 100644 index 4461c2b79..000000000 --- a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/bean/AggregateWithBeanTest.java +++ /dev/null @@ -1,110 +0,0 @@ -/* - * Copyright 2016-2019 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 - * - * https://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.stream.test.aggregate.bean; - -import java.util.concurrent.TimeUnit; - -import org.junit.Ignore; -import org.junit.Test; -import org.junit.runner.RunWith; - -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.cloud.stream.aggregate.AggregateApplication; -import org.springframework.cloud.stream.aggregate.AggregateApplicationBuilder; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.messaging.Processor; -import org.springframework.cloud.stream.test.binder.MessageCollector; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; -import org.springframework.integration.annotation.Transformer; -import org.springframework.integration.support.MessageBuilder; -import org.springframework.messaging.Message; -import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; - -import static org.assertj.core.api.Assertions.assertThat; - -/** - * @author Marius Bogoevici - */ -@RunWith(SpringJUnit4ClassRunner.class) -@SpringBootTest(classes = AggregateWithBeanTest.ChainedProcessors.class, properties = { - "server.port=-1", "--spring.cloud.stream.bindings.input.contentType=text/plain", - "--spring.cloud.stream.bindings.output.contentType=text/plain" }) -@Ignore -public class AggregateWithBeanTest { - - @Autowired - public MessageCollector messageCollector; - - @Autowired - public AggregateApplication aggregateApplication; - - @Test - @SuppressWarnings("unchecked") - public void testAggregateApplication() throws InterruptedException { - Processor uppercaseProcessor = this.aggregateApplication - .getBinding(Processor.class, "upper"); - Processor suffixProcessor = this.aggregateApplication.getBinding(Processor.class, - "suffix"); - uppercaseProcessor.input().send(MessageBuilder.withPayload("Hello").build()); - Message receivedMessage = (Message) this.messageCollector - .forChannel(suffixProcessor.output()).poll(1, TimeUnit.SECONDS); - assertThat(receivedMessage).isNotNull(); - assertThat(receivedMessage.getPayload()).isEqualTo("HELLO WORLD!"); - } - - @SpringBootApplication - @EnableBinding - public static class ChainedProcessors { - - @Bean - public AggregateApplication aggregateApplication() { - return new AggregateApplicationBuilder().from(UppercaseProcessor.class) - .namespace("upper").to(SuffixProcessor.class).namespace("suffix") - .build(); - } - - } - - @Configuration - @EnableBinding(Processor.class) - @EnableAutoConfiguration - public static class UppercaseProcessor { - - @Transformer(inputChannel = Processor.INPUT, outputChannel = Processor.OUTPUT) - public String transform(String in) { - return in.toUpperCase(); - } - - } - - @Configuration - @EnableBinding(Processor.class) - @EnableAutoConfiguration - public static class SuffixProcessor { - - @Transformer(inputChannel = Processor.INPUT, outputChannel = Processor.OUTPUT) - public String transform(String in) { - return in + " WORLD!"; - } - - } - -} diff --git a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/main/AggregateWithMainTest.java b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/main/AggregateWithMainTest.java deleted file mode 100644 index b61256f47..000000000 --- a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/aggregate/main/AggregateWithMainTest.java +++ /dev/null @@ -1,105 +0,0 @@ -/* - * Copyright 2016-2018 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 - * - * https://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.stream.test.aggregate.main; - -import java.util.concurrent.TimeUnit; - -import org.junit.Ignore; -import org.junit.Test; - -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringBootConfiguration; -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.cloud.stream.aggregate.AggregateApplication; -import org.springframework.cloud.stream.aggregate.AggregateApplicationBuilder; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.messaging.Processor; -import org.springframework.cloud.stream.test.binder.MessageCollector; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Configuration; -import org.springframework.integration.annotation.Transformer; -import org.springframework.integration.support.MessageBuilder; -import org.springframework.messaging.Message; - -import static org.assertj.core.api.Assertions.assertThat; - -/** - * @author Marius Bogoevici - * @author Artem Bilan - */ -@Ignore -public class AggregateWithMainTest { - - @SuppressWarnings("unchecked") - @Test - public void testAggregateApplication() throws InterruptedException { - // emulate a main method - ConfigurableApplicationContext context = new AggregateApplicationBuilder( - MainConfiguration.class).web(false).from(UppercaseProcessor.class) - .namespace("upper").to(SuffixProcessor.class).namespace("suffix") - .run("--spring.cloud.stream.bindings.input.contentType=text/plain", - "--spring.cloud.stream.bindings.output.contentType=text/plain"); - - AggregateApplication aggregateAccessor = context - .getBean(AggregateApplication.class); - MessageCollector messageCollector = context.getBean(MessageCollector.class); - Processor uppercaseProcessor = aggregateAccessor.getBinding(Processor.class, - "upper"); - Processor suffixProcessor = aggregateAccessor.getBinding(Processor.class, - "suffix"); - uppercaseProcessor.input().send(MessageBuilder.withPayload("Hello").build()); - Message receivedMessage = (Message) messageCollector - .forChannel(suffixProcessor.output()).poll(1, TimeUnit.SECONDS); - assertThat(receivedMessage).isNotNull(); - assertThat(receivedMessage.getPayload()).isEqualTo("HELLO WORLD!"); - context.close(); - } - - @SpringBootConfiguration - @EnableAutoConfiguration - public static class MainConfiguration { - - } - - @Configuration - @EnableBinding(Processor.class) - @EnableAutoConfiguration - static class UppercaseProcessor { - - @Autowired - Processor processor; - - @Transformer(inputChannel = Processor.INPUT, outputChannel = Processor.OUTPUT) - public String transform(String in) { - return in.toUpperCase(); - } - - } - - @Configuration - @EnableBinding(Processor.class) - @EnableAutoConfiguration - static class SuffixProcessor { - - @Transformer(inputChannel = Processor.INPUT, outputChannel = Processor.OUTPUT) - public String transform(String in) { - return in + " WORLD!"; - } - - } - -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java deleted file mode 100644 index f2374ba57..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java +++ /dev/null @@ -1,40 +0,0 @@ -/* - * Copyright 2016-2017 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 - * - * https://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.stream.aggregate; - -/** - * Handle to an aggregate application, providing access to the underlying components of - * the aggregate (e.g. bindable instances). - * - * @author Marius Bogoevici - */ -public interface AggregateApplication { - - /** - * Retrieves the bindable proxy instance (e.g. - * {@link org.springframework.cloud.stream.messaging.Processor}, - * {@link org.springframework.cloud.stream.messaging.Source}, - * {@link org.springframework.cloud.stream.messaging.Sink} or custom interface) from - * the given namespace. - * @param bindableType the bindable type - * @param namespace the namespace - * @param parameterized bindable type - * @return binding - */ - T getBinding(Class bindableType, String namespace); - -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationBuilder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationBuilder.java deleted file mode 100644 index 4f369572b..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationBuilder.java +++ /dev/null @@ -1,531 +0,0 @@ -/* - * Copyright 2015-2018 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 - * - * https://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.stream.aggregate; - -import java.util.ArrayList; -import java.util.Arrays; -import java.util.HashMap; -import java.util.LinkedHashMap; -import java.util.LinkedHashSet; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.regex.Pattern; - -import org.springframework.beans.BeansException; -import org.springframework.beans.factory.BeanFactoryUtils; -import org.springframework.beans.factory.SmartInitializingSingleton; -import org.springframework.boot.actuate.autoconfigure.endpoint.EndpointAutoConfiguration; -import org.springframework.boot.autoconfigure.ImportAutoConfiguration; -import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; -import org.springframework.boot.autoconfigure.web.servlet.ServletWebServerFactoryAutoConfiguration; -import org.springframework.boot.builder.SpringApplicationBuilder; -import org.springframework.boot.context.properties.bind.BindResult; -import org.springframework.boot.context.properties.bind.Bindable; -import org.springframework.boot.context.properties.bind.Binder; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.binding.BindableProxyFactory; -import org.springframework.cloud.stream.config.ChannelBindingAutoConfiguration; -import org.springframework.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.core.env.Environment; -import org.springframework.util.Assert; -import org.springframework.util.ClassUtils; -import org.springframework.util.StringUtils; - -/** - * Application builder for {@link AggregateApplication}. - * - * @author Dave Syer - * @author Ilayaperumal Gopinathan - * @author Marius Bogoevici - * @author Venil Noronha - * @author Janne Valkealahti - * @author Vinicius Carvalho - * @author Oleg Zhurakousky - */ -@EnableBinding -public class AggregateApplicationBuilder implements AggregateApplication, - ApplicationContextAware, SmartInitializingSingleton { - - private static final String CHILD_CONTEXT_SUFFIX = ".spring.cloud.stream.context"; - - private static final Bindable> STRING_STRING_MAP = Bindable - .mapOf(String.class, String.class); - - private static final Pattern DOLLAR_ESCAPE_PATTERN = Pattern.compile("\\$"); - - private SourceConfigurer sourceConfigurer; - - private SinkConfigurer sinkConfigurer; - - private List processorConfigurers = new ArrayList<>(); - - private AggregateApplicationBuilder applicationBuilder = this; - - private ConfigurableApplicationContext parentContext; - - private List parentSources = new ArrayList<>(); - - private List parentArgs = new ArrayList<>(); - - private boolean headless = true; - - private boolean webEnvironment = true; - - public AggregateApplicationBuilder(String... args) { - this(new Object[] { ParentConfiguration.class }, args); - } - - public AggregateApplicationBuilder(Object source, String... args) { - this(new Object[] { source }, args); - } - - public AggregateApplicationBuilder(Object[] sources, String[] args) { - addParentSources(sources); - this.parentArgs.addAll(Arrays.asList(args)); - } - - /** - * Adding auto configuration classes to parent sources excluding the configuration - * classes related to binder/binding. - * @param sources sources to which parent sources will be added - */ - private void addParentSources(Object[] sources) { - if (!this.parentSources.contains(ParentConfiguration.class)) { - this.parentSources.add(ParentConfiguration.class); - if (ClassUtils.isPresent( - "org.springframework.boot.actuate.autoconfigure.endpoint.EndpointAutoConfiguration", - null)) { - this.parentSources.add(ParentActuatorConfiguration.class); - } - } - this.parentSources.addAll(Arrays.asList(sources)); - } - - public AggregateApplicationBuilder parent(Object source, String... args) { - return parent(new Object[] { source }, args); - } - - public AggregateApplicationBuilder parent(Object[] sources, String[] args) { - addParentSources(sources); - this.parentArgs.addAll(Arrays.asList(args)); - return this; - } - - /** - * Flag to explicitly request a web or non-web environment. - * @param webEnvironment true if the application has a web environment - * @return the AggregateApplicationBuilder being constructed - * @see SpringApplicationBuilder#web(boolean) - */ - public AggregateApplicationBuilder web(boolean webEnvironment) { - this.webEnvironment = webEnvironment; - return this; - } - - /** - * Configures the headless attribute of the build application. - * @param headless true if the application is headless - * @return the AggregateApplicationBuilder being constructed - * @see SpringApplicationBuilder#headless(boolean) - */ - public AggregateApplicationBuilder headless(boolean headless) { - this.headless = headless; - return this; - } - - @Override - public void afterSingletonsInstantiated() { - this.run(); - } - - @Override - public void setApplicationContext(ApplicationContext applicationContext) - throws BeansException { - this.parentContext = (ConfigurableApplicationContext) applicationContext; - } - - @Override - public T getBinding(Class bindableType, String namespace) { - if (this.parentContext == null) { - throw new IllegalStateException( - "The aggregate application has not been started yet"); - } - try { - ChildContextHolder contextHolder = this.parentContext - .getBean(namespace + CHILD_CONTEXT_SUFFIX, ChildContextHolder.class); - return contextHolder.getChildContext().getBean(bindableType); - } - catch (BeansException e) { - throw new IllegalStateException("Binding not found for '" - + bindableType.getName() + "' into namespace " + namespace); - } - } - - public SourceConfigurer from(Class app) { - SourceConfigurer sourceConfigurer = new SourceConfigurer(app); - this.sourceConfigurer = sourceConfigurer; - return sourceConfigurer; - } - - public ConfigurableApplicationContext run(String... parentArgs) { - this.parentArgs.addAll(Arrays.asList(parentArgs)); - List> apps = new ArrayList<>(); - if (this.sourceConfigurer != null) { - apps.add(this.sourceConfigurer); - } - if (!this.processorConfigurers.isEmpty()) { - for (ProcessorConfigurer processorConfigurer : this.processorConfigurers) { - apps.add(processorConfigurer); - } - } - if (this.sinkConfigurer != null) { - apps.add(this.sinkConfigurer); - } - LinkedHashMap, String> appsToEmbed = new LinkedHashMap<>(); - LinkedHashMap, String> appConfigurers = new LinkedHashMap<>(); - for (int i = 0; i < apps.size(); i++) { - AppConfigurer appConfigurer = apps.get(i); - Class appToEmbed = appConfigurer.getApp(); - // Always update namespace before preparing SharedChannelRegistry - if (appConfigurer.namespace == null) { - // to remove illegal characters for new properties - // binder - // org.springframework.cloud.stream.aggregation.AggregationTest$TestSource - appConfigurer.namespace = AggregateApplicationUtils.getDefaultNamespace( - DOLLAR_ESCAPE_PATTERN.matcher(appConfigurer.getApp().getName()) - .replaceAll("."), - i); - } - appsToEmbed.put(appToEmbed, appConfigurer.namespace); - appConfigurers.put(appConfigurer, appConfigurer.namespace); - } - if (this.parentContext == null) { - if (Boolean.TRUE.equals(this.webEnvironment)) { - Assert.isTrue( - ClassUtils.isPresent("javax.servlet.ServletRequest", - ClassUtils.getDefaultClassLoader()), - "'webEnvironment' is set to 'true' but 'javax.servlet.*' does not appear to be available in " - + "the classpath. Consider adding `org.springframework.boot:spring-boot-starter-web"); - this.addParentSources( - new Object[] { ServletWebServerFactoryAutoConfiguration.class }); - } - this.parentContext = AggregateApplicationUtils.createParentContext( - this.parentSources.toArray(new Class[0]), - this.parentArgs.toArray(new String[0]), selfContained(), - this.webEnvironment, this.headless); - } - else { - if (BeanFactoryUtils.beansOfTypeIncludingAncestors(this.parentContext, - SharedBindingTargetRegistry.class).size() == 0) { - SharedBindingTargetRegistry sharedBindingTargetRegistry = new SharedBindingTargetRegistry(); - this.parentContext.getBeanFactory().registerSingleton( - "sharedBindingTargetRegistry", sharedBindingTargetRegistry); - } - } - SharedBindingTargetRegistry sharedBindingTargetRegistry = this.parentContext - .getBean(SharedBindingTargetRegistry.class); - AggregateApplicationUtils.prepareSharedBindingTargetRegistry( - sharedBindingTargetRegistry, appsToEmbed); - for (Map.Entry, String> appConfigurerEntry : appConfigurers - .entrySet()) { - - AppConfigurer appConfigurer = appConfigurerEntry.getKey(); - if (appConfigurerEntry.getValue() == null) { - continue; - } - String namespace = appConfigurerEntry.getValue().toLowerCase(); - Set argsToUpdate = new LinkedHashSet<>(); - Set argKeys = new LinkedHashSet<>(); - Map target = bindProperties(namespace, - this.parentContext.getEnvironment()); - - if (!target.isEmpty()) { - for (Map.Entry entry : target.entrySet()) { - String key = entry.getKey(); - argKeys.add(key); - argsToUpdate.add("--" + key + "=" + entry.getValue()); - } - } - - if (!argsToUpdate.isEmpty()) { - appConfigurer.args(argsToUpdate.toArray(new String[0])); - } - } - for (int i = apps.size() - 1; i >= 0; i--) { - AppConfigurer appConfigurer = apps.get(i); - appConfigurer.embed(); - } - if (BeanFactoryUtils.beansOfTypeIncludingAncestors(this.parentContext, - AggregateApplication.class).size() == 0) { - this.parentContext.getBeanFactory() - .registerSingleton("aggregateApplicationAccessor", this); - } - return this.parentContext; - } - - private boolean selfContained() { - return (this.sourceConfigurer != null) && (this.sinkConfigurer != null); - } - - private ChildContextBuilder childContext(Class app, - ConfigurableApplicationContext parentContext, String namespace) { - return new ChildContextBuilder( - AggregateApplicationUtils.embedApp(parentContext, namespace, app)); - } - - private Map bindProperties(String namepace, Environment environment) { - Map target; - BindResult> bindResult = Binder.get(environment) - .bind(namepace, STRING_STRING_MAP); - if (bindResult.isBound()) { - target = bindResult.get(); - } - else { - target = new HashMap<>(); - } - return target; - } - - private static class ChildContextHolder { - - private final ConfigurableApplicationContext childContext; - - ChildContextHolder(ConfigurableApplicationContext childContext) { - Assert.notNull(childContext, "cannot be null"); - this.childContext = childContext; - } - - public ConfigurableApplicationContext getChildContext() { - return this.childContext; - } - - } - - /** - * Auto configuration for {@link SharedBindingTargetRegistry}. - */ - @ImportAutoConfiguration(ChannelBindingAutoConfiguration.class) - @EnableBinding - public static class ParentConfiguration { - - @Bean - @ConditionalOnMissingBean(SharedBindingTargetRegistry.class) - public SharedBindingTargetRegistry sharedBindingTargetRegistry() { - return new SharedBindingTargetRegistry(); - } - - } - - /** - * Auto configuration for {@link EndpointAutoConfiguration}. - */ - @ImportAutoConfiguration(EndpointAutoConfiguration.class) - public static class ParentActuatorConfiguration { - - } - - /** - * Source configurer. - */ - public class SourceConfigurer extends AppConfigurer { - - public SourceConfigurer(Class app) { - this.app = app; - AggregateApplicationBuilder.this.sourceConfigurer = this; - } - - public SinkConfigurer to(Class sink) { - return new SinkConfigurer(sink); - } - - public ProcessorConfigurer via(Class processor) { - return new ProcessorConfigurer(processor); - } - - } - - /** - * Sink configurer. - */ - public class SinkConfigurer extends AppConfigurer { - - public SinkConfigurer(Class app) { - this.app = app; - AggregateApplicationBuilder.this.sinkConfigurer = this; - } - - } - - /** - * Processor configurer. - */ - public class ProcessorConfigurer extends AppConfigurer { - - public ProcessorConfigurer(Class app) { - this.app = app; - AggregateApplicationBuilder.this.processorConfigurers.add(this); - } - - public SinkConfigurer to(Class sink) { - return new SinkConfigurer(sink); - } - - public ProcessorConfigurer via(Class processor) { - return new ProcessorConfigurer(processor); - } - - } - - /** - * Abstraction over configuration of an applciation. - * - * @param type of a configurer - */ - public abstract class AppConfigurer> { - - Class app; - - String[] args; - - String[] names; - - String[] profiles; - - String namespace; - - Class getApp() { - return this.app; - } - - public T as(String... names) { - this.names = names; - return getConfigurer(); - } - - public T args(String... args) { - this.args = args; - return getConfigurer(); - } - - public T profiles(String... profiles) { - this.profiles = profiles; - return getConfigurer(); - } - - @SuppressWarnings("unchecked") - private T getConfigurer() { - return (T) this; - } - - public T namespace(String namespace) { - this.namespace = namespace; - return getConfigurer(); - } - - public ConfigurableApplicationContext run(String... args) { - return AggregateApplicationBuilder.this.applicationBuilder.run(args); - } - - void embed() { - final ConfigurableApplicationContext childContext = childContext(this.app, - AggregateApplicationBuilder.this.parentContext, this.namespace) - .args(this.args).config(this.names).profiles(this.profiles) - .run(); - // Register bindable proxies as beans so they can be queried for later - Map bindableProxies = BeanFactoryUtils - .beansOfTypeIncludingAncestors(childContext.getBeanFactory(), - BindableProxyFactory.class); - for (String bindableProxyName : bindableProxies.keySet()) { - try { - AggregateApplicationBuilder.this.parentContext.getBeanFactory() - .registerSingleton(this.getNamespace() + CHILD_CONTEXT_SUFFIX, - new ChildContextHolder(childContext)); - } - catch (Exception e) { - throw new IllegalStateException( - "Error while trying to register the aggregate bound interface '" - + bindableProxyName + "' into namespace '" - + this.getNamespace() + "'", - e); - } - } - - } - - public AggregateApplication build() { - return AggregateApplicationBuilder.this.applicationBuilder; - } - - public String[] getArgs() { - return this.args; - } - - public String getNamespace() { - return this.namespace; - } - - } - - private final class ChildContextBuilder { - - private SpringApplicationBuilder builder; - - private String configName; - - private String[] args; - - private ChildContextBuilder(SpringApplicationBuilder builder) { - this.builder = builder; - } - - public ChildContextBuilder profiles(String... profiles) { - if (profiles != null) { - this.builder.profiles(profiles); - } - return this; - } - - public ChildContextBuilder config(String... configs) { - if (configs != null) { - this.configName = StringUtils.arrayToCommaDelimitedString(configs); - } - return this; - } - - public ChildContextBuilder args(String... args) { - this.args = args; - return this; - } - - public ConfigurableApplicationContext run() { - List args = new ArrayList(); - if (this.args != null) { - args.addAll(Arrays.asList(this.args)); - } - if (this.configName != null) { - args.add("--spring.config.name=" + this.configName); - } - return this.builder.run(args.toArray(new String[0])); - } - - } - -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationUtils.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationUtils.java deleted file mode 100644 index dd7a521c6..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplicationUtils.java +++ /dev/null @@ -1,91 +0,0 @@ -/* - * Copyright 2015-2017 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 - * - * https://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.stream.aggregate; - -import java.util.LinkedHashMap; -import java.util.Map.Entry; - -import org.springframework.boot.Banner.Mode; -import org.springframework.boot.WebApplicationType; -import org.springframework.boot.builder.SpringApplicationBuilder; -import org.springframework.cloud.stream.internal.InternalPropertyNames; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.integration.channel.DirectChannel; -import org.springframework.messaging.SubscribableChannel; - -/** - * Utilities for embedding applications in aggregates. - * - * @author Marius Bogoevici - * @author Ilayaperumal Gopinathan - * @author Venil Noronha - * @author Janne Valkealahti - */ -abstract class AggregateApplicationUtils { - - public static final String INPUT_BINDING_NAME = "input"; - - public static final String OUTPUT_BINDING_NAME = "output"; - - static ConfigurableApplicationContext createParentContext(Class[] sources, - String[] args, final boolean selfContained, boolean webEnvironment, - boolean headless) { - SpringApplicationBuilder aggregatorParentConfiguration = new SpringApplicationBuilder(); - aggregatorParentConfiguration.sources(sources).web(WebApplicationType.NONE) - .headless(headless) - .properties("spring.jmx.default-domain=" - + AggregateApplicationBuilder.ParentConfiguration.class.getName(), - InternalPropertyNames.SELF_CONTAINED_APP_PROPERTY_NAME + "=" - + selfContained); - return aggregatorParentConfiguration.run(args); - } - - static String getDefaultNamespace(String appClassName, int index) { - return appClassName + "-" + index; - } - - protected static SpringApplicationBuilder embedApp( - ConfigurableApplicationContext parentContext, String namespace, - Class app) { - return new SpringApplicationBuilder(app).web(WebApplicationType.NONE).main(app) - .bannerMode(Mode.OFF).properties("spring.jmx.default-domain=" + namespace) - .properties( - InternalPropertyNames.NAMESPACE_PROPERTY_NAME + "=" + namespace) - .registerShutdownHook(false).parent(parentContext); - } - - static void prepareSharedBindingTargetRegistry( - SharedBindingTargetRegistry sharedBindingTargetRegistry, - LinkedHashMap, String> appsWithNamespace) { - int i = 0; - SubscribableChannel sharedChannel = null; - for (Entry, String> appEntry : appsWithNamespace.entrySet()) { - String namespace = appEntry.getValue(); - if (i > 0) { - sharedBindingTargetRegistry.register(namespace + "." + INPUT_BINDING_NAME, - sharedChannel); - } - sharedChannel = new DirectChannel(); - if (i < appsWithNamespace.size() - 1) { - sharedBindingTargetRegistry - .register(namespace + "." + OUTPUT_BINDING_NAME, sharedChannel); - } - i++; - } - } - -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/SharedBindingTargetRegistry.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/SharedBindingTargetRegistry.java deleted file mode 100644 index 1c6fffc26..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/SharedBindingTargetRegistry.java +++ /dev/null @@ -1,58 +0,0 @@ -/* - * Copyright 2016-2017 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 - * - * https://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.stream.aggregate; - -import java.util.Collections; -import java.util.Map; -import java.util.concurrent.ConcurrentSkipListMap; - -/** - * Stores binding targets shared by the components of an aggregate application. - * - * @author Marius Bogoevici - * @since 1.1.1 - */ -public class SharedBindingTargetRegistry { - - private Map sharedBindingTargets = new ConcurrentSkipListMap<>( - String.CASE_INSENSITIVE_ORDER); - - @SuppressWarnings("unchecked") - public T get(String id, Class bindingTargetType) { - Object sharedBindingTarget = this.sharedBindingTargets.get(id); - if (sharedBindingTarget == null) { - return null; - } - if (!bindingTargetType.isAssignableFrom(sharedBindingTarget.getClass())) { - throw new IllegalArgumentException("A shared " + bindingTargetType.getName() - + " was requested, " + "but the existing shared target with id '" + id - + "' is a " + sharedBindingTarget.getClass()); - } - else { - return (T) sharedBindingTarget; - } - } - - public void register(String id, Object bindingTarget) { - this.sharedBindingTargets.put(id, bindingTarget); - } - - public Map getAll() { - return Collections.unmodifiableMap(this.sharedBindingTargets); - } - -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index 559364a2a..1d78efb86 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -71,6 +71,7 @@ import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.support.ChannelInterceptor; +import org.springframework.messaging.support.InterceptableChannel; import org.springframework.retry.RecoveryCallback; import org.springframework.util.Assert; import org.springframework.util.StringUtils; @@ -968,10 +969,10 @@ public abstract class AbstractMessageChannelBinder bindingTargetFactories; @@ -126,17 +122,10 @@ public class BindableProxyFactory String name = BindingBeanDefinitionRegistryUtils .getBindingTargetName(input, method); Class returnType = method.getReturnType(); - Object sharedBindingTarget = locateSharedBindingTarget(name, - returnType); - if (sharedBindingTarget != null) { - BindableProxyFactory.this.inputHolders.put(name, - new BoundTargetHolder(sharedBindingTarget, false)); - } - else { - BindableProxyFactory.this.inputHolders.put(name, - new BoundTargetHolder(getBindingTargetFactory(returnType) - .createInput(name), true)); - } + + BindableProxyFactory.this.inputHolders.put(name, + new BoundTargetHolder(getBindingTargetFactory(returnType) + .createInput(name), true)); } } }); @@ -148,17 +137,10 @@ public class BindableProxyFactory String name = BindingBeanDefinitionRegistryUtils .getBindingTargetName(output, method); Class returnType = method.getReturnType(); - Object sharedBindingTarget = locateSharedBindingTarget(name, - returnType); - if (sharedBindingTarget != null) { - BindableProxyFactory.this.outputHolders.put(name, - new BoundTargetHolder(sharedBindingTarget, false)); - } - else { - BindableProxyFactory.this.outputHolders.put(name, - new BoundTargetHolder(getBindingTargetFactory(returnType) - .createOutput(name), true)); - } + + BindableProxyFactory.this.outputHolders.put(name, + new BoundTargetHolder(getBindingTargetFactory(returnType) + .createOutput(name), true)); } } }); @@ -195,17 +177,6 @@ public class BindableProxyFactory } } - private T locateSharedBindingTarget(String name, Class bindingTargetType) { - return this.sharedBindingTargetRegistry != null - ? this.sharedBindingTargetRegistry.get( - getNamespacePrefixedBindingTargetName(name), bindingTargetType) - : null; - } - - private String getNamespacePrefixedBindingTargetName(String name) { - return this.namespace + "." + name; - } - @Override public synchronized Object getObject() throws Exception { if (this.proxy == null) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java index d075006b6..4d2edc0fc 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java @@ -27,7 +27,7 @@ import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Headers; import org.springframework.messaging.handler.annotation.Payload; import org.springframework.messaging.handler.annotation.support.MethodArgumentNotValidException; -import org.springframework.messaging.handler.annotation.support.PayloadArgumentResolver; +import org.springframework.messaging.handler.annotation.support.PayloadMethodArgumentResolver; import org.springframework.util.ClassUtils; import org.springframework.util.StringUtils; import org.springframework.validation.BeanPropertyBindingResult; @@ -39,7 +39,7 @@ import org.springframework.validation.Validator; * @author Oleg Zhurakousky * @author Gary Russell */ -class SmartPayloadArgumentResolver extends PayloadArgumentResolver { +class SmartPayloadArgumentResolver extends PayloadMethodArgumentResolver { private final MessageConverter messageConverter; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/AggregationTest.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/AggregationTest.java deleted file mode 100644 index 5afd37b1a..000000000 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/AggregationTest.java +++ /dev/null @@ -1,458 +0,0 @@ -/* - * Copyright 2015-2018 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 - * - * https://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.stream.aggregation; - -import java.lang.reflect.Field; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.List; -import java.util.Map; - -import org.junit.After; -import org.junit.Ignore; -import org.junit.Test; - -import org.springframework.beans.DirectFieldAccessor; -import org.springframework.boot.WebApplicationType; -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.builder.SpringApplicationBuilder; -import org.springframework.cloud.stream.aggregate.AggregateApplicationBuilder; -import org.springframework.cloud.stream.aggregate.AggregateApplicationBuilder.SourceConfigurer; -import org.springframework.cloud.stream.aggregate.SharedBindingTargetRegistry; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Output; -import org.springframework.cloud.stream.binding.BindableProxyFactory; -import org.springframework.cloud.stream.binding.BindingTargetFactory; -import org.springframework.cloud.stream.binding.SubscribableChannelBindingTargetFactory; -import org.springframework.cloud.stream.messaging.Processor; -import org.springframework.cloud.stream.messaging.Source; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Configuration; -import org.springframework.messaging.MessageChannel; -import org.springframework.util.ReflectionUtils; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.fail; - -/** - * @author Marius Bogoevici - * @author Ilayaperumal Gopinathan - * @author Artem Bilan - * @author Janne Valkealahti - * @author Gary Russell - * @author Soby Chacko - */ -@Ignore -public class AggregationTest { - - private ConfigurableApplicationContext aggregatedApplicationContext; - - @After - public void closeContext() { - System.clearProperty("a.foo-value"); - System.clearProperty("c.fooValue"); - System.clearProperty("a.foo.value"); - System.clearProperty("c.foo.value"); - if (this.aggregatedApplicationContext != null) { - this.aggregatedApplicationContext.close(); - } - } - - @Test - public void aggregation() { - this.aggregatedApplicationContext = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock") - .web(false).from(TestSource.class).to(TestProcessor.class).run(); - SharedBindingTargetRegistry sharedBindingTargetRegistry = this.aggregatedApplicationContext - .getBean(SharedBindingTargetRegistry.class); - BindingTargetFactory channelFactory = this.aggregatedApplicationContext - .getBean(SubscribableChannelBindingTargetFactory.class); - assertThat(channelFactory).isNotNull(); - assertThat(sharedBindingTargetRegistry.getAll().keySet()).hasSize(2); - this.aggregatedApplicationContext.close(); - } - - @Test - public void testModuleAggregationUsingSharedChannelRegistry() { - // test backward compatibility - this.aggregatedApplicationContext = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock") - .web(false).from(TestSource.class).to(TestProcessor.class).run(); - SharedBindingTargetRegistry sharedChannelRegistry = this.aggregatedApplicationContext - .getBean(SharedBindingTargetRegistry.class); - BindingTargetFactory channelFactory = this.aggregatedApplicationContext - .getBean(SubscribableChannelBindingTargetFactory.class); - assertThat(channelFactory).isNotNull(); - assertThat(sharedChannelRegistry.getAll().keySet()).hasSize(2); - this.aggregatedApplicationContext.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testParentArgsAndSources() { - - List argsToVerify = new ArrayList<>(); - argsToVerify.add("--foo1=bar1"); - argsToVerify.add("--foo2=bar2"); - argsToVerify.add("--foo3=bar3"); - argsToVerify.add("--spring.cloud.stream.default-binder=mock"); - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--foo1=bar1"); - final ConfigurableApplicationContext context = aggregateApplicationBuilder - .parent(DummyConfig.class, "--foo2=bar2").web(false) - .from(TestSource.class).namespace("foo").to(TestProcessor.class) - .namespace("bar") - .run("--foo3=bar3", "--spring.cloud.stream.default-binder=mock"); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - final List parentArgs = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("parentArgs"); - assertThat(parentArgs).containsExactlyInAnyOrder( - argsToVerify.toArray(new String[argsToVerify.size()])); - context.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testParentArgsAndSourcesWithWebDisabled() { - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--foo1=bar1"); - final ConfigurableApplicationContext context = aggregateApplicationBuilder - .parent(DummyConfig.class, "--foo2=bar2").web(false) - .from(TestSource.class).namespace("foo").to(TestProcessor.class) - .namespace("bar").run("--spring.cloud.stream.default-binder=mock"); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - List sources = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("parentSources"); - assertThat(sources).containsExactlyInAnyOrder( - AggregateApplicationBuilder.ParentConfiguration.class, - AggregateApplicationBuilder.ParentActuatorConfiguration.class, - AggregationAppConfig.class, DummyConfig.class); - context.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testNamespacePrefixesFromCmdLine() { - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock"); - this.aggregatedApplicationContext = aggregateApplicationBuilder - .parent(DummyConfig.class).web(false).from(TestSource.class) - .namespace("a").via(TestProcessor.class).namespace("b") - .via(TestProcessor.class).namespace("c") - .run("--a.foo1=bar1", "--b.foo1=bar2", "--c.foo1=bar3"); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - assertThat(Arrays - .asList(((SourceConfigurer) aggregateApplicationBuilderAccessor - .getPropertyValue("sourceConfigurer")).getArgs()) - .contains("--foo1=bar1")).isTrue(); - final List processorConfigurers; - processorConfigurers = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("processorConfigurers"); - for (AggregateApplicationBuilder.ProcessorConfigurer processorConfigurer : processorConfigurers) { - if (processorConfigurer.getNamespace().equals("b")) { - assertThat(Arrays.equals(processorConfigurer.getArgs(), - new String[] { "--foo1=bar2" })).isTrue(); - } - if (processorConfigurer.getNamespace().equals("c")) { - assertThat(Arrays.asList(processorConfigurer.getArgs()) - .contains("--foo1=bar3")).isTrue(); - } - } - this.aggregatedApplicationContext.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testNamespacePrefixesFromCmdLineVsArgs() { - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock"); - this.aggregatedApplicationContext = aggregateApplicationBuilder - .parent(DummyConfig.class).web(false).from(TestSource.class) - .namespace("a").args("--fooValue=bar").via(TestProcessor.class) - .namespace("b").args("--foo1=argbarb").via(TestProcessor.class) - .namespace("c").run("--a.fooValue=bara", "--c.foo1=barc"); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - assertThat(Arrays - .asList(((SourceConfigurer) aggregateApplicationBuilderAccessor - .getPropertyValue("sourceConfigurer")).getArgs()) - .contains("--fooValue=bara")).isTrue(); - final List processorConfigurers; - processorConfigurers = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("processorConfigurers"); - for (AggregateApplicationBuilder.ProcessorConfigurer processorConfigurer : processorConfigurers) { - if (processorConfigurer.getNamespace().equals("b")) { - assertThat(Arrays.equals(processorConfigurer.getArgs(), - new String[] { "--foo1=argbarb" })).isTrue(); - } - if (processorConfigurer.getNamespace().equals("c")) { - assertThat(Arrays - .asList(((SourceConfigurer) aggregateApplicationBuilderAccessor - .getPropertyValue("sourceConfigurer")).getArgs()) - .contains("--fooValue=bara")).isTrue(); - } - } - this.aggregatedApplicationContext.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testNamespacePrefixesFromCmdLineWithRelaxedNames() { - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock"); - this.aggregatedApplicationContext = aggregateApplicationBuilder - .parent(DummyConfig.class).web(false).from(TestSource.class) - .namespace("a").args("--foo-value=bar").via(TestProcessor.class) - .namespace("b").args("--fooValue=argbarb").via(TestProcessor.class) - .namespace("c") - .run("--a.fooValue=bara", "--b.foo-value=barb", "--c.foo1=barc"); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - assertThat(Arrays - .asList(((SourceConfigurer) aggregateApplicationBuilderAccessor - .getPropertyValue("sourceConfigurer")).getArgs()) - .contains("--fooValue=bara")).isTrue(); - final List processorConfigurers; - processorConfigurers = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("processorConfigurers"); - for (AggregateApplicationBuilder.ProcessorConfigurer processorConfigurer : processorConfigurers) { - if (processorConfigurer.getNamespace().equals("b")) { - assertThat(Arrays.equals(processorConfigurer.getArgs(), - new String[] { "--foo-value=barb" })).isTrue(); - } - if (processorConfigurer.getNamespace().equals("c")) { - assertThat(Arrays.asList(processorConfigurer.getArgs()) - .contains("--foo1=barc")).isTrue(); - } - } - this.aggregatedApplicationContext.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testNamespacePrefixesFromCmdLineWithRelaxedNamesAndMorePropertySources() { - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock"); - System.setProperty("a.foo-value", "sysbara"); - System.setProperty("c.fooValue", "sysbarc"); - this.aggregatedApplicationContext = aggregateApplicationBuilder - .parent(DummyConfig.class).web(false).from(TestSource.class) - .namespace("a").args("--foo-value=bar").via(TestProcessor.class) - .namespace("b").args("--fooValue=argbarb").via(TestProcessor.class) - .namespace("c").args("--foo-value=argbarc").run("--a.fooValue=bara"); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - SourceConfigurer sourceConfigurer = (SourceConfigurer) aggregateApplicationBuilderAccessor - .getPropertyValue("sourceConfigurer"); - assertThat(Arrays.asList(sourceConfigurer.getArgs()).contains("--fooValue=bara")) - .isTrue(); - assertThat(Arrays.asList(sourceConfigurer.getArgs()).contains("--foo-value=bara")) - .isTrue(); - final List processorConfigurers; - processorConfigurers = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("processorConfigurers"); - for (AggregateApplicationBuilder.ProcessorConfigurer processorConfigurer : processorConfigurers) { - if (processorConfigurer.getNamespace().equals("b")) { - assertThat(Arrays.equals(processorConfigurer.getArgs(), - new String[] { "--fooValue=argbarb" })).isTrue(); - } - if (processorConfigurer.getNamespace().equals("c")) { - assertThat(Arrays.asList(processorConfigurer.getArgs()) - .contains("--fooValue=sysbarc")).isTrue(); - } - } - this.aggregatedApplicationContext.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testNamespacePrefixesWithoutCmdLinePropertySource() { - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock"); - System.setProperty("a.foo-value", "sysbara"); - System.setProperty("c.fooValue", "sysbarc"); - this.aggregatedApplicationContext = aggregateApplicationBuilder - .parent(DummyConfig.class).web(false).from(TestSource.class) - .namespace("a").args("--foo-value=bar").via(TestProcessor.class) - .namespace("b").args("--fooValue=argbarb").via(TestProcessor.class) - .namespace("c").args("--foo-value=argbarc").run(); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - SourceConfigurer sourceConfigurer = (SourceConfigurer) aggregateApplicationBuilderAccessor - .getPropertyValue("sourceConfigurer"); - assertThat( - Arrays.asList(sourceConfigurer.getArgs()).contains("--foo-value=sysbara")) - .isTrue(); - List configurers; - configurers = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("processorConfigurers"); - for (AggregateApplicationBuilder.ProcessorConfigurer processorConfigurer : configurers) { - if (processorConfigurer.getNamespace().equals("b")) { - assertThat(Arrays.equals(processorConfigurer.getArgs(), - new String[] { "--fooValue=argbarb" })).isTrue(); - } - if (processorConfigurer.getNamespace().equals("c")) { - assertThat(Arrays.asList(processorConfigurer.getArgs()) - .contains("--fooValue=sysbarc")).isTrue(); - } - } - this.aggregatedApplicationContext.close(); - } - - @Test - @SuppressWarnings("unchecked") - public void testNamespacePrefixesWithCAPSProperties() { - AggregateApplicationBuilder aggregateApplicationBuilder = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock"); - System.setProperty("a.fooValue", "sysbara"); - System.setProperty("c.fooValue", "sysbarc"); - this.aggregatedApplicationContext = aggregateApplicationBuilder - .parent(DummyConfig.class).web(false).from(TestSource.class) - .namespace("a").args("--foo-value=bar").via(TestProcessor.class) - .namespace("b").args("--fooValue=argbarb").via(TestProcessor.class) - .namespace("c").args("--foo-value=argbarc").run("--a.fooValue=highest"); - DirectFieldAccessor aggregateApplicationBuilderAccessor = new DirectFieldAccessor( - aggregateApplicationBuilder); - String[] configurers = ((SourceConfigurer) aggregateApplicationBuilderAccessor - .getPropertyValue("sourceConfigurer")).getArgs(); - assertThat(configurers).contains(new String[] { "--fooValue=highest" }); - final List processorConfigurers; - processorConfigurers = (List) aggregateApplicationBuilderAccessor - .getPropertyValue("processorConfigurers"); - for (AggregateApplicationBuilder.ProcessorConfigurer processorConfigurer : processorConfigurers) { - if (processorConfigurer.getNamespace().equals("b")) { - assertThat(Arrays.equals(processorConfigurer.getArgs(), - new String[] { "--fooValue=argbarb" })).isTrue(); - } - if (processorConfigurer.getNamespace().equals("c")) { - assertThat(Arrays.asList(processorConfigurer.getArgs()) - .contains("--fooValue=sysbarc")).isTrue(); - } - } - this.aggregatedApplicationContext.close(); - } - - @Test - public void testNamespaces() { - this.aggregatedApplicationContext = new AggregateApplicationBuilder( - AggregationAppConfig.class, "--spring.cloud.stream.default-binder=mock") - .web(false).from(TestSource.class).namespace("foo") - .to(TestProcessor.class).namespace("bar").run(); - SharedBindingTargetRegistry sharedChannelRegistry = this.aggregatedApplicationContext - .getBean(SharedBindingTargetRegistry.class); - BindingTargetFactory channelFactory = this.aggregatedApplicationContext - .getBean(SubscribableChannelBindingTargetFactory.class); - MessageChannel fooOutput = sharedChannelRegistry.get("foo.output", - MessageChannel.class); - assertThat(fooOutput).isNotNull(); - Object barInput = sharedChannelRegistry.get("bar.input", MessageChannel.class); - assertThat(barInput).isNotNull(); - assertThat(channelFactory).isNotNull(); - assertThat(sharedChannelRegistry.getAll().keySet()).hasSize(2); - this.aggregatedApplicationContext.close(); - } - - @Test - public void testBindableProxyFactoryCaching() { - ConfigurableApplicationContext context = new SpringApplicationBuilder( - TestSource2.class, TestProcessor.class).web(WebApplicationType.NONE) - .run("--spring.cloud.stream.default-binder=mock"); - - Map factories = context - .getBeansOfType(BindableProxyFactory.class); - assertThat(factories).hasSize(2); - - Map sources = context.getBeansOfType(Source.class); - assertThat(sources).hasSize(1); - for (Source source : sources.values()) { - source.output(); - } - - Map fooSources = context.getBeansOfType(FooSource.class); - assertThat(fooSources).hasSize(1); - for (FooSource source : fooSources.values()) { - source.output(); - } - - Map processors = context.getBeansOfType(Processor.class); - assertThat(processors).hasSize(1); - for (Processor processor : processors.values()) { - processor.input(); - processor.output(); - } - - for (BindableProxyFactory factory : factories.values()) { - Field field = ReflectionUtils.findField(BindableProxyFactory.class, - "targetCache"); - ReflectionUtils.makeAccessible(field); - Map targetCache = (Map) ReflectionUtils.getField(field, factory); - if (factory.getObjectType() == Source.class) { - assertThat(targetCache).hasSize(1); - } - if (factory.getObjectType() == FooSource.class) { - assertThat(targetCache).hasSize(1); - } - else if (factory.getObjectType() == Processor.class) { - assertThat(targetCache).hasSize(2); - } - else { - fail("Found unexpected type"); - } - } - context.close(); - } - - public interface FooSource { - - @Output("fooOutput") - MessageChannel output(); - - } - - @EnableBinding(Source.class) - @EnableAutoConfiguration - public static class TestSource { - - } - - @EnableBinding(Processor.class) - @EnableAutoConfiguration - public static class TestProcessor { - - } - - @EnableBinding(FooSource.class) - @EnableAutoConfiguration - public static class TestSource2 { - - } - - @Configuration - public static class DummyConfig { - - } - - @Configuration - @EnableAutoConfiguration - public static class AggregationAppConfig { - - } - -} diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java index 4f9be83da..2c67b4bcc 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2018 the original author or authors. + * Copyright 2013-2019 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. @@ -41,7 +41,6 @@ import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; -import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.interceptor.GlobalChannelInterceptorWrapper; import org.springframework.integration.support.MessageBuilder; @@ -51,6 +50,7 @@ import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessagingException; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.support.ImmutableMessageChannelInterceptor; +import org.springframework.messaging.support.InterceptableChannel; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.fail; @@ -107,9 +107,9 @@ public class BinderAwareChannelResolverTests { assertThat(bindable.getOutputs().size()).isEqualTo(0); // consumer } MessageChannel registered = this.resolver.resolveDestination("foo"); - assertThat(((AbstractMessageChannel) registered).getChannelInterceptors().size()) + assertThat(((InterceptableChannel) registered).getInterceptors().size()) .isEqualTo(2); - assertThat(((AbstractMessageChannel) registered).getChannelInterceptors() + assertThat(((InterceptableChannel) registered).getInterceptors() .get(1) instanceof ImmutableMessageChannelInterceptor).isTrue(); bindables = this.context.getBeansOfType(Bindable.class); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/BindingLifecycleTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/BindingLifecycleTests.java index 5b40978a8..2976a0a25 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/BindingLifecycleTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/BindingLifecycleTests.java @@ -33,6 +33,7 @@ import static org.mockito.Mockito.mock; * @author Oleg Zhurakousky * */ +@SuppressWarnings("unchecked") public class BindingLifecycleTests { @Test diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/CustomPartitionedProducerTest.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/CustomPartitionedProducerTest.java index e52282b66..93d74ff15 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/CustomPartitionedProducerTest.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/CustomPartitionedProducerTest.java @@ -62,7 +62,7 @@ public class CustomPartitionedProducerTest { Source testSource = context.getBean(Source.class); DirectChannel messageChannel = (DirectChannel) testSource.output(); for (ChannelInterceptor channelInterceptor : messageChannel - .getChannelInterceptors()) { + .getInterceptors()) { if (channelInterceptor instanceof MessageConverterConfigurer.PartitioningInterceptor) { Field partitionHandlerField = ReflectionUtils.findField( MessageConverterConfigurer.PartitioningInterceptor.class, @@ -97,7 +97,7 @@ public class CustomPartitionedProducerTest { Source testSource = context.getBean(Source.class); DirectChannel messageChannel = (DirectChannel) testSource.output(); for (ChannelInterceptor channelInterceptor : messageChannel - .getChannelInterceptors()) { + .getInterceptors()) { if (channelInterceptor instanceof MessageConverterConfigurer.PartitioningInterceptor) { Field partitionHandlerField = ReflectionUtils.findField( MessageConverterConfigurer.PartitioningInterceptor.class, @@ -130,7 +130,7 @@ public class CustomPartitionedProducerTest { Source testSource = context.getBean(Source.class); DirectChannel messageChannel = (DirectChannel) testSource.output(); for (ChannelInterceptor channelInterceptor : messageChannel - .getChannelInterceptors()) { + .getInterceptors()) { if (channelInterceptor instanceof MessageConverterConfigurer.PartitioningInterceptor) { Field partitionHandlerField = ReflectionUtils.findField( MessageConverterConfigurer.PartitioningInterceptor.class, @@ -164,7 +164,7 @@ public class CustomPartitionedProducerTest { Source testSource = context.getBean(Source.class); DirectChannel messageChannel = (DirectChannel) testSource.output(); for (ChannelInterceptor channelInterceptor : messageChannel - .getChannelInterceptors()) { + .getInterceptors()) { if (channelInterceptor instanceof MessageConverterConfigurer.PartitioningInterceptor) { Field partitionHandlerField = ReflectionUtils.findField( MessageConverterConfigurer.PartitioningInterceptor.class, diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java index 6a3cf7012..1af9fb66e 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java @@ -333,6 +333,7 @@ public class RoutingFunctionTests { } } + @SuppressWarnings("unused") private static class Person { private String name;