INT-3573: Fix EEvalSplitter for Iterator

JIRA: https://jira.spring.io/browse/INT-3573

Fix `seen` `Queue` `NPE` for `(S)FtpInboundRemoteFileSystemSynchronizerTests`:
https://build.spring.io/browse/INT-B41-JOB1-164/test/case/155357717

**Cherry-pick to 4.0.x**
This commit is contained in:
Artem Bilan
2014-12-09 22:47:57 +02:00
parent 26dae7bae7
commit 41dcd8a453
5 changed files with 41 additions and 17 deletions

View File

@@ -16,8 +16,6 @@
package org.springframework.integration.splitter;
import java.util.List;
import org.springframework.expression.Expression;
import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor;
@@ -35,7 +33,7 @@ public class ExpressionEvaluatingSplitter extends AbstractMessageProcessingSplit
@SuppressWarnings({"unchecked", "rawtypes"})
public ExpressionEvaluatingSplitter(Expression expression) {
super(new ExpressionEvaluatingMessageProcessor(expression, List.class));
super(new ExpressionEvaluatingMessageProcessor(expression));
}
}

View File

@@ -21,4 +21,6 @@
<beans:bean id="testBean" class="org.springframework.integration.splitter.SpelSplitterIntegrationTests$TestBean"/>
<splitter input-channel="spelIteratorInput" expression="@testBean.splitIterator(payload)" output-channel="output"/>
</beans:beans>

View File

