diff --git a/spring-integration-java-dsl/build.gradle b/spring-integration-java-dsl/build.gradle index a10be01..e020027 100644 --- a/spring-integration-java-dsl/build.gradle +++ b/spring-integration-java-dsl/build.gradle @@ -28,12 +28,13 @@ ext { apacheSshdVersion = '0.10.1' embedMongoVersion = '1.46.0' ftpServerVersion = '1.0.6' + hsqldbVersion = '2.3.2' jacocoVersion = '0.7.1.201405082137' jmsApiVersion = '1.1-rev-1' mailVersion = '1.4.7' slf4jVersion = '1.7.7' springIntegrationVersion = '4.0.4.RELEASE' - springBootVersion = '1.1.7.RELEASE' + springBootVersion = '1.1.8.RELEASE' linkHomepage = 'https://github.com/spring-projects/spring-integration-extensions' linkCi = 'https://build.spring.io/browse/INTEXT' @@ -96,6 +97,7 @@ dependencies { testRuntime "com.sun.mail:smtp:$mailVersion" testRuntime "com.sun.mail:pop3:$mailVersion" testRuntime "com.sun.mail:imap:$mailVersion" + testRuntime "org.hsqldb:hsqldb:$hsqldbVersion" jacoco "org.jacoco:org.jacoco.agent:$jacocoVersion:runtime" } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/Adapters.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/Adapters.java index 102af56..c88f9a9 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/Adapters.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/Adapters.java @@ -38,11 +38,13 @@ import org.springframework.integration.dsl.mail.MailSendingMessageHandlerSpec; import org.springframework.integration.dsl.sftp.Sftp; import org.springframework.integration.dsl.sftp.SftpMessageHandlerSpec; import org.springframework.integration.dsl.sftp.SftpOutboundGatewaySpec; +import org.springframework.integration.dsl.support.Function; import org.springframework.integration.file.remote.RemoteFileTemplate; import org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway; import org.springframework.integration.file.remote.session.SessionFactory; import org.springframework.integration.file.support.FileExistsMode; import org.springframework.jms.core.JmsTemplate; +import org.springframework.messaging.Message; import com.jcraft.jsch.ChannelSftp; @@ -67,6 +69,10 @@ public class Adapters { return Files.outboundAdapter(directoryExpression); } + public

FileWritingMessageHandlerSpec file(Function, ?> directoryFunction) { + return Files.outboundAdapter(directoryFunction); + } + public FileWritingMessageHandlerSpec fileGateway(File destinationDirectory) { return Files.outboundGateway(destinationDirectory); } @@ -75,6 +81,10 @@ public class Adapters { return Files.outboundGateway(directoryExpression); } + public

FileWritingMessageHandlerSpec fileGateway(Function, ?> directoryFunction) { + return Files.outboundGateway(directoryFunction); + } + public FtpMessageHandlerSpec ftp(SessionFactory sessionFactory) { return Ftp.outboundAdapter(sessionFactory); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/CorrelationHandlerSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/CorrelationHandlerSpec.java index 16685e9..065a0a7 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/CorrelationHandlerSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/CorrelationHandlerSpec.java @@ -25,7 +25,10 @@ import org.springframework.integration.aggregator.ReleaseStrategy; import org.springframework.integration.config.CorrelationStrategyFactoryBean; import org.springframework.integration.config.ReleaseStrategyFactoryBean; import org.springframework.integration.dsl.core.MessageHandlerSpec; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.expression.ValueExpression; +import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.MessageGroupStore; import org.springframework.messaging.MessageChannel; import org.springframework.scheduling.TaskScheduler; @@ -76,8 +79,13 @@ public abstract class return _this(); } - public S groupTimeoutExpression(String expression) { - this.groupTimeoutExpression = PARSER.parseExpression(expression); + public S groupTimeoutExpression(String groupTimeoutExpression) { + this.groupTimeoutExpression = PARSER.parseExpression(groupTimeoutExpression); + return _this(); + } + + public S groupTimeout(Function groupTimeoutFunction) { + this.groupTimeoutExpression = new FunctionExpression(groupTimeoutFunction); return _this(); } @@ -98,7 +106,7 @@ public abstract class public S processor(Object target, String methodName) { try { - return this.correlationStrategy(new CorrelationStrategyFactoryBean(target, methodName).getObject()) + return correlationStrategy(new CorrelationStrategyFactoryBean(target, methodName).getObject()) .releaseStrategy(new ReleaseStrategyFactoryBean(target, methodName).getObject()); } catch (Exception e) { @@ -107,7 +115,7 @@ public abstract class } public S correlationExpression(String correlationExpression) { - return this.correlationStrategy(new ExpressionEvaluatingCorrelationStrategy(correlationExpression)); + return correlationStrategy(new ExpressionEvaluatingCorrelationStrategy(correlationExpression)); } public S correlationStrategy(Object target, String methodName) { @@ -125,7 +133,7 @@ public abstract class } public S releaseExpression(String releaseExpression) { - return this.releaseStrategy(new ExpressionEvaluatingReleaseStrategy(releaseExpression)); + return releaseStrategy(new ExpressionEvaluatingReleaseStrategy(releaseExpression)); } public S releaseStrategy(Object target, String methodName) { diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/EnricherSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/EnricherSpec.java index 52dafba..ce0a622 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/EnricherSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/EnricherSpec.java @@ -20,14 +20,16 @@ import java.util.HashMap; import java.util.Map; import org.springframework.expression.Expression; -import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.dsl.core.MessageHandlerSpec; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.expression.ValueExpression; import org.springframework.integration.transformer.ContentEnricher; import org.springframework.integration.transformer.support.AbstractHeaderValueMessageProcessor; import org.springframework.integration.transformer.support.ExpressionEvaluatingHeaderValueMessageProcessor; import org.springframework.integration.transformer.support.HeaderValueMessageProcessor; import org.springframework.integration.transformer.support.StaticHeaderValueMessageProcessor; +import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.util.Assert; @@ -37,13 +39,12 @@ import org.springframework.util.Assert; */ public class EnricherSpec extends MessageHandlerSpec { - private final static SpelExpressionParser PARSER = new SpelExpressionParser(); - private final ContentEnricher enricher = new ContentEnricher(); private final Map propertyExpressions = new HashMap(); - private final Map> headerExpressions = new HashMap>(); + private final Map> headerExpressions = + new HashMap>(); EnricherSpec() { } @@ -83,6 +84,11 @@ public class EnricherSpec extends MessageHandlerSpec EnricherSpec requestPayload(Function, ?> requestPayloadFunction) { + this.enricher.setRequestPayloadExpression(new FunctionExpression>(requestPayloadFunction)); + return _this(); + } + public EnricherSpec shouldClonePayload(boolean shouldClonePayload) { this.enricher.setShouldClonePayload(shouldClonePayload); return _this(); @@ -99,6 +105,12 @@ public class EnricherSpec extends MessageHandlerSpec EnricherSpec propertyFunction(String key, Function, Object> function) { + Assert.notNull(key); + this.propertyExpressions.put(key, new FunctionExpression>(function)); + return _this(); + } + public EnricherSpec header(String name, V value) { return this.header(name, value, null); } @@ -107,22 +119,36 @@ public class EnricherSpec extends MessageHandlerSpec headerValueMessageProcessor = new StaticHeaderValueMessageProcessor(value); headerValueMessageProcessor.setOverwrite(overwrite); - return this.header(name, headerValueMessageProcessor); + return header(name, headerValueMessageProcessor); } public EnricherSpec headerExpression(String name, String expression) { - return this.headerExpression(name, expression, null); + return headerExpression(name, expression, null); } public EnricherSpec headerExpression(String name, String expression, Boolean overwrite) { + Assert.hasText(expression); + return headerExpression(name, PARSER.parseExpression(expression), overwrite); + } + + public

EnricherSpec headerFunction(String name, Function, Object> function) { + return headerFunction(name, function, null); + } + + public

EnricherSpec headerFunction(String name, Function, Object> function, Boolean overwrite) { + Assert.notNull(function); + return headerExpression(name, new FunctionExpression>(function), overwrite); + } + + private EnricherSpec headerExpression(String name, Expression expression, Boolean overwrite) { AbstractHeaderValueMessageProcessor headerValueMessageProcessor = new ExpressionEvaluatingHeaderValueMessageProcessor(expression, null); headerValueMessageProcessor.setOverwrite(overwrite); - return this.header(name, headerValueMessageProcessor); + return header(name, headerValueMessageProcessor); } public EnricherSpec header(String name, HeaderValueMessageProcessor headerValueMessageProcessor) { - Assert.notNull(name); + Assert.hasText(name); this.headerExpressions.put(name, headerValueMessageProcessor); return _this(); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/HeaderEnricherSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/HeaderEnricherSpec.java index e321f91..08a0d71 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/HeaderEnricherSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/HeaderEnricherSpec.java @@ -21,12 +21,13 @@ import java.util.Map; import java.util.Map.Entry; import org.springframework.expression.Expression; -import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.dsl.core.IntegrationComponentSpec; import org.springframework.integration.dsl.support.BeanNameMessageProcessor; +import org.springframework.integration.dsl.support.Consumer; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.dsl.support.MapBuilder; -import org.springframework.integration.dsl.support.MapBuilder.MapBuilderConfigurer; import org.springframework.integration.dsl.support.StringStringMapBuilder; import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor; import org.springframework.integration.handler.MessageProcessor; @@ -35,6 +36,7 @@ import org.springframework.integration.transformer.support.AbstractHeaderValueMe import org.springframework.integration.transformer.support.ExpressionEvaluatingHeaderValueMessageProcessor; import org.springframework.integration.transformer.support.HeaderValueMessageProcessor; import org.springframework.integration.transformer.support.StaticHeaderValueMessageProcessor; +import org.springframework.messaging.Message; import org.springframework.util.Assert; /** @@ -43,8 +45,6 @@ import org.springframework.util.Assert; */ public class HeaderEnricherSpec extends IntegrationComponentSpec { - private final static SpelExpressionParser PARSER = new SpelExpressionParser(); - private final Map> headerToAdd = new HashMap>(); private final HeaderEnricher headerEnricher = new HeaderEnricher(headerToAdd); @@ -62,19 +62,17 @@ public class HeaderEnricherSpec extends IntegrationComponentSpec messageProcessor) { this.headerEnricher.setMessageProcessor(messageProcessor); return _this(); } public HeaderEnricherSpec messageProcessor(String expression) { - return this.messageProcessor(new ExpressionEvaluatingMessageProcessor( - PARSER.parseExpression(expression))); + return messageProcessor(new ExpressionEvaluatingMessageProcessor(PARSER.parseExpression(expression))); } public HeaderEnricherSpec messageProcessor(String beanName, String methodName) { - return this.messageProcessor(new BeanNameMessageProcessor(beanName, methodName)); + return messageProcessor(new BeanNameMessageProcessor(beanName, methodName)); } public HeaderEnricherSpec headers(MapBuilder headers) { @@ -97,13 +95,14 @@ public class HeaderEnricherSpec extends IntegrationComponentSpec headers) { + Assert.notNull(headers); return headerExpressions(headers.get()); } - public HeaderEnricherSpec headerExpressions( - MapBuilderConfigurer configurer) { + public HeaderEnricherSpec headerExpressions(Consumer configurer) { + Assert.notNull(configurer); StringStringMapBuilder builder = new StringStringMapBuilder(); - configurer.configure(builder); + configurer.accept(builder); return headerExpressions(builder.get()); } @@ -116,29 +115,44 @@ public class HeaderEnricherSpec extends IntegrationComponentSpec HeaderEnricherSpec header(String name, V value) { - return this.header(name, value, null); + return header(name, value, null); } public HeaderEnricherSpec header(String name, V value, Boolean overwrite) { AbstractHeaderValueMessageProcessor headerValueMessageProcessor = new StaticHeaderValueMessageProcessor(value); headerValueMessageProcessor.setOverwrite(overwrite); - return this.header(name, headerValueMessageProcessor); + return header(name, headerValueMessageProcessor); } public HeaderEnricherSpec headerExpression(String name, String expression) { - return this.headerExpression(name, expression, null); + return headerExpression(name, expression, null); } public HeaderEnricherSpec headerExpression(String name, String expression, Boolean overwrite) { + Assert.hasText(expression); + return headerExpression(name, PARSER.parseExpression(expression), overwrite); + } + + public

HeaderEnricherSpec headerFunction(String name, Function, Object> function) { + return headerFunction(name, function, null); + } + + public

HeaderEnricherSpec headerFunction(String name, Function, Object> function, + Boolean overwrite) { + Assert.notNull(function); + return headerExpression(name, new FunctionExpression>(function), overwrite); + } + + private HeaderEnricherSpec headerExpression(String name, Expression expression, Boolean overwrite) { AbstractHeaderValueMessageProcessor headerValueMessageProcessor = new ExpressionEvaluatingHeaderValueMessageProcessor(expression, null); headerValueMessageProcessor.setOverwrite(overwrite); - return this.header(name, headerValueMessageProcessor); + return header(name, headerValueMessageProcessor); } public HeaderEnricherSpec header(String name, HeaderValueMessageProcessor headerValueMessageProcessor) { - Assert.notNull(name); + Assert.hasText(name); this.headerToAdd.put(name, headerValueMessageProcessor); return _this(); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlow.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlow.java index a0e0458..81d6bcd 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlow.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlow.java @@ -16,11 +16,10 @@ package org.springframework.integration.dsl; +import org.springframework.integration.dsl.support.Consumer; + /** * @author Artem Bilan */ -public interface IntegrationFlow { - - void define(IntegrationFlowDefinition flow); - +public interface IntegrationFlow extends Consumer> { } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java index d23ef49..abacec6 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java @@ -62,7 +62,7 @@ public final class IntegrationFlowBuilder extends IntegrationFlowDefinition flow) { + public void accept(IntegrationFlowDefinition flow) { throw new UnsupportedOperationException(); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java index aba01b8..2d1a8a7 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java @@ -24,6 +24,7 @@ import java.util.Set; import org.springframework.aop.framework.Advised; import org.springframework.aop.support.AopUtils; import org.springframework.beans.factory.BeanCreationException; +import org.springframework.expression.Expression; import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.aggregator.AbstractCorrelatingMessageHandler; import org.springframework.integration.aggregator.AggregatingMessageHandler; @@ -40,9 +41,8 @@ import org.springframework.integration.dsl.support.BeanNameMessageProcessor; import org.springframework.integration.dsl.support.Consumer; import org.springframework.integration.dsl.support.FixedSubscriberChannelPrototype; import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.dsl.support.GenericHandler; -import org.springframework.integration.dsl.support.GenericRouter; -import org.springframework.integration.dsl.support.GenericSplitter; import org.springframework.integration.dsl.support.MapBuilder; import org.springframework.integration.dsl.support.MessageChannelReference; import org.springframework.integration.expression.ControlBusMethodFilter; @@ -73,6 +73,7 @@ import org.springframework.integration.transformer.HeaderFilter; import org.springframework.integration.transformer.MessageTransformingHandler; import org.springframework.integration.transformer.MethodInvokingTransformer; import org.springframework.integration.transformer.Transformer; +import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.util.Assert; @@ -151,7 +152,7 @@ public abstract class IntegrationFlowDefinition> endpointConfigurer) { return this.handle(new ServiceActivatingHandler(new ExpressionCommandMessageProcessor( - new ControlBusMethodFilter())), endpointConfigurer); + new ControlBusMethodFilter())), endpointConfigurer); } public B transform(String expression) { @@ -241,7 +242,7 @@ public abstract class IntegrationFlowDefinition B handle(GenericHandler

handler) { - return this.handle(null, handler); + return handle(null, handler); } public

B handle(GenericHandler

handler, @@ -281,15 +282,38 @@ public abstract class IntegrationFlowDefinition(new BridgeHandler()), endpointConfigurer); } + public B delay(String groupId) { + return this.delay(groupId, (String) null); + } + + public B delay(String groupId, Consumer endpointConfigurer) { + return this.delay(groupId, (String) null, endpointConfigurer); + } + public B delay(String groupId, String expression) { return this.delay(groupId, expression, null); } - public B delay(String groupId, String expression, + public

B delay(String groupId, Function, Object> function) { + return this.delay(groupId, function, null); + } + + public

B delay(String groupId, Function, Object> function, Consumer endpointConfigurer) { + Assert.notNull(function); + return this.delay(groupId, new FunctionExpression>(function), endpointConfigurer); + } + + public B delay(String groupId, String expression, Consumer endpointConfigurer) { + return delay(groupId, + StringUtils.hasText(expression) ? PARSER.parseExpression(expression) : null, + endpointConfigurer); + } + + private B delay(String groupId, Expression expression, Consumer endpointConfigurer) { DelayHandler delayHandler = new DelayHandler(groupId); - if (StringUtils.hasText(expression)) { - delayHandler.setDelayExpression(PARSER.parseExpression(expression)); + if (expression != null) { + delayHandler.setDelayExpression(expression); } return this.register(new DelayerEndpointSpec(delayHandler), endpointConfigurer); } @@ -351,6 +375,10 @@ public abstract class IntegrationFlowDefinition>) null); + } + public B split(Consumer> endpointConfigurer) { return this.split(new DefaultMessageSplitter(), endpointConfigurer); } @@ -366,24 +394,24 @@ public abstract class IntegrationFlowDefinition> endpointConfigurer) { - return this.split(new MethodInvokingSplitter(new BeanNameMessageProcessor>(beanName, methodName)), + return this.split(new MethodInvokingSplitter(new BeanNameMessageProcessor(beanName, methodName)), endpointConfigurer); } - public

B split(Class

payloadType, GenericSplitter

splitter) { - return this.split(payloadType, splitter, null); + public

B split(Class

payloadType, Function splitter) { + return split(payloadType, splitter, null); } - public B split(GenericSplitter splitter, + public

B split(Function splitter, Consumer> endpointConfigurer) { return split(null, splitter, endpointConfigurer); } - public

B split(Class

payloadType, GenericSplitter

splitter, + public

B split(Class

payloadType, Function splitter, Consumer> endpointConfigurer) { MethodInvokingSplitter split = isLambda(splitter) ? new MethodInvokingSplitter(new LambdaMessageProcessor(splitter, payloadType)) - : new MethodInvokingSplitter(splitter, "split"); + : new MethodInvokingSplitter(splitter); return this.split(split, endpointConfigurer); } @@ -509,31 +537,31 @@ public abstract class IntegrationFlowDefinition B route(GenericRouter router) { + public B route(Function router) { return this.route(null, router); } - public B route(GenericRouter router, + public B route(Function router, Consumer> routerConfigurer) { return this.route(null, router, routerConfigurer); } - public B route(Class

payloadType, GenericRouter router) { + public B route(Class

payloadType, Function router) { return this.route(payloadType, router, null, null); } - public B route(Class

payloadType, GenericRouter router, + public B route(Class

payloadType, Function router, Consumer> routerConfigurer) { return this.route(payloadType, router, routerConfigurer, null); } - public B route(GenericRouter router, + public B route(Function router, Consumer> routerConfigurer, Consumer> endpointConfigurer) { return route(null, router, routerConfigurer, endpointConfigurer); } - public B route(Class

payloadType, GenericRouter router, + public B route(Class

payloadType, Function router, Consumer> routerConfigurer, Consumer> endpointConfigurer) { MethodInvokingRouter methodInvokingRouter = isLambda(router) diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/core/IntegrationFlowBeanPostProcessor.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/core/IntegrationFlowBeanPostProcessor.java index 636b9a9..dcc66da 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/core/IntegrationFlowBeanPostProcessor.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/core/IntegrationFlowBeanPostProcessor.java @@ -177,7 +177,7 @@ public class IntegrationFlowBeanPostProcessor implements BeanPostProcessor, Bean private Object processIntegrationFlowImpl(IntegrationFlow flow, String beanName) { IntegrationFlowBuilder flowBuilder = IntegrationFlows.from(beanName + ".input"); - flow.define(flowBuilder); + flow.accept(flowBuilder); return processStandardIntegrationFlow(flowBuilder.get(), beanName); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileTransferringMessageHandlerSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileTransferringMessageHandlerSpec.java index 9c83694..90984f5 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileTransferringMessageHandlerSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileTransferringMessageHandlerSpec.java @@ -24,12 +24,15 @@ import java.util.Collections; import org.springframework.expression.common.LiteralExpression; import org.springframework.integration.dsl.core.ComponentsRegistration; import org.springframework.integration.dsl.core.MessageHandlerSpec; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.file.DefaultFileNameGenerator; import org.springframework.integration.file.FileNameGenerator; import org.springframework.integration.file.remote.RemoteFileTemplate; import org.springframework.integration.file.remote.handler.FileTransferringMessageHandler; import org.springframework.integration.file.remote.session.SessionFactory; import org.springframework.integration.file.support.FileExistsMode; +import org.springframework.messaging.Message; import org.springframework.util.Assert; import org.springframework.util.ClassUtils; @@ -95,6 +98,11 @@ public abstract class FileTransferringMessageHandlerSpec S remoteDirectory(Function, String> remoteDirectoryFunction) { + this.target.setRemoteDirectoryExpression(new FunctionExpression>(remoteDirectoryFunction)); + return _this(); + } + public S temporaryRemoteDirectory(String temporaryRemoteDirectory) { this.target.setTemporaryRemoteDirectoryExpression(new LiteralExpression(temporaryRemoteDirectory)); return _this(); @@ -105,17 +113,24 @@ public abstract class FileTransferringMessageHandlerSpec S temporaryRemoteDirectory(Function, String> temporaryRemoteDirectoryFunction) { + this.target.setTemporaryRemoteDirectoryExpression( + new FunctionExpression>(temporaryRemoteDirectoryFunction)); + return _this(); + } + public S useTemporaryFileName(boolean useTemporaryFileName) { this.target.setUseTemporaryFileName(useTemporaryFileName); return _this(); } public S fileNameGenerator(FileNameGenerator fileNameGenerator) { + this.fileNameGenerator = fileNameGenerator; this.target.setFileNameGenerator(fileNameGenerator); return _this(); } - public S fileNameGeneratorExpression(String fileNameGeneratorExpression) { + public S fileNameExpression(String fileNameGeneratorExpression) { Assert.isNull(this.fileNameGenerator, "'fileNameGenerator' and 'fileNameGeneratorExpression' are mutually exclusive."); this.defaultFileNameGenerator = new DefaultFileNameGenerator(); diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileWritingMessageHandlerSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileWritingMessageHandlerSpec.java index 2555478..fdb4f22 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileWritingMessageHandlerSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/FileWritingMessageHandlerSpec.java @@ -16,15 +16,19 @@ package org.springframework.integration.dsl.file; +import java.io.File; import java.util.Collection; import java.util.Collections; import org.springframework.integration.dsl.core.ComponentsRegistration; import org.springframework.integration.dsl.core.MessageHandlerSpec; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.file.DefaultFileNameGenerator; import org.springframework.integration.file.FileNameGenerator; import org.springframework.integration.file.FileWritingMessageHandler; import org.springframework.integration.file.support.FileExistsMode; +import org.springframework.messaging.Message; import org.springframework.util.Assert; /** @@ -38,7 +42,7 @@ public class FileWritingMessageHandlerSpec private DefaultFileNameGenerator defaultFileNameGenerator; - FileWritingMessageHandlerSpec(java.io.File destinationDirectory) { + FileWritingMessageHandlerSpec(File destinationDirectory) { this.target = new FileWritingMessageHandler(destinationDirectory); } @@ -46,6 +50,10 @@ public class FileWritingMessageHandlerSpec this.target = new FileWritingMessageHandler(PARSER.parseExpression(directoryExpression)); } +

FileWritingMessageHandlerSpec(Function, ?> directoryFunction) { + this.target = new FileWritingMessageHandler(new FunctionExpression>(directoryFunction)); + } + FileWritingMessageHandlerSpec expectReply(boolean expectReply) { target.setExpectReply(expectReply); return _this(); @@ -72,11 +80,11 @@ public class FileWritingMessageHandlerSpec return _this(); } - public FileWritingMessageHandlerSpec fileNameGeneratorExpression(String fileNameGeneratorExpression) { + public FileWritingMessageHandlerSpec fileNameExpression(String fileNameExpression) { Assert.isNull(this.fileNameGenerator, "'fileNameGenerator' and 'fileNameGeneratorExpression' are mutually exclusive."); this.defaultFileNameGenerator = new DefaultFileNameGenerator(); - this.defaultFileNameGenerator.setExpression(fileNameGeneratorExpression); + this.defaultFileNameGenerator.setExpression(fileNameExpression); return fileNameGenerator(this.defaultFileNameGenerator); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/Files.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/Files.java index 25e5613..6dfac9a 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/Files.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/Files.java @@ -19,6 +19,9 @@ package org.springframework.integration.dsl.file; import java.io.File; import java.util.Comparator; +import org.springframework.integration.dsl.support.Function; +import org.springframework.messaging.Message; + /** * @author Artem Bilan */ @@ -41,6 +44,10 @@ public abstract class Files { return new FileWritingMessageHandlerSpec(directoryExpression).expectReply(false); } + public static

FileWritingMessageHandlerSpec outboundAdapter(Function, ?> directoryFunction) { + return new FileWritingMessageHandlerSpec(directoryFunction).expectReply(false); + } + public static FileWritingMessageHandlerSpec outboundGateway(File destinationDirectory) { return new FileWritingMessageHandlerSpec(destinationDirectory).expectReply(true); } @@ -49,6 +56,10 @@ public abstract class Files { return new FileWritingMessageHandlerSpec(directoryExpression).expectReply(true); } + public static

FileWritingMessageHandlerSpec outboundGateway(Function, ?> directoryFunction) { + return new FileWritingMessageHandlerSpec(directoryFunction).expectReply(true); + } + public static TailAdapterSpec tailAdapter(File file) { return new TailAdapterSpec().file(file); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileInboundChannelAdapterSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileInboundChannelAdapterSpec.java index 5757073..1324c63 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileInboundChannelAdapterSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileInboundChannelAdapterSpec.java @@ -22,6 +22,8 @@ import java.util.Collections; import org.springframework.integration.dsl.core.ComponentsRegistration; import org.springframework.integration.dsl.core.MessageSourceSpec; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.file.filters.FileListFilter; import org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizer; import org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizingMessageSource; @@ -62,8 +64,13 @@ public abstract class RemoteFileInboundChannelAdapterSpec localFilenameFunction) { + this.synchronizer.setLocalFilenameGeneratorExpression(new FunctionExpression(localFilenameFunction)); return _this(); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileOutboundGatewaySpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileOutboundGatewaySpec.java index 4851215..d03ea8d 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileOutboundGatewaySpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/RemoteFileOutboundGatewaySpec.java @@ -19,10 +19,13 @@ package org.springframework.integration.dsl.file; import java.io.File; import org.springframework.integration.dsl.core.MessageHandlerSpec; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.file.filters.FileListFilter; import org.springframework.integration.file.filters.RegexPatternFileListFilter; import org.springframework.integration.file.filters.SimplePatternFileListFilter; import org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway; +import org.springframework.messaging.Message; import org.springframework.util.Assert; /** @@ -59,7 +62,7 @@ public abstract class RemoteFileOutboundGatewaySpec S localDirectory(Function, String> localDirectoryFunction) { + this.target.setLocalDirectoryExpression(new FunctionExpression>(localDirectoryFunction)); + return _this(); + } + public S autoCreateLocalDirectory(boolean autoCreateLocalDirectory) { this.target.setAutoCreateLocalDirectory(autoCreateLocalDirectory); return _this(); @@ -114,8 +122,13 @@ public abstract class RemoteFileOutboundGatewaySpec S localFilename(Function, String> localFilenameFunction) { + this.target.setLocalFilenameGeneratorExpression(new FunctionExpression>(localFilenameFunction)); return _this(); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/TailAdapterSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/TailAdapterSpec.java index 4cd816c..c83b900 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/TailAdapterSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/file/TailAdapterSpec.java @@ -16,6 +16,8 @@ package org.springframework.integration.dsl.file; +import java.io.File; + import org.springframework.beans.factory.support.DefaultListableBeanFactory; import org.springframework.core.task.TaskExecutor; import org.springframework.integration.channel.NullChannel; @@ -42,7 +44,7 @@ public class TailAdapterSpec extends MessageProducerSpec S destination(Function, ?> destinationFunction) { + this.target.setDestinationExpression(new FunctionExpression>(destinationFunction)); + return _this(); + } + @Override protected JmsSendingMessageHandler doGet() { return null; diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/jms/JmsOutboundGatewaySpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/jms/JmsOutboundGatewaySpec.java index 36931bc..699ef5c 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/jms/JmsOutboundGatewaySpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/jms/JmsOutboundGatewaySpec.java @@ -24,10 +24,13 @@ import javax.jms.Destination; import org.springframework.integration.dsl.core.IntegrationComponentSpec; import org.springframework.integration.dsl.core.MessageHandlerSpec; import org.springframework.integration.dsl.support.Consumer; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.jms.JmsHeaderMapper; import org.springframework.integration.jms.JmsOutboundGateway; import org.springframework.jms.support.converter.MessageConverter; import org.springframework.jms.support.destination.DestinationResolver; +import org.springframework.messaging.Message; import org.springframework.util.Assert; /** @@ -70,6 +73,11 @@ public class JmsOutboundGatewaySpec extends MessageHandlerSpec JmsOutboundGatewaySpec requestDestination(Function, ?> destinationFunction) { + this.target.setRequestDestinationExpression(new FunctionExpression>(destinationFunction)); + return _this(); + } + public JmsOutboundGatewaySpec replyDestination(Destination destination) { this.target.setReplyDestination(destination); return _this(); @@ -85,6 +93,11 @@ public class JmsOutboundGatewaySpec extends MessageHandlerSpec JmsOutboundGatewaySpec replyDestination(Function, ?> destinationFunction) { + this.target.setReplyDestinationExpression(new FunctionExpression>(destinationFunction)); + return _this(); + } + public JmsOutboundGatewaySpec destinationResolver(DestinationResolver destinationResolver) { this.target.setDestinationResolver(destinationResolver); return _this(); diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/ImapIdleChannelAdapterSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/ImapIdleChannelAdapterSpec.java index 79d0391..9bab06f 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/ImapIdleChannelAdapterSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/ImapIdleChannelAdapterSpec.java @@ -23,13 +23,16 @@ import java.util.concurrent.Executor; import javax.mail.Authenticator; import javax.mail.Session; +import javax.mail.internet.MimeMessage; import org.aopalliance.aop.Advice; import org.springframework.integration.dsl.core.ComponentsRegistration; import org.springframework.integration.dsl.core.MessageProducerSpec; +import org.springframework.integration.dsl.support.Consumer; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.dsl.support.PropertiesBuilder; -import org.springframework.integration.dsl.support.PropertiesBuilder.PropertiesConfigurer; import org.springframework.integration.mail.ImapIdleChannelAdapter; import org.springframework.integration.mail.ImapMailReceiver; import org.springframework.integration.mail.SearchTermStrategy; @@ -37,6 +40,7 @@ import org.springframework.integration.transaction.TransactionSynchronizationFac /** * @author Gary Russell + * @author Artem Bilan */ public class ImapIdleChannelAdapterSpec extends MessageProducerSpec @@ -54,6 +58,11 @@ public class ImapIdleChannelAdapterSpec return this; } + public ImapIdleChannelAdapterSpec selector(Function selectorFunction) { + this.receiver.setSelectorExpression(new FunctionExpression(selectorFunction)); + return this; + } + public ImapIdleChannelAdapterSpec session(Session session) { this.receiver.setSession(session); return this; @@ -64,9 +73,9 @@ public class ImapIdleChannelAdapterSpec return this; } - public ImapIdleChannelAdapterSpec javaMailProperties(PropertiesConfigurer configurer) { + public ImapIdleChannelAdapterSpec javaMailProperties(Consumer configurer) { PropertiesBuilder properties = new PropertiesBuilder(); - configurer.configure(properties); + configurer.accept(properties); return javaMailProperties(properties.get()); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailHeadersBuilder.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailHeadersBuilder.java index 8696e52..4a7ec82 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailHeadersBuilder.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailHeadersBuilder.java @@ -16,8 +16,11 @@ package org.springframework.integration.dsl.mail; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.dsl.support.MapBuilder; import org.springframework.integration.mail.MailHeaders; +import org.springframework.messaging.Message; /** * @author Artem Bilan @@ -25,9 +28,6 @@ import org.springframework.integration.mail.MailHeaders; */ public class MailHeadersBuilder extends MapBuilder { - MailHeadersBuilder() { - } - public MailHeadersBuilder subject(String subject) { return put(MailHeaders.SUBJECT, subject); } @@ -36,7 +36,11 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder subjectFunction(Function, String> subject) { + return put(MailHeaders.SUBJECT, new FunctionExpression>(subject)); + } + + public MailHeadersBuilder to(String... to) { return put(MailHeaders.TO, to); } @@ -44,7 +48,11 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder toFunction(Function, String[]> to) { + return put(MailHeaders.TO, new FunctionExpression>(to)); + } + + public MailHeadersBuilder cc(String... cc) { return put(MailHeaders.CC, cc); } @@ -52,7 +60,11 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder ccFunction(Function, String[]> cc) { + return put(MailHeaders.CC, new FunctionExpression>(cc)); + } + + public MailHeadersBuilder bcc(String... bcc) { return put(MailHeaders.BCC, bcc); } @@ -60,6 +72,10 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder bccFunction(Function, String[]> bcc) { + return put(MailHeaders.BCC, new FunctionExpression>(bcc)); + } + public MailHeadersBuilder from(String from) { return put(MailHeaders.FROM, from); } @@ -68,6 +84,10 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder fromFunction(Function, String> from) { + return put(MailHeaders.FROM, new FunctionExpression>(from)); + } + public MailHeadersBuilder replyTo(String replyTo) { return put(MailHeaders.REPLY_TO, replyTo); } @@ -76,6 +96,10 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder replyToFunction(Function, String> replyTo) { + return put(MailHeaders.REPLY_TO, new FunctionExpression>(replyTo)); + } + /** * @param multipartMode header value * @return this @@ -89,6 +113,10 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder multipartModeFunction(Function, Integer> multipartMode) { + return put(MailHeaders.MULTIPART_MODE, new FunctionExpression>(multipartMode)); + } + public MailHeadersBuilder attachmentFilename(String attachmentFilename) { return put(MailHeaders.ATTACHMENT_FILENAME, attachmentFilename); } @@ -97,6 +125,10 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder attachmentFilenameFunction(Function, String> attachmentFilename) { + return put(MailHeaders.ATTACHMENT_FILENAME, new FunctionExpression>(attachmentFilename)); + } + public MailHeadersBuilder contentType(String contentType) { return put(MailHeaders.CONTENT_TYPE, contentType); } @@ -105,8 +137,15 @@ public class MailHeadersBuilder extends MapBuilder MailHeadersBuilder contentTypeFunction(Function, String> contentType) { + return put(MailHeaders.CONTENT_TYPE, new FunctionExpression>(contentType)); + } + private MailHeadersBuilder putExpression(String key, String expression) { return put(key, PARSER.parseExpression(expression)); } + MailHeadersBuilder() { + } + } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailInboundChannelAdapterSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailInboundChannelAdapterSpec.java index 3a02519..bb0c7c4 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailInboundChannelAdapterSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailInboundChannelAdapterSpec.java @@ -21,17 +21,20 @@ import java.util.Properties; import javax.mail.Authenticator; import javax.mail.Session; +import javax.mail.internet.MimeMessage; import org.springframework.integration.dsl.core.ComponentsRegistration; import org.springframework.integration.dsl.core.MessageSourceSpec; +import org.springframework.integration.dsl.support.Consumer; +import org.springframework.integration.dsl.support.Function; +import org.springframework.integration.dsl.support.FunctionExpression; import org.springframework.integration.dsl.support.PropertiesBuilder; -import org.springframework.integration.dsl.support.PropertiesBuilder.PropertiesConfigurer; import org.springframework.integration.mail.AbstractMailReceiver; import org.springframework.integration.mail.MailReceivingMessageSource; /** * @author Gary Russell - * + * @author Artem Bilan */ public abstract class MailInboundChannelAdapterSpec, R extends AbstractMailReceiver> @@ -45,6 +48,11 @@ public abstract class MailInboundChannelAdapterSpec selectorFunction) { + this.receiver.setSelectorExpression(new FunctionExpression(selectorFunction)); + return _this(); + } + public S session(Session session) { this.receiver.setSession(session); return _this(); @@ -55,9 +63,9 @@ public abstract class MailInboundChannelAdapterSpec configurer) { PropertiesBuilder properties = new PropertiesBuilder(); - configurer.configure(properties); + configurer.accept(properties); return javaMailProperties(properties.get()); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailSendingMessageHandlerSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailSendingMessageHandlerSpec.java index 37514e6..6762b65 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailSendingMessageHandlerSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/mail/MailSendingMessageHandlerSpec.java @@ -20,8 +20,8 @@ import java.util.Properties; import javax.activation.FileTypeMap; import org.springframework.integration.dsl.core.MessageHandlerSpec; +import org.springframework.integration.dsl.support.Consumer; import org.springframework.integration.dsl.support.PropertiesBuilder; -import org.springframework.integration.dsl.support.PropertiesBuilder.PropertiesConfigurer; import org.springframework.integration.mail.MailSendingMessageHandler; import org.springframework.mail.javamail.JavaMailSenderImpl; @@ -44,9 +44,9 @@ public class MailSendingMessageHandlerSpec return this; } - public MailSendingMessageHandlerSpec javaMailProperties(PropertiesConfigurer propertiesConfigurer) { + public MailSendingMessageHandlerSpec javaMailProperties(Consumer propertiesConfigurer) { PropertiesBuilder properties = new PropertiesBuilder(); - propertiesConfigurer.configure(properties); + propertiesConfigurer.accept(properties); return javaMailProperties(properties.get()); } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Consumer.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Consumer.java index 01089fe..be8719e 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Consumer.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Consumer.java @@ -19,6 +19,8 @@ package org.springframework.integration.dsl.support; /** * Implementations accept a given value and perform work on the argument. * + *

This is a copy of Java 8 {@code Consumer} interface. + * * @param the type of values to accept * * @author Jon Brisbin diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Function.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Function.java index b263598..64f841d 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Function.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Function.java @@ -20,6 +20,8 @@ package org.springframework.integration.dsl.support; * Implementations of this class perform work on the given parameter * and return a result of an optionally different type. * + *

This is a copy of Java 8 {@code Function} interface. + * * @param The type of the input to the apply operation * @param The type of the result of the apply operation * diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/FunctionExpression.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/FunctionExpression.java new file mode 100644 index 0000000..d88b509 --- /dev/null +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/FunctionExpression.java @@ -0,0 +1,177 @@ +/* + * Copyright 2014 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.dsl.support; + +import org.springframework.core.convert.TypeDescriptor; +import org.springframework.expression.EvaluationContext; +import org.springframework.expression.EvaluationException; +import org.springframework.expression.Expression; +import org.springframework.expression.TypedValue; +import org.springframework.expression.common.ExpressionUtils; +import org.springframework.expression.spel.support.StandardEvaluationContext; + +/** + * An {@link Expression} that simply invokes {@link Function#apply(Object)} on its + * provided {@link Function}. + *

+ * This is a powerful alternative to the SpEL, when Java 8 and its Lambda support is in use. + *

+ * If the target component has support for an {@link Expression} property, + * a {@link FunctionExpression} can be specified instead of a + * {@link org.springframework.expression.spel.standard.SpelExpression} + * as an alternative to evaluate the value from the Lambda, rather than runtime SpEL resolution. + *

+ * The {@link FunctionExpression} is 'read-only', hence only {@link #getValue} operations + * are allowed. + * Any {@link #setValue} operations and {@link #getValueType} related operations + * throw {@link EvaluationException}. + * + * @author Artem Bilan + */ +public class FunctionExpression implements Expression { + + private final Function function; + + private final EvaluationContext defaultContext = new StandardEvaluationContext(); + + private final EvaluationException readOnlyException; + + public FunctionExpression(Function function) { + this.function = function; + this.readOnlyException = new EvaluationException(getExpressionString(), + "FunctionExpression is a 'read only' Expression implementation"); + } + + @Override + public Object getValue() throws EvaluationException { + return this.function.apply(null); + } + + @Override + @SuppressWarnings("unchecked") + public Object getValue(Object rootObject) throws EvaluationException { + return this.function.apply((S) rootObject); + } + + @Override + public T getValue(Class desiredResultType) throws EvaluationException { + return getValue(this.defaultContext, desiredResultType); + } + + @Override + public T getValue(Object rootObject, Class desiredResultType) throws EvaluationException { + return getValue(this.defaultContext, rootObject, desiredResultType); + } + + @Override + public Object getValue(EvaluationContext context) throws EvaluationException { + return getValue(); + } + + @Override + public Object getValue(EvaluationContext context, Object rootObject) throws EvaluationException { + return getValue(rootObject); + } + + @Override + public T getValue(EvaluationContext context, Class desiredResultType) throws EvaluationException { + return ExpressionUtils.convertTypedValue(context, new TypedValue(getValue()), desiredResultType); + } + + @Override + public T getValue(EvaluationContext context, Object rootObject, Class desiredResultType) + throws EvaluationException { + return ExpressionUtils.convertTypedValue(context, new TypedValue(getValue(rootObject)), desiredResultType); + } + + @Override + public Class getValueType() throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public Class getValueType(Object rootObject) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public Class getValueType(EvaluationContext context) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public Class getValueType(EvaluationContext context, Object rootObject) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public TypeDescriptor getValueTypeDescriptor() throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public TypeDescriptor getValueTypeDescriptor(Object rootObject) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public TypeDescriptor getValueTypeDescriptor(EvaluationContext context) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public TypeDescriptor getValueTypeDescriptor(EvaluationContext context, Object rootObject) + throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public void setValue(EvaluationContext context, Object value) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public void setValue(Object rootObject, Object value) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public void setValue(EvaluationContext context, Object rootObject, Object value) throws EvaluationException { + throw this.readOnlyException; + } + + @Override + public boolean isWritable(EvaluationContext context) throws EvaluationException { + return false; + } + + @Override + public boolean isWritable(EvaluationContext context, Object rootObject) throws EvaluationException { + return false; + } + + @Override + public boolean isWritable(Object rootObject) throws EvaluationException { + return false; + } + + @Override + public String getExpressionString() { + return this.function.toString(); + } + +} diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericRouter.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericRouter.java deleted file mode 100644 index a9be8ee..0000000 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericRouter.java +++ /dev/null @@ -1,26 +0,0 @@ -/* - * Copyright 2014 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.dsl.support; - -/** - * @author Artem Bilan - */ -public interface GenericRouter { - - T route(S source); - -} diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericSplitter.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericSplitter.java deleted file mode 100644 index 18a2b6e..0000000 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericSplitter.java +++ /dev/null @@ -1,28 +0,0 @@ -/* - * Copyright 2014 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.dsl.support; - -import java.util.Collection; - -/** - * @author Artem Bilan - */ -public interface GenericSplitter { - - Collection split(T target); - -} diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/MapBuilder.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/MapBuilder.java index 01715c1..981b0e9 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/MapBuilder.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/MapBuilder.java @@ -44,10 +44,4 @@ public class MapBuilder, K, V> { return (B) this; } - public interface MapBuilderConfigurer, K, V> { - - void configure(MapBuilder builder); - - } - } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/PropertiesBuilder.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/PropertiesBuilder.java index 2af207e..17eb6e2 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/PropertiesBuilder.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/PropertiesBuilder.java @@ -36,10 +36,4 @@ public class PropertiesBuilder { return this.properties; } - public interface PropertiesConfigurer { - - void configure(PropertiesBuilder propertiesBuilder); - - } - } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Transformers.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Transformers.java index bf66a6d..5225d2b 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Transformers.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/Transformers.java @@ -26,7 +26,6 @@ import org.springframework.core.io.Resource; import org.springframework.core.serializer.Deserializer; import org.springframework.core.serializer.Serializer; import org.springframework.expression.Expression; -import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.dsl.support.tuple.Tuple2; import org.springframework.integration.file.transformer.FileToByteArrayTransformer; import org.springframework.integration.file.transformer.FileToStringTransformer; @@ -60,8 +59,6 @@ import org.springframework.xml.xpath.NodeMapper; */ public abstract class Transformers { - private final static SpelExpressionParser PARSER = new SpelExpressionParser(); - public static ObjectToStringTransformer objectToString() { return objectToString(null); } @@ -325,12 +322,14 @@ public abstract class Transformers { return transformer; } - public static XsltPayloadTransformer xslt(Resource xsltTemplate, Tuple2... xslParameterMappings) { + @SuppressWarnings("unchecked") + public static XsltPayloadTransformer xslt(Resource xsltTemplate, + Tuple2... xslParameterMappings) { XsltPayloadTransformer transformer = new XsltPayloadTransformer(xsltTemplate); if (xslParameterMappings != null) { Map params = new HashMap(xslParameterMappings.length); - for (Tuple2 mapping : xslParameterMappings) { - params.put(mapping.getT1(), PARSER.parseExpression(mapping.getT2())); + for (Tuple2 mapping : xslParameterMappings) { + params.put(mapping.getT1(), mapping.getT2()); } transformer.setXslParameterMappings(params); } diff --git a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java index 81d76c7..5cf2fba 100644 --- a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java +++ b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java @@ -110,11 +110,10 @@ import org.springframework.integration.dsl.MessagingGateways; import org.springframework.integration.dsl.amqp.Amqp; import org.springframework.integration.dsl.channel.DirectChannelSpec; import org.springframework.integration.dsl.channel.MessageChannels; -import org.springframework.integration.dsl.file.Files; +import org.springframework.integration.dsl.core.Pollers; import org.springframework.integration.dsl.ftp.Ftp; import org.springframework.integration.dsl.jms.Jms; import org.springframework.integration.dsl.sftp.Sftp; -import org.springframework.integration.dsl.core.Pollers; import org.springframework.integration.dsl.support.Transformers; import org.springframework.integration.dsl.test.TestFtpServer; import org.springframework.integration.dsl.test.TestSftpServer; @@ -1422,7 +1421,7 @@ public class IntegrationFlowTests { .preserveTimestamp(true) .remoteDirectory("ftpSource") .regexFilter(".*\\.txt$") - .localFilenameGeneratorExpression("#this.toUpperCase() + '.a'") + .localFilename(f -> f.toUpperCase() + ".a") .localDirectory(this.ftpServer.getTargetLocalDirectory()), e -> e.id("ftpInboundAdapter")) .channel(MessageChannels.queue("ftpInboundResultChannel")) @@ -1436,7 +1435,7 @@ public class IntegrationFlowTests { .preserveTimestamp(true) .remoteDirectory("sftpSource") .regexFilter(".*\\.txt$") - .localFilenameGeneratorExpression("#this.toUpperCase() + '.a'") + .localFilenameExpression("#this.toUpperCase() + '.a'") .localDirectory(this.sftpServer.getTargetLocalDirectory()), e -> e.id("sftpInboundAdapter")) .channel(MessageChannels.queue("sftpInboundResultChannel")) @@ -1473,8 +1472,8 @@ public class IntegrationFlowTests { .options(AbstractRemoteFileOutboundGateway.Option.RECURSIVE) .regexFileNameFilter("(subFtpSource|.*1.txt)") .localDirectoryExpression("@ftpServer.targetLocalDirectoryName + #remoteDirectory") - .localFilenameGeneratorExpression("#remoteFileName.replaceFirst('ftpSource', " + - "'localTarget')").get(); + .localFilenameExpression("#remoteFileName.replaceFirst('ftpSource', 'localTarget')") + .get(); } @Bean @@ -1494,8 +1493,7 @@ public class IntegrationFlowTests { .options(AbstractRemoteFileOutboundGateway.Option.RECURSIVE) .regexFileNameFilter("(subSftpSource|.*1.txt)") .localDirectoryExpression("@sftpServer.targetLocalDirectoryName + #remoteDirectory") - .localFilenameGeneratorExpression( - "#remoteFileName.replaceFirst('sftpSource', 'localTarget')")) + .localFilenameExpression("#remoteFileName.replaceFirst('sftpSource', 'localTarget')")) .channel(remoteFileOutputChannel()) .get(); } @@ -1733,7 +1731,7 @@ public class IntegrationFlowTests { .requestPayloadExpression("payload") .shouldClonePayload(false) .propertyExpression("name", "payload['name']") - .propertyExpression("date", "new java.util.Date()") + .propertyFunction("date", m -> new Date()) .headerExpression("foo", "payload['name']") ) .get(); @@ -1757,10 +1755,9 @@ public class IntegrationFlowTests { public IntegrationFlow enricherFlow3() { return IntegrationFlows.from("enricherInput3", true) .enrich(e -> e.requestChannel("enrichChannel") - .requestPayloadExpression("payload") - .shouldClonePayload(false) - .headerExpression("foo", "payload['name']") - ) + .requestPayload(Message::getPayload) + .shouldClonePayload(false) + .>headerFunction("foo", m -> m.getPayload().get("name"))) .get(); } @@ -1799,7 +1796,8 @@ public class IntegrationFlowTests { .split(s -> s.applySequence(false).get().getT2().setDelimiters(",")) .channel(c -> c.executor(this.taskExecutor())) .transform(Integer::parseInt) - .enrichHeaders(s -> s.headerExpression(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, "payload")) + .enrichHeaders(h -> + h.headerFunction(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, Message::getPayload)) .resequence(r -> r.releasePartialSequences(true).correlationExpression("'foo'"), null) .headerFilter("foo", false); } @@ -1807,7 +1805,7 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow splitAggregateFlow() { return IntegrationFlows.from("splitAggregateInput", true) - .split(null) + .split() .channel(MessageChannels.executor(this.taskExecutor())) .resequence() .aggregate() @@ -1967,7 +1965,7 @@ public class IntegrationFlowTests { e -> e.poller(Pollers.fixedDelay(100))) .transform(Transformers.fileToString()) .aggregate(a -> a.correlationExpression("1") - .releaseExpression("size() == 25"), null) + .releaseStrategy(g -> g.size() == 25), null) .channel(MessageChannels.queue("fileReadingResultChannel")) .get(); } @@ -1977,7 +1975,7 @@ public class IntegrationFlowTests { return IntegrationFlows.from("fileWritingInput") .enrichHeaders(h -> h.header(FileHeaders.FILENAME, "foo.sitest") .header("directory", new File(tmpDir, "fileWritingFlow"))) - .handle(Files.outboundGateway("headers[directory]")) + .handleWithAdapter(a -> a.fileGateway(m -> m.getHeaders().get("directory"))) .channel(MessageChannels.queue("fileWritingResultChannel")) .get(); } diff --git a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/jdbc/JdbcTests.java b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/jdbc/JdbcTests.java new file mode 100644 index 0000000..b63537b --- /dev/null +++ b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/jdbc/JdbcTests.java @@ -0,0 +1,139 @@ +/* + * Copyright 2014 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.dsl.test.jdbc; + +import static org.junit.Assert.assertNotNull; + +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.util.Iterator; + +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.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.jdbc.InvalidResultSetAccessException; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.core.RowMapper; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.PollableChannel; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +/** + * @author Artem Bilan + */ +@ContextConfiguration +@RunWith(SpringJUnit4ClassRunner.class) +@DirtiesContext +public class JdbcTests { + + @Autowired + @Qualifier("jdbcSplitter.input") + private MessageChannel jdbcSplitterChannel; + + @Autowired + private PollableChannel splitResultsChannel; + + @Test + public void testJdbcSplitter() { + this.jdbcSplitterChannel.send(new GenericMessage<>("foo")); + for (int i = 0; i < 10; i++) { + Message result = this.splitResultsChannel.receive(1000); + assertNotNull(result); + } + } + + @Configuration + @EnableAutoConfiguration + public static class ContextConfiguration { + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Bean + public IntegrationFlow jdbcSplitter() { + return f -> + f.split(p -> + jdbcTemplate.execute("SELECT * from FOO", + (PreparedStatement ps) -> + new ResultSetIterator(ps.executeQuery(), + (rs1, rowNum) -> + new Foo(rs1.getInt(1), rs1.getString(2)))) + , null) + .channel(c -> c.queue("splitResultsChannel")); + } + + } + + private static class Foo { + + private final int id; + + private final String name; + + private Foo(int id, String name) { + this.id = id; + this.name = name; + } + + } + + private static class ResultSetIterator implements Iterator { + + private final ResultSet rs; + + private final RowMapper rowMapper; + + private ResultSetIterator(ResultSet rs, RowMapper rowMapper) { + this.rs = rs; + this.rowMapper = rowMapper; + } + + @Override + public boolean hasNext() { + try { + return !this.rs.isLast(); + } + catch (SQLException e) { + throw new InvalidResultSetAccessException(e); + } + } + + @Override + public T next() { + try { + this.rs.next(); + return this.rowMapper.mapRow(this.rs, this.rs.getRow()); + } + catch (SQLException e) { + throw new InvalidResultSetAccessException(e); + } + } + + } + +} diff --git a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/mail/MailTests.java b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/mail/MailTests.java index 94f87d6..fe39f42 100644 --- a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/mail/MailTests.java +++ b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/mail/MailTests.java @@ -46,12 +46,12 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.dsl.HeaderEnricherSpec; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.dsl.MessageProducers; import org.springframework.integration.dsl.channel.MessageChannels; import org.springframework.integration.dsl.mail.Mail; -import org.springframework.integration.dsl.core.Pollers; import org.springframework.integration.dsl.test.mail.PoorMansMailServer.ImapServer; import org.springframework.integration.dsl.test.mail.PoorMansMailServer.Pop3Server; import org.springframework.integration.dsl.test.mail.PoorMansMailServer.SmtpServer; @@ -197,7 +197,10 @@ public class MailTests { @Bean public IntegrationFlow sendMailFlow() { return IntegrationFlows.from("sendMailChannel") - .enrichHeaders(Mail.headers().subject("foo").from("foo@bar").to("bar@baz")) + .enrichHeaders(Mail.headers() + .subjectFunction(m -> "foo") + .from("foo@bar") + .toFunction(m -> new String[] {"bar@baz"})) .handleWithAdapter(h -> h.mail("localhost") .port(smtpPort) .credentials("user", "pw") diff --git a/spring-integration-java-dsl/src/test/resources/data.sql b/spring-integration-java-dsl/src/test/resources/data.sql new file mode 100644 index 0000000..715c1c4 --- /dev/null +++ b/spring-integration-java-dsl/src/test/resources/data.sql @@ -0,0 +1,10 @@ +INSERT INTO FOO VALUES (1, 'foo1'); +INSERT INTO FOO VALUES (2, 'foo2'); +INSERT INTO FOO VALUES (3, 'foo3'); +INSERT INTO FOO VALUES (4, 'foo4'); +INSERT INTO FOO VALUES (5, 'foo5'); +INSERT INTO FOO VALUES (6, 'foo6'); +INSERT INTO FOO VALUES (7, 'foo7'); +INSERT INTO FOO VALUES (8, 'foo8'); +INSERT INTO FOO VALUES (9, 'foo9'); +INSERT INTO FOO VALUES (10, 'foo10'); diff --git a/spring-integration-java-dsl/src/test/resources/schema.sql b/spring-integration-java-dsl/src/test/resources/schema.sql new file mode 100644 index 0000000..38de881 --- /dev/null +++ b/spring-integration-java-dsl/src/test/resources/schema.sql @@ -0,0 +1,4 @@ +CREATE TABLE FOO ( + id INTEGER IDENTITY PRIMARY KEY, + name VARCHAR(30), +); \ No newline at end of file