From ea89682368a3b7946f5378fbbfba82e692894021 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 16 May 2017 17:33:02 -0400 Subject: [PATCH] INT-4252 IntegrationFlow: Allow Custom Bean Names JIRA: https://jira.spring.io/browse/INT-4252 To allow to provide arbitrary bean names for the intermediate components in the `IntegrationFlow` change `ComponentsRegistration` to return a `Map` instead of raw `Collection` * Use the value from that map in the `IntegrationFlowBeanPostProcessor` as a fallback option before walking into bean name generation * Provide `containerId` option for the AMQP DSL components * Revert `@AfterClass` in the `AmqpTests` * Expose `id()` for the `JmsDestinationAccessorSpec` * Register `ListenerContainer` and `JmsTemplate` as bean via `ComponentRegistration` in the particular JMS Java DSL components * Fix `Jms` factory to populate the `JmsDefaultListenerContainerSpec` for the `messageDrivenChannelAdapter()` without class specified * Refactor AMQP DSL Inbound components to deal with the new ContainerSpec API * Remove `containerId()` in favor of `id()` in the ContainerSpec --- .../AmqpInboundChannelAdapterDMLCSpec.java | 10 +- .../AmqpInboundChannelAdapterSMLCSpec.java | 10 +- .../dsl/AmqpInboundChannelAdapterSpec.java | 15 +- .../amqp/dsl/AmqpInboundGatewayDMLCSpec.java | 8 +- .../amqp/dsl/AmqpInboundGatewaySMLCSpec.java | 8 +- .../amqp/dsl/AmqpInboundGatewaySpec.java | 34 +++-- .../integration/amqp/dsl/AmqpTests.java | 16 +- .../dsl/IntegrationFlowBeanPostProcessor.java | 140 +++++++++++------- .../integration/dsl/AbstractRouterSpec.java | 2 +- .../dsl/ComponentsRegistration.java | 6 +- .../integration/dsl/ConsumerEndpointSpec.java | 2 +- .../integration/dsl/EndpointSpec.java | 12 +- .../integration/dsl/EnricherSpec.java | 2 +- .../integration/dsl/FilterEndpointSpec.java | 2 +- .../integration/dsl/HeaderEnricherSpec.java | 2 +- .../dsl/IntegrationFlowAdapter.java | 4 +- .../dsl/IntegrationFlowDefinition.java | 25 ++-- .../integration/dsl/PollerSpec.java | 16 +- .../integration/dsl/PublishSubscribeSpec.java | 17 +-- .../dsl/PublisherIntegrationFlow.java | 4 +- .../dsl/RecipientListRouterSpec.java | 6 +- .../integration/dsl/RouterSpec.java | 7 +- .../dsl/StandardIntegrationFlow.java | 21 +-- .../dsl/channel/MessageChannelSpec.java | 9 +- .../integration/dsl/channel/WireTapSpec.java | 6 +- .../context/IntegrationFlowRegistration.java | 5 +- .../dsl/FileInboundChannelAdapterSpec.java | 6 +- .../FileTransferringMessageHandlerSpec.java | 8 +- .../dsl/FileWritingMessageHandlerSpec.java | 6 +- .../RemoteFileInboundChannelAdapterSpec.java | 13 +- .../dsl/RemoteFileOutboundGatewaySpec.java | 13 +- ...ileStreamingInboundChannelAdapterSpec.java | 6 +- .../http/dsl/BaseHttpInboundEndpointSpec.java | 7 +- .../http/dsl/BaseHttpMessageHandlerSpec.java | 6 +- .../ip/dsl/TcpInboundChannelAdapterSpec.java | 13 +- .../ip/dsl/TcpInboundGatewaySpec.java | 13 +- .../ip/dsl/TcpOutboundChannelAdapterSpec.java | 13 +- .../ip/dsl/TcpOutboundGatewaySpec.java | 13 +- .../integration/jms/DynamicJmsTemplate.java | 8 +- .../integration/jms/dsl/Jms.java | 13 +- .../jms/dsl/JmsDestinationAccessorSpec.java | 11 +- .../jms/dsl/JmsInboundChannelAdapterSpec.java | 14 +- .../JmsMessageDrivenChannelAdapterSpec.java | 13 +- .../dsl/JmsOutboundChannelAdapterSpec.java | 13 +- .../integration/jms/dsl/JmsTests.java | 38 +++-- .../jpa/dsl/JpaBaseOutboundEndpointSpec.java | 8 +- .../jpa/dsl/JpaInboundChannelAdapterSpec.java | 8 +- .../mail/dsl/ImapIdleChannelAdapterSpec.java | 14 +- .../dsl/MailInboundChannelAdapterSpec.java | 10 +- .../dsl/ScriptMessageSourceSpec.java | 7 +- 50 files changed, 396 insertions(+), 267 deletions(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterDMLCSpec.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterDMLCSpec.java index 6004e57068..a930d610af 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterDMLCSpec.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterDMLCSpec.java @@ -24,18 +24,20 @@ import org.springframework.amqp.rabbit.listener.DirectMessageListenerContainer; * Spec for an inbound channel adapter with a {@link DirectMessageListenerContainer}. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ -public class AmqpInboundChannelAdapterDMLCSpec extends AmqpInboundChannelAdapterSpec { +public class AmqpInboundChannelAdapterDMLCSpec + extends AmqpInboundChannelAdapterSpec { AmqpInboundChannelAdapterDMLCSpec(DirectMessageListenerContainer listenerContainer) { - super(listenerContainer); + super(new DirectMessageListenerContainerSpec(listenerContainer)); } AmqpInboundChannelAdapterDMLCSpec configureContainer(Consumer configurer) { - configurer.accept(new DirectMessageListenerContainerSpec(this.listenerContainer)); + configurer.accept((DirectMessageListenerContainerSpec) this.listenerContainerSpec); return this; } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSMLCSpec.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSMLCSpec.java index ad12069cf5..8f4c4c2289 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSMLCSpec.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSMLCSpec.java @@ -24,18 +24,20 @@ import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; * Spec for an inbound channel adapter with a {@link SimpleMessageListenerContainer}. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ -public class AmqpInboundChannelAdapterSMLCSpec extends AmqpInboundChannelAdapterSpec { +public class AmqpInboundChannelAdapterSMLCSpec + extends AmqpInboundChannelAdapterSpec { AmqpInboundChannelAdapterSMLCSpec(SimpleMessageListenerContainer listenerContainer) { - super(listenerContainer); + super(new SimpleMessageListenerContainerSpec(listenerContainer)); } AmqpInboundChannelAdapterSMLCSpec configureContainer(Consumer configurer) { - configurer.accept(new SimpleMessageListenerContainerSpec(this.listenerContainer)); + configurer.accept((SimpleMessageListenerContainerSpec) this.listenerContainerSpec); return this; } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSpec.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSpec.java index dc092be0f2..e3c9eae19c 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSpec.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundChannelAdapterSpec.java @@ -16,8 +16,8 @@ package org.springframework.integration.amqp.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; @@ -31,6 +31,7 @@ import org.springframework.integration.dsl.MessageProducerSpec; * @param the container type. * * @author Artem Bilan + * * @since 5.0 */ public abstract class AmqpInboundChannelAdapterSpec @@ -38,16 +39,16 @@ public abstract class AmqpInboundChannelAdapterSpec extends AmqpBaseInboundChannelAdapterSpec implements ComponentsRegistration { - protected final C listenerContainer; + protected final AbstractMessageListenerContainerSpec listenerContainerSpec; - AmqpInboundChannelAdapterSpec(C listenerContainer) { - super(new AmqpInboundChannelAdapter(listenerContainer)); - this.listenerContainer = listenerContainer; + AmqpInboundChannelAdapterSpec(AbstractMessageListenerContainerSpec listenerContainerSpec) { + super(new AmqpInboundChannelAdapter(listenerContainerSpec.get())); + this.listenerContainerSpec = listenerContainerSpec; } @Override - public Collection getComponentsToRegister() { - return Collections.singleton(this.listenerContainer); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.listenerContainerSpec.get(), this.listenerContainerSpec.getId()); } } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewayDMLCSpec.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewayDMLCSpec.java index 5257a6e8ba..c7cfabc370 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewayDMLCSpec.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewayDMLCSpec.java @@ -25,6 +25,8 @@ import org.springframework.amqp.rabbit.listener.DirectMessageListenerContainer; * Spec for a gateway with a {@link DirectMessageListenerContainer}. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ @@ -32,15 +34,15 @@ public class AmqpInboundGatewayDMLCSpec extends AmqpInboundGatewaySpec { AmqpInboundGatewayDMLCSpec(DirectMessageListenerContainer listenerContainer, AmqpTemplate amqpTemplate) { - super(listenerContainer, amqpTemplate); + super(new DirectMessageListenerContainerSpec(listenerContainer), amqpTemplate); } AmqpInboundGatewayDMLCSpec(DirectMessageListenerContainer listenerContainer) { - super(listenerContainer); + super(new DirectMessageListenerContainerSpec(listenerContainer)); } public AmqpInboundGatewayDMLCSpec configureContainer(Consumer configurer) { - configurer.accept(new DirectMessageListenerContainerSpec(this.listenerContainer)); + configurer.accept((DirectMessageListenerContainerSpec) this.listenerContainerSpec); return this; } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySMLCSpec.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySMLCSpec.java index 8eef0c103b..39716aa58e 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySMLCSpec.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySMLCSpec.java @@ -25,6 +25,8 @@ import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; * Spec for a gateway with a {@link SimpleMessageListenerContainer}. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ @@ -32,15 +34,15 @@ public class AmqpInboundGatewaySMLCSpec extends AmqpInboundGatewaySpec { AmqpInboundGatewaySMLCSpec(SimpleMessageListenerContainer listenerContainer, AmqpTemplate amqpTemplate) { - super(listenerContainer, amqpTemplate); + super(new SimpleMessageListenerContainerSpec(listenerContainer), amqpTemplate); } AmqpInboundGatewaySMLCSpec(SimpleMessageListenerContainer listenerContainer) { - super(listenerContainer); + super(new SimpleMessageListenerContainerSpec(listenerContainer)); } public AmqpInboundGatewaySMLCSpec configureContainer(Consumer configurer) { - configurer.accept(new SimpleMessageListenerContainerSpec(this.listenerContainer)); + configurer.accept((SimpleMessageListenerContainerSpec) this.listenerContainerSpec); return this; } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySpec.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySpec.java index 8fe7e64319..ecd3263453 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySpec.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpInboundGatewaySpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2015 the original author or authors. + * Copyright 2014-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. @@ -16,12 +16,11 @@ package org.springframework.integration.amqp.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import org.springframework.amqp.core.AmqpTemplate; import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; -import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.integration.amqp.inbound.AmqpInboundGateway; import org.springframework.integration.dsl.ComponentsRegistration; @@ -29,34 +28,41 @@ import org.springframework.integration.dsl.ComponentsRegistration; * An {@link AmqpBaseInboundGatewaySpec} implementation for a {@link AmqpInboundGateway}. * Allows to provide {@link AbstractMessageListenerContainer} options. * + * @param the spec type. + * @param the container type. + * * @author Artem Bilan + * * @since 5.0 */ public abstract class AmqpInboundGatewaySpec , C extends AbstractMessageListenerContainer> - extends AmqpBaseInboundGatewaySpec implements ComponentsRegistration { + extends AmqpBaseInboundGatewaySpec + implements ComponentsRegistration { - protected final C listenerContainer; + protected final AbstractMessageListenerContainerSpec listenerContainerSpec; - AmqpInboundGatewaySpec(C listenerContainer) { - super(new AmqpInboundGateway(listenerContainer)); - this.listenerContainer = listenerContainer; + AmqpInboundGatewaySpec(AbstractMessageListenerContainerSpec listenerContainerSpec) { + super(new AmqpInboundGateway(listenerContainerSpec.get())); + this.listenerContainerSpec = listenerContainerSpec; } /** * Instantiate {@link AmqpInboundGateway} based on the provided {@link AbstractMessageListenerContainer} * and {@link AmqpTemplate}. - * @param listenerContainer the {@link SimpleMessageListenerContainer} to use. + * @param listenerContainerSpec the {@link AbstractMessageListenerContainerSpec} to use. * @param amqpTemplate the {@link AmqpTemplate} to use. */ - AmqpInboundGatewaySpec(C listenerContainer, AmqpTemplate amqpTemplate) { - super(new AmqpInboundGateway(listenerContainer, amqpTemplate)); - this.listenerContainer = listenerContainer; + AmqpInboundGatewaySpec( + AbstractMessageListenerContainerSpec listenerContainerSpec, + AmqpTemplate amqpTemplate) { + super(new AmqpInboundGateway(listenerContainerSpec.get(), amqpTemplate)); + this.listenerContainerSpec = listenerContainerSpec; } @Override - public Collection getComponentsToRegister() { - return Collections.singleton(this.listenerContainer); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.listenerContainerSpec.get(), this.listenerContainerSpec.getId()); } } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java index 4b3eaeeb92..4b2545349b 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java @@ -37,6 +37,7 @@ import org.springframework.amqp.rabbit.core.RabbitAdmin; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.junit.BrokerRunning; import org.springframework.amqp.rabbit.listener.DirectMessageListenerContainer; +import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; @@ -61,6 +62,7 @@ import org.springframework.test.context.junit4.SpringRunner; /** * @author Artem Bilan * @author Gary Russell + * * @since 5.0 */ @RunWith(SpringRunner.class) @@ -68,7 +70,9 @@ import org.springframework.test.context.junit4.SpringRunner; public class AmqpTests { @ClassRule - public static BrokerRunning brokerRunning = BrokerRunning.isRunning(); + public static BrokerRunning brokerRunning = + BrokerRunning.isRunningWithEmptyQueues("amqpOutboundInput", "amqpReplyChannel", "asyncReplies", + "defaultReplyTo", "si.dsl.test", "testTemplateChannelTransacted"); @Autowired private ConnectionFactory rabbitConnectionFactory; @@ -83,10 +87,13 @@ public class AmqpTests { @Autowired private AmqpInboundGateway amqpInboundGateway; + @Autowired(required = false) + @Qualifier("amqpInboundGatewayContainer") + private SimpleMessageListenerContainer amqpInboundGatewayContainer; + @AfterClass public static void tearDown() { - brokerRunning.removeTestQueues("amqpOutboundInput", "amqpReplyChannel", "asyncReplies", "defaultReplyTo", - "si.dsl.test", "testTemplateChannelTransacted"); + brokerRunning.removeTestQueues(); } @Test @@ -103,6 +110,8 @@ public class AmqpTests { result = this.amqpTemplate.receiveAndConvert("defaultReplyTo"); assertEquals("HELLO WORLD", result); assertSame(this.amqpTemplate, TestUtils.getPropertyValue(this.amqpInboundGateway, "amqpTemplate")); + + assertNotNull(this.amqpInboundGatewayContainer); } @Autowired @@ -213,6 +222,7 @@ public class AmqpTests { .from(Amqp.inboundGateway(rabbitConnectionFactory, amqpTemplate, queue()) .id("amqpInboundGateway") .configureContainer(c -> c + .id("amqpInboundGatewayContainer") .recoveryInterval(5000) .concurrentConsumers(1)) .defaultReplyTo(defaultReplyTo().getName())) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/dsl/IntegrationFlowBeanPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/dsl/IntegrationFlowBeanPostProcessor.java index fbece4390f..76497a98b6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/dsl/IntegrationFlowBeanPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/dsl/IntegrationFlowBeanPostProcessor.java @@ -16,10 +16,10 @@ package org.springframework.integration.config.dsl; -import java.util.ArrayList; import java.util.Collection; import java.util.HashSet; -import java.util.List; +import java.util.LinkedHashMap; +import java.util.Map; import java.util.Set; import org.springframework.beans.BeansException; @@ -130,45 +130,44 @@ public class IntegrationFlowBeanPostProcessor implements BeanPostProcessor, Bean } } - private Object processStandardIntegrationFlow(StandardIntegrationFlow flow, String beanName) { - String flowNamePrefix = beanName + "."; + private Object processStandardIntegrationFlow(StandardIntegrationFlow flow, String flowBeanName) { + String flowNamePrefix = flowBeanName + "."; int subFlowNameIndex = 0; int channelNameIndex = 0; boolean registerSingleton = flow.isRegisterComponents(); - List integrationComponents = new ArrayList<>(flow.getIntegrationComponents()); - for (int i = 0; i < integrationComponents.size(); i++) { - Object component = integrationComponents.get(i); + Map integrationComponents = flow.getIntegrationComponents(); + Map targetIntegrationComponents = new LinkedHashMap<>(integrationComponents.size()); + + for (Map.Entry entry : integrationComponents.entrySet()) { + Object component = entry.getKey(); if (component instanceof ConsumerEndpointSpec) { ConsumerEndpointSpec endpointSpec = (ConsumerEndpointSpec) component; MessageHandler messageHandler = endpointSpec.get().getT2(); ConsumerEndpointFactoryBean endpoint = endpointSpec.get().getT1(); String id = endpointSpec.getId(); - Collection messageHandlers = this.beanFactory.getBeansOfType(messageHandler.getClass(), false, - false).values(); + if (id == null) { + id = generateBeanName(endpoint, entry.getValue()); + } + + Collection messageHandlers = + this.beanFactory.getBeansOfType(messageHandler.getClass(), false, false) + .values(); if (!messageHandlers.contains(messageHandler)) { String handlerBeanName = generateBeanName(messageHandler); - String[] handlerAlias = id != null - ? new String[] { id + IntegrationConfigUtils.HANDLER_ALIAS_SUFFIX } - : null; + String[] handlerAlias = new String[] { id + IntegrationConfigUtils.HANDLER_ALIAS_SUFFIX }; - registerComponent(messageHandler, handlerBeanName, beanName, registerSingleton); - if (handlerAlias != null) { - for (String alias : handlerAlias) { - this.beanFactory.registerAlias(handlerBeanName, alias); - } + registerComponent(messageHandler, handlerBeanName, flowBeanName, registerSingleton); + for (String alias : handlerAlias) { + this.beanFactory.registerAlias(handlerBeanName, alias); } } - String endpointBeanName = id; - if (endpointBeanName == null) { - endpointBeanName = generateBeanName(endpoint); - } - registerComponent(endpoint, endpointBeanName, beanName, registerSingleton); - integrationComponents.set(i, endpoint); + registerComponent(endpoint, id, flowBeanName, registerSingleton); + targetIntegrationComponents.put(endpoint, id); } else { Collection values = this.beanFactory.getBeansOfType(component.getClass(), false, false).values(); @@ -176,17 +175,21 @@ public class IntegrationFlowBeanPostProcessor implements BeanPostProcessor, Bean if (component instanceof AbstractMessageChannel) { String channelBeanName = ((AbstractMessageChannel) component).getComponentName(); if (channelBeanName == null) { - channelBeanName = flowNamePrefix + "channel" + - BeanFactoryUtils.GENERATED_BEAN_NAME_SEPARATOR + channelNameIndex++; + channelBeanName = entry.getValue(); + if (channelBeanName == null) { + channelBeanName = flowNamePrefix + "channel" + + BeanFactoryUtils.GENERATED_BEAN_NAME_SEPARATOR + channelNameIndex++; + } } - registerComponent(component, channelBeanName, beanName, registerSingleton); + registerComponent(component, channelBeanName, flowBeanName, registerSingleton); + targetIntegrationComponents.put(component, channelBeanName); } else if (component instanceof MessageChannelReference) { String channelBeanName = ((MessageChannelReference) component).getName(); if (!this.beanFactory.containsBean(channelBeanName)) { DirectChannel directChannel = new DirectChannel(); - registerComponent(directChannel, channelBeanName, beanName, registerSingleton); - integrationComponents.set(i, directChannel); + registerComponent(directChannel, channelBeanName, flowBeanName, registerSingleton); + targetIntegrationComponents.put(directChannel, channelBeanName); } } else if (component instanceof FixedSubscriberChannel) { @@ -196,26 +199,29 @@ public class IntegrationFlowBeanPostProcessor implements BeanPostProcessor, Bean channelBeanName = flowNamePrefix + "channel" + BeanFactoryUtils.GENERATED_BEAN_NAME_SEPARATOR + channelNameIndex++; } - registerComponent(component, channelBeanName, beanName, registerSingleton); + registerComponent(component, channelBeanName, flowBeanName, registerSingleton); + targetIntegrationComponents.put(component, channelBeanName); } else if (component instanceof SourcePollingChannelAdapterSpec) { SourcePollingChannelAdapterSpec spec = (SourcePollingChannelAdapterSpec) component; - Collection componentsToRegister = spec.getComponentsToRegister(); + Map componentsToRegister = spec.getComponentsToRegister(); if (!CollectionUtils.isEmpty(componentsToRegister)) { - componentsToRegister.stream() - .filter(o -> !this.beanFactory.getBeansOfType(o.getClass(), false, false) - .values() - .contains(o)) - .forEach(o -> registerComponent(o, generateBeanName(o)) - ); + componentsToRegister.entrySet() + .stream() + .filter(o -> + !this.beanFactory.getBeansOfType(o.getKey().getClass(), false, false) + .values() + .contains(o.getKey())) + .forEach(o -> + registerComponent(o.getKey(), generateBeanName(o.getKey(), o.getValue()))); } SourcePollingChannelAdapterFactoryBean pollingChannelAdapterFactoryBean = spec.get().getT1(); String id = spec.getId(); if (!StringUtils.hasText(id)) { - id = generateBeanName(pollingChannelAdapterFactoryBean); + id = generateBeanName(pollingChannelAdapterFactoryBean, entry.getValue()); } - registerComponent(pollingChannelAdapterFactoryBean, id, beanName, registerSingleton); - integrationComponents.set(i, pollingChannelAdapterFactoryBean); + registerComponent(pollingChannelAdapterFactoryBean, id, flowBeanName, registerSingleton); + targetIntegrationComponents.put(pollingChannelAdapterFactoryBean, id); MessageSource messageSource = spec.get().getT2(); if (!this.beanFactory.getBeansOfType(messageSource.getClass(), false, false) @@ -226,25 +232,37 @@ public class IntegrationFlowBeanPostProcessor implements BeanPostProcessor, Bean && ((NamedComponent) messageSource).getComponentName() != null) { messageSourceId = ((NamedComponent) messageSource).getComponentName(); } - registerComponent(messageSource, messageSourceId, beanName, registerSingleton); + registerComponent(messageSource, messageSourceId, flowBeanName, registerSingleton); } } else if (component instanceof StandardIntegrationFlow) { - String subFlowBeanName = flowNamePrefix + "subFlow" + - BeanFactoryUtils.GENERATED_BEAN_NAME_SEPARATOR + subFlowNameIndex++; - registerComponent(component, subFlowBeanName, beanName, registerSingleton); + String subFlowBeanName = + entry.getValue() != null + ? entry.getValue() + : flowNamePrefix + "subFlow" + + BeanFactoryUtils.GENERATED_BEAN_NAME_SEPARATOR + subFlowNameIndex++; + registerComponent(component, subFlowBeanName, flowBeanName, registerSingleton); + targetIntegrationComponents.put(component, subFlowBeanName); } else if (component instanceof AnnotationGatewayProxyFactoryBean) { - registerComponent(component, flowNamePrefix + "gateway", beanName, registerSingleton); + String gatewayId = entry.getValue() != null + ? entry.getValue() + : flowNamePrefix + "gateway"; + registerComponent(component, gatewayId, flowBeanName, registerSingleton); + targetIntegrationComponents.put(component, gatewayId); } else { - String generateBeanName = generateBeanName(component); - registerComponent(component, generateBeanName, beanName, registerSingleton); + String generateBeanName = generateBeanName(component, entry.getValue()); + registerComponent(component, generateBeanName, flowBeanName, registerSingleton); } } + else { + targetIntegrationComponents.put(entry.getKey(), entry.getValue()); + } } } - flow.setIntegrationComponents(integrationComponents); + + flow.setIntegrationComponents(targetIntegrationComponents); return flow; } @@ -256,15 +274,21 @@ public class IntegrationFlowBeanPostProcessor implements BeanPostProcessor, Bean } private void processIntegrationComponentSpec(IntegrationComponentSpec bean) { - registerComponent(bean.get(), generateBeanName(bean.get()), null, false); + registerComponent(bean.get(), generateBeanName(bean.get(), bean.getId()), null, false); if (bean instanceof ComponentsRegistration) { - Collection componentsToRegister = ((ComponentsRegistration) bean).getComponentsToRegister(); + Map componentsToRegister = ((ComponentsRegistration) bean).getComponentsToRegister(); if (!CollectionUtils.isEmpty(componentsToRegister)) { - componentsToRegister.stream() - .filter(component -> !this.beanFactory.getBeansOfType(component.getClass(), false, false) - .values() - .contains(component)) - .forEach(component -> registerComponent(component, generateBeanName(component))); + + componentsToRegister.entrySet() + .stream() + .filter(component -> + !this.beanFactory.getBeansOfType(component.getKey().getClass(), false, false) + .values() + .contains(component.getKey())) + .forEach(component -> + registerComponent(component.getKey(), + generateBeanName(component.getKey(), component.getValue()))); + } } } @@ -293,9 +317,17 @@ public class IntegrationFlowBeanPostProcessor implements BeanPostProcessor, Bean } private String generateBeanName(Object instance) { + return generateBeanName(instance, null); + } + + private String generateBeanName(Object instance, String fallbackId) { if (instance instanceof NamedComponent && ((NamedComponent) instance).getComponentName() != null) { return ((NamedComponent) instance).getComponentName(); } + else if (fallbackId != null) { + return fallbackId; + } + String generatedBeanName = instance.getClass().getName(); String id = generatedBeanName; int counter = -1; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/AbstractRouterSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/AbstractRouterSpec.java index 6da90343af..2515f13333 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/AbstractRouterSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/AbstractRouterSpec.java @@ -96,7 +96,7 @@ public class AbstractRouterSpec, R extends Ab IntegrationFlowBuilder flowBuilder = IntegrationFlows.from(channel); subFlow.configure(flowBuilder); - this.componentsToRegister.add(flowBuilder); + this.componentsToRegister.put(flowBuilder, null); return defaultOutputChannel(channel); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/ComponentsRegistration.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/ComponentsRegistration.java index e02cb07f98..8385731c02 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/ComponentsRegistration.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/ComponentsRegistration.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,7 +16,7 @@ package org.springframework.integration.dsl; -import java.util.Collection; +import java.util.Map; /** * The marker interface for the {@link IntegrationComponentSpec} implementation, @@ -33,6 +33,6 @@ import java.util.Collection; @FunctionalInterface public interface ComponentsRegistration { - Collection getComponentsToRegister(); + Map getComponentsToRegister(); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java index 91107c7456..15a06a170d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java @@ -149,7 +149,7 @@ public abstract class ConsumerEndpointSpec, */ public S transactional(boolean handleMessageAdvice) { TransactionInterceptor transactionInterceptor = new TransactionInterceptorBuilder(handleMessageAdvice).build(); - this.componentsToRegister.add(transactionInterceptor); + this.componentsToRegister.put(transactionInterceptor, null); return transactional(transactionInterceptor); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/EndpointSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/EndpointSpec.java index 0fc0764f89..3b3012d3fa 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/EndpointSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/EndpointSpec.java @@ -16,8 +16,8 @@ package org.springframework.integration.dsl; -import java.util.ArrayList; -import java.util.Collection; +import java.util.LinkedHashMap; +import java.util.Map; import java.util.function.Function; import org.springframework.beans.factory.BeanNameAware; @@ -46,7 +46,7 @@ public abstract class EndpointSpec, F extends Be extends IntegrationComponentSpec> implements ComponentsRegistration { - protected final Collection componentsToRegister = new ArrayList<>(); + protected final Map componentsToRegister = new LinkedHashMap<>(); protected H handler; @@ -87,9 +87,9 @@ public abstract class EndpointSpec, F extends Be * @see PollerSpec */ public S poller(PollerSpec pollerMetadataSpec) { - Collection componentsToRegister = pollerMetadataSpec.getComponentsToRegister(); + Map componentsToRegister = pollerMetadataSpec.getComponentsToRegister(); if (componentsToRegister != null) { - this.componentsToRegister.addAll(componentsToRegister); + this.componentsToRegister.putAll(componentsToRegister); } return poller(pollerMetadataSpec.get()); } @@ -116,7 +116,7 @@ public abstract class EndpointSpec, F extends Be public abstract S autoStartup(boolean autoStartup); @Override - public Collection getComponentsToRegister() { + public Map getComponentsToRegister() { return this.componentsToRegister.isEmpty() ? null : this.componentsToRegister; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java index 76598bbf57..0bcbcf7365 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java @@ -150,7 +150,7 @@ public class EnricherSpec extends ConsumerEndpointSpec targetFlow = buildFlow(); Assert.state(targetFlow != null, "the 'buildFlow()' must not return null"); flow.integrationComponents.clear(); - flow.integrationComponents.addAll(targetFlow.integrationComponents); + flow.integrationComponents.putAll(targetFlow.integrationComponents); this.targetIntegrationFlow = flow.get(); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java index e958181bca..50395b93a3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java @@ -16,9 +16,8 @@ package org.springframework.integration.dsl; -import java.util.Collection; import java.util.HashSet; -import java.util.LinkedHashSet; +import java.util.LinkedHashMap; import java.util.Map; import java.util.Optional; import java.util.Set; @@ -120,7 +119,7 @@ public abstract class IntegrationFlowDefinition REFERENCED_REPLY_PRODUCERS = new HashSet<>(); - protected final Set integrationComponents = new LinkedHashSet<>(); + protected final Map integrationComponents = new LinkedHashMap<>(); protected MessageChannel currentMessageChannel; @@ -134,13 +133,13 @@ public abstract class IntegrationFlowDefinition components) { + B addComponents(Map components) { if (components != null) { - this.integrationComponents.addAll(components); + this.integrationComponents.putAll(components); } return _this(); } @@ -1974,9 +1973,10 @@ public abstract class IntegrationFlowDefinition componentsToRegister = routerSpec.getComponentsToRegister(); + Map componentsToRegister = routerSpec.getComponentsToRegister(); if (!CollectionUtils.isEmpty(componentsToRegister)) { - for (Object component : componentsToRegister) { + for (Map.Entry entry : componentsToRegister.entrySet()) { + Object component = entry.getKey(); if (component instanceof IntegrationFlowDefinition) { IntegrationFlowDefinition flowBuilder = (IntegrationFlowDefinition) component; if (flowBuilder.isOutputChannelRequired()) { @@ -1986,7 +1986,7 @@ public abstract class IntegrationFlowDefinition lastComponent = this.integrationComponents.stream().reduce((first, second) -> second); + Optional lastComponent = + this.integrationComponents.keySet() + .stream() + .reduce((first, second) -> second); if (lastComponent.get() instanceof WireTapSpec) { channel(IntegrationContextUtils.NULL_CHANNEL_BEAN_NAME); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PollerSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PollerSpec.java index afc0ae1228..70a04d31a7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PollerSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PollerSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,11 +16,11 @@ package org.springframework.integration.dsl; -import java.util.ArrayList; import java.util.Arrays; -import java.util.Collection; +import java.util.LinkedHashMap; import java.util.LinkedList; import java.util.List; +import java.util.Map; import java.util.concurrent.Executor; import org.aopalliance.aop.Advice; @@ -49,7 +49,7 @@ public final class PollerSpec extends IntegrationComponentSpec adviceChain = new LinkedList<>(); - private final Collection componentsToRegister = new ArrayList<>(); + private final Map componentsToRegister = new LinkedHashMap<>(); PollerSpec(Trigger trigger) { this.target = new PollerMetadata(); @@ -92,7 +92,7 @@ public final class PollerSpec extends IntegrationComponentSpec getComponentsToRegister() { + public Map getComponentsToRegister() { return this.componentsToRegister; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java index 56bc88cfae..da516a6686 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java @@ -16,9 +16,8 @@ package org.springframework.integration.dsl; -import java.util.ArrayList; -import java.util.Collection; -import java.util.List; +import java.util.LinkedHashMap; +import java.util.Map; import java.util.concurrent.Executor; import org.springframework.integration.dsl.channel.PublishSubscribeChannelSpec; @@ -30,7 +29,7 @@ import org.springframework.integration.dsl.channel.PublishSubscribeChannelSpec; */ public class PublishSubscribeSpec extends PublishSubscribeChannelSpec { - private final List subscriberFlows = new ArrayList<>(); + private final Map subscriberFlows = new LinkedHashMap<>(); PublishSubscribeSpec() { super(); @@ -50,15 +49,15 @@ public class PublishSubscribeSpec extends PublishSubscribeChannelSpec getComponentsToRegister() { - List objects = new ArrayList(); - objects.addAll(super.getComponentsToRegister()); - objects.addAll(this.subscriberFlows); + public Map getComponentsToRegister() { + Map objects = new LinkedHashMap<>(); + objects.putAll(super.getComponentsToRegister()); + objects.putAll(this.subscriberFlows); return objects; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java index 9bd5f01f94..f91389ff5d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java @@ -16,7 +16,7 @@ package org.springframework.integration.dsl; -import java.util.Set; +import java.util.Map; import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; @@ -35,7 +35,7 @@ class PublisherIntegrationFlow extends StandardIntegrationFlow implements Pub private final Publisher> delegate; - PublisherIntegrationFlow(Set integrationComponents, Publisher> publisher) { + PublisherIntegrationFlow(Map integrationComponents, Publisher> publisher) { super(integrationComponents); this.delegate = publisher; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/RecipientListRouterSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/RecipientListRouterSpec.java index 5fe72efb77..95f9c17392 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/RecipientListRouterSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/RecipientListRouterSpec.java @@ -76,7 +76,7 @@ public class RecipientListRouterSpec extends AbstractRouterSpec IntegrationFlowBuilder flowBuilder = IntegrationFlows.from(channel); subFlow.configure(flowBuilder); - this.componentsToRegister.add(flowBuilder); + this.componentsToRegister.put(flowBuilder, null); this.mappingProvider.addMapping(key, channel); return _this(); } @Override - public Collection getComponentsToRegister() { + public Map getComponentsToRegister() { // The 'mappingProvider' must be added to the 'componentsToRegister' in the end to // let all other components to be registered before the 'RouterMappingProvider.onInit()' logic. if (!this.mappingProviderRegistered) { if (!this.mappingProvider.mapping.isEmpty()) { - this.componentsToRegister.add(this.mappingProvider); + this.componentsToRegister.put(this.mappingProvider, null); } this.mappingProviderRegistered = true; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/StandardIntegrationFlow.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/StandardIntegrationFlow.java index 2960c04c88..57dc69783d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/StandardIntegrationFlow.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/StandardIntegrationFlow.java @@ -16,10 +16,12 @@ package org.springframework.integration.dsl; +import java.util.Collections; +import java.util.LinkedHashMap; import java.util.LinkedList; import java.util.List; import java.util.ListIterator; -import java.util.Set; +import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; import org.springframework.context.SmartLifecycle; @@ -63,7 +65,7 @@ import org.springframework.context.SmartLifecycle; */ public class StandardIntegrationFlow implements IntegrationFlow, SmartLifecycle { - private final List integrationComponents; + private final Map integrationComponents; private final List lifecycles = new LinkedList<>(); @@ -71,8 +73,8 @@ public class StandardIntegrationFlow implements IntegrationFlow, SmartLifecycle private boolean running; - StandardIntegrationFlow(Set integrationComponents) { - this.integrationComponents = new LinkedList<>(integrationComponents); + StandardIntegrationFlow(Map integrationComponents) { + this.integrationComponents = new LinkedHashMap<>(integrationComponents); } //TODO Figure out some custom DestinationResolver when we don't register singletons - remove NOSONAR above when done @@ -84,13 +86,13 @@ public class StandardIntegrationFlow implements IntegrationFlow, SmartLifecycle return this.registerComponents; } - public void setIntegrationComponents(List integrationComponents) { + public void setIntegrationComponents(Map integrationComponents) { this.integrationComponents.clear(); - this.integrationComponents.addAll(integrationComponents); + this.integrationComponents.putAll(integrationComponents); } - public List getIntegrationComponents() { - return this.integrationComponents; + public Map getIntegrationComponents() { + return Collections.unmodifiableMap(this.integrationComponents); } @Override @@ -101,7 +103,8 @@ public class StandardIntegrationFlow implements IntegrationFlow, SmartLifecycle @Override public void start() { if (!this.running) { - ListIterator iterator = this.integrationComponents.listIterator(this.integrationComponents.size()); + List components = new LinkedList<>(this.integrationComponents.keySet()); + ListIterator iterator = components.listIterator(this.integrationComponents.size()); this.lifecycles.clear(); while (iterator.hasPrevious()) { Object component = iterator.previous(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/MessageChannelSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/MessageChannelSpec.java index 85010e75f0..c904046ec7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/MessageChannelSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/MessageChannelSpec.java @@ -18,9 +18,10 @@ package org.springframework.integration.dsl.channel; import java.util.ArrayList; import java.util.Arrays; -import java.util.Collection; +import java.util.LinkedHashMap; import java.util.LinkedList; import java.util.List; +import java.util.Map; import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.channel.interceptor.WireTap; @@ -44,7 +45,7 @@ public abstract class MessageChannelSpec, C e extends IntegrationComponentSpec implements ComponentsRegistration { - private final List componentsToRegister = new ArrayList<>(); + private final Map componentsToRegister = new LinkedHashMap<>(); private final List> datatypes = new ArrayList<>(); @@ -108,7 +109,7 @@ public abstract class MessageChannelSpec, C e */ public S wireTap(WireTapSpec wireTapSpec) { WireTap interceptor = wireTapSpec.get(); - this.componentsToRegister.add(interceptor); + this.componentsToRegister.put(interceptor, null); return interceptor(interceptor); } @@ -118,7 +119,7 @@ public abstract class MessageChannelSpec, C e } @Override - public Collection getComponentsToRegister() { + public Map getComponentsToRegister() { return this.componentsToRegister; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java index cbb297c21b..53ef6bf291 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/channel/WireTapSpec.java @@ -16,8 +16,8 @@ package org.springframework.integration.dsl.channel; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import org.springframework.expression.Expression; import org.springframework.integration.channel.interceptor.WireTap; @@ -101,9 +101,9 @@ public class WireTapSpec extends IntegrationComponentSpec } @Override - public Collection getComponentsToRegister() { + public Map getComponentsToRegister() { if (this.selector != null) { - return Collections.singleton(this.selector); + return Collections.singletonMap(this.selector, null); } else { return null; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/IntegrationFlowRegistration.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/IntegrationFlowRegistration.java index 9be6e7daf6..014c694a04 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/IntegrationFlowRegistration.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/IntegrationFlowRegistration.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -30,6 +30,7 @@ import org.springframework.messaging.MessageChannel; * and provide an API for some useful {@link IntegrationFlow} options and its lifecycle. * * @author Artem Bilan + * * @since 5.0 * * @see IntegrationFlowContext @@ -80,7 +81,7 @@ public class IntegrationFlowRegistration { if (this.inputChannel == null) { if (this.integrationFlow instanceof StandardIntegrationFlow) { StandardIntegrationFlow integrationFlow = (StandardIntegrationFlow) this.integrationFlow; - Object next = integrationFlow.getIntegrationComponents().iterator().next(); + Object next = integrationFlow.getIntegrationComponents().keySet().iterator().next(); if (next instanceof MessageChannel) { this.inputChannel = (MessageChannel) next; } diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileInboundChannelAdapterSpec.java b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileInboundChannelAdapterSpec.java index 3356c171f0..786c8a2e2d 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileInboundChannelAdapterSpec.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileInboundChannelAdapterSpec.java @@ -17,9 +17,9 @@ package org.springframework.integration.file.dsl; import java.io.File; -import java.util.Collection; import java.util.Collections; import java.util.Comparator; +import java.util.Map; import java.util.function.Function; import org.springframework.beans.factory.BeanCreationException; @@ -260,9 +260,9 @@ public class FileInboundChannelAdapterSpec } @Override - public Collection getComponentsToRegister() { + public Map getComponentsToRegister() { if (this.expressionFileListFilter != null) { - return Collections.singleton(this.expressionFileListFilter); + return Collections.singletonMap(this.expressionFileListFilter, null); } else { return null; diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileTransferringMessageHandlerSpec.java b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileTransferringMessageHandlerSpec.java index 0934c6350b..b325bd4cf3 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileTransferringMessageHandlerSpec.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileTransferringMessageHandlerSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -17,8 +17,8 @@ package org.springframework.integration.file.dsl; import java.nio.charset.Charset; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import java.util.function.Function; import org.springframework.expression.common.LiteralExpression; @@ -228,9 +228,9 @@ public abstract class FileTransferringMessageHandlerSpec getComponentsToRegister() { + public Map getComponentsToRegister() { if (this.defaultFileNameGenerator != null) { - return Collections.singletonList(this.defaultFileNameGenerator); + return Collections.singletonMap(this.defaultFileNameGenerator, null); } return null; } diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileWritingMessageHandlerSpec.java b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileWritingMessageHandlerSpec.java index 78418aa711..ea826c2a06 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileWritingMessageHandlerSpec.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/FileWritingMessageHandlerSpec.java @@ -17,8 +17,8 @@ package org.springframework.integration.file.dsl; import java.io.File; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import java.util.function.Function; import org.springframework.expression.Expression; @@ -261,9 +261,9 @@ public class FileWritingMessageHandlerSpec } @Override - public Collection getComponentsToRegister() { + public Map getComponentsToRegister() { if (this.defaultFileNameGenerator != null) { - return Collections.singletonList(this.defaultFileNameGenerator); + return Collections.singletonMap(this.defaultFileNameGenerator, null); } return null; } diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileInboundChannelAdapterSpec.java b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileInboundChannelAdapterSpec.java index 4b3dd2dab5..1dc46e1cb3 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileInboundChannelAdapterSpec.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileInboundChannelAdapterSpec.java @@ -17,9 +17,8 @@ package org.springframework.integration.file.dsl; import java.io.File; -import java.util.ArrayList; -import java.util.Collection; -import java.util.List; +import java.util.LinkedHashMap; +import java.util.Map; import java.util.function.Function; import org.springframework.expression.Expression; @@ -234,12 +233,12 @@ public abstract class RemoteFileInboundChannelAdapterSpec getComponentsToRegister() { - List componentsToRegister = new ArrayList<>(); - componentsToRegister.add(this.synchronizer); + public Map getComponentsToRegister() { + Map componentsToRegister = new LinkedHashMap<>(); + componentsToRegister.put(this.synchronizer, null); if (this.expressionFileListFilter != null) { - componentsToRegister.add(this.expressionFileListFilter); + componentsToRegister.put(this.expressionFileListFilter, null); } return componentsToRegister; diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileOutboundGatewaySpec.java b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileOutboundGatewaySpec.java index 7e62c011b8..e80512e78e 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileOutboundGatewaySpec.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileOutboundGatewaySpec.java @@ -17,9 +17,8 @@ package org.springframework.integration.file.dsl; import java.io.File; -import java.util.ArrayList; -import java.util.Collection; -import java.util.List; +import java.util.LinkedHashMap; +import java.util.Map; import java.util.function.Function; import org.springframework.expression.Expression; @@ -340,13 +339,13 @@ public abstract class RemoteFileOutboundGatewaySpec getComponentsToRegister() { - List componentsToRegister = new ArrayList<>(); + public Map getComponentsToRegister() { + Map componentsToRegister = new LinkedHashMap<>(); if (this.expressionFileListFilter != null) { - componentsToRegister.add(this.expressionFileListFilter); + componentsToRegister.put(this.expressionFileListFilter, null); } if (this.mputExpressionFileListFilter != null) { - componentsToRegister.add(this.expressionFileListFilter); + componentsToRegister.put(this.mputExpressionFileListFilter, null); } return componentsToRegister; diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileStreamingInboundChannelAdapterSpec.java b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileStreamingInboundChannelAdapterSpec.java index 64b5214888..0dca643858 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileStreamingInboundChannelAdapterSpec.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/dsl/RemoteFileStreamingInboundChannelAdapterSpec.java @@ -16,8 +16,8 @@ package org.springframework.integration.file.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import java.util.function.Function; import org.springframework.expression.Expression; @@ -127,9 +127,9 @@ public abstract class RemoteFileStreamingInboundChannelAdapterSpec getComponentsToRegister() { + public Map getComponentsToRegister() { if (this.expressionFileListFilter != null) { - return Collections.singleton(this.expressionFileListFilter); + return Collections.singletonMap(this.expressionFileListFilter, null); } else { return null; diff --git a/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpInboundEndpointSpec.java b/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpInboundEndpointSpec.java index 3d60580ec1..0c9235b549 100644 --- a/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpInboundEndpointSpec.java +++ b/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpInboundEndpointSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -17,7 +17,6 @@ package org.springframework.integration.http.dsl; import java.util.Arrays; -import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -306,10 +305,10 @@ public abstract class BaseHttpInboundEndpointSpec getComponentsToRegister() { + public Map getComponentsToRegister() { HeaderMapper headerMapperToRegister = (this.explicitHeaderMapper != null ? this.explicitHeaderMapper : this.headerMapper); - return Collections.singletonList(headerMapperToRegister); + return Collections.singletonMap(headerMapperToRegister, null); } /** diff --git a/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpMessageHandlerSpec.java b/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpMessageHandlerSpec.java index bbde39d722..ee53735bc9 100644 --- a/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpMessageHandlerSpec.java +++ b/spring-integration-http/src/main/java/org/springframework/integration/http/dsl/BaseHttpMessageHandlerSpec.java @@ -16,7 +16,6 @@ package org.springframework.integration.http.dsl; -import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -310,10 +309,11 @@ public abstract class BaseHttpMessageHandlerSpec getComponentsToRegister() { + public Map getComponentsToRegister() { this.target.setUriVariableExpressions(this.uriVariableExpressions); - return Collections.singletonList(this.headerMapper); + return Collections.singletonMap(this.headerMapper, null); } protected abstract boolean isClientSet(); + } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundChannelAdapterSpec.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundChannelAdapterSpec.java index bb7d928350..81fffa1681 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundChannelAdapterSpec.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundChannelAdapterSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,8 +16,8 @@ package org.springframework.integration.ip.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import org.springframework.integration.dsl.ComponentsRegistration; import org.springframework.integration.dsl.MessageProducerSpec; @@ -29,6 +29,8 @@ import org.springframework.scheduling.TaskScheduler; * A {@link MessageProducerSpec} for {@link TcpReceivingChannelAdapter}s. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ @@ -89,9 +91,10 @@ public class TcpInboundChannelAdapterSpec } @Override - public Collection getComponentsToRegister() { - return this.connectionFactory == null ? Collections.emptyList() - : Collections.singletonList(this.connectionFactory); + public Map getComponentsToRegister() { + return this.connectionFactory != null + ? Collections.singletonMap(this.connectionFactory, this.connectionFactory.getComponentName()) + : null; } } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundGatewaySpec.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundGatewaySpec.java index f632352c97..c5493c00f4 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundGatewaySpec.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpInboundGatewaySpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,8 +16,8 @@ package org.springframework.integration.ip.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import org.springframework.integration.dsl.ComponentsRegistration; import org.springframework.integration.dsl.MessagingGatewaySpec; @@ -29,6 +29,8 @@ import org.springframework.scheduling.TaskScheduler; * A {@link MessagingGatewaySpec} for {@link TcpInboundGateway}s. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ @@ -88,9 +90,10 @@ public class TcpInboundGatewaySpec extends MessagingGatewaySpec getComponentsToRegister() { - return this.connectionFactory == null ? Collections.emptyList() - : Collections.singletonList(this.connectionFactory); + public Map getComponentsToRegister() { + return this.connectionFactory != null + ? Collections.singletonMap(this.connectionFactory, this.connectionFactory.getComponentName()) + : null; } } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundChannelAdapterSpec.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundChannelAdapterSpec.java index 17df57766a..02b8682c19 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundChannelAdapterSpec.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundChannelAdapterSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,8 +16,8 @@ package org.springframework.integration.ip.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import org.springframework.integration.dsl.ComponentsRegistration; import org.springframework.integration.dsl.MessageHandlerSpec; @@ -29,6 +29,8 @@ import org.springframework.scheduling.TaskScheduler; * A {@link MessageHandlerSpec} for {@link TcpSendingMessageHandler}s. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ @@ -89,9 +91,10 @@ public class TcpOutboundChannelAdapterSpec } @Override - public Collection getComponentsToRegister() { - return this.connectionFactory == null ? Collections.emptyList() - : Collections.singletonList(this.connectionFactory); + public Map getComponentsToRegister() { + return this.connectionFactory != null + ? Collections.singletonMap(this.connectionFactory, this.connectionFactory.getComponentName()) + : null; } } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundGatewaySpec.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundGatewaySpec.java index bbb5d6c2e9..4671480a94 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundGatewaySpec.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/dsl/TcpOutboundGatewaySpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,8 +16,8 @@ package org.springframework.integration.ip.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import java.util.function.Function; import org.springframework.integration.dsl.ComponentsRegistration; @@ -31,6 +31,8 @@ import org.springframework.messaging.Message; * A {@link MessageHandlerSpec} for {@link TcpOutboundGateway}s. * * @author Gary Russell + * @author Artem Bilan + * * @since 5.0 * */ @@ -88,9 +90,10 @@ public class TcpOutboundGatewaySpec extends MessageHandlerSpec getComponentsToRegister() { - return this.connectionFactory == null ? Collections.emptyList() - : Collections.singletonList(this.connectionFactory); + public Map getComponentsToRegister() { + return this.connectionFactory != null + ? Collections.singletonMap(this.connectionFactory, this.connectionFactory.getComponentName()) + : null; } } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/DynamicJmsTemplate.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/DynamicJmsTemplate.java index 8b5251744f..98ec5fcf02 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/DynamicJmsTemplate.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/DynamicJmsTemplate.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-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. @@ -21,6 +21,8 @@ import org.springframework.util.Assert; /** * @author Mark Fisher + * @author Artem Bilan + * * @since 2.0.2 */ public class DynamicJmsTemplate extends JmsTemplate { @@ -32,13 +34,13 @@ public class DynamicJmsTemplate extends JmsTemplate { return super.getPriority(); } Assert.isTrue(priority >= 0 && priority <= 9, "JMS priority must be in the range of 0-9"); - return priority.intValue(); + return priority; } @Override public long getReceiveTimeout() { Long receiveTimeout = DynamicJmsTemplateProperties.getReceiveTimeout(); - return (receiveTimeout != null) ? receiveTimeout.longValue() : super.getReceiveTimeout(); + return (receiveTimeout != null) ? receiveTimeout : super.getReceiveTimeout(); } } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/Jms.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/Jms.java index 49b2ff0d82..93e56f81c9 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/Jms.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/Jms.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2016 the original author or authors. + * Copyright 2014-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. @@ -211,10 +211,15 @@ public final class Jms { * @param connectionFactory the JMS ConnectionFactory to build on * @return the {@link JmsMessageDrivenChannelAdapterSpec} instance */ - public static JmsMessageDrivenChannelAdapterSpec - .JmsMessageDrivenChannelAdapterListenerContainerSpec + public static JmsMessageDrivenChannelAdapterSpec.JmsMessageDrivenChannelAdapterListenerContainerSpec messageDrivenChannelAdapter(ConnectionFactory connectionFactory) { - return messageDrivenChannelAdapter(connectionFactory, DefaultMessageListenerContainer.class); + try { + return new JmsMessageDrivenChannelAdapterSpec.JmsMessageDrivenChannelAdapterListenerContainerSpec<>( + new JmsDefaultListenerContainerSpec().connectionFactory(connectionFactory)); + } + catch (Exception e) { + throw new IllegalStateException(e); + } } /** diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsDestinationAccessorSpec.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsDestinationAccessorSpec.java index ed3e7c98a1..9be56602bc 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsDestinationAccessorSpec.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsDestinationAccessorSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -29,10 +29,10 @@ import org.springframework.jms.support.destination.JmsDestinationAccessor; * @param the target {@link JmsDestinationAccessor} implementation type. * * @author Artem Bilan + * * @since 5.0 */ -public abstract class - JmsDestinationAccessorSpec, A extends JmsDestinationAccessor> +public abstract class JmsDestinationAccessorSpec, A extends JmsDestinationAccessor> extends IntegrationComponentSpec { protected JmsDestinationAccessorSpec(A accessor) { @@ -44,6 +44,11 @@ public abstract class return _this(); } + @Override + public S id(String id) { + return super.id(id); + } + /** * A {@link DestinationResolver} to use. * @param destinationResolver the {@link DestinationResolver} to use. diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsInboundChannelAdapterSpec.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsInboundChannelAdapterSpec.java index fe106f6bfc..f91b128666 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsInboundChannelAdapterSpec.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsInboundChannelAdapterSpec.java @@ -16,11 +16,14 @@ package org.springframework.integration.jms.dsl; +import java.util.Collections; +import java.util.Map; import java.util.function.Consumer; import javax.jms.ConnectionFactory; import javax.jms.Destination; +import org.springframework.integration.dsl.ComponentsRegistration; import org.springframework.integration.dsl.MessageSourceSpec; import org.springframework.integration.jms.JmsDestinationPollingSource; import org.springframework.integration.jms.JmsHeaderMapper; @@ -33,6 +36,7 @@ import org.springframework.util.Assert; * @param the target {@link JmsInboundChannelAdapterSpec} implementation type. * * @author Artem Bilan + * * @since 5.0 */ public class JmsInboundChannelAdapterSpec> @@ -92,8 +96,9 @@ public class JmsInboundChannelAdapterSpec { + public static class JmsInboundChannelSpecTemplateAware + extends JmsInboundChannelAdapterSpec + implements ComponentsRegistration { JmsInboundChannelSpecTemplateAware(ConnectionFactory connectionFactory) { super(connectionFactory); @@ -111,6 +116,11 @@ public class JmsInboundChannelAdapterSpec getComponentsToRegister() { + return Collections.singletonMap(this.jmsTemplateSpec.get(), this.jmsTemplateSpec.getId()); + } + } } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsMessageDrivenChannelAdapterSpec.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsMessageDrivenChannelAdapterSpec.java index 511550d4b9..27e1cb06bf 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsMessageDrivenChannelAdapterSpec.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsMessageDrivenChannelAdapterSpec.java @@ -16,10 +16,13 @@ package org.springframework.integration.jms.dsl; +import java.util.Collections; +import java.util.Map; import java.util.function.Consumer; import javax.jms.Destination; +import org.springframework.integration.dsl.ComponentsRegistration; import org.springframework.integration.dsl.MessageProducerSpec; import org.springframework.integration.jms.ChannelPublishingJmsMessageListener; import org.springframework.integration.jms.JmsHeaderMapper; @@ -34,6 +37,7 @@ import org.springframework.util.Assert; * @param the target {@link JmsMessageDrivenChannelAdapterSpec} implementation type. * * @author Artem Bilan + * * @since 5.0 */ public class JmsMessageDrivenChannelAdapterSpec> @@ -81,7 +85,8 @@ public class JmsMessageDrivenChannelAdapterSpec, C extends AbstractMessageListenerContainer> - extends JmsMessageDrivenChannelAdapterSpec> { + extends JmsMessageDrivenChannelAdapterSpec> + implements ComponentsRegistration { private final JmsListenerContainerSpec spec; @@ -89,6 +94,7 @@ public class JmsMessageDrivenChannelAdapterSpec getComponentsToRegister() { + return Collections.singletonMap(this.spec.get(), this.spec.getId()); + } + } } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundChannelAdapterSpec.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundChannelAdapterSpec.java index be40ad2176..368b6c631c 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundChannelAdapterSpec.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundChannelAdapterSpec.java @@ -16,12 +16,15 @@ package org.springframework.integration.jms.dsl; +import java.util.Collections; +import java.util.Map; import java.util.function.Consumer; import java.util.function.Function; import javax.jms.ConnectionFactory; import javax.jms.Destination; +import org.springframework.integration.dsl.ComponentsRegistration; import org.springframework.integration.dsl.MessageHandlerSpec; import org.springframework.integration.expression.FunctionExpression; import org.springframework.integration.jms.JmsHeaderMapper; @@ -127,8 +130,9 @@ public class JmsOutboundChannelAdapterSpec { + public static class JmsOutboundChannelSpecTemplateAware + extends JmsOutboundChannelAdapterSpec + implements ComponentsRegistration { JmsOutboundChannelSpecTemplateAware(ConnectionFactory connectionFactory) { super(connectionFactory); @@ -140,6 +144,11 @@ public class JmsOutboundChannelAdapterSpec getComponentsToRegister() { + return Collections.singletonMap(this.jmsTemplateSpec.get(), this.jmsTemplateSpec.getId()); + } + } } diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/dsl/JmsTests.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/dsl/JmsTests.java index 65223cc805..7888aecab9 100644 --- a/spring-integration-jms/src/test/java/org/springframework/integration/jms/dsl/JmsTests.java +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/dsl/JmsTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -64,6 +64,7 @@ import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.test.util.TestUtils; import org.springframework.jms.connection.CachingConnectionFactory; import org.springframework.jms.core.JmsTemplate; +import org.springframework.jms.listener.MessageListenerContainer; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; @@ -127,6 +128,17 @@ public class JmsTests { @Autowired private AtomicBoolean jmsInboundGatewayChannelCalled; + @Autowired(required = false) + @Qualifier("jmsOutboundFlowTemplate") + private JmsTemplate jmsOutboundFlowTemplate; + + @Autowired(required = false) + @Qualifier("jmsMessageDrivenRedeliveryFlowContainer") + private MessageListenerContainer jmsMessageDrivenRedeliveryFlowContainer; + + @Autowired + private CountDownLatch redeliveryLatch; + @Test public void testPollingFlow() { this.controlBus.send("@'jmsTests.ContextConfiguration.integerMessageSource.inboundChannelAdapter'.start()"); @@ -174,6 +186,8 @@ public class JmsTests { assertNotNull(receive); assertEquals("foo", receive.getPayload()); + + assertNotNull(this.jmsOutboundFlowTemplate); } @Test @@ -206,9 +220,6 @@ public class JmsTests { assertEquals("foo", received.getPayload()); } - @Autowired - private CountDownLatch redeliveryLatch; - @Test public void testJmsRedeliveryFlow() throws InterruptedException { this.jmsOutboundInboundChannel.send(MessageBuilder.withPayload("foo") @@ -216,6 +227,8 @@ public class JmsTests { .build()); assertTrue(this.redeliveryLatch.await(10, TimeUnit.SECONDS)); + + assertNotNull(this.jmsMessageDrivenRedeliveryFlowContainer); } @MessagingGateway(defaultRequestChannel = "controlBus.input") @@ -279,8 +292,10 @@ public class JmsTests { @Bean public IntegrationFlow jmsOutboundFlow() { - return f -> f.handle(Jms.outboundAdapter(jmsConnectionFactory()) - .destinationExpression("headers." + SimpMessageHeaderAccessor.DESTINATION_HEADER)); + return f -> f + .handle(Jms.outboundAdapter(jmsConnectionFactory()) + .destinationExpression("headers." + SimpMessageHeaderAccessor.DESTINATION_HEADER) + .configureJmsTemplate(t -> t.id("jmsOutboundFlowTemplate"))); } @Bean @@ -353,16 +368,16 @@ public class JmsTests { @Bean public IntegrationFlow jmsOutboundGatewayFlow() { return f -> f.handle(Jms.outboundGateway(jmsConnectionFactory()) - .replyContainer(c -> c.idleReplyContainerTimeout(10)) - .requestDestination("jmsPipelineTest"), + .replyContainer(c -> c.idleReplyContainerTimeout(10)) + .requestDestination("jmsPipelineTest"), e -> e.id("jmsOutboundGateway")); } @Bean public IntegrationFlow jmsInboundGatewayFlow() { return IntegrationFlows.from(Jms.inboundGateway(jmsConnectionFactory()) - .requestChannel(jmsInboundGatewayInputChannel()) - .destination("jmsPipelineTest")) + .requestChannel(jmsInboundGatewayInputChannel()) + .destination("jmsPipelineTest")) .transform(String::toUpperCase) .get(); } @@ -392,7 +407,8 @@ public class JmsTests { return IntegrationFlows .from(Jms.messageDrivenChannelAdapter(jmsConnectionFactory()) .errorChannel(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME) - .destination("jmsMessageDrivenRedelivery")) + .destination("jmsMessageDrivenRedelivery") + .configureListenerContainer(c -> c.id("jmsMessageDrivenRedeliveryFlowContainer"))) .transform(p -> { throw new RuntimeException("intentional"); }) diff --git a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaBaseOutboundEndpointSpec.java b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaBaseOutboundEndpointSpec.java index e80f1fded8..38961fb30c 100644 --- a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaBaseOutboundEndpointSpec.java +++ b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaBaseOutboundEndpointSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,10 +16,10 @@ package org.springframework.integration.jpa.dsl; -import java.util.Collection; import java.util.Collections; import java.util.LinkedList; import java.util.List; +import java.util.Map; import org.springframework.integration.dsl.ComponentsRegistration; import org.springframework.integration.dsl.MessageHandlerSpec; @@ -159,8 +159,8 @@ public abstract class JpaBaseOutboundEndpointSpec getComponentsToRegister() { - return Collections.singletonList(this.jpaExecutor); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.jpaExecutor, null); } } diff --git a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaInboundChannelAdapterSpec.java b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaInboundChannelAdapterSpec.java index 23c91e08d9..d45ca16e09 100644 --- a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaInboundChannelAdapterSpec.java +++ b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/dsl/JpaInboundChannelAdapterSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,8 +16,8 @@ package org.springframework.integration.jpa.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import org.springframework.expression.Expression; import org.springframework.integration.dsl.ComponentsRegistration; @@ -183,8 +183,8 @@ public class JpaInboundChannelAdapterSpec } @Override - public Collection getComponentsToRegister() { - return Collections.singletonList(this.jpaExecutor); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.jpaExecutor, null); } } diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/ImapIdleChannelAdapterSpec.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/ImapIdleChannelAdapterSpec.java index 674ad08f12..a0d7b71160 100644 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/ImapIdleChannelAdapterSpec.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/ImapIdleChannelAdapterSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2016 the original author or authors. + * Copyright 2014-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. @@ -16,11 +16,11 @@ package org.springframework.integration.mail.dsl; -import java.util.ArrayList; import java.util.Arrays; -import java.util.Collection; +import java.util.LinkedHashMap; import java.util.LinkedList; import java.util.List; +import java.util.Map; import java.util.Properties; import java.util.concurrent.Executor; import java.util.function.Consumer; @@ -64,7 +64,7 @@ public class ImapIdleChannelAdapterSpec private final ImapMailReceiver receiver; - private final Collection componentsToRegister = new ArrayList(); + private final Map componentsToRegister = new LinkedHashMap<>(); private final List adviceChain = new LinkedList<>(); @@ -80,7 +80,7 @@ public class ImapIdleChannelAdapterSpec super(new ImapIdleChannelAdapter(receiver)); this.target.setAdviceChain(this.adviceChain); this.receiver = receiver; - this.componentsToRegister.add(receiver); + this.componentsToRegister.put(receiver, receiver.getComponentName()); this.externalReceiver = externalReceiver; } @@ -326,7 +326,7 @@ public class ImapIdleChannelAdapterSpec */ public ImapIdleChannelAdapterSpec transactional() { TransactionInterceptor transactionInterceptor = new TransactionInterceptorBuilder(false).build(); - this.componentsToRegister.add(transactionInterceptor); + this.componentsToRegister.put(transactionInterceptor, null); return transactional(transactionInterceptor); } @@ -352,7 +352,7 @@ public class ImapIdleChannelAdapterSpec } @Override - public Collection getComponentsToRegister() { + public Map getComponentsToRegister() { return this.componentsToRegister; } diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/MailInboundChannelAdapterSpec.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/MailInboundChannelAdapterSpec.java index 63e41d6cf1..bfe1544d0e 100644 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/MailInboundChannelAdapterSpec.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/dsl/MailInboundChannelAdapterSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2016 the original author or authors. + * Copyright 2014-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. @@ -16,8 +16,8 @@ package org.springframework.integration.mail.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import java.util.Properties; import java.util.function.Consumer; import java.util.function.Function; @@ -50,7 +50,7 @@ import org.springframework.util.Assert; * @since 5.0 */ public abstract class - MailInboundChannelAdapterSpec, R extends AbstractMailReceiver> +MailInboundChannelAdapterSpec, R extends AbstractMailReceiver> extends MessageSourceSpec implements ComponentsRegistration { @@ -253,8 +253,8 @@ public abstract class } @Override - public Collection getComponentsToRegister() { - return Collections.singletonList(this.receiver); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.receiver, this.receiver.getComponentName()); } @Override diff --git a/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/dsl/ScriptMessageSourceSpec.java b/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/dsl/ScriptMessageSourceSpec.java index 4d16b0b057..7ec7710bf3 100644 --- a/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/dsl/ScriptMessageSourceSpec.java +++ b/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/dsl/ScriptMessageSourceSpec.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * 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. @@ -16,7 +16,6 @@ package org.springframework.integration.scripting.dsl; -import java.util.Collection; import java.util.Collections; import java.util.Map; @@ -126,8 +125,8 @@ public class ScriptMessageSourceSpec extends MessageSourceSpec getComponentsToRegister() { - return Collections.singletonList(this.delegate.get()); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.delegate.get(), this.delegate.getId()); } }