diff --git a/spring-integration-java-dsl/build.gradle b/spring-integration-java-dsl/build.gradle index 0d7cbc5..659a4d4 100644 --- a/spring-integration-java-dsl/build.gradle +++ b/spring-integration-java-dsl/build.gradle @@ -24,7 +24,7 @@ compileTestJava { } ext { - embedMongoVersion = '1.43' + embedMongoVersion = '1.45' jacocoVersion = '0.7.0.201403182114' log4jVersion = '1.2.17' springIntegrationVersion = '4.0.0.RELEASE' diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlows.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlows.java index 604fedb..92a0543 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlows.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlows.java @@ -16,12 +16,15 @@ package org.springframework.integration.dsl; +import org.springframework.beans.DirectFieldAccessor; +import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.config.SourcePollingChannelAdapterFactoryBean; import org.springframework.integration.core.MessageSource; import org.springframework.integration.dsl.channel.MessageChannelSpec; +import org.springframework.integration.dsl.support.EndpointConfigurer; import org.springframework.integration.dsl.support.FixedSubscriberChannelPrototype; import org.springframework.integration.dsl.support.MessageChannelReference; -import org.springframework.integration.dsl.support.EndpointConfigurer; +import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.messaging.MessageChannel; /** @@ -74,6 +77,16 @@ public final class IntegrationFlows { .currentComponent(sourcePollingChannelAdapterFactoryBean); } + public static IntegrationFlowBuilder from(MessageProducerSupport messageProducer) { + DirectFieldAccessor dfa = new DirectFieldAccessor(messageProducer); + MessageChannel outputChannel = (MessageChannel) dfa.getPropertyValue("outputChannel"); + if (outputChannel == null) { + outputChannel = new DirectChannel(); + messageProducer.setOutputChannel(outputChannel); + } + return from(outputChannel).addComponent(messageProducer); + } + /*public static IntegrationFlowBuilder from(AbstractEndpoint endpoint) { return new IntegrationFlowBuilder(); }*/ diff --git a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java index 5ddbb78..4c032ac 100644 --- a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java +++ b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java @@ -19,6 +19,7 @@ package org.springframework.integration.dsl.test; import static org.junit.Assert.*; import java.io.File; +import java.io.FileOutputStream; import java.io.IOException; import java.util.ArrayList; import java.util.Arrays; @@ -77,7 +78,6 @@ import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; import org.springframework.integration.config.GlobalChannelInterceptor; import org.springframework.integration.core.MessageSource; -import org.springframework.integration.dsl.AggregatorSpec; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.dsl.ResequencerSpec; @@ -92,6 +92,7 @@ import org.springframework.integration.event.outbound.ApplicationEventPublishing import org.springframework.integration.file.DefaultFileNameGenerator; import org.springframework.integration.file.FileHeaders; import org.springframework.integration.file.FileWritingMessageHandler; +import org.springframework.integration.file.tail.ApacheCommonsFileTailingMessageProducer; import org.springframework.integration.handler.advice.ExpressionEvaluatingRequestHandlerAdvice; import org.springframework.integration.mongodb.store.MongoDbChannelMessageStore; import org.springframework.integration.router.MethodInvokingRouter; @@ -130,6 +131,8 @@ public class IntegrationFlowTests { private static final File tmpDir = new File(System.getProperty("java.io.tmpdir")); + private static int mongoPort; + private static MongodExecutable mongodExe; private static MongodProcess mongod; @@ -277,12 +280,17 @@ public class IntegrationFlowTests { @Qualifier("lamdasInput") private MessageChannel lamdasInput; + @Autowired + @Qualifier("tailChannel") + private PollableChannel tailChannel; + @BeforeClass public static void setup() throws IOException { + mongoPort = Network.getFreeServerPort(); mongodExe = MongodStarter.getDefaultInstance() .prepare(new MongodConfigBuilder() .version(Version.Main.PRODUCTION) - .net(new Net(12345, Network.localhostIsIPv6())) + .net(new Net(mongoPort, Network.localhostIsIPv6())) .build()); mongod = mongodExe.start(); } @@ -817,6 +825,23 @@ public class IntegrationFlowTests { assertEquals("none", receive.getPayload()); } + @Test + public void testMessageProducerFlow() throws Exception { + FileOutputStream file = new FileOutputStream(new File(tmpDir, "TailTest")); + for (int i = 0; i < 50; i++) { + file.write((i + "\n").getBytes()); + } + file.flush(); + file.close(); + + for (int i = 0; i < 50; i++) { + Message message = this.tailChannel.receive(5000); + assertNotNull(message); + assertEquals("hello " + i, message.getPayload()); + } + assertNull(this.tailChannel.receive(10)); + } + @MessagingGateway(defaultRequestChannel = "controlBus") private static interface ControlBusGateway { @@ -907,7 +932,7 @@ public class IntegrationFlowTests { @Bean public MongoDbFactory mongoDbFactory() throws Exception { - return new SimpleMongoDbFactory(new MongoClient("localhost",12345), "local"); + return new SimpleMongoDbFactory(new MongoClient("localhost",mongoPort), "local"); } @Bean @@ -1216,12 +1241,23 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow routeMultiMethodInvocationFlow() { return IntegrationFlows.from("routerMultiInput") - .route(p -> p.equals("foo") || p.equals("bar") ? new String[]{"foo", "bar"} : null, + .route(p -> p.equals("foo") || p.equals("bar") ? new String[] {"foo", "bar"} : null, s -> s.suffix("-channel") ) .get(); } + @Bean + public IntegrationFlow tailFlow() { + ApacheCommonsFileTailingMessageProducer adapter = new ApacheCommonsFileTailingMessageProducer(); + adapter.setFile(new File(tmpDir, "TailTest")); + + return IntegrationFlows.from(adapter) + .transform("hello "::concat) + .channel(MessageChannels.queue("tailChannel")) + .get(); + } + } private static class RoutingTestBean {