commit 04c563788afe7ddd11c604f1cbe3584360a1f521 Author: David Turanski Date: Fri Jul 1 16:40:56 2011 -0400 initial import diff --git a/.classpath b/.classpath new file mode 100644 index 0000000..98cf23a --- /dev/null +++ b/.classpath @@ -0,0 +1,21 @@ + + + + + + + + + + + + + + + + + + + + + diff --git a/.project b/.project new file mode 100644 index 0000000..a4c1e0d --- /dev/null +++ b/.project @@ -0,0 +1,38 @@ + + + spring-integration-flow + + + spring-integration-core + + + + org.eclipse.wst.common.project.facet.core.builder + + + + + org.eclipse.jdt.core.javabuilder + + + + + org.eclipse.wst.validation.validationbuilder + + + + + org.maven.ide.eclipse.maven2Builder + + + + + + org.eclipse.jem.workbench.JavaEMFNature + org.eclipse.wst.common.modulecore.ModuleCoreNature + org.eclipse.jdt.groovy.core.groovyNature + org.maven.ide.eclipse.maven2Nature + org.eclipse.jdt.core.javanature + org.eclipse.wst.common.project.facet.core.nature + + diff --git a/.settings/com.springsource.sts.config.flow.prefs b/.settings/com.springsource.sts.config.flow.prefs new file mode 100644 index 0000000..6ad5ae5 --- /dev/null +++ b/.settings/com.springsource.sts.config.flow.prefs @@ -0,0 +1,30 @@ +#Thu Jun 30 08:13:21 EDT 2011 +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/cic-subflows/src/test/resources/FlowConfigTest-context.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/cic-subflows/src/test/resources/META-INF/spring/subflows/test-subflow-1-config.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/cic-subflows/src/test/resources/META-INF/spring/subflows/test-subflow-2-config.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/cic-subflows/src/test/resources/META-INF/spring/subflows/test-subflow-3-config.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/cic-subflows/src/test/resources/SubFlow1-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/cic-subflows/src/test/resources/SubFlowConfigTest-context.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/cic-subflows/src/test/resources/SubFlowTest-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-cic/src/test/resources/META-INF/spring/subflows/test-subflow-config.xml=\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-cic/src/test/resources/NamespaceTest-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-cic/src/test/resources/SubFlow1-context.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-cic/src/test/resources/SubFlow2-context.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-cic/src/test/resources/SubFlowTest-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/FlowClientNamespace2Test-context.xml=\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/FlowClientNamespaceTest-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/FlowConfigTest-context.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/FlowWithErrorTest-context.xml=\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/FlowWithMissingReferencesTest-context.xml=\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/FlowWithReferencedBeanTest-context.xml=\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/FlowWithReferencesTest-context.xml=\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/META-INF/spring/integration/flows/subflow1-bean-config.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/META-INF/spring/integration/flows/subflow1-context-new.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/META-INF/spring/integration/flows/subflow1-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/META-INF/spring/integration/flows/subflow2-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/META-INF/spring/integration/flows/subflow3-context.xml=\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/META-INF/spring/integration/flows/subflow4-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/META-INF/spring/integration/flows/subflow5-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/MultiBridgeTest-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-integration-flow/src/test/resources/NamespaceTest-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n +eclipse.preferences.version=1 diff --git a/.settings/org.eclipse.jdt.core.prefs b/.settings/org.eclipse.jdt.core.prefs new file mode 100644 index 0000000..cc49a65 --- /dev/null +++ b/.settings/org.eclipse.jdt.core.prefs @@ -0,0 +1,77 @@ +#Wed Jun 29 17:51:26 EDT 2011 +eclipse.preferences.version=1 +org.eclipse.jdt.core.compiler.codegen.inlineJsrBytecode=enabled +org.eclipse.jdt.core.compiler.codegen.targetPlatform=1.6 +org.eclipse.jdt.core.compiler.codegen.unusedLocal=preserve +org.eclipse.jdt.core.compiler.compliance=1.6 +org.eclipse.jdt.core.compiler.debug.lineNumber=generate +org.eclipse.jdt.core.compiler.debug.localVariable=generate +org.eclipse.jdt.core.compiler.debug.sourceFile=generate +org.eclipse.jdt.core.compiler.problem.annotationSuperInterface=warning +org.eclipse.jdt.core.compiler.problem.assertIdentifier=error +org.eclipse.jdt.core.compiler.problem.autoboxing=ignore +org.eclipse.jdt.core.compiler.problem.comparingIdentical=warning +org.eclipse.jdt.core.compiler.problem.deadCode=warning +org.eclipse.jdt.core.compiler.problem.deprecation=warning +org.eclipse.jdt.core.compiler.problem.deprecationInDeprecatedCode=disabled +org.eclipse.jdt.core.compiler.problem.deprecationWhenOverridingDeprecatedMethod=disabled +org.eclipse.jdt.core.compiler.problem.discouragedReference=warning +org.eclipse.jdt.core.compiler.problem.emptyStatement=ignore +org.eclipse.jdt.core.compiler.problem.enumIdentifier=error +org.eclipse.jdt.core.compiler.problem.fallthroughCase=ignore +org.eclipse.jdt.core.compiler.problem.fatalOptionalError=disabled +org.eclipse.jdt.core.compiler.problem.fieldHiding=ignore +org.eclipse.jdt.core.compiler.problem.finalParameterBound=warning +org.eclipse.jdt.core.compiler.problem.finallyBlockNotCompletingNormally=warning +org.eclipse.jdt.core.compiler.problem.forbiddenReference=error +org.eclipse.jdt.core.compiler.problem.hiddenCatchBlock=warning +org.eclipse.jdt.core.compiler.problem.incompatibleNonInheritedInterfaceMethod=warning +org.eclipse.jdt.core.compiler.problem.incompleteEnumSwitch=ignore +org.eclipse.jdt.core.compiler.problem.indirectStaticAccess=ignore +org.eclipse.jdt.core.compiler.problem.localVariableHiding=ignore +org.eclipse.jdt.core.compiler.problem.methodWithConstructorName=warning +org.eclipse.jdt.core.compiler.problem.missingDeprecatedAnnotation=ignore +org.eclipse.jdt.core.compiler.problem.missingHashCodeMethod=ignore +org.eclipse.jdt.core.compiler.problem.missingOverrideAnnotation=ignore +org.eclipse.jdt.core.compiler.problem.missingOverrideAnnotationForInterfaceMethodImplementation=enabled +org.eclipse.jdt.core.compiler.problem.missingSerialVersion=warning +org.eclipse.jdt.core.compiler.problem.missingSynchronizedOnInheritedMethod=ignore +org.eclipse.jdt.core.compiler.problem.noEffectAssignment=warning +org.eclipse.jdt.core.compiler.problem.noImplicitStringConversion=warning +org.eclipse.jdt.core.compiler.problem.nonExternalizedStringLiteral=ignore +org.eclipse.jdt.core.compiler.problem.nullReference=warning +org.eclipse.jdt.core.compiler.problem.overridingPackageDefaultMethod=warning +org.eclipse.jdt.core.compiler.problem.parameterAssignment=ignore +org.eclipse.jdt.core.compiler.problem.possibleAccidentalBooleanAssignment=ignore +org.eclipse.jdt.core.compiler.problem.potentialNullReference=ignore +org.eclipse.jdt.core.compiler.problem.rawTypeReference=warning +org.eclipse.jdt.core.compiler.problem.redundantNullCheck=ignore +org.eclipse.jdt.core.compiler.problem.redundantSuperinterface=ignore +org.eclipse.jdt.core.compiler.problem.specialParameterHidingField=disabled +org.eclipse.jdt.core.compiler.problem.staticAccessReceiver=warning +org.eclipse.jdt.core.compiler.problem.suppressOptionalErrors=disabled +org.eclipse.jdt.core.compiler.problem.suppressWarnings=enabled +org.eclipse.jdt.core.compiler.problem.syntheticAccessEmulation=ignore +org.eclipse.jdt.core.compiler.problem.typeParameterHiding=warning +org.eclipse.jdt.core.compiler.problem.uncheckedTypeOperation=warning +org.eclipse.jdt.core.compiler.problem.undocumentedEmptyBlock=ignore +org.eclipse.jdt.core.compiler.problem.unhandledWarningToken=warning +org.eclipse.jdt.core.compiler.problem.unnecessaryElse=ignore +org.eclipse.jdt.core.compiler.problem.unnecessaryTypeCheck=ignore +org.eclipse.jdt.core.compiler.problem.unqualifiedFieldAccess=ignore +org.eclipse.jdt.core.compiler.problem.unusedDeclaredThrownException=ignore +org.eclipse.jdt.core.compiler.problem.unusedDeclaredThrownExceptionExemptExceptionAndThrowable=enabled +org.eclipse.jdt.core.compiler.problem.unusedDeclaredThrownExceptionIncludeDocCommentReference=enabled +org.eclipse.jdt.core.compiler.problem.unusedDeclaredThrownExceptionWhenOverriding=disabled +org.eclipse.jdt.core.compiler.problem.unusedImport=warning +org.eclipse.jdt.core.compiler.problem.unusedLabel=warning +org.eclipse.jdt.core.compiler.problem.unusedLocal=warning +org.eclipse.jdt.core.compiler.problem.unusedObjectAllocation=ignore +org.eclipse.jdt.core.compiler.problem.unusedParameter=ignore +org.eclipse.jdt.core.compiler.problem.unusedParameterIncludeDocCommentReference=enabled +org.eclipse.jdt.core.compiler.problem.unusedParameterWhenImplementingAbstract=disabled +org.eclipse.jdt.core.compiler.problem.unusedParameterWhenOverridingConcrete=disabled +org.eclipse.jdt.core.compiler.problem.unusedPrivateMember=warning +org.eclipse.jdt.core.compiler.problem.unusedWarningToken=warning +org.eclipse.jdt.core.compiler.problem.varargsArgumentNeedCast=warning +org.eclipse.jdt.core.compiler.source=1.6 diff --git a/.settings/org.eclipse.wst.common.component b/.settings/org.eclipse.wst.common.component new file mode 100644 index 0000000..81c01f3 --- /dev/null +++ b/.settings/org.eclipse.wst.common.component @@ -0,0 +1,7 @@ + + + + + + + diff --git a/.settings/org.eclipse.wst.common.project.facet.core.xml b/.settings/org.eclipse.wst.common.project.facet.core.xml new file mode 100644 index 0000000..5c9bd75 --- /dev/null +++ b/.settings/org.eclipse.wst.common.project.facet.core.xml @@ -0,0 +1,5 @@ + + + + + diff --git a/.settings/org.maven.ide.eclipse.prefs b/.settings/org.maven.ide.eclipse.prefs new file mode 100644 index 0000000..7374e14 --- /dev/null +++ b/.settings/org.maven.ide.eclipse.prefs @@ -0,0 +1,8 @@ +#Tue May 24 12:20:41 EDT 2011 +activeProfiles= +eclipse.preferences.version=1 +fullBuildGoals=process-test-resources +resolveWorkspaceProjects=true +resourceFilterGoals=process-resources resources\:testResources +skipCompilerPlugin=true +version=1 diff --git a/.springBeans b/.springBeans new file mode 100644 index 0000000..38aad46 --- /dev/null +++ b/.springBeans @@ -0,0 +1,17 @@ + + + 1 + + + + + + + src/test/resources/META-INF/spring/integration/flows/subflow1-context.xml + src/test/resources/FlowConfigNamespaceTest-context.xml + src/test/resources/FlowClientNamespaceTest-context.xml + src/test/resources/ref-bean-config.xml + + + + diff --git a/README b/README new file mode 100644 index 0000000..208b58a --- /dev/null +++ b/README @@ -0,0 +1,44 @@ +Goals +--------- +This component is exploring mechanisms to encapsulate a referenced Spring Integration message flow as a component. +A flow is a Spring Integration message flow intended for reuse. A flow is accessed via logical "ports" which map +to internal channels. + +A flow may expose multiple inputs and multiple outputs. A port mapping is defined for each input and has multiple +outputs associated with it. See src/test/resources/META-INF/spring/integration/flows/subflow1-context.xml, for example. + +Generally, a message flow may behave like a router. For example, a flow may define +a primary output and a discard output. Additionally, it may act like a delayer, providing no immediate response. Or it +could act as an outbound channel adapter, providing no output. + +The goal is to support these, and possibly other semantics. Additional goals are: +- Encapsulation: the flow channels, and component should not be included the consumer application context. The consumer need +not know how the flow is configured (i.e, it is contained in a jar). +- Configuration: The flow may define optional or required properties to be provided by the consumer. The flow may define optional +or required bean definitions provided by the consumer (e.g, a generic XML processor may require an OXM marshaller) + - The flow should be self describing, in terms of ports, properties, and beans it exposes + - It should be easy to implement or wrap an existing configuration as a flow and provide a better option than simple importing + a spring configuration file and sending a message to one of its input channels + +Usage +------- +The flow consumer instantiates a flow and defines one or more flow outbound gateways for each input port: + + + + + +A message sent on the input-channel is delegated to the flow. The message on the output-channel is a response from one of the output +ports. The output port name is contained in the response header 'flow.output.port' + +Implementation +---------------- +The flow element creates a Flow instance (eventually, configured with properties, referenced beans, etc.) The flow id is used to derive the flow's +spring bean definition file by convention. This bean definition file will be used to create a standalone application context which must provide a +FlowConfiguration containing the metadata about the exposed ports, beans, and properties. + +Currently, all exposed outputs are bridged to a PublishSubscribeChannel which acts as a single 'flowOutputChannel'. +Each flow outbound gateway instance is backed by a FlowMessageHandler that bridges the 'flow output channel to its own QueueChannel. This is +analogous to a JMS topic. Each flow message handler sends the request message to the flow input channel corresponding to the input port and checks +its queue for a response. A correlation id is used to match the response to the request. \ No newline at end of file diff --git a/pom.xml b/pom.xml new file mode 100644 index 0000000..fe93ef9 --- /dev/null +++ b/pom.xml @@ -0,0 +1,154 @@ + + + + 4.0.0 + org.springframework.integration + spring-integration-flow + 2.0.5.BUILD-SNAPSHOT + Spring Integration Flow Support + + + The Apache Software License, Version 2.0 + http://www.apache.org/licenses/LICENSE-2.0.txt + repo + + + + + + + log4j + log4j + 1.2.12 + test + + + + org.springframework.integration + spring-integration-test + 2.0.5.BUILD-SNAPSHOT + test + + + + org.springframework + spring-test + ${spring.framework.version} + test + + + + org.springframework.integration + spring-integration-core + 2.0.5.BUILD-SNAPSHOT + + + + org.springframework + spring-context-support + ${spring.framework.version} + compile + + + + org.springframework + spring-context + ${spring.framework.version} + compile + + + + org.springframework + spring-core + ${spring.framework.version} + compile + + + org.springframework + spring-asm + ${spring.framework.version} + compile + + + + org.springframework + spring-expression + ${spring.framework.version} + compile + + + + commons-lang + commons-lang + 2.6 + + + + org.codehaus.groovy + groovy-all + 1.8.0 + + + + junit + junit + 4.8.2 + test + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + 2.3.2 + + 1.5 + 1.5 + + + + maven-antrun-plugin + + + test-compile + test-compile + + + + + + + + + + + + run + + + + + + + + + UTF8 + 3.1.0.M2 + true + + + + + http://maven.springframework.org/milestone + http://maven.springframework.org/milestone + + + \ No newline at end of file diff --git a/src/main/java/META-INF/MANIFEST.MF b/src/main/java/META-INF/MANIFEST.MF new file mode 100644 index 0000000..5e94951 --- /dev/null +++ b/src/main/java/META-INF/MANIFEST.MF @@ -0,0 +1,3 @@ +Manifest-Version: 1.0 +Class-Path: + diff --git a/src/main/java/org/springframework/integration/flow/ChannelNamePortConfiguration.java b/src/main/java/org/springframework/integration/flow/ChannelNamePortConfiguration.java new file mode 100644 index 0000000..9ab2df9 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/ChannelNamePortConfiguration.java @@ -0,0 +1,84 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +/** + * @author David Turanski + */ +package org.springframework.integration.flow; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +public class ChannelNamePortConfiguration extends NamedResourceConfiguration implements FlowProviderPortConfiguration { + + private PortMetadata inputPortMetadata; + + public ChannelNamePortConfiguration(PortMetadata inputPortMetadata, List outputPortMetadataList) { + super(outputPortMetadataList); + this.inputPortMetadata = inputPortMetadata; + } + + @Override + public String getInputPortName() { + return this.inputPortMetadata.getName(); + } + + @Override + public String getInputChannel() { + return this.inputPortMetadata.getChannelName(); + } + + @Override + public String getInputPortDescription() { + return this.inputPortMetadata.getDescription(); + } + + + @Override + public String getOutputChannel(String portName) { + PortMetadata portMetadata = (PortMetadata) find(portName); + if (portMetadata != null) { + return portMetadata.getChannelName(); + } + return null; + } + + @Override + public List getOutputPortNames() { + List results = new ArrayList(); + for (NamedResourceMetadata resourceMetadata : getConfiguredResources()) { + PortMetadata portMetadata = (PortMetadata) resourceMetadata; + results.add(portMetadata.getName()); + } + /** + * consistent with ClientPortConfiguration impl. + */ + + if (results.isEmpty()) { + return null; + } + + return results; + } + + /* (non-Javadoc) + * @see org.springframework.integration.flow.FlowProviderPortConfiguration#getOutputPortMetadata() + */ + @Override + public List getOutputPortMetadata() { + return Collections.unmodifiableList(super.getConfiguredResources()); + } +} diff --git a/src/main/java/org/springframework/integration/flow/Flow.java b/src/main/java/org/springframework/integration/flow/Flow.java new file mode 100644 index 0000000..d0cbfb0 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/Flow.java @@ -0,0 +1,232 @@ +package org.springframework.integration.flow; + +import java.util.Properties; + +import org.apache.commons.lang.ArrayUtils; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.BeanNameAware; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.beans.factory.support.BeanDefinitionRegistryPostProcessor; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.core.env.MutablePropertySources; +import org.springframework.core.env.PropertiesPropertySource; +import org.springframework.core.env.PropertySource; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.channel.AbstractMessageChannel; +import org.springframework.integration.channel.PublishSubscribeChannel; +import org.springframework.integration.flow.config.FlowUtils; +import org.springframework.integration.flow.interceptor.FlowInterceptor; +import org.springframework.integration.support.channel.BeanFactoryChannelResolver; +import org.springframework.integration.support.channel.ChannelResolver; +import org.springframework.util.Assert; +import org.springframework.util.StringUtils; + +/** + * Encapsulates a message flow with inputs and outputs exposed via a message + * port + * + * @author David Turanski + * + */ +public class Flow implements InitializingBean, BeanNameAware, ChannelResolver, BeanDefinitionRegistryPostProcessor { + + protected Log logger = LogFactory.getLog(getClass()); + + private volatile ConfigurableApplicationContext flowContext; + + private volatile FlowConfiguration flowConfiguration; + + private volatile String[] configLocations; + + private volatile String[] referencedBeanLocations; + + private volatile Properties flowProperties; + + private volatile String name; + + private volatile ChannelResolver flowChannelResolver; + + private volatile PublishSubscribeChannel flowOutputChannel; + + private volatile boolean help; + + public Flow() { + + } + + public Flow(String[] configLocations) { + this.configLocations = configLocations; + } + + @Override + public void afterPropertiesSet() { + + if (configLocations == null) { + configLocations = new String[] { String.format( + "classpath*:META-INF/spring/integration/flows/%s-context.xml", this.name) }; + } + + if (referencedBeanLocations != null) { + configLocations = (String[]) ArrayUtils.addAll(configLocations, referencedBeanLocations); + } + + logger.debug("instantiating flow context from configLocations [" + + StringUtils.arrayToCommaDelimitedString(configLocations) + "]"); + + Assert.notEmpty(configLocations, "configLocations cannot be empty"); + + flowContext = new ClassPathXmlApplicationContext(configLocations); + + this.flowConfiguration = flowContext.getBean(FlowConfiguration.class); + Assert.notNull(flowConfiguration, "flow context does not contain a flow configuration"); + + if (help) { + System.out.println(displayFlowConfiguration()); + } + + validatePortMapping(); + + this.flowChannelResolver = new BeanFactoryChannelResolver(flowContext); + + addReferencedProperties(); + + } + + public FlowConfiguration getFlowConfiguration() { + return this.flowConfiguration; + } + + @Override + public void setBeanName(String name) { + this.name = name; + + } + + public String getBeanName() { + return this.name; + } + + public void setReferencedBeanLocations(String[] referencedBeanLocations) { + this.referencedBeanLocations = referencedBeanLocations; + } + + public void setProperties(Properties flowProperties) { + this.flowProperties = flowProperties; + } + + public void setHelp(boolean help) { + this.help = help; + } + + public PublishSubscribeChannel getFlowOutputChannel() { + return flowOutputChannel; + } + + public void setFlowOutputChannel(PublishSubscribeChannel flowOutputChannel) { + this.flowOutputChannel = flowOutputChannel; + } + + @Override + public MessageChannel resolveChannelName(String channelName) { + return flowChannelResolver.resolveChannelName(channelName); + } + + @Override + public void postProcessBeanFactory(ConfigurableListableBeanFactory arg0) throws BeansException { + // TODO Auto-generated method stub + + } + + @Override + public void postProcessBeanDefinitionRegistry(BeanDefinitionRegistry registry) throws BeansException { + bridgeMessagingPorts(registry); + + } + + /** + * + */ + private String displayFlowConfiguration() { + StringBuilder sb = new StringBuilder(); + sb.append("\nFlow configuration for [").append(this.getBeanName()).append("]:\n"); + sb.append("Port configuration:\n"); + for (FlowProviderPortConfiguration portConfiguration : this.getFlowConfiguration().getPortConfigurations()) { + sb.append("\tinput port:").append(portConfiguration.getInputPortName()).append("\n\n\t") + .append(portConfiguration.getInputPortDescription()).append("\n\n").append("\toutput ports:\n"); + + for (NamedResourceMetadata metadata : portConfiguration.getOutputPortMetadata()) { + sb.append("\t\t").append(metadata.getName()).append("\t") + .append(metadata.getDescription()).append("\n"); + } + + NamedResourceConfiguration referencedBeansConfig = this.getFlowConfiguration() + .getReferencedBeansConfiguration(); + if (referencedBeansConfig != null && !referencedBeansConfig.getConfiguredResources().isEmpty()) { + sb.append("\nReferenced beans:\n"); + for (NamedResourceMetadata metadata : referencedBeansConfig.getConfiguredResources()) { + sb.append("\t").append(metadata).append("\n"); + } + } + + NamedResourceConfiguration propertiesConfig = this.getFlowConfiguration().getPropertiesConfiguration(); + if (propertiesConfig != null && !propertiesConfig.getConfiguredResources().isEmpty()) { + sb.append("\nFlow properties:\n"); + for (NamedResourceMetadata metadata : propertiesConfig.getConfiguredResources()) { + sb.append("\t").append(metadata).append("\n"); + } + } + + } + + return sb.toString(); + } + + private void addReferencedProperties() { + if (flowProperties != null) { + + PropertySource propertySource = new PropertiesPropertySource("flowProperties", flowProperties); + NamedResourceConfiguration propertiesConfiguration = this.getFlowConfiguration() + .getPropertiesConfiguration(); + if (propertiesConfiguration != null) { + for (NamedResourceMetadata resource : propertiesConfiguration.getRequiredResources()) { + Assert.isTrue(propertySource.containsProperty(resource.getName()), "Flow [" + this.name + + "] is missing required property [" + resource.getName() + "]"); + } + } + MutablePropertySources propertySources = flowContext.getEnvironment().getPropertySources(); + propertySources.addLast(propertySource); + flowContext.refresh(); + } + + } + + private void validatePortMapping() { + Assert.notEmpty(this.flowConfiguration.getPortConfigurations(), + "flow configuration contains no port configurations"); + } + + private void bridgeMessagingPorts(BeanDefinitionRegistry registry) { + + /* + * create a bridge for each target output port to the flow outputChannel + */ + for (FlowProviderPortConfiguration targetPortConfiguration : this.getFlowConfiguration() + .getPortConfigurations()) { + for (String outputPort : targetPortConfiguration.getOutputPortNames()) { + String targetOutputChannelName = (String) targetPortConfiguration.getOutputChannel(outputPort); + AbstractMessageChannel inputChannel = (AbstractMessageChannel) resolveChannelName(targetOutputChannelName); + + inputChannel.addInterceptor(new FlowInterceptor(outputPort)); + + logger.debug("creating output bridge on [" + outputPort + "] inputChannelName = [" + + targetOutputChannelName + "] outputChannel = [" + this.flowOutputChannel + "]"); + FlowUtils.createBridge(inputChannel, this.flowOutputChannel, registry); + } + } + } +} diff --git a/src/main/java/org/springframework/integration/flow/FlowConfiguration.java b/src/main/java/org/springframework/integration/flow/FlowConfiguration.java new file mode 100644 index 0000000..29ac726 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/FlowConfiguration.java @@ -0,0 +1,70 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow; + +import java.util.List; + +/** + * + * @author David Turanski + * + */ +public class FlowConfiguration { + + private final List portConfigurations; + + private volatile NamedResourceConfiguration propertiesConfiguration; + + private volatile NamedResourceConfiguration referencedBeansConfiguration; + + public FlowConfiguration(List portConfigurations) { + this.portConfigurations = portConfigurations; + } + + public FlowProviderPortConfiguration getConfigurationForInputPort(String inputPortName) { + for (FlowProviderPortConfiguration pc : portConfigurations) { + if (pc.getInputPortName().equals(inputPortName)) { + return pc; + } + } + return null; + } + + public void setPropertiesConfiguration(NamedResourceConfiguration propertiesConfiguration) { + this.propertiesConfiguration = propertiesConfiguration; + } + + public void setReferenceedBeansConfiguration(NamedResourceConfiguration referencedBeansConfiguration) { + this.setReferencedBeansConfiguration(referencedBeansConfiguration); + } + + public List getPortConfigurations() { + return portConfigurations; + } + + public NamedResourceConfiguration getPropertiesConfiguration() { + return propertiesConfiguration; + } + + public void setReferencedBeansConfiguration(NamedResourceConfiguration referencedBeansConfiguration) { + this.referencedBeansConfiguration = referencedBeansConfiguration; + } + + public NamedResourceConfiguration getReferencedBeansConfiguration() { + return referencedBeansConfiguration; + } + +} diff --git a/src/main/java/org/springframework/integration/flow/FlowProviderPortConfiguration.java b/src/main/java/org/springframework/integration/flow/FlowProviderPortConfiguration.java new file mode 100644 index 0000000..b5b53af --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/FlowProviderPortConfiguration.java @@ -0,0 +1,28 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow; + +import java.util.List; + +/** + * + * @author David Turanski + * + */ +public interface FlowProviderPortConfiguration extends PortConfiguration { + public abstract String getInputPortDescription(); + public abstract List getOutputPortMetadata(); +} \ No newline at end of file diff --git a/src/main/java/org/springframework/integration/flow/NamedResourceConfiguration.java b/src/main/java/org/springframework/integration/flow/NamedResourceConfiguration.java new file mode 100644 index 0000000..72809f1 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/NamedResourceConfiguration.java @@ -0,0 +1,65 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * + * @author David Turanski + * + */ +public class NamedResourceConfiguration { + private final List resourceMetadata; + + public NamedResourceConfiguration(List resourceMetadata) { + this.resourceMetadata = resourceMetadata; + } + + public boolean isRequired(String namedResource) { + for (NamedResourceMetadata metadata : resourceMetadata) { + if (metadata.getName().equals(namedResource)) { + return metadata.isRequired(); + } + } + return false; + } + + public List getConfiguredResources() { + return Collections.unmodifiableList(resourceMetadata); + } + + public List getRequiredResources() { + List requiredResources = new ArrayList(); + for (NamedResourceMetadata metadata : resourceMetadata) { + if (metadata.isRequired()) { + requiredResources.add(metadata); + } + } + return requiredResources; + } + + public NamedResourceMetadata find(String name) { + for (NamedResourceMetadata metadata : resourceMetadata) { + if (metadata.getName().equals(name)) { + return metadata; + } + } + return null; + } +} diff --git a/src/main/java/org/springframework/integration/flow/NamedResourceMetadata.java b/src/main/java/org/springframework/integration/flow/NamedResourceMetadata.java new file mode 100644 index 0000000..1819e20 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/NamedResourceMetadata.java @@ -0,0 +1,60 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow; + +import org.springframework.util.Assert; + +/** + * + * @author David Turanski + * + */ +public class NamedResourceMetadata { + private final String name; + + private final String description; + + private final boolean required; + + public NamedResourceMetadata(String name, String description, boolean required) { + Assert.hasText(name, "name is required"); + this.name = name; + this.description = description; + this.required = required; + } + + public String getName() { + return name; + } + + public String getDescription() { + return description; + } + + public boolean isRequired() { + return required; + } + + public String toString() { + StringBuilder sb = new StringBuilder(); + sb.append(getName()).append("\t") + .append(required? "required" : "optional") + .append("\t") + .append(getDescription()); + return sb.toString(); + } + +} diff --git a/src/main/java/org/springframework/integration/flow/PortConfiguration.java b/src/main/java/org/springframework/integration/flow/PortConfiguration.java new file mode 100644 index 0000000..813a6ea --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/PortConfiguration.java @@ -0,0 +1,34 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow; + +import java.util.Collection; + +/** + * + * @author David Turanski + * + */ +public interface PortConfiguration { + public String getInputPortName(); + + public Object getInputChannel(); + + public Collection getOutputPortNames(); + + public Object getOutputChannel(String portName); + // TODO: Do we need error channel? +} diff --git a/src/main/java/org/springframework/integration/flow/PortMetadata.java b/src/main/java/org/springframework/integration/flow/PortMetadata.java new file mode 100644 index 0000000..3eebfca --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/PortMetadata.java @@ -0,0 +1,39 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow; + +/** + * + * @author David Turanski + * + */ +public class PortMetadata extends NamedResourceMetadata { + + private final String channelName; + + public PortMetadata(String portName, String channelName) { + this(portName, "", channelName); + } + + public PortMetadata(String portName, String description, String channelName) { + super(portName, description, true); + this.channelName = channelName; + } + + public String getChannelName() { + return channelName; + } +} diff --git a/src/main/java/org/springframework/integration/flow/config/FlowMessageHandlerFactoryBean.java b/src/main/java/org/springframework/integration/flow/config/FlowMessageHandlerFactoryBean.java new file mode 100644 index 0000000..0967f7e --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/config/FlowMessageHandlerFactoryBean.java @@ -0,0 +1,112 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.beans.factory.support.BeanDefinitionRegistryPostProcessor; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.config.AbstractSimpleMessageHandlerFactoryBean; +import org.springframework.integration.core.MessageHandler; +import org.springframework.integration.core.PollableChannel; +import org.springframework.integration.flow.Flow; +import org.springframework.integration.flow.FlowProviderPortConfiguration; +import org.springframework.integration.flow.handler.FlowMessageHandler; +import org.springframework.util.Assert; + +/** + * + * @author David Turanski + * + */ +public class FlowMessageHandlerFactoryBean extends AbstractSimpleMessageHandlerFactoryBean implements + InitializingBean, BeanDefinitionRegistryPostProcessor { + + private volatile Flow flow; + + private volatile String inputPortName; + + private volatile long timeout; + + private volatile PollableChannel flowReceiveChannel; + + private volatile FlowProviderPortConfiguration flowConfiguration; + + @Override + protected MessageHandler createHandler() { + + MessageChannel flowInputChannel = flow.resolveChannelName((String) flowConfiguration.getInputChannel()); + + FlowMessageHandler flowMessageHandler = new FlowMessageHandler(flowInputChannel, flowReceiveChannel, timeout); + + return flowMessageHandler; + } + + public void setFlow(Flow flow) { + this.flow = flow; + } + + public void setInputPortName(String inputPortName) { + this.inputPortName = inputPortName; + } + + public void setTimeout(long timeout) { + this.timeout = timeout; + } + + public void setFlowOutputChannel(PollableChannel flowOutputChannel) { + this.flowReceiveChannel = flowOutputChannel; + } + + @Override + public void postProcessBeanFactory(ConfigurableListableBeanFactory beanFactory) throws BeansException { + // TODO Auto-generated method stub + + } + + @Override + public void postProcessBeanDefinitionRegistry(BeanDefinitionRegistry registry) throws BeansException { + bridgeMessagingPorts(registry); + + } + + + + private void bridgeMessagingPorts(BeanDefinitionRegistry registry) { + FlowUtils.createBridge(this.flow.getFlowOutputChannel(), this.flowReceiveChannel, registry); + } + + @Override + public void afterPropertiesSet() throws Exception { + this.flowConfiguration = null; + if (this.inputPortName == null){ + Assert.isTrue(!(this.flow.getFlowConfiguration().getPortConfigurations().size() > 1), + "flow [" + this.flow.getBeanName() +"] exposes multiple port configurations. Must specify an input port"); + + this.flowConfiguration = this.flow.getFlowConfiguration().getPortConfigurations().get(0); + this.inputPortName = this.flowConfiguration.getInputPortName(); + } + else { + this.flowConfiguration = this.flow.getFlowConfiguration().getConfigurationForInputPort( + this.inputPortName); + } + Assert.notEmpty(this.flowConfiguration.getOutputPortNames(), "flow [" + this.flow.getBeanName() + + "] has no configured output ports"); + } + +} diff --git a/src/main/java/org/springframework/integration/flow/config/FlowUtils.java b/src/main/java/org/springframework/integration/flow/config/FlowUtils.java new file mode 100644 index 0000000..bc704e4 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/config/FlowUtils.java @@ -0,0 +1,64 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +package org.springframework.integration.flow.config; + +import org.apache.commons.lang.StringUtils; +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.support.AbstractBeanDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionReaderUtils; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.config.ConsumerEndpointFactoryBean; +import org.springframework.integration.handler.BridgeHandler; + +/** + * @author David Turanski + * + */ +public class FlowUtils { + /* + * create a bridge + */ + public static void createBridge(MessageChannel inputChannel, MessageChannel outputChannel, + BeanDefinitionRegistry registry) { + + BeanDefinitionBuilder handlerBuilder = BeanDefinitionBuilder + .genericBeanDefinition(BridgeHandler.class); + + handlerBuilder.addPropertyValue("outputChannel", outputChannel); + + AbstractBeanDefinition handlerBeanDefinition = handlerBuilder.getBeanDefinition(); + + BeanDefinitionBuilder consumerEndpointBuilder = BeanDefinitionBuilder + .genericBeanDefinition(ConsumerEndpointFactoryBean.class); + + String handlerBeanName = registerBeanDefinition(handlerBeanDefinition, registry); + consumerEndpointBuilder.addPropertyReference("handler", handlerBeanName); + consumerEndpointBuilder.addPropertyValue("inputChannel", inputChannel); + registerBeanDefinition(consumerEndpointBuilder.getBeanDefinition(), registry); + } + + public static String registerBeanDefinition(BeanDefinition beanDefinition, BeanDefinitionRegistry registry){ + String beanName = BeanDefinitionReaderUtils.generateBeanName(beanDefinition, registry); + beanName = "flow."+ beanName; + String strIndex = StringUtils.substringAfter(beanName,"#"); + int index = Integer.valueOf(strIndex); + while (registry.isBeanNameInUse(beanName)){ + index++; + beanName = beanName.replaceAll("#\\d$","#"+ (index)); + } + registry.registerBeanDefinition(beanName, beanDefinition); + return beanName; + } +} diff --git a/src/main/java/org/springframework/integration/flow/config/xml/FlowConfigurationParser.java b/src/main/java/org/springframework/integration/flow/config/xml/FlowConfigurationParser.java new file mode 100644 index 0000000..b13e1bd --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/config/xml/FlowConfigurationParser.java @@ -0,0 +1,140 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import java.util.List; +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionReaderUtils; +import org.springframework.beans.factory.support.ManagedList; +import org.springframework.beans.factory.xml.BeanDefinitionParser; +import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.integration.flow.ChannelNamePortConfiguration; +import org.springframework.integration.flow.FlowConfiguration; +import org.springframework.integration.flow.NamedResourceConfiguration; +import org.springframework.integration.flow.NamedResourceMetadata; +import org.springframework.integration.flow.PortMetadata; +import org.springframework.util.xml.DomUtils; +import org.w3c.dom.Element; + +/** + * + * @author David Turanski + * + */ +public class FlowConfigurationParser implements BeanDefinitionParser { + + @Override + public BeanDefinition parse(Element element, ParserContext parserContext) { + + List portMappings = DomUtils.getChildElementsByTagName(element, "port-mapping"); + List portMappingRefs = DomUtils.getChildElementsByTagName(element, "port-mapping-ref"); + + BeanDefinitionBuilder flowConfigurationBuilder = BeanDefinitionBuilder + .genericBeanDefinition(FlowConfiguration.class); + + ManagedList portConfigList = new ManagedList(); + + for (Element el : portMappings) { + + BeanDefinition portConfiguration = buildFlowProviderPortConfiguration(el, parserContext); + portConfigList.add(portConfiguration); + } + + flowConfigurationBuilder.addConstructorArgValue(portConfigList); + + + BeanDefinition referencedProperties = this.buildNamedResourceConfiguration(element, "referenced-property"); + + flowConfigurationBuilder.addPropertyValue("propertiesConfiguration", referencedProperties); + + BeanDefinition referencedBeans = this.buildNamedResourceConfiguration(element, "referenced-bean"); + + flowConfigurationBuilder.addPropertyValue("referencedBeansConfiguration", referencedBeans); + + BeanDefinitionReaderUtils.registerWithGeneratedName(flowConfigurationBuilder.getBeanDefinition(), + parserContext.getRegistry()); + + + + + return null; + } + + private BeanDefinition buildFlowProviderPortConfiguration(Element el, ParserContext parserContext) { + Element inputPortEl = DomUtils.getChildElementByTagName(el, "input-port"); + + BeanDefinitionBuilder portConfigurationBuilder = BeanDefinitionBuilder + .genericBeanDefinition(ChannelNamePortConfiguration.class); + + BeanDefinition portMetadata = this.buildPortMetadata(el, inputPortEl); + + portConfigurationBuilder.addConstructorArgValue(portMetadata); + + List outputPortElements = DomUtils.getChildElementsByTagName(el, "output-port"); + ManagedList outputList = null; + + if (outputPortElements != null) { + + outputList = new ManagedList(); + for (Element outputPortEl : outputPortElements) { + portMetadata = this.buildPortMetadata(el, outputPortEl); + outputList.add(portMetadata); + } + } + + portConfigurationBuilder.addConstructorArgValue(outputList); + + return portConfigurationBuilder.getBeanDefinition(); + } + + private BeanDefinition buildNamedResourceConfiguration(Element parent, String elementName) { + ManagedList namedResourceList = new ManagedList(); + List namedResources = DomUtils.getChildElementsByTagName(parent, elementName); + for (Element el : namedResources ) { + BeanDefinitionBuilder namedResourceBuilder = BeanDefinitionBuilder.genericBeanDefinition(NamedResourceMetadata.class); + boolean required = ("true".equals(el.getAttribute("required"))); + String name = el.getAttribute("id"); + String description = getChildElementText(el, "description", ""); + namedResourceBuilder.addConstructorArgValue(name); + namedResourceBuilder.addConstructorArgValue(description); + namedResourceBuilder.addConstructorArgValue(required); + namedResourceList.add(namedResourceBuilder.getBeanDefinition()); + } + BeanDefinitionBuilder namedResourceConfigurationBuilder = BeanDefinitionBuilder.genericBeanDefinition(NamedResourceConfiguration.class); + namedResourceConfigurationBuilder.addConstructorArgValue(namedResourceList); + return namedResourceConfigurationBuilder.getBeanDefinition(); + } + + private String getChildElementText(Element parent, String elementName, String defaultValue ){ + String value = defaultValue; + Element child = DomUtils.getChildElementByTagName(parent, elementName); + if (child != null ) { + value = child.getTextContent(); + } + return value; + } + + private BeanDefinition buildPortMetadata(Element element, Element portElement) { + BeanDefinitionBuilder portMetadataBuilder = BeanDefinitionBuilder.genericBeanDefinition(PortMetadata.class); + portMetadataBuilder.addConstructorArgValue(portElement.getAttribute("name")); + String description = getChildElementText(portElement, "description", ""); + portMetadataBuilder.addConstructorArgValue(description); + portMetadataBuilder.addConstructorArgValue(portElement.getAttribute("channel")); + return portMetadataBuilder.getBeanDefinition(); + } + +} diff --git a/src/main/java/org/springframework/integration/flow/config/xml/FlowNamespaceHandler.java b/src/main/java/org/springframework/integration/flow/config/xml/FlowNamespaceHandler.java new file mode 100644 index 0000000..8804094 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/config/xml/FlowNamespaceHandler.java @@ -0,0 +1,35 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import org.springframework.integration.config.xml.AbstractIntegrationNamespaceHandler; + +/** + * + * @author David Turanski + * + */ +public class FlowNamespaceHandler extends AbstractIntegrationNamespaceHandler { + + @Override + public void init() { + registerBeanDefinitionParser("flow", new FlowParser()); + registerBeanDefinitionParser("outbound-gateway", new FlowOutboundGatewayParser()); + registerBeanDefinitionParser("flow-configuration", new FlowConfigurationParser()); + registerBeanDefinitionParser("port-mapping", new FlowProviderPortConfigurationParser()); + } + +} diff --git a/src/main/java/org/springframework/integration/flow/config/xml/FlowOutboundGatewayParser.java b/src/main/java/org/springframework/integration/flow/config/xml/FlowOutboundGatewayParser.java new file mode 100644 index 0000000..a059130 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/config/xml/FlowOutboundGatewayParser.java @@ -0,0 +1,53 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.config.xml.AbstractConsumerEndpointParser; +import org.springframework.integration.config.xml.IntegrationNamespaceUtils; +import org.springframework.integration.flow.config.FlowMessageHandlerFactoryBean; +import org.springframework.integration.flow.config.FlowUtils; +import org.w3c.dom.Element; + +/** + * + * @author David Turanski + * + */ +public class FlowOutboundGatewayParser extends AbstractConsumerEndpointParser { + + @Override + protected BeanDefinitionBuilder parseHandler(Element element, ParserContext parserContext) { + + BeanDefinitionBuilder flowHandlerBuilder = BeanDefinitionBuilder + .genericBeanDefinition(FlowMessageHandlerFactoryBean.class); + String flowName = element.getAttribute("flow"); + flowHandlerBuilder.addPropertyReference("flow", flowName); + + IntegrationNamespaceUtils.setValueIfAttributeDefined(flowHandlerBuilder, element, "input-port","inputPortName"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(flowHandlerBuilder, element, "timeout"); + + BeanDefinitionBuilder flowOutputChannelBuilder = BeanDefinitionBuilder.genericBeanDefinition(QueueChannel.class); + String flowOutputChannelName = + FlowUtils.registerBeanDefinition(flowOutputChannelBuilder.getBeanDefinition(), parserContext.getRegistry()); + + flowHandlerBuilder.addPropertyReference("flowOutputChannel", flowOutputChannelName); + + return flowHandlerBuilder; + } +} diff --git a/src/main/java/org/springframework/integration/flow/config/xml/FlowParser.java b/src/main/java/org/springframework/integration/flow/config/xml/FlowParser.java new file mode 100644 index 0000000..caf2cc7 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/config/xml/FlowParser.java @@ -0,0 +1,54 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.xml.BeanDefinitionParser; +import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.integration.channel.PublishSubscribeChannel; +import org.springframework.integration.config.xml.IntegrationNamespaceUtils; +import org.springframework.integration.flow.Flow; +import org.springframework.integration.flow.config.FlowUtils; +import org.w3c.dom.Element; + +/** + * + * @author David Turanski + * + */ +public class FlowParser implements BeanDefinitionParser { + + @Override + public BeanDefinition parse(Element element, ParserContext parserContext) { + BeanDefinitionBuilder flowBuilder = BeanDefinitionBuilder.genericBeanDefinition(Flow.class); + String id = element.getAttribute("id"); + BeanDefinitionBuilder flowOutputChannelBuilder = BeanDefinitionBuilder + .genericBeanDefinition(PublishSubscribeChannel.class); + String beanName = FlowUtils.registerBeanDefinition(flowOutputChannelBuilder.getBeanDefinition(), + parserContext.getRegistry()); + flowBuilder.addPropertyReference("flowOutputChannel", beanName); + + IntegrationNamespaceUtils.setValueIfAttributeDefined(flowBuilder, element, "referenced-bean-locations"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(flowBuilder, element, "properties"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(flowBuilder, element, "help"); + + BeanDefinition beanDefinition = flowBuilder.getBeanDefinition(); + parserContext.getRegistry().registerBeanDefinition(id, beanDefinition); + + return beanDefinition; + } +} diff --git a/src/main/java/org/springframework/integration/flow/config/xml/FlowProviderPortConfigurationParser.java b/src/main/java/org/springframework/integration/flow/config/xml/FlowProviderPortConfigurationParser.java new file mode 100644 index 0000000..568e583 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/config/xml/FlowProviderPortConfigurationParser.java @@ -0,0 +1,36 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.xml.BeanDefinitionParser; +import org.springframework.beans.factory.xml.ParserContext; +import org.w3c.dom.Element; + +/** + * + * @author David Turanski + * + */ +public class FlowProviderPortConfigurationParser implements BeanDefinitionParser { + + @Override + public BeanDefinition parse(Element arg0, ParserContext arg1) { + // TODO Auto-generated method stub + return null; + } + +} diff --git a/src/main/java/org/springframework/integration/flow/handler/FlowMessageHandler.java b/src/main/java/org/springframework/integration/flow/handler/FlowMessageHandler.java new file mode 100644 index 0000000..8f27511 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/handler/FlowMessageHandler.java @@ -0,0 +1,95 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.handler; + +import java.util.Collections; +import java.util.Map; +import java.util.UUID; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.MessagingException; +import org.springframework.integration.core.PollableChannel; +import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; +import org.springframework.integration.message.ErrorMessage; +import org.springframework.integration.support.MessageBuilder; + +/** + * + * @author David Turanski + * + */ +public class FlowMessageHandler extends AbstractReplyProducingMessageHandler { + + private static Log log = LogFactory.getLog(FlowMessageHandler.class); + + private final MessageChannel flowInputChannel; + + private final PollableChannel flowOutputChannel; + + private final long timeout; + + public FlowMessageHandler(MessageChannel flowInputChannel, PollableChannel flowOutputChannel, long timeout) { + this.flowInputChannel = flowInputChannel; + this.flowOutputChannel = flowOutputChannel; + this.timeout = timeout; + } + + @Override + protected Object handleRequestMessage(Message requestMessage) { + + UUID conversationId = requestMessage.getHeaders().getId(); + Map flowConversationIdHeader = Collections.singletonMap("flow.conversation.id", + (Object) conversationId); + + Message message = MessageBuilder + .fromMessage(requestMessage) + .copyHeadersIfAbsent(flowConversationIdHeader) + + .build(); + Message response = null; + + try { + + flowInputChannel.send(message); + + while ((response = flowOutputChannel.receive(timeout)) != null) { + if (conversationId.equals(response.getHeaders().get("flow.conversation.id"))) { + return response; + } else { + + if (response.getPayload() instanceof MessagingException) { + MessagingException me = (MessagingException) response.getPayload(); + log.debug("failed message: " + me.getFailedMessage()); + if (conversationId.equals(me.getFailedMessage().getHeaders().get("flow.conversation.id"))) { + return response; + } + + } + } + } + } catch (MessagingException me) { + log.debug("caught exception - failed message: " + me.getFailedMessage()); + if (conversationId.equals(me.getFailedMessage().getHeaders().get("flow.conversation.id"))) { + return new ErrorMessage(me); + } + } + return null; + } + +} diff --git a/src/main/java/org/springframework/integration/flow/interceptor/FlowInterceptor.java b/src/main/java/org/springframework/integration/flow/interceptor/FlowInterceptor.java new file mode 100644 index 0000000..4fd29f2 --- /dev/null +++ b/src/main/java/org/springframework/integration/flow/interceptor/FlowInterceptor.java @@ -0,0 +1,52 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.interceptor; + +import java.util.Collections; +import java.util.Map; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.channel.interceptor.ChannelInterceptorAdapter; +import org.springframework.integration.support.MessageBuilder; + +/** + * + * @author David Turanski + * + */ +public class FlowInterceptor extends ChannelInterceptorAdapter { + private static Log log = LogFactory.getLog(FlowInterceptor.class); + + private final String portName; + + public FlowInterceptor(String portName) { + this.portName = portName; + + } + + @Override + public Message preSend(Message message, MessageChannel channel) { + + log.debug("flow interceptor " + this.hashCode() + " received a message from port " + portName + " on channel " + + channel); + Map headersToCopy = Collections.singletonMap("flow.output.port", (Object) portName); + return MessageBuilder.fromMessage(message).copyHeadersIfAbsent(headersToCopy).build(); + + } +} diff --git a/src/main/resources/META-INF/spring.handlers b/src/main/resources/META-INF/spring.handlers new file mode 100644 index 0000000..75e92b4 --- /dev/null +++ b/src/main/resources/META-INF/spring.handlers @@ -0,0 +1 @@ +http\://www.springframework.org/schema/integration/flow=org.springframework.integration.flow.config.xml.FlowNamespaceHandler \ No newline at end of file diff --git a/src/main/resources/META-INF/spring.schemas b/src/main/resources/META-INF/spring.schemas new file mode 100644 index 0000000..eeab0be --- /dev/null +++ b/src/main/resources/META-INF/spring.schemas @@ -0,0 +1,2 @@ +http\://www.springframework.org/schema/integration/flow/spring-integration-flow-2.0.xsd=org/springframework/integration/flow/config/spring-integration-flow-2.0.xsd +http\://www.springframework.org/schema/integration/flow/spring-integration-flow.xsd=org/springframework/integration/flow/config/spring-integration-flow-2.0.xsd \ No newline at end of file diff --git a/src/main/resources/META-INF/spring.tooling b/src/main/resources/META-INF/spring.tooling new file mode 100644 index 0000000..8b6b973 --- /dev/null +++ b/src/main/resources/META-INF/spring.tooling @@ -0,0 +1,4 @@ +# Tooling related information for the integration flow namespace +http\://www.springframework.org/schema/integration/flow@name=integration flow Namespace +http\://www.springframework.org/schema/integration/flow@prefix=int-flow +http\://www.springframework.org/schema/integration/flow@icon=org/springframework/integration/flow/config/spring-integration-flow.gif diff --git a/src/main/resources/org/springframework/integration/flow/config/spring-integration-flow-2.0.xsd b/src/main/resources/org/springframework/integration/flow/config/spring-integration-flow-2.0.xsd new file mode 100644 index 0000000..65150e4 --- /dev/null +++ b/src/main/resources/org/springframework/integration/flow/config/spring-integration-flow-2.0.xsd @@ -0,0 +1,180 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/src/main/resources/org/springframework/integration/flow/config/spring-integration-flow.gif b/src/main/resources/org/springframework/integration/flow/config/spring-integration-flow.gif new file mode 100644 index 0000000..a2ba674 Binary files /dev/null and b/src/main/resources/org/springframework/integration/flow/config/spring-integration-flow.gif differ diff --git a/src/test/java/org/springframework/integration/flow/NamedResourcesConfigurationTest.groovy b/src/test/java/org/springframework/integration/flow/NamedResourcesConfigurationTest.groovy new file mode 100644 index 0000000..be65799 --- /dev/null +++ b/src/test/java/org/springframework/integration/flow/NamedResourcesConfigurationTest.groovy @@ -0,0 +1,43 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow +import org.junit.Test +/** + * + * @author David Turanski + * + */ +class NamedResourcesConfigurationTest { + def resourceMetadataList = [ + new NamedResourceMetadata('resource1','describes resource1',true), + new NamedResourceMetadata('resource2','describes resource2',false), + new NamedResourceMetadata('resource3','describes resource3',true) + ] + + @Test + public void testGetRequired(){ + def namedResourceConfiguration = new NamedResourceConfiguration(resourceMetadataList); + def required = namedResourceConfiguration.getRequiredResources() + assert required.size() == 2 + required.each { assert it.required; assert namedResourceConfiguration.isRequired(it.name) } + } + + @Test + public void testGetAll(){ + def namedResourceConfiguration = new NamedResourceConfiguration(resourceMetadataList); + assert namedResourceConfiguration.getConfiguredResources().size() == 3 + } +} diff --git a/src/test/java/org/springframework/integration/flow/config/xml/ExceptionBean.java b/src/test/java/org/springframework/integration/flow/config/xml/ExceptionBean.java new file mode 100644 index 0000000..038dd13 --- /dev/null +++ b/src/test/java/org/springframework/integration/flow/config/xml/ExceptionBean.java @@ -0,0 +1,25 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +package org.springframework.integration.flow.config.xml; + +/** + * @author David Turanski + * + */ +public class ExceptionBean { + + public void exception() { + throw new RuntimeException("exception"); + } + +} diff --git a/src/test/java/org/springframework/integration/flow/config/xml/FlowClientNamespaceTest.java b/src/test/java/org/springframework/integration/flow/config/xml/FlowClientNamespaceTest.java new file mode 100644 index 0000000..bc2ced2 --- /dev/null +++ b/src/test/java/org/springframework/integration/flow/config/xml/FlowClientNamespaceTest.java @@ -0,0 +1,95 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.core.PollableChannel; +import org.springframework.integration.message.GenericMessage; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +/** + * + * @author David Turanski + * + */ +@RunWith(SpringJUnit4ClassRunner.class) +@ContextConfiguration("classpath:/FlowClientNamespaceTest-context.xml") +public class FlowClientNamespaceTest { + + @Autowired + @Qualifier("another-input") + MessageChannel gatewayInput; + + @Autowired + @Qualifier("another-output") + PollableChannel gatewayOutput; + + + @Test + public void testGateway(){ + Message msg = new GenericMessage("hello"); + gatewayInput.send(msg); + Message reply = gatewayOutput.receive(); + assertNotNull(reply); + } + + @Autowired + @Qualifier("another-input2") + MessageChannel gatewayInput2; + + @Autowired + @Qualifier("another-output2") + PollableChannel gatewayOutput2; + + + @Test + public void testOutboundGateway(){ + Message msg1 = new GenericMessage("hello"); + Message msg2 = new GenericMessage("world"); + Message reply = null; + + gatewayInput2.send(msg1); + reply = gatewayOutput2.receive(); + assertNotNull(reply); + assertEquals("gateway-output",reply.getHeaders().get("flow.output.port")); + assertEquals("yeah!",reply.getHeaders().get("gateway")); + + gatewayInput2.send(msg2); + reply = gatewayOutput2.receive(); + assertNotNull(reply); + assertEquals("gateway-discard",reply.getHeaders().get("flow.output.port")); + assertEquals("yeah!",reply.getHeaders().get("gateway")); + + gatewayInput2.send(msg1); + reply = gatewayOutput2.receive(); + assertNotNull(reply); + assertEquals("gateway-output",reply.getHeaders().get("flow.output.port")); + assertEquals("yeah!",reply.getHeaders().get("gateway")); + + } + + + +} diff --git a/src/test/java/org/springframework/integration/flow/config/xml/FlowConfigNamespaceTest.java b/src/test/java/org/springframework/integration/flow/config/xml/FlowConfigNamespaceTest.java new file mode 100644 index 0000000..23be2da --- /dev/null +++ b/src/test/java/org/springframework/integration/flow/config/xml/FlowConfigNamespaceTest.java @@ -0,0 +1,51 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import static org.junit.Assert.*; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.integration.flow.FlowConfiguration; +import org.springframework.integration.flow.FlowProviderPortConfiguration; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +/** + * + * @author David Turanski + * + */ +@RunWith(SpringJUnit4ClassRunner.class) +@ContextConfiguration("classpath:/FlowConfigNamespaceTest-context.xml") +public class FlowConfigNamespaceTest { + @Autowired + FlowConfiguration flowConfiguration; + + @Test + public void test() { + assertNotNull(flowConfiguration.getPortConfigurations()); + assertEquals(2, flowConfiguration.getPortConfigurations().size()); + FlowProviderPortConfiguration pc0 = flowConfiguration.getPortConfigurations().get(0); + assertEquals("input", pc0.getInputPortName()); + assertEquals("subflow-input", pc0.getInputChannel()); + assertEquals("", pc0.getInputPortDescription()); + assertEquals("subflow-output", pc0.getOutputChannel("output")); + assertEquals(1, pc0.getOutputPortNames().size()); + } + +} diff --git a/src/test/java/org/springframework/integration/flow/config/xml/FlowWithErrorTest.java b/src/test/java/org/springframework/integration/flow/config/xml/FlowWithErrorTest.java new file mode 100644 index 0000000..c9b120a --- /dev/null +++ b/src/test/java/org/springframework/integration/flow/config/xml/FlowWithErrorTest.java @@ -0,0 +1,90 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; + +import org.junit.Ignore; +import org.junit.Test; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.MessagingException; +import org.springframework.integration.core.MessageHandler; +import org.springframework.integration.core.PollableChannel; +import org.springframework.integration.core.SubscribableChannel; +import org.springframework.integration.message.GenericMessage; + + +/** + * + * @author David Turanski + * + */ + +public class FlowWithErrorTest { + + @Test + public void testExceptionThrown(){ + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/FlowWithErrorTest-context.xml"); + MessageChannel input = applicationContext.getBean("inputC",MessageChannel.class); + PollableChannel output = applicationContext.getBean("outputC",PollableChannel.class); + Message msg = new GenericMessage("hello"); + input.send(msg); + + Message reply = output.receive(100); + assertNotNull(reply); + assertTrue(reply.getPayload() instanceof MessagingException); + } + + @Test + public void testDirectCallWithErrorChannel(){ + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/META-INF/spring/integration/flows/subflow5-context.xml"); + MessageChannel input = applicationContext.getBean("subflow-input",MessageChannel.class); + SubscribableChannel error = applicationContext.getBean("errorChannel",SubscribableChannel.class); + + error.subscribe(new MessageHandler() { + + @Override + public void handleMessage(Message message) throws MessagingException { + assertTrue( message.getPayload() instanceof MessagingException ); + System.out.println("got error message"); + } + }); + + Message msg = new GenericMessage("hello"); + assertTrue(input.send(msg)); + } + + + @Test + public void testWithErrorChannel(){ + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/FlowWithErrorTest-context.xml"); + MessageChannel input = applicationContext.getBean("inputC1",MessageChannel.class); + PollableChannel output = applicationContext.getBean("outputC1",PollableChannel.class); + Message msg = new GenericMessage("hello"); + input.send(msg); + + Message reply = output.receive(100); + assertNotNull(reply); + assertTrue(reply.getPayload() instanceof MessagingException); + } + + + +} diff --git a/src/test/java/org/springframework/integration/flow/config/xml/FlowWithReferencesTest.java b/src/test/java/org/springframework/integration/flow/config/xml/FlowWithReferencesTest.java new file mode 100644 index 0000000..671a0d1 --- /dev/null +++ b/src/test/java/org/springframework/integration/flow/config/xml/FlowWithReferencesTest.java @@ -0,0 +1,65 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package org.springframework.integration.flow.config.xml; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.fail; + +import org.junit.Test; +import org.springframework.beans.factory.BeanCreationException; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.core.PollableChannel; +import org.springframework.integration.message.GenericMessage; + +/** + * + * @author David Turanski + * + */ + +public class FlowWithReferencesTest { + + @Test + public void testReferencedBeanConfig(){ + ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:/FlowWithReferencesTest-context.xml"); + MessageChannel input = applicationContext.getBean("inputC",MessageChannel.class); + PollableChannel output = applicationContext.getBean("outputC",PollableChannel.class); + Message msg = new GenericMessage("hello"); + input.send(msg); + Message reply = output.receive(); + assertNotNull(reply); + assertEquals("it works!",reply.getHeaders().get("refbean.value")); + assertEquals("val1",reply.getHeaders().get("property.value.1")); + assertEquals("undefined",reply.getHeaders().get("property.value.2")); + } + + @Test + public void testMissingPropertyReference() { + try { + new ClassPathXmlApplicationContext("classpath:/FlowWithMissingReferencesTest-context.xml"); + fail("should throw exception"); + } catch (BeanCreationException e) { + if(!("Flow [subflow3] is missing required property [key1]".equals(e.getCause().getMessage() ) ) ) { + e.printStackTrace(); + fail(e.getMessage()); + } + } + } +} diff --git a/src/test/java/org/springframework/integration/flow/config/xml/RefBean.java b/src/test/java/org/springframework/integration/flow/config/xml/RefBean.java new file mode 100644 index 0000000..323a7dd --- /dev/null +++ b/src/test/java/org/springframework/integration/flow/config/xml/RefBean.java @@ -0,0 +1,30 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +package org.springframework.integration.flow.config.xml; + +/** + * @author David Turanski + * + */ +public class RefBean { + private String value; + + public void setValue(String value) { + this.value = value; + } + + public String getValue() { + return value; + } + +} diff --git a/src/test/resources/FlowClientNamespaceTest-context.xml b/src/test/resources/FlowClientNamespaceTest-context.xml new file mode 100644 index 0000000..a09588f --- /dev/null +++ b/src/test/resources/FlowClientNamespaceTest-context.xml @@ -0,0 +1,28 @@ + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/test/resources/FlowConfigNamespaceTest-context.xml b/src/test/resources/FlowConfigNamespaceTest-context.xml new file mode 100644 index 0000000..8f20764 --- /dev/null +++ b/src/test/resources/FlowConfigNamespaceTest-context.xml @@ -0,0 +1,37 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/test/resources/FlowWithErrorTest-context.xml b/src/test/resources/FlowWithErrorTest-context.xml new file mode 100644 index 0000000..4f9814c --- /dev/null +++ b/src/test/resources/FlowWithErrorTest-context.xml @@ -0,0 +1,40 @@ + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/test/resources/FlowWithMissingReferencesTest-context.xml b/src/test/resources/FlowWithMissingReferencesTest-context.xml new file mode 100644 index 0000000..7756002 --- /dev/null +++ b/src/test/resources/FlowWithMissingReferencesTest-context.xml @@ -0,0 +1,19 @@ + + + + + val1 + + + + + + diff --git a/src/test/resources/FlowWithReferencesTest-context.xml b/src/test/resources/FlowWithReferencesTest-context.xml new file mode 100644 index 0000000..6c22a69 --- /dev/null +++ b/src/test/resources/FlowWithReferencesTest-context.xml @@ -0,0 +1,29 @@ + + + + + val1 + + + + + + + + + + + + + diff --git a/src/test/resources/META-INF/spring/integration/flows/subflow1-bean-config.xml b/src/test/resources/META-INF/spring/integration/flows/subflow1-bean-config.xml new file mode 100644 index 0000000..e1cc01d --- /dev/null +++ b/src/test/resources/META-INF/spring/integration/flows/subflow1-bean-config.xml @@ -0,0 +1,75 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/test/resources/META-INF/spring/integration/flows/subflow1-context.xml b/src/test/resources/META-INF/spring/integration/flows/subflow1-context.xml new file mode 100644 index 0000000..a274a15 --- /dev/null +++ b/src/test/resources/META-INF/spring/integration/flows/subflow1-context.xml @@ -0,0 +1,41 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/test/resources/META-INF/spring/integration/flows/subflow2-context.xml b/src/test/resources/META-INF/spring/integration/flows/subflow2-context.xml new file mode 100644 index 0000000..794df5a --- /dev/null +++ b/src/test/resources/META-INF/spring/integration/flows/subflow2-context.xml @@ -0,0 +1,34 @@ + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/test/resources/META-INF/spring/integration/flows/subflow3-context.xml b/src/test/resources/META-INF/spring/integration/flows/subflow3-context.xml new file mode 100644 index 0000000..b74f256 --- /dev/null +++ b/src/test/resources/META-INF/spring/integration/flows/subflow3-context.xml @@ -0,0 +1,37 @@ + + + + + + + + + + This optionally describes the flow + + + This describes the purpose/ necessary conditions for this output + + + + + This is an example of a required property + + + This is an example of an optional property + + + + + + diff --git a/src/test/resources/META-INF/spring/integration/flows/subflow4-context.xml b/src/test/resources/META-INF/spring/integration/flows/subflow4-context.xml new file mode 100644 index 0000000..1451b8b --- /dev/null +++ b/src/test/resources/META-INF/spring/integration/flows/subflow4-context.xml @@ -0,0 +1,34 @@ + + + + + + + This optionally describes the flow + + + This describes the purpose/ necessary conditions for this output + + + + + + + + + + + + + + diff --git a/src/test/resources/META-INF/spring/integration/flows/subflow5-context.xml b/src/test/resources/META-INF/spring/integration/flows/subflow5-context.xml new file mode 100644 index 0000000..c2b8a1f --- /dev/null +++ b/src/test/resources/META-INF/spring/integration/flows/subflow5-context.xml @@ -0,0 +1,47 @@ + + + + + + + This optionally describes the flow + + + This describes the purpose/ necessary conditions for this output + + + + This describes the purpose/ necessary conditions for this output + + + + + + + + + + + + + + + + diff --git a/src/test/resources/log4j.xml b/src/test/resources/log4j.xml new file mode 100644 index 0000000..75e0f19 --- /dev/null +++ b/src/test/resources/log4j.xml @@ -0,0 +1,28 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/src/test/resources/ref-bean-config.xml b/src/test/resources/ref-bean-config.xml new file mode 100644 index 0000000..88efa47 --- /dev/null +++ b/src/test/resources/ref-bean-config.xml @@ -0,0 +1,10 @@ + + + + + + + +