This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 60615fa0ba2c CAMEL-25016: camel-file - Retry a file whose download
failed also with noop or idempotent consumers (#26881)
60615fa0ba2c is described below
commit 60615fa0ba2c5dd50ae4c2cfdf2d769ce269549f
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 13:44:59 2026 +0530
CAMEL-25016: camel-file - Retry a file whose download failed also with noop
or idempotent consumers (#26881)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../camel/component/file/GenericFileConsumer.java | 69 ++++++-
.../file/FileConsumerRetrieveFailureTest.java | 203 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 13 ++
3 files changed, 284 insertions(+), 1 deletion(-)
diff --git
a/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileConsumer.java
b/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileConsumer.java
index f42276871293..6a6039e0b85b 100644
---
a/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileConsumer.java
+++
b/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileConsumer.java
@@ -468,7 +468,7 @@ public abstract class GenericFileConsumer<T> extends
ScheduledBatchPollingConsum
final String name = target.getAbsoluteFilePath();
try {
if (isRetrieveFile()) {
- if (!tryRetrievingFile(exchange, name, target,
absoluteFileName, file)) {
+ if (!retrieveFileOrAbort(exchange, name, target,
absoluteFileName, file)) {
return false;
}
} else {
@@ -513,6 +513,73 @@ public abstract class GenericFileConsumer<T> extends
ScheduledBatchPollingConsum
return true;
}
+ /**
+ * Retrieves the file, and aborts if the file could not be retrieved (or
cannot be retrieved and is ignored).
+ * <p/>
+ * The read lock has already been acquired by the begin strategy, and no
{@link GenericFileOnCompletion} is
+ * registered yet, so the abort must release the read lock here. The
idempotent key, when it was added eagerly, is
+ * removed here as well, for the original file: the begin strategy may
have pre moved the file and bound the pre
+ * moved file to the exchange (preMove), so the key cannot be derived from
the exchange file later. Returning
+ * <tt>false</tt> marks the file as not started, so the file can be
retried on a later poll (the same as when the
+ * read lock could not be acquired).
+ *
+ * @return <tt>true</tt> if the file was retrieved, <tt>false</tt> if the
file was not retrieved and processing was
+ * aborted
+ */
+ private boolean retrieveFileOrAbort(
+ Exchange exchange, String name, GenericFile<T> target, String
absoluteFileName, GenericFile<T> file) {
+ Exception retrieveCause = null;
+ boolean retrieved = false;
+ try {
+ retrieved = tryRetrievingFile(exchange, name, target,
absoluteFileName, file);
+ } catch (Exception e) {
+ retrieveCause = e;
+ }
+ if (retrieved) {
+ return true;
+ }
+
+ LOG.debug("{} cannot retrieve file: {}", endpoint, target);
+ Exception abortCause = null;
+ try {
+ processStrategy.abort(operations, endpoint, exchange, target);
+ } catch (Exception e) {
+ abortCause = e;
+ } finally {
+ // the file is no longer in progress
+ endpoint.getInProgressRepository().remove(absoluteFileName);
+ removeEagerIdempotentKey(exchange, absoluteFileName);
+ }
+ if (retrieveCause != null) {
+ String msg = "Error processing file " + file + " due to " +
retrieveCause.getMessage();
+ handleException(msg, exchange, retrieveCause);
+ }
+ if (abortCause != null) {
+ String msg2 = endpoint + " cannot abort processing file: " +
target + " due to: " + abortCause.getMessage();
+ handleException(msg2, exchange, abortCause);
+ }
+ return false;
+ }
+
+ /**
+ * Removes the idempotent key that was added eagerly while polling the
file, so the file can be consumed again.
+ * <p/>
+ * Uses the key of the original file (as {@link GenericFileOnCompletion}
does on rollback), and not the file bound
+ * to the exchange, which is the pre moved file when using preMove.
+ */
+ private void removeEagerIdempotentKey(Exchange exchange, String
absoluteFileName) {
+ if (Boolean.TRUE.equals(endpoint.isIdempotent()) &&
endpoint.isIdempotentEager()
+ && endpoint.getIdempotentRepository() != null) {
+ String key = absoluteFileName;
+ if (endpoint.getIdempotentKey() != null) {
+ key = exchange.getProperty(Exchange.FILE_IDEMPOTENT_KEY,
absoluteFileName, String.class);
+ }
+ if (key != null) {
+ endpoint.getIdempotentRepository().remove(key);
+ }
+ }
+ }
+
boolean tryRetrievingFile(
Exchange exchange, String name, GenericFile<T> target, String
absoluteFileName, GenericFile<T> file)
throws Exception {
diff --git
a/core/camel-core/src/test/java/org/apache/camel/component/file/FileConsumerRetrieveFailureTest.java
b/core/camel-core/src/test/java/org/apache/camel/component/file/FileConsumerRetrieveFailureTest.java
new file mode 100644
index 000000000000..130d86f58ddb
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/component/file/FileConsumerRetrieveFailureTest.java
@@ -0,0 +1,203 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You 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.apache.camel.component.file;
+
+import java.io.File;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.Component;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.spi.ExceptionHandler;
+import
org.apache.camel.support.processor.idempotent.MemoryIdempotentRepository;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * When retrieving a file fails after the read lock was acquired (e.g. a FTP,
SFTP or SMB download fails), the eagerly
+ * added idempotent key and the read lock must be released, so the file is
retried on a later poll.
+ * <p/>
+ * The retrieve of the file component itself does not fail, so a {@link
FileOperations} whose first retrieve fails is
+ * plugged into the real {@link FileConsumer}, which shares {@link
GenericFileConsumer} with the remote components.
+ */
+class FileConsumerRetrieveFailureTest extends ContextTestSupport {
+
+ private final AtomicInteger failuresLeft = new AtomicInteger(1);
+ private final AtomicInteger retrieveCalls = new AtomicInteger();
+ private final AtomicBoolean ignoreCannotRetrieve = new AtomicBoolean();
+ private final AtomicInteger handledExceptions = new AtomicInteger();
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Test
+ void testNoopRetriedAfterRetrieveFailure() throws Exception {
+ FileEndpoint endpoint = createEndpoint();
+ endpoint.setNoop(true);
+
+ assertFileRetried(endpoint);
+ assertEquals(1, handledExceptions.get(), "The retrieve failure should
be reported");
+ }
+
+ @Test
+ void testIdempotentPreMoveKeyRemovedAfterRetrieveFailure() throws
Exception {
+ MemoryIdempotentRepository repo = new MemoryIdempotentRepository();
+ FileEndpoint endpoint = createEndpoint();
+ endpoint.setIdempotent(true);
+ endpoint.setIdempotentRepository(repo);
+ endpoint.setPreMove("inprogress");
+
+ context.start();
+ template.sendBodyAndHeader(fileUri(), "Hello World",
Exchange.FILE_NAME, "hello.txt");
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("Hello World");
+
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from(endpoint).convertBodyTo(String.class).to("mock:result");
+ }
+ });
+
+ // the first retrieve fails, after the file was pre moved
+ Path preMoved = testFile("inprogress/hello.txt");
+ await().atMost(5, TimeUnit.SECONDS).until(() ->
handledExceptions.get() == 1);
+ assertTrue(Files.exists(preMoved), "The file should be left in the pre
move directory");
+ // the key was added for the original file, and not for the pre moved
file bound to the exchange
+ assertEquals(0, repo.getCacheSize(), "The eager idempotent key of the
original file should be removed");
+
+ // move the file back, such as an operator would do, and the file must
be consumed
+ Files.move(preMoved, testFile("hello.txt"));
+
+ mock.assertIsSatisfied(5000);
+ assertEquals(2, retrieveCalls.get(), "The file should be retrieved
again after the failure");
+ }
+
+ @Test
+ void testIdempotentReadLockRetriedAfterRetrieveFailure() throws Exception {
+ MemoryIdempotentRepository repo = new MemoryIdempotentRepository();
+ FileEndpoint endpoint = createEndpoint();
+ endpoint.setReadLock("idempotent");
+ endpoint.setIdempotentRepository(repo);
+
+ assertFileRetried(endpoint);
+ assertEquals(1, handledExceptions.get(), "The retrieve failure should
be reported");
+ }
+
+ @Test
+ void testIdempotentReadLockRetriedAfterIgnoredRetrieveFailure() throws
Exception {
+ ignoreCannotRetrieve.set(true);
+ MemoryIdempotentRepository repo = new MemoryIdempotentRepository();
+ FileEndpoint endpoint = createEndpoint();
+ endpoint.setReadLock("idempotent");
+ endpoint.setIdempotentRepository(repo);
+
+ assertFileRetried(endpoint);
+ assertEquals(0, handledExceptions.get(), "An ignored retrieve failure
should not be reported");
+ }
+
+ private void assertFileRetried(FileEndpoint endpoint) throws Exception {
+ context.start();
+ template.sendBodyAndHeader(fileUri(), "Hello World",
Exchange.FILE_NAME, "hello.txt");
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("Hello World");
+
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from(endpoint).convertBodyTo(String.class).to("mock:result");
+ }
+ });
+
+ // the first retrieve fails, and the next poll must retrieve the file
again
+ mock.assertIsSatisfied(5000);
+ assertEquals(2, retrieveCalls.get(), "The file should be retrieved
again after the failure");
+ }
+
+ private FileEndpoint createEndpoint() {
+ FileEndpoint endpoint = new FlakyFileEndpoint(fileUri(),
context.getComponent("file"));
+ endpoint.setCamelContext(context);
+ endpoint.setFile(testDirectory().toFile());
+ endpoint.setInitialDelay(0);
+ endpoint.setDelay(10);
+ endpoint.setExceptionHandler(new CountingExceptionHandler());
+ return endpoint;
+ }
+
+ private final class FlakyFileEndpoint extends FileEndpoint {
+
+ FlakyFileEndpoint(String endpointUri, Component component) {
+ super(endpointUri, component);
+ }
+
+ @Override
+ protected FileConsumer newFileConsumer(Processor processor,
GenericFileOperations<File> operations) {
+ FileOperations flaky = new FileOperations(this) {
+ @Override
+ public boolean retrieveFile(String name, Exchange exchange,
long size) {
+ retrieveCalls.incrementAndGet();
+ if (failuresLeft.getAndDecrement() > 0) {
+ if (ignoreCannotRetrieve.get()) {
+ return false;
+ }
+ // such as a FTP download that fails
+ throw new
GenericFileOperationFailedException("Simulated connection reset while
retrieving " + name);
+ }
+ return super.retrieveFile(name, exchange, size);
+ }
+ };
+ return new FileConsumer(this, processor, flaky,
createGenericFileStrategy()) {
+ @Override
+ protected boolean ignoreCannotRetrieveFile(String name,
Exchange exchange, Exception cause) {
+ return ignoreCannotRetrieve.get();
+ }
+ };
+ }
+ }
+
+ private final class CountingExceptionHandler implements ExceptionHandler {
+
+ @Override
+ public void handleException(Throwable exception) {
+ handledExceptions.incrementAndGet();
+ }
+
+ @Override
+ public void handleException(String message, Throwable exception) {
+ handledExceptions.incrementAndGet();
+ }
+
+ @Override
+ public void handleException(String message, Exchange exchange,
Throwable exception) {
+ handledExceptions.incrementAndGet();
+ }
+ }
+}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index e16f9a2740bc..45581cbb9a5c 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -2587,6 +2587,19 @@ Downloads to a configured `localWorkDirectory` now
resolve existing filesystem p
checking that the destination remains inside that directory. Downloads through
a symbolic link that
resolves outside the `localWorkDirectory` are rejected. Valid nested download
paths continue to work.
+=== camel-file, camel-ftp, camel-smb, camel-mina-sftp and camel-azure-files -
a file whose download fails is retried by idempotent consumers
+
+When the consumer cannot download a file (for example a FTP, SFTP or SMB
transfer fails with a socket timeout or a
+connection reset), the file is now retried on the next poll also when the
consumer is idempotent, which includes every
+consumer with `noop=true`. Consumers that are not idempotent already retried
such a file. Previously the idempotent key,
+which is added when the file is polled (`idempotentEager=true`, the default),
was kept, so the file was skipped until the
+application was restarted, or for good with a persistent idempotent
repository. The exclusive read lock acquired for the
+file is now released as well.
+
+A poll in which the download of every file failed now counts as an idle poll,
the same as a poll in which no read lock
+could be acquired. This affects `backoffIdleThreshold`,
`sendEmptyMessageWhenIdle` (an empty message is sent) and
+`greedy` (the consumer does not poll again immediately).
+
=== camel-infinispan - the aggregation repository keeps completed exchanges
for recovery
`InfinispanAggregationRepository` implements
`RecoverableAggregationRepository`, but it had no recovery