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());
+ }
+}