@@ -56,6 +56,9 @@ public class SpelSplitterIntegrationTests {
@Autowired
private MessageChannel iteratorInput;
@Autowired
private MessageChannel spelIteratorInput;
@Autowired
private PollableChannel output;
@@ -127,6 +130,28 @@ public class SpelSplitterIntegrationTests {
assertNull(output.receive(0));
}
@Test
public void spelIteratorSplitter() {
this.spelIteratorInput.send(new GenericMessage<String>("a,b,c,d"));
Message<?> a = output.receive(0);
Message<?> b = output.receive(0);
Message<?> c = output.receive(0);
Message<?> d = output.receive(0);
assertEquals("a", a.getPayload());
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(a).getSequenceNumber());
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(a).getSequenceSize());
assertEquals("b", b.getPayload());
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(b).getSequenceNumber());
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(b).getSequenceSize());
assertEquals("c", c.getPayload());
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(c).getSequenceNumber());
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(c).getSequenceSize());
assertEquals("d", d.getPayload());
assertEquals(new Integer(4), new IntegrationMessageHeaderAccessor(d).getSequenceNumber());
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(d).getSequenceSize());
assertNull(output.receive(0));
}
static class TestBean {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2013 the original author or authors.
* Copyright 2002-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.
@@ -34,7 +34,6 @@ import java.util.ArrayList;
import java.util.Calendar;
import java.util.Collection;
import java.util.List;
import java.util.Queue;
import org.apache.commons.net.ftp.FTPClient;
import org.apache.commons.net.ftp.FTPFile;
@@ -112,8 +111,7 @@ public class FtpInboundRemoteFileSystemSynchronizerTests {
Expression expression = expressionParser.parseExpression("#this.toUpperCase() + '.a'");
synchronizer.setLocalFilenameGeneratorExpression(expression);
FtpInboundFileSynchronizingMessageSource ms =
new FtpInboundFileSynchronizingMessageSource(synchronizer);
FtpInboundFileSynchronizingMessageSource ms = new FtpInboundFileSynchronizingMessageSource(synchronizer);
ms.setAutoCreateLocalDirectory(true);
@@ -141,7 +139,7 @@ public class FtpInboundRemoteFileSystemSynchronizerTests {
assertTrue(new File("test/A.TEST.a").exists());
assertTrue(new File("test/B.TEST.a").exists());
TestUtils.getPropertyValue(ms, "localFileListFilter.seen", Queue.class).clear();
TestUtils.getPropertyValue(ms, "localFileListFilter.seenSet", Collection.class).clear();
new File("test/A.TEST.a").delete();
new File("test/B.TEST.a").delete();
@@ -154,7 +152,7 @@ public class FtpInboundRemoteFileSystemSynchronizerTests {
public static class TestFtpSessionFactory extends AbstractFtpSessionFactory<FTPClient> {
private final Collection<Object> ftpFiles = new ArrayList<Object>();
private final Collection<FTPFile> ftpFiles = new ArrayList<FTPFile>();
private void init() {
String[] files = new File("remote-test-dir").list();
@@ -183,9 +181,10 @@ public class FtpInboundRemoteFileSystemSynchronizerTests {
String[] files = new File("remote-test-dir").list();
for (String fileName : files) {
when(ftpClient.retrieveFile(Mockito.eq("remote-test-dir/" + fileName) , Mockito.any(OutputStream.class))).thenReturn(true);
when(ftpClient.retrieveFile(Mockito.eq("remote-test-dir/" + fileName),
Mockito.any(OutputStream.class))).thenReturn(true);
}
when(ftpClient.listFiles("remote-test-dir")).thenReturn(ftpFiles.toArray(new FTPFile[]{}));
when(ftpClient.listFiles("remote-test-dir")).thenReturn(ftpFiles.toArray(new FTPFile[ftpFiles.size()]));
when(ftpClient.deleteFile(Mockito.anyString())).thenReturn(true);
return ftpClient;
} catch (Exception e) {
@@ -193,4 +192,5 @@ public class FtpInboundRemoteFileSystemSynchronizerTests {
}
}
}
}

View File

@@ -32,8 +32,8 @@ import java.io.File;
import java.io.FileInputStream;
import java.util.ArrayList;
import java.util.Calendar;
import java.util.Collection;
import java.util.List;
import java.util.Queue;
import java.util.Vector;
import org.hamcrest.Matchers;
@@ -84,7 +84,6 @@ public class SftpInboundRemoteFileSystemSynchronizerTests {
@Test
public void testCopyFileToLocalDir() throws Exception {
this.cleanup();
File localDirectoy = new File("test");
assertFalse(localDirectoy.exists());
@@ -109,8 +108,7 @@ public class SftpInboundRemoteFileSystemSynchronizerTests {
synchronizer.setFilter(filter);
synchronizer.setIntegrationEvaluationContext(ExpressionUtils.createStandardEvaluationContext());
SftpInboundFileSynchronizingMessageSource ms =
new SftpInboundFileSynchronizingMessageSource(synchronizer);
SftpInboundFileSynchronizingMessageSource ms = new SftpInboundFileSynchronizingMessageSource(synchronizer);
ms.setAutoCreateLocalDirectory(true);
ms.setLocalDirectory(localDirectoy);
ms.setBeanFactory(mock(BeanFactory.class));
@@ -136,7 +134,7 @@ public class SftpInboundRemoteFileSystemSynchronizerTests {
assertTrue(new File("test/a.test").exists());
assertTrue(new File("test/b.test").exists());
TestUtils.getPropertyValue(ms, "localFileListFilter.seen", Queue.class).clear();
TestUtils.getPropertyValue(ms, "localFileListFilter.seenSet", Collection.class).clear();
new File("test/a.test").delete();
new File("test/b.test").delete();
@@ -177,7 +175,8 @@ public class SftpInboundRemoteFileSystemSynchronizerTests {
String[] files = new File("remote-test-dir").list();
for (String fileName : files) {
when(channel.get("remote-test-dir/"+fileName)).thenReturn(new FileInputStream("remote-test-dir/" + fileName));
when(channel.get("remote-test-dir/"+fileName))
.thenReturn(new FileInputStream("remote-test-dir/" + fileName));
}
when(channel.ls("remote-test-dir")).thenReturn(sftpEntries);