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 657adf3c2e0e CAMEL-25289: camel-iggy - a failed send or poll must not 
keep its pooled client, and stopped producers and consumers must close their 
clients (#27321)
657adf3c2e0e is described below

commit 657adf3c2e0e20e415e981855e953bf005697b3c
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:33:50 2026 +0530

    CAMEL-25289: camel-iggy - a failed send or poll must not keep its pooled 
client, and stopped producers and consumers must close their clients (#27321)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 components/camel-iggy/pom.xml                      |   6 +
 .../camel/component/iggy/IggyConfiguration.java    |   2 +-
 .../apache/camel/component/iggy/IggyConsumer.java  |  13 +-
 .../camel/component/iggy/IggyFetchRecords.java     |  91 ++++++----
 .../apache/camel/component/iggy/IggyProducer.java  |  55 +++++--
 .../iggy/client/IggyClientConnectionPool.java      |  22 +++
 .../component/iggy/client/IggyClientFactory.java   |   8 +
 .../camel/component/iggy/IggyMockClientTest.java   | 183 +++++++++++++++++++++
 8 files changed, 327 insertions(+), 53 deletions(-)

diff --git a/components/camel-iggy/pom.xml b/components/camel-iggy/pom.xml
index 6c4164a81e1f..50aa8996a656 100644
--- a/components/camel-iggy/pom.xml
+++ b/components/camel-iggy/pom.xml
@@ -83,6 +83,12 @@
             <groupId>org.assertj</groupId>
             <artifactId>assertj-core</artifactId>
         </dependency>
+        <dependency>
+            <groupId>org.mockito</groupId>
+            <artifactId>mockito-core</artifactId>
+            <version>${mockito-version}</version>
+            <scope>test</scope>
+        </dependency>
     </dependencies>
 
 </project>
diff --git 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConfiguration.java
 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConfiguration.java
index a09f17147fe1..59318e4897b6 100644
--- 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConfiguration.java
+++ 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConfiguration.java
@@ -82,7 +82,7 @@ public class IggyConfiguration implements Cloneable {
     @UriParam(label = "consumer", defaultValue = "0",
               description = "Defines the initial message offset position when 
autoCommit is disabled. " +
                             "Use 0 to start from the beginning of the stream, 
or specify a custom offset to resume from a particular point")
-    private Long startingOffset;
+    private Long startingOffset = 0L;
     @UriParam(label = "security", defaultValue = "false",
               description = "Whether to enable TLS for the connection to the 
Iggy server")
     private boolean tlsEnabled;
diff --git 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConsumer.java
 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConsumer.java
index b837e0a0336f..e331f9625a84 100644
--- 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConsumer.java
+++ 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyConsumer.java
@@ -59,9 +59,12 @@ public class IggyConsumer extends DefaultConsumer {
                 endpoint.getConfiguration().getSslContextParameters());
 
         IggyBaseClient client = iggyClientConnectionPool.borrowObject();
-        endpoint.initializeTopic(client);
-        endpoint.initializeConsumerGroup(client);
-        iggyClientConnectionPool.returnClient(client);
+        try {
+            endpoint.initializeTopic(client);
+            endpoint.initializeConsumerGroup(client);
+        } finally {
+            iggyClientConnectionPool.returnClient(client);
+        }
 
         executor = endpoint.createExecutor();
         BridgeExceptionHandlerToErrorHandler bridge = new 
BridgeExceptionHandlerToErrorHandler(this);
@@ -110,6 +113,10 @@ public class IggyConsumer extends DefaultConsumer {
         }
         tasks.clear();
         executor = null;
+        if (iggyClientConnectionPool != null) {
+            iggyClientConnectionPool.close();
+            iggyClientConnectionPool = null;
+        }
 
         super.doStop();
     }
diff --git 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyFetchRecords.java
 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyFetchRecords.java
index 1cf8f413dcf6..4192edcf2712 100644
--- 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyFetchRecords.java
+++ 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyFetchRecords.java
@@ -71,17 +71,7 @@ public class IggyFetchRecords implements Runnable {
         while (running) {
             if (iggyConsumer.isSuspending() || iggyConsumer.isSuspended()) {
                 LOG.trace("Consumer is suspended. Skipping message polling.");
-                // Use Camel's task API to avoid busy-waiting instead of 
Thread.sleep()
-                // We use initialDelay for the actual delay, and 
maxIterations(1) to run once
-                Tasks.foregroundTask()
-                        .withBudget(Budgets.iterationBudget()
-                                .withMaxIterations(1)
-                                .withInitialDelay(Duration.ofSeconds(1))
-                                .withInterval(Duration.ZERO)
-                                .build())
-                        .withName("IggySuspendedDelay")
-                        .build()
-                        .run(endpoint.getCamelContext(), () -> true);
+                delay("IggySuspendedDelay");
                 continue;
             }
 
@@ -105,31 +95,19 @@ public class IggyFetchRecords implements Runnable {
 
             PolledMessages polledMessages;
             IggyBaseClient client = iggyClientConnectionPool.borrowObject();
-            if (configuration.isAutoCommit()) {
-                polledMessages = client.messages()
-                        .pollMessages(streamId,
-                                topicId,
-                                
Optional.ofNullable(configuration.getPartitionId()),
-                                Consumer.group(consumerId),
-                                resolvePollingStrategy(),
-                                configuration.getPollBatchSize(),
-                                configuration.isAutoCommit());
-            } else {
-                polledMessages = client.messages()
-                        .pollMessages(streamId,
-                                topicId,
-                                
Optional.ofNullable(configuration.getPartitionId()),
-                                Consumer.group(consumerId),
-                                PollingStrategy.offset(offset),
-                                configuration.getPollBatchSize(),
-                                false);
-
-                // Update offset
-                offset = 
offset.add(BigInteger.valueOf(polledMessages.count()));
+            boolean polled = false;
+            try {
+                polledMessages = poll(client, streamId, topicId, consumerId);
+                polled = true;
+            } finally {
+                if (polled) {
+                    iggyClientConnectionPool.returnClient(client);
+                } else {
+                    // the client of a failed poll may be broken (connection 
lost): do not hand it out again
+                    iggyClientConnectionPool.invalidateClient(client);
+                }
             }
 
-            iggyClientConnectionPool.returnClient(client);
-
             LOG.debug("Fetched {} messages from partition {}, current offset 
{}",
                     polledMessages.count(),
                     polledMessages.partitionId(),
@@ -145,7 +123,52 @@ public class IggyFetchRecords implements Runnable {
             }
         } catch (Exception e) {
             bridgeExceptionHandlerToErrorHandler.handleException("Error 
polling messages from Iggy", e);
+            if (running) {
+                // do not poll a server that fails again in a tight loop
+                delay("IggyPollErrorDelay");
+            }
+        }
+    }
+
+    private void delay(String name) {
+        // Use Camel's task API to avoid busy-waiting instead of Thread.sleep()
+        // We use initialDelay for the actual delay, and maxIterations(1) to 
run once
+        Tasks.foregroundTask()
+                .withBudget(Budgets.iterationBudget()
+                        .withMaxIterations(1)
+                        .withInitialDelay(Duration.ofSeconds(1))
+                        .withInterval(Duration.ZERO)
+                        .build())
+                .withName(name)
+                .build()
+                .run(endpoint.getCamelContext(), () -> true);
+    }
+
+    private PolledMessages poll(IggyBaseClient client, StreamId streamId, 
TopicId topicId, ConsumerId consumerId) {
+        PolledMessages polledMessages;
+        if (configuration.isAutoCommit()) {
+            polledMessages = client.messages()
+                    .pollMessages(streamId,
+                            topicId,
+                            
Optional.ofNullable(configuration.getPartitionId()),
+                            Consumer.group(consumerId),
+                            resolvePollingStrategy(),
+                            configuration.getPollBatchSize(),
+                            configuration.isAutoCommit());
+        } else {
+            polledMessages = client.messages()
+                    .pollMessages(streamId,
+                            topicId,
+                            
Optional.ofNullable(configuration.getPartitionId()),
+                            Consumer.group(consumerId),
+                            PollingStrategy.offset(offset),
+                            configuration.getPollBatchSize(),
+                            false);
+
+            // Update offset
+            offset = offset.add(BigInteger.valueOf(polledMessages.count()));
         }
+        return polledMessages;
     }
 
     private PollingStrategy resolvePollingStrategy() {
diff --git 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyProducer.java
 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyProducer.java
index 2ba2ee9c4ae9..becae90f49bf 100644
--- 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyProducer.java
+++ 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/IggyProducer.java
@@ -58,8 +58,20 @@ public class IggyProducer extends DefaultAsyncProducer {
                 endpoint.getConfiguration().getSslContextParameters());
 
         IggyBaseClient client = iggyClientConnectionPool.borrowObject();
-        endpoint.initializeTopic(client);
-        iggyClientConnectionPool.returnClient(client);
+        try {
+            endpoint.initializeTopic(client);
+        } finally {
+            iggyClientConnectionPool.returnClient(client);
+        }
+    }
+
+    @Override
+    protected void doStop() throws Exception {
+        if (iggyClientConnectionPool != null) {
+            iggyClientConnectionPool.close();
+            iggyClientConnectionPool = null;
+        }
+        super.doStop();
     }
 
     @Override
@@ -99,8 +111,6 @@ public class IggyProducer extends DefaultAsyncProducer {
                 .unwrap();
              */
 
-            IggyBaseClient client = iggyClientConnectionPool.borrowObject();
-
             Optional<String> topicOverride
                     = 
Optional.ofNullable(exchange.getMessage().getHeader(IggyConstants.TOPIC_OVERRIDE,
 String.class));
             Optional<String> streamOverride
@@ -108,19 +118,25 @@ public class IggyProducer extends DefaultAsyncProducer {
 
             String topic = topicOverride.orElse(endpoint.getTopicName());
             String stream = 
streamOverride.orElse(iggyConfiguration.getStreamName());
-            if (topicOverride.isPresent() || streamOverride.isPresent()) {
-                endpoint.initializeTopic(client,
-                        topic,
-                        stream);
-            }
 
-            client.messages().sendMessages(
-                    StreamId.of(stream),
-                    TopicId.of(topic),
-                    iggyConfiguration.getPartitioning(),
-                    messages);
+            IggyBaseClient client = iggyClientConnectionPool.borrowObject();
+            boolean sent = false;
+            try {
+                if (topicOverride.isPresent() || streamOverride.isPresent()) {
+                    endpoint.initializeTopic(client,
+                            topic,
+                            stream);
+                }
 
-            iggyClientConnectionPool.returnClient(client);
+                client.messages().sendMessages(
+                        StreamId.of(stream),
+                        TopicId.of(topic),
+                        iggyConfiguration.getPartitioning(),
+                        messages);
+                sent = true;
+            } finally {
+                releaseClient(client, sent);
+            }
         } catch (Exception e) {
             exchange.setException(e);
         }
@@ -128,6 +144,15 @@ public class IggyProducer extends DefaultAsyncProducer {
         return true;
     }
 
+    private void releaseClient(IggyBaseClient client, boolean success) {
+        if (success) {
+            iggyClientConnectionPool.returnClient(client);
+        } else {
+            // the client of a failed request may be broken (connection lost): 
do not hand it out again
+            iggyClientConnectionPool.invalidateClient(client);
+        }
+    }
+
     private boolean isListOfStrings(List<?> list) {
         return list != null &&
                 list.stream().allMatch(item -> item == null || item instanceof 
String);
diff --git 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientConnectionPool.java
 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientConnectionPool.java
index fb386b1620b0..ef54c29489de 100644
--- 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientConnectionPool.java
+++ 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientConnectionPool.java
@@ -19,9 +19,13 @@ package org.apache.camel.component.iggy.client;
 import org.apache.camel.support.jsse.SSLContextParameters;
 import org.apache.commons.pool2.impl.GenericObjectPool;
 import org.apache.iggy.client.blocking.IggyBaseClient;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 public class IggyClientConnectionPool {
 
+    private static final Logger LOG = 
LoggerFactory.getLogger(IggyClientConnectionPool.class);
+
     private final GenericObjectPool<IggyBaseClient> pool;
 
     public IggyClientConnectionPool(String host, int port, String username, 
String password, String transport,
@@ -41,6 +45,24 @@ public class IggyClientConnectionPool {
         pool.returnObject(client);
     }
 
+    /**
+     * Removes a client whose request failed from the pool, and closes it.
+     */
+    public void invalidateClient(IggyBaseClient client) {
+        try {
+            pool.invalidateObject(client);
+        } catch (Exception e) {
+            LOG.debug("Error closing Iggy client: {}", e.getMessage(), e);
+        }
+    }
+
+    /**
+     * Closes the pool and the clients that are not in use.
+     */
+    public void close() {
+        pool.close();
+    }
+
     public int getNumActive() {
         return pool.getNumActive();
     }
diff --git 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientFactory.java
 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientFactory.java
index bcda20421537..665780fddb78 100644
--- 
a/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientFactory.java
+++ 
b/components/camel-iggy/src/main/java/org/apache/camel/component/iggy/client/IggyClientFactory.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.iggy.client;
 
+import java.io.Closeable;
 import java.util.function.Consumer;
 
 import org.apache.camel.support.jsse.SSLContextParameters;
@@ -106,4 +107,11 @@ public class IggyClientFactory extends 
BasePooledObjectFactory<IggyBaseClient> {
         return new DefaultPooledObject<>(iggyBaseClient);
     }
 
+    @Override
+    public void destroyObject(PooledObject<IggyBaseClient> pooledObject) 
throws Exception {
+        if (pooledObject.getObject() instanceof Closeable closeable) {
+            closeable.close();
+        }
+    }
+
 }
diff --git 
a/components/camel-iggy/src/test/java/org/apache/camel/component/iggy/IggyMockClientTest.java
 
b/components/camel-iggy/src/test/java/org/apache/camel/component/iggy/IggyMockClientTest.java
new file mode 100644
index 000000000000..62e541b1dbb5
--- /dev/null
+++ 
b/components/camel-iggy/src/test/java/org/apache/camel/component/iggy/IggyMockClientTest.java
@@ -0,0 +1,183 @@
+/*
+ * 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.iggy;
+
+import java.io.Closeable;
+import java.math.BigInteger;
+import java.time.Duration;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Consumer;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.Producer;
+import org.apache.camel.component.iggy.client.IggyClientFactory;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.iggy.client.blocking.IggyBaseClient;
+import org.apache.iggy.identifier.StreamId;
+import org.apache.iggy.identifier.TopicId;
+import org.apache.iggy.message.Partitioning;
+import org.apache.iggy.message.PollingStrategy;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyList;
+import static org.mockito.Mockito.CALLS_REAL_METHODS;
+import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.when;
+import static org.mockito.Mockito.withSettings;
+
+/**
+ * The producer and the consumer must give back or discard the pooled client 
of a failed request, and close their
+ * clients when they stop; without autoCommit the consumer starts at the 
default starting offset 0. The Iggy clients are
+ * mocks, no Iggy server is needed.
+ */
+public class IggyMockClientTest {
+
+    private static final String URI
+            = 
"iggy:topic?streamName=stream&autoCreateStream=false&autoCreateTopic=false&consumerGroupName=group";
+
+    private final AtomicInteger closed = new AtomicInteger();
+
+    @Test
+    void testFailedSendsDoNotExhaustThePool() throws Exception {
+        IggyBaseClient client = newClient();
+        when(client.messages().sendMessages(any(StreamId.class), 
any(TopicId.class), any(Partitioning.class), anyList()))
+                .thenThrow(new IllegalStateException("Connection reset"));
+
+        try (MockedConstruction<IggyClientFactory> factories = 
mockFactories(client);
+             CamelContext context = new DefaultCamelContext()) {
+            context.start();
+            Endpoint endpoint = context.getEndpoint(URI);
+            // created and started in this thread, where the client factory is 
mocked
+            Producer producer = endpoint.createProducer();
+            producer.start();
+
+            // the pool holds at most 8 clients: if each failed send keeps its 
client, the 9th send waits forever
+            assertTimeoutPreemptively(Duration.ofSeconds(30), () -> {
+                for (int i = 0; i < 10; i++) {
+                    Exchange exchange = endpoint.createExchange();
+                    exchange.getIn().setBody("hello");
+                    producer.process(exchange);
+                    assertInstanceOf(IllegalStateException.class, 
exchange.getException());
+                }
+            });
+            // the clients of the failed sends are not reused but closed
+            assertEquals(10, closed.get());
+
+            producer.stop();
+        }
+    }
+
+    @Test
+    void testProducerStopClosesTheClients() throws Exception {
+        IggyBaseClient client = newClient();
+
+        try (MockedConstruction<IggyClientFactory> factories = 
mockFactories(client);
+             CamelContext context = new DefaultCamelContext()) {
+            context.start();
+            Producer producer = context.getEndpoint(URI).createProducer();
+            producer.start();
+            assertEquals(0, closed.get());
+
+            producer.stop();
+            assertEquals(1, closed.get());
+        }
+    }
+
+    @Test
+    void testFailedPollDiscardsTheClient() throws Exception {
+        IggyBaseClient client = newClient();
+        AtomicInteger polls = new AtomicInteger();
+        AtomicInteger closedBeforeSecondPoll = new AtomicInteger(-1);
+        CountDownLatch secondPoll = new CountDownLatch(1);
+        when(client.messages().pollMessages(any(StreamId.class), 
any(TopicId.class), any(), any(), any(), any(),
+                anyBoolean())).thenAnswer(invocation -> {
+                    if (polls.incrementAndGet() == 2) {
+                        closedBeforeSecondPoll.set(closed.get());
+                        secondPoll.countDown();
+                    }
+                    throw new IllegalStateException("Connection reset");
+                });
+
+        try (MockedConstruction<IggyClientFactory> factories = 
mockFactories(client);
+             CamelContext context = new DefaultCamelContext()) {
+            context.start();
+            Consumer consumer = 
context.getEndpoint(URI).createConsumer(exchange -> {
+            });
+            consumer.start();
+
+            assertTrue(secondPoll.await(30, TimeUnit.SECONDS));
+            // the client of the failed poll was closed (not kept out of the 
pool) before the next poll
+            assertEquals(1, closedBeforeSecondPoll.get());
+
+            consumer.stop();
+        }
+    }
+
+    @Test
+    void testManualCommitStartsAtOffsetZeroByDefault() throws Exception {
+        IggyBaseClient client = newClient();
+        AtomicReference<PollingStrategy> strategy = new AtomicReference<>();
+        CountDownLatch polled = new CountDownLatch(1);
+        when(client.messages().pollMessages(any(StreamId.class), 
any(TopicId.class), any(), any(), any(), any(),
+                anyBoolean())).thenAnswer(invocation -> {
+                    strategy.compareAndSet(null, invocation.getArgument(4));
+                    polled.countDown();
+                    throw new IllegalStateException("Stop here");
+                });
+
+        try (MockedConstruction<IggyClientFactory> factories = 
mockFactories(client);
+             CamelContext context = new DefaultCamelContext()) {
+            context.start();
+            Consumer consumer = context.getEndpoint(URI + 
"&autoCommit=false").createConsumer(exchange -> {
+            });
+            consumer.start();
+
+            assertTrue(polled.await(30, TimeUnit.SECONDS));
+            assertEquals(PollingStrategy.offset(BigInteger.ZERO), 
strategy.get());
+
+            consumer.stop();
+        }
+    }
+
+    private IggyBaseClient newClient() throws Exception {
+        IggyBaseClient client = mock(IggyBaseClient.class,
+                
withSettings().extraInterfaces(Closeable.class).defaultAnswer(RETURNS_DEEP_STUBS));
+        doAnswer(invocation -> closed.incrementAndGet()).when((Closeable) 
client).close();
+        return client;
+    }
+
+    private static MockedConstruction<IggyClientFactory> 
mockFactories(IggyBaseClient client) {
+        return mockConstruction(IggyClientFactory.class, 
withSettings().defaultAnswer(CALLS_REAL_METHODS),
+                (factory, context) -> doReturn(client).when(factory).create());
+    }
+}

Reply via email to