(S)FTP Lambdas
AMQP Lambdas Event Lambdas Feed/File Lambdas Gemfile/Groovy Lambdas * Ensure that JavaDocs are checked via ` check.dependsOn javadoc` * Add `@return` tag to the `MessageBuilder.readOnlyHeaders()`
This commit is contained in:
committed by
Artem Bilan
parent
3402a86fb5
commit
e5ae7d886d
@@ -42,7 +42,6 @@ import org.springframework.integration.file.filters.FileListFilter;
|
||||
import org.springframework.integration.file.remote.AbstractFileInfo;
|
||||
import org.springframework.integration.file.remote.MessageSessionCallback;
|
||||
import org.springframework.integration.file.remote.RemoteFileTemplate;
|
||||
import org.springframework.integration.file.remote.SessionCallback;
|
||||
import org.springframework.integration.file.remote.session.Session;
|
||||
import org.springframework.integration.file.remote.session.SessionFactory;
|
||||
import org.springframework.integration.file.support.FileExistsMode;
|
||||
@@ -520,15 +519,8 @@ public abstract class AbstractRemoteFileOutboundGateway<F> extends AbstractReply
|
||||
return doMput(requestMessage);
|
||||
}
|
||||
}
|
||||
return this.remoteFileTemplate.execute(new SessionCallback<F, Object>() {
|
||||
|
||||
@Override
|
||||
public Object doInSession(Session<F> session) throws IOException {
|
||||
return AbstractRemoteFileOutboundGateway.this.messageSessionCallback.doInSession(session,
|
||||
requestMessage);
|
||||
}
|
||||
|
||||
});
|
||||
return this.remoteFileTemplate.execute(session ->
|
||||
AbstractRemoteFileOutboundGateway.this.messageSessionCallback.doInSession(session, requestMessage));
|
||||
}
|
||||
|
||||
private Object doLs(Message<?> requestMessage) {
|
||||
@@ -537,13 +529,8 @@ public abstract class AbstractRemoteFileOutboundGateway<F> extends AbstractReply
|
||||
dir += this.remoteFileTemplate.getRemoteFileSeparator();
|
||||
}
|
||||
final String fullDir = dir;
|
||||
List<?> payload = this.remoteFileTemplate.execute(new SessionCallback<F, List<?>>() {
|
||||
|
||||
@Override
|
||||
public List<?> doInSession(Session<F> session) throws IOException {
|
||||
return AbstractRemoteFileOutboundGateway.this.ls(session, fullDir);
|
||||
}
|
||||
});
|
||||
List<?> payload = this.remoteFileTemplate.execute(session ->
|
||||
AbstractRemoteFileOutboundGateway.this.ls(session, fullDir));
|
||||
return this.getMessageBuilderFactory().withPayload(payload)
|
||||
.setHeader(FileHeaders.REMOTE_DIRECTORY, dir)
|
||||
.build();
|
||||
@@ -567,14 +554,8 @@ public abstract class AbstractRemoteFileOutboundGateway<F> extends AbstractReply
|
||||
}
|
||||
}
|
||||
else {
|
||||
payload = this.remoteFileTemplate.execute(new SessionCallback<F, File>() {
|
||||
|
||||
@Override
|
||||
public File doInSession(Session<F> session) throws IOException {
|
||||
return get(requestMessage, session, remoteDir, remoteFilePath, remoteFilename, true);
|
||||
|
||||
}
|
||||
});
|
||||
payload = this.remoteFileTemplate.execute(session1 ->
|
||||
get(requestMessage, session1, remoteDir, remoteFilePath, remoteFilename, true));
|
||||
}
|
||||
return getMessageBuilderFactory().withPayload(payload)
|
||||
.setHeader(FileHeaders.REMOTE_DIRECTORY, remoteDir)
|
||||
@@ -588,13 +569,8 @@ public abstract class AbstractRemoteFileOutboundGateway<F> extends AbstractReply
|
||||
final String remoteFilePath = this.fileNameProcessor.processMessage(requestMessage);
|
||||
final String remoteFilename = getRemoteFilename(remoteFilePath);
|
||||
final String remoteDir = getRemoteDirectory(remoteFilePath, remoteFilename);
|
||||
List<File> payload = this.remoteFileTemplate.execute(new SessionCallback<F, List<File>>() {
|
||||
|
||||
@Override
|
||||
public List<File> doInSession(Session<F> session) throws IOException {
|
||||
return mGet(requestMessage, session, remoteDir, remoteFilename);
|
||||
}
|
||||
});
|
||||
List<File> payload = this.remoteFileTemplate.execute(session ->
|
||||
mGet(requestMessage, session, remoteDir, remoteFilename));
|
||||
return this.getMessageBuilderFactory().withPayload(payload)
|
||||
.setHeader(FileHeaders.REMOTE_DIRECTORY, remoteDir)
|
||||
.setHeader(FileHeaders.REMOTE_FILE, remoteFilename)
|
||||
|
||||
@@ -40,7 +40,6 @@ import org.springframework.integration.expression.ExpressionUtils;
|
||||
import org.springframework.integration.file.filters.FileListFilter;
|
||||
import org.springframework.integration.file.filters.ReversibleFileListFilter;
|
||||
import org.springframework.integration.file.remote.RemoteFileTemplate;
|
||||
import org.springframework.integration.file.remote.SessionCallback;
|
||||
import org.springframework.integration.file.remote.session.Session;
|
||||
import org.springframework.integration.file.remote.session.SessionFactory;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
@@ -241,52 +240,40 @@ public abstract class AbstractInboundFileSynchronizer<F>
|
||||
}
|
||||
final String remoteDirectory = this.remoteDirectoryExpression.getValue(this.evaluationContext, String.class);
|
||||
try {
|
||||
int transferred = this.remoteFileTemplate.execute(new SessionCallback<F, Integer>() {
|
||||
|
||||
@Override
|
||||
public Integer doInSession(Session<F> session) throws IOException {
|
||||
F[] files = session.list(remoteDirectory);
|
||||
if (!ObjectUtils.isEmpty(files)) {
|
||||
List<F> filteredFiles = filterFiles(files);
|
||||
if (maxFetchSize >= 0 && filteredFiles.size() > maxFetchSize) {
|
||||
rollbackFromFileToListEnd(filteredFiles, filteredFiles.get(maxFetchSize));
|
||||
List<F> newList = new ArrayList<>(maxFetchSize);
|
||||
for (int i = 0; i < maxFetchSize; i++) {
|
||||
newList.add(filteredFiles.get(i));
|
||||
}
|
||||
filteredFiles = newList;
|
||||
int transferred = this.remoteFileTemplate.execute(session -> {
|
||||
F[] files = session.list(remoteDirectory);
|
||||
if (!ObjectUtils.isEmpty(files)) {
|
||||
List<F> filteredFiles = filterFiles(files);
|
||||
if (maxFetchSize >= 0 && filteredFiles.size() > maxFetchSize) {
|
||||
rollbackFromFileToListEnd(filteredFiles, filteredFiles.get(maxFetchSize));
|
||||
List<F> newList = new ArrayList<>(maxFetchSize);
|
||||
for (int i = 0; i < maxFetchSize; i++) {
|
||||
newList.add(filteredFiles.get(i));
|
||||
}
|
||||
for (F file : filteredFiles) {
|
||||
try {
|
||||
if (file != null) {
|
||||
copyFileToLocalDirectory(
|
||||
remoteDirectory, file, localDirectory,
|
||||
session);
|
||||
}
|
||||
}
|
||||
catch (RuntimeException e) {
|
||||
rollbackFromFileToListEnd(filteredFiles, file);
|
||||
throw e;
|
||||
}
|
||||
catch (IOException e) {
|
||||
rollbackFromFileToListEnd(filteredFiles, file);
|
||||
throw e;
|
||||
filteredFiles = newList;
|
||||
}
|
||||
for (F file : filteredFiles) {
|
||||
try {
|
||||
if (file != null) {
|
||||
copyFileToLocalDirectory(
|
||||
remoteDirectory, file, localDirectory,
|
||||
session);
|
||||
}
|
||||
}
|
||||
return filteredFiles.size();
|
||||
}
|
||||
else {
|
||||
return 0;
|
||||
catch (RuntimeException e1) {
|
||||
rollbackFromFileToListEnd(filteredFiles, file);
|
||||
throw e1;
|
||||
}
|
||||
catch (IOException e2) {
|
||||
rollbackFromFileToListEnd(filteredFiles, file);
|
||||
throw e2;
|
||||
}
|
||||
}
|
||||
return filteredFiles.size();
|
||||
}
|
||||
|
||||
public void rollbackFromFileToListEnd(List<F> filteredFiles, F file) {
|
||||
if (AbstractInboundFileSynchronizer.this.filter instanceof ReversibleFileListFilter) {
|
||||
((ReversibleFileListFilter<F>) AbstractInboundFileSynchronizer.this.filter)
|
||||
.rollback(file, filteredFiles);
|
||||
}
|
||||
else {
|
||||
return 0;
|
||||
}
|
||||
|
||||
});
|
||||
if (this.logger.isDebugEnabled()) {
|
||||
this.logger.debug(transferred + " files transferred");
|
||||
@@ -297,6 +284,13 @@ public abstract class AbstractInboundFileSynchronizer<F>
|
||||
}
|
||||
}
|
||||
|
||||
protected void rollbackFromFileToListEnd(List<F> filteredFiles, F file) {
|
||||
if (this.filter instanceof ReversibleFileListFilter) {
|
||||
((ReversibleFileListFilter<F>) this.filter)
|
||||
.rollback(file, filteredFiles);
|
||||
}
|
||||
}
|
||||
|
||||
protected void copyFileToLocalDirectory(String remoteDirectoryPath, F remoteFile, File localDirectory,
|
||||
Session<F> session) throws IOException {
|
||||
String remoteFileName = this.getFilename(remoteFile);
|
||||
|
||||
@@ -84,13 +84,7 @@ public class OSDelegatingFileTailingMessageProducer extends FileTailingMessagePr
|
||||
super.doStart();
|
||||
destroyProcess();
|
||||
this.command = "tail " + this.options + " " + this.getFile().getAbsolutePath();
|
||||
this.getTaskExecutor().execute(new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
runExec();
|
||||
}
|
||||
});
|
||||
this.getTaskExecutor().execute(() -> runExec());
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -145,48 +139,39 @@ public class OSDelegatingFileTailingMessageProducer extends FileTailingMessagePr
|
||||
* Runs a thread that waits for the Process result.
|
||||
*/
|
||||
private void startProcessMonitor() {
|
||||
this.getTaskExecutor().execute(new Runnable() {
|
||||
this.getTaskExecutor().execute(() -> {
|
||||
Process process = OSDelegatingFileTailingMessageProducer.this.process;
|
||||
if (process == null) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Process destroyed before starting process monitor");
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
Process process = OSDelegatingFileTailingMessageProducer.this.process;
|
||||
if (process == null) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Process destroyed before starting process monitor");
|
||||
}
|
||||
return;
|
||||
int result = Integer.MIN_VALUE;
|
||||
try {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Monitoring process " + process);
|
||||
}
|
||||
|
||||
int result = Integer.MIN_VALUE;
|
||||
try {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Monitoring process " + process);
|
||||
}
|
||||
result = process.waitFor();
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("tail process terminated with value " + result);
|
||||
}
|
||||
result = process.waitFor();
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("tail process terminated with value " + result);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
logger.error("Interrupted - stopping adapter", e);
|
||||
stop();
|
||||
}
|
||||
finally {
|
||||
destroyProcess();
|
||||
}
|
||||
if (isRunning()) {
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Restarting tail process in " + getMissingFileDelay() + " milliseconds");
|
||||
}
|
||||
getRequiredTaskScheduler().schedule(new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
runExec();
|
||||
}
|
||||
}, new Date(System.currentTimeMillis() + getMissingFileDelay()));
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
logger.error("Interrupted - stopping adapter", e);
|
||||
stop();
|
||||
}
|
||||
finally {
|
||||
destroyProcess();
|
||||
}
|
||||
if (isRunning()) {
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Restarting tail process in " + getMissingFileDelay() + " milliseconds");
|
||||
}
|
||||
getRequiredTaskScheduler().schedule((Runnable) () -> runExec(),
|
||||
new Date(System.currentTimeMillis() + getMissingFileDelay()));
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -204,35 +189,31 @@ public class OSDelegatingFileTailingMessageProducer extends FileTailingMessagePr
|
||||
return;
|
||||
}
|
||||
final BufferedReader errorReader = new BufferedReader(new InputStreamReader(process.getErrorStream()));
|
||||
this.getTaskExecutor().execute(new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
String statusMessage;
|
||||
this.getTaskExecutor().execute(() -> {
|
||||
String statusMessage;
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Reading stderr");
|
||||
}
|
||||
try {
|
||||
while ((statusMessage = errorReader.readLine()) != null) {
|
||||
publish(statusMessage);
|
||||
if (logger.isTraceEnabled()) {
|
||||
logger.trace(statusMessage);
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (IOException e1) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Reading stderr");
|
||||
logger.debug("Exception on tail error reader", e1);
|
||||
}
|
||||
}
|
||||
finally {
|
||||
try {
|
||||
while ((statusMessage = errorReader.readLine()) != null) {
|
||||
publish(statusMessage);
|
||||
if (logger.isTraceEnabled()) {
|
||||
logger.trace(statusMessage);
|
||||
}
|
||||
}
|
||||
errorReader.close();
|
||||
}
|
||||
catch (IOException e) {
|
||||
catch (IOException e2) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Exception on tail error reader", e);
|
||||
}
|
||||
}
|
||||
finally {
|
||||
try {
|
||||
errorReader.close();
|
||||
}
|
||||
catch (IOException e) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Exception while closing stderr", e);
|
||||
}
|
||||
logger.debug("Exception while closing stderr", e2);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,7 +36,6 @@ import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.integration.endpoint.SourcePollingChannelAdapter;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
import org.springframework.messaging.PollableChannel;
|
||||
import org.springframework.messaging.SubscribableChannel;
|
||||
@@ -87,15 +86,11 @@ public class FileInboundTransactionTests {
|
||||
public void testNoTx() throws Exception {
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean crash = new AtomicBoolean();
|
||||
input.subscribe(new MessageHandler() {
|
||||
|
||||
@Override
|
||||
public void handleMessage(Message<?> message) throws MessagingException {
|
||||
if (crash.get()) {
|
||||
throw new MessagingException("eek");
|
||||
}
|
||||
latch.countDown();
|
||||
input.subscribe(message -> {
|
||||
if (crash.get()) {
|
||||
throw new MessagingException("eek");
|
||||
}
|
||||
latch.countDown();
|
||||
});
|
||||
pseudoTx.start();
|
||||
File file = new File(tmpDir.getRoot(), "si-test1/foo");
|
||||
@@ -123,15 +118,11 @@ public class FileInboundTransactionTests {
|
||||
public void testTx() throws Exception {
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean crash = new AtomicBoolean();
|
||||
txInput.subscribe(new MessageHandler() {
|
||||
|
||||
@Override
|
||||
public void handleMessage(Message<?> message) throws MessagingException {
|
||||
if (crash.get()) {
|
||||
throw new MessagingException("eek");
|
||||
}
|
||||
latch.countDown();
|
||||
txInput.subscribe(message -> {
|
||||
if (crash.get()) {
|
||||
throw new MessagingException("eek");
|
||||
}
|
||||
latch.countDown();
|
||||
});
|
||||
realTx.start();
|
||||
File file = new File(tmpDir.getRoot(), "si-test2/baz");
|
||||
|
||||
@@ -142,28 +142,18 @@ public class FileReadingMessageSourceIntegrationTests {
|
||||
@Repeat(5)
|
||||
public void concurrentProcessing() throws Exception {
|
||||
CountDownLatch go = new CountDownLatch(1);
|
||||
Runnable successfulConsumer = new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
Message<File> received = pollableFileSource.receive();
|
||||
while (received == null) {
|
||||
Thread.yield();
|
||||
received = pollableFileSource.receive();
|
||||
}
|
||||
Runnable successfulConsumer = () -> {
|
||||
Message<File> received = pollableFileSource.receive();
|
||||
while (received == null) {
|
||||
Thread.yield();
|
||||
received = pollableFileSource.receive();
|
||||
}
|
||||
|
||||
};
|
||||
Runnable failingConsumer = new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
Message<File> received = pollableFileSource.receive();
|
||||
if (received != null) {
|
||||
pollableFileSource.onFailure(received);
|
||||
}
|
||||
Runnable failingConsumer = () -> {
|
||||
Message<File> received = pollableFileSource.receive();
|
||||
if (received != null) {
|
||||
pollableFileSource.onFailure(received);
|
||||
}
|
||||
|
||||
};
|
||||
CountDownLatch successfulDone = doConcurrently(3, successfulConsumer, go);
|
||||
CountDownLatch failingDone = doConcurrently(10, failingConsumer, go);
|
||||
|
||||
@@ -50,8 +50,6 @@ import org.junit.rules.TemporaryFolder;
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.integration.channel.NullChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.file.FileWritingMessageHandler.FlushPredicate;
|
||||
import org.springframework.integration.file.FileWritingMessageHandler.MessageFlushPredicate;
|
||||
import org.springframework.integration.file.support.FileExistsMode;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
@@ -381,13 +379,7 @@ public class FileWritingMessageHandlerTests {
|
||||
final String anyFilename = "fooBar.test";
|
||||
QueueChannel output = new QueueChannel();
|
||||
handler.setOutputChannel(output);
|
||||
handler.setFileNameGenerator(new FileNameGenerator() {
|
||||
|
||||
@Override
|
||||
public String generateFileName(Message<?> message) {
|
||||
return anyFilename;
|
||||
}
|
||||
});
|
||||
handler.setFileNameGenerator(message -> anyFilename);
|
||||
Message<?> message = MessageBuilder.withPayload("test").build();
|
||||
handler.handleMessage(message);
|
||||
File result = (File) output.receive(0).getPayload();
|
||||
@@ -449,13 +441,7 @@ public class FileWritingMessageHandlerTests {
|
||||
File tempFolder = this.temp.newFolder();
|
||||
FileWritingMessageHandler handler = new FileWritingMessageHandler(tempFolder);
|
||||
handler.setFileExistsMode(FileExistsMode.APPEND_NO_FLUSH);
|
||||
handler.setFileNameGenerator(new FileNameGenerator() {
|
||||
|
||||
@Override
|
||||
public String generateFileName(Message<?> message) {
|
||||
return "foo.txt";
|
||||
}
|
||||
});
|
||||
handler.setFileNameGenerator(message -> "foo.txt");
|
||||
ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler();
|
||||
taskScheduler.afterPropertiesSet();
|
||||
handler.setTaskScheduler(taskScheduler);
|
||||
@@ -487,14 +473,9 @@ public class FileWritingMessageHandlerTests {
|
||||
|
||||
handler.setFlushInterval(30000);
|
||||
final AtomicBoolean called = new AtomicBoolean();
|
||||
handler.setFlushPredicate(new MessageFlushPredicate() {
|
||||
|
||||
@Override
|
||||
public boolean shouldFlush(String fileAbsolutePath, long lastWrite, Message<?> triggerMessage) {
|
||||
called.set(true);
|
||||
return true;
|
||||
}
|
||||
|
||||
handler.setFlushPredicate((fileAbsolutePath, lastWrite, triggerMessage) -> {
|
||||
called.set(true);
|
||||
return true;
|
||||
});
|
||||
handler.handleMessage(new GenericMessage<InputStream>(new ByteArrayInputStream("box".getBytes())));
|
||||
handler.trigger(new GenericMessage<String>("foo"));
|
||||
@@ -503,14 +484,9 @@ public class FileWritingMessageHandlerTests {
|
||||
|
||||
handler.handleMessage(new GenericMessage<InputStream>(new ByteArrayInputStream("bux".getBytes())));
|
||||
called.set(false);
|
||||
handler.flushIfNeeded(new FlushPredicate() {
|
||||
|
||||
@Override
|
||||
public boolean shouldFlush(String fileAbsolutePath, long lastWrite) {
|
||||
called.set(true);
|
||||
return true;
|
||||
}
|
||||
|
||||
handler.flushIfNeeded((fileAbsolutePath, lastWrite) -> {
|
||||
called.set(true);
|
||||
return true;
|
||||
});
|
||||
assertThat(file.length(), equalTo(24L));
|
||||
assertTrue(called.get());
|
||||
|
||||
@@ -20,7 +20,6 @@ import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.io.File;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
@@ -29,8 +28,6 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.mockito.Mockito;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.stubbing.Answer;
|
||||
|
||||
import org.springframework.data.redis.core.RedisTemplate;
|
||||
import org.springframework.data.redis.serializer.StringRedisSerializer;
|
||||
@@ -81,21 +78,17 @@ public class PersistentAcceptOnceFileListFilterExternalStoreTests extends RedisA
|
||||
|
||||
store = Mockito.spy(store);
|
||||
|
||||
Mockito.doAnswer(new Answer<Object>() {
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
if (suspend.get()) {
|
||||
latch2.countDown();
|
||||
try {
|
||||
latch1.await(10, TimeUnit.SECONDS);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
Mockito.doAnswer(invocation -> {
|
||||
if (suspend.get()) {
|
||||
latch2.countDown();
|
||||
try {
|
||||
latch1.await(10, TimeUnit.SECONDS);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
return invocation.callRealMethod();
|
||||
}
|
||||
return invocation.callRealMethod();
|
||||
}).when(store).replace(Mockito.anyString(), Mockito.anyString(), Mockito.anyString());
|
||||
|
||||
final FileSystemPersistentAcceptOnceFileListFilter filter =
|
||||
@@ -114,13 +107,8 @@ public class PersistentAcceptOnceFileListFilterExternalStoreTests extends RedisA
|
||||
suspend.set(true);
|
||||
file.setLastModified(file.lastModified() + 5000L);
|
||||
|
||||
Future<Integer> result = Executors.newSingleThreadExecutor().submit(new Callable<Integer>() {
|
||||
|
||||
@Override
|
||||
public Integer call() throws Exception {
|
||||
return filter.filterFiles(new File[] {file}).size();
|
||||
}
|
||||
});
|
||||
Future<Integer> result = Executors.newSingleThreadExecutor()
|
||||
.submit(() -> filter.filterFiles(new File[] { file }).size());
|
||||
assertTrue(latch2.await(10, TimeUnit.SECONDS));
|
||||
store.put("foo:" + file.getAbsolutePath(), "43");
|
||||
latch1.countDown();
|
||||
|
||||
@@ -54,8 +54,6 @@ import org.junit.Rule;
|
||||
import org.junit.Test;
|
||||
import org.junit.rules.TemporaryFolder;
|
||||
import org.mockito.ArgumentCaptor;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.stubbing.Answer;
|
||||
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.expression.common.LiteralExpression;
|
||||
@@ -260,14 +258,10 @@ public class RemoteFileOutboundGatewayTests {
|
||||
gw.afterPropertiesSet();
|
||||
Session<?> session = mock(Session.class);
|
||||
final AtomicReference<String> args = new AtomicReference<String>();
|
||||
doAnswer(new Answer<Object>() {
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
Object[] arguments = invocation.getArguments();
|
||||
args.set((String) arguments[0] + (String) arguments[1]);
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
Object[] arguments = invocation.getArguments();
|
||||
args.set((String) arguments[0] + (String) arguments[1]);
|
||||
return null;
|
||||
}).when(session).rename(anyString(), anyString());
|
||||
when(sessionFactory.getSession()).thenReturn(session);
|
||||
Message<String> requestMessage = MessageBuilder.withPayload("foo")
|
||||
@@ -287,14 +281,10 @@ public class RemoteFileOutboundGatewayTests {
|
||||
gw.afterPropertiesSet();
|
||||
Session<?> session = mock(Session.class);
|
||||
final AtomicReference<String> args = new AtomicReference<String>();
|
||||
doAnswer(new Answer<Object>() {
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
Object[] arguments = invocation.getArguments();
|
||||
args.set((String) arguments[0] + (String) arguments[1]);
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
Object[] arguments = invocation.getArguments();
|
||||
args.set((String) arguments[0] + (String) arguments[1]);
|
||||
return null;
|
||||
}).when(session).rename(anyString(), anyString());
|
||||
when(sessionFactory.getSession()).thenReturn(session);
|
||||
Message<?> out = (Message<?>) gw.handleRequestMessage(new GenericMessage<String>("foo"));
|
||||
@@ -312,23 +302,15 @@ public class RemoteFileOutboundGatewayTests {
|
||||
gw.afterPropertiesSet();
|
||||
Session<?> session = mock(Session.class);
|
||||
final AtomicReference<String> args = new AtomicReference<String>();
|
||||
doAnswer(new Answer<Object>() {
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
Object[] arguments = invocation.getArguments();
|
||||
args.set((String) arguments[0] + (String) arguments[1]);
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
Object[] arguments = invocation.getArguments();
|
||||
args.set((String) arguments[0] + (String) arguments[1]);
|
||||
return null;
|
||||
}).when(session).rename(anyString(), anyString());
|
||||
final List<String> madeDirs = new ArrayList<String>();
|
||||
doAnswer(new Answer<Object>() {
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
madeDirs.add((String) invocation.getArguments()[0]);
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
madeDirs.add((String) invocation.getArguments()[0]);
|
||||
return null;
|
||||
}).when(session).mkdir(anyString());
|
||||
when(sessionFactory.getSession()).thenReturn(session);
|
||||
Message<String> requestMessage = MessageBuilder.withPayload("foo")
|
||||
@@ -926,13 +908,9 @@ public class RemoteFileOutboundGatewayTests {
|
||||
gw.afterPropertiesSet();
|
||||
when(sessionFactory.getSession()).thenReturn(session);
|
||||
final AtomicReference<String> written = new AtomicReference<String>();
|
||||
doAnswer(new Answer<Object>() {
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
written.set((String) invocation.getArguments()[1]);
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
written.set((String) invocation.getArguments()[1]);
|
||||
return null;
|
||||
}).when(session).write(any(InputStream.class), anyString());
|
||||
tempFolder.newFile("baz.txt");
|
||||
tempFolder.newFile("qux.txt");
|
||||
@@ -964,13 +942,9 @@ public class RemoteFileOutboundGatewayTests {
|
||||
gw.afterPropertiesSet();
|
||||
when(sessionFactory.getSession()).thenReturn(session);
|
||||
final AtomicReference<String> written = new AtomicReference<String>();
|
||||
doAnswer(new Answer<Object>() {
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
written.set((String) invocation.getArguments()[1]);
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
written.set((String) invocation.getArguments()[1]);
|
||||
return null;
|
||||
}).when(session).write(any(InputStream.class), anyString());
|
||||
tempFolder.newFile("baz.txt");
|
||||
tempFolder.newFile("qux.txt");
|
||||
|
||||
@@ -35,8 +35,6 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.mockito.Mockito;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.stubbing.Answer;
|
||||
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.expression.ExpressionParser;
|
||||
@@ -45,11 +43,11 @@ import org.springframework.expression.spel.standard.SpelExpressionParser;
|
||||
import org.springframework.integration.file.remote.session.CachingSessionFactory;
|
||||
import org.springframework.integration.file.remote.session.Session;
|
||||
import org.springframework.integration.file.remote.session.SessionFactory;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.integration.util.SimplePool;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
/**
|
||||
* @author Oleg Zhurakousky
|
||||
@@ -65,12 +63,10 @@ public class FileTransferringMessageHandlerTests {
|
||||
Session<F> session = mock(Session.class);
|
||||
|
||||
when(sf.getSession()).thenReturn(session);
|
||||
doAnswer(new Answer<Object>() {
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
String path = (String) invocation.getArguments()[1];
|
||||
assertFalse(path.startsWith("/"));
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
String path = (String) invocation.getArguments()[1];
|
||||
assertFalse(path.startsWith("/"));
|
||||
return null;
|
||||
}).when(session).rename(Mockito.anyString(), Mockito.anyString());
|
||||
ExpressionParser parser = new SpelExpressionParser();
|
||||
FileTransferringMessageHandler<F> handler = new FileTransferringMessageHandler<F>(sf);
|
||||
@@ -90,12 +86,10 @@ public class FileTransferringMessageHandlerTests {
|
||||
final AtomicReference<String> temporaryPath = new AtomicReference<String>();
|
||||
final AtomicReference<String> finalPath = new AtomicReference<String>();
|
||||
when(sf.getSession()).thenReturn(session);
|
||||
doAnswer(new Answer<Object>() {
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
temporaryPath.set((String) invocation.getArguments()[0]);
|
||||
finalPath.set((String) invocation.getArguments()[1]);
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
temporaryPath.set((String) invocation.getArguments()[0]);
|
||||
finalPath.set((String) invocation.getArguments()[1]);
|
||||
return null;
|
||||
}).when(session).rename(Mockito.anyString(), Mockito.anyString());
|
||||
FileTransferringMessageHandler<F> handler = new FileTransferringMessageHandler<F>(sf);
|
||||
handler.setRemoteDirectoryExpression(new LiteralExpression("foo"));
|
||||
@@ -115,12 +109,10 @@ public class FileTransferringMessageHandlerTests {
|
||||
Session<F> session = mock(Session.class);
|
||||
|
||||
when(sf.getSession()).thenReturn(session);
|
||||
doAnswer(new Answer<Object>() {
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
String path = (String) invocation.getArguments()[1];
|
||||
assertFalse(path.startsWith("/"));
|
||||
return null;
|
||||
}
|
||||
doAnswer(invocation -> {
|
||||
String path = (String) invocation.getArguments()[1];
|
||||
assertFalse(path.startsWith("/"));
|
||||
return null;
|
||||
}).when(session).rename(Mockito.anyString(), Mockito.anyString());
|
||||
ExpressionParser parser = new SpelExpressionParser();
|
||||
FileTransferringMessageHandler<F> handler = new FileTransferringMessageHandler<F>(sf);
|
||||
|
||||
@@ -92,12 +92,8 @@ public class CachingSessionFactoryTests {
|
||||
template.setBeanFactory(mock(BeanFactory.class));
|
||||
template.afterPropertiesSet();
|
||||
try {
|
||||
template.get(new GenericMessage<String>("foo"), new InputStreamCallback() {
|
||||
|
||||
@Override
|
||||
public void doWithInputStream(InputStream stream) throws IOException {
|
||||
throw new RuntimeException("bar");
|
||||
}
|
||||
template.get(new GenericMessage<String>("foo"), (InputStreamCallback) stream -> {
|
||||
throw new RuntimeException("bar");
|
||||
});
|
||||
fail("Expected exception");
|
||||
}
|
||||
|
||||
@@ -36,8 +36,6 @@ import org.junit.Test;
|
||||
|
||||
import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.context.ApplicationEvent;
|
||||
import org.springframework.context.ApplicationEventPublisher;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.file.tail.FileTailingMessageProducerSupport.FileTailingEvent;
|
||||
import org.springframework.messaging.Message;
|
||||
@@ -128,20 +126,10 @@ public class FileTailingMessageProducerTests {
|
||||
throws Exception {
|
||||
this.adapter = adapter;
|
||||
final List<FileTailingEvent> events = new ArrayList<FileTailingEvent>();
|
||||
adapter.setApplicationEventPublisher(new ApplicationEventPublisher() {
|
||||
|
||||
@Override
|
||||
public void publishEvent(ApplicationEvent event) {
|
||||
FileTailingEvent tailEvent = (FileTailingEvent) event;
|
||||
logger.warn(event);
|
||||
events.add(tailEvent);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publishEvent(Object event) {
|
||||
|
||||
}
|
||||
|
||||
adapter.setApplicationEventPublisher(event -> {
|
||||
FileTailingEvent tailEvent = (FileTailingEvent) event;
|
||||
logger.warn(event);
|
||||
events.add(tailEvent);
|
||||
});
|
||||
adapter.setFile(new File(testDir, "foo"));
|
||||
QueueChannel outputChannel = new QueueChannel();
|
||||
|
||||
@@ -22,7 +22,6 @@ import java.io.File;
|
||||
import java.io.FileOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.OutputStream;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Executors;
|
||||
@@ -86,26 +85,18 @@ public class TailRule extends TestWatcher {
|
||||
fos.close();
|
||||
final AtomicReference<Integer> c = new AtomicReference<Integer>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
Future<Process> future = Executors.newSingleThreadExecutor().submit(new Callable<Process>() {
|
||||
|
||||
@Override
|
||||
public Process call() throws Exception {
|
||||
final Process process = Runtime.getRuntime().exec(commandToTest + " " + file.getAbsolutePath());
|
||||
Executors.newSingleThreadExecutor().execute(new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
c.set(process.getInputStream().read());
|
||||
latch.countDown();
|
||||
}
|
||||
catch (IOException e) {
|
||||
logger.error("Error reading test stream", e);
|
||||
}
|
||||
}
|
||||
});
|
||||
return process;
|
||||
}
|
||||
Future<Process> future = Executors.newSingleThreadExecutor().submit(() -> {
|
||||
final Process process = Runtime.getRuntime().exec(commandToTest + " " + file.getAbsolutePath());
|
||||
Executors.newSingleThreadExecutor().execute(() -> {
|
||||
try {
|
||||
c.set(process.getInputStream().read());
|
||||
latch.countDown();
|
||||
}
|
||||
catch (IOException e) {
|
||||
logger.error("Error reading test stream", e);
|
||||
}
|
||||
});
|
||||
return process;
|
||||
});
|
||||
try {
|
||||
Process process = future.get(10, TimeUnit.SECONDS);
|
||||
|
||||
Reference in New Issue
Block a user