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 9729f27567bf CAMEL-25246: camel-hazelcast - queue and instance
consumers must remove their listener on stop, and the queue poll must survive
an error (#27243)
9729f27567bf is described below
commit 9729f27567bf31d537cd47d17f2926608ddc28dd
Author: allthingssecurity <[email protected]>
AuthorDate: Fri Oct 2 13:23:16 2026 +0530
CAMEL-25246: camel-hazelcast - queue and instance consumers must remove
their listener on stop, and the queue poll must survive an error (#27243)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../instance/HazelcastInstanceConsumer.java | 26 ++++-
.../hazelcast/queue/HazelcastQueueConsumer.java | 103 +++++++++++++------
.../hazelcast/HazelcastInstanceConsumerTest.java | 14 ++-
.../HazelcastQueueConsumerPollErrorTest.java | 89 ++++++++++++++++
.../HazelcastQueueConsumerRestartTest.java | 112 +++++++++++++++++++++
.../hazelcast/HazelcastQueueConsumerTest.java | 28 +++---
6 files changed, 324 insertions(+), 48 deletions(-)
diff --git
a/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/instance/HazelcastInstanceConsumer.java
b/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/instance/HazelcastInstanceConsumer.java
index 6a9b39253fc6..8a7395b0c3b5 100644
---
a/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/instance/HazelcastInstanceConsumer.java
+++
b/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/instance/HazelcastInstanceConsumer.java
@@ -17,7 +17,9 @@
package org.apache.camel.component.hazelcast.instance;
import java.net.InetSocketAddress;
+import java.util.UUID;
+import com.hazelcast.cluster.Cluster;
import com.hazelcast.cluster.MembershipEvent;
import com.hazelcast.cluster.MembershipListener;
import com.hazelcast.core.HazelcastInstance;
@@ -30,10 +32,32 @@ import org.apache.camel.support.DefaultEndpoint;
public class HazelcastInstanceConsumer extends DefaultConsumer {
+ private final HazelcastInstance hazelcastInstance;
+ private Cluster cluster;
+ private UUID listener;
+
public HazelcastInstanceConsumer(HazelcastInstance hazelcastInstance,
DefaultEndpoint endpoint, Processor processor) {
super(endpoint, processor);
+ this.hazelcastInstance = hazelcastInstance;
+ }
+
+ @Override
+ protected void doStart() throws Exception {
+ super.doStart();
+
+ // register the listener here, so that doStop can remove it
(CAMEL-15899)
+ cluster = hazelcastInstance.getCluster();
+ listener = cluster.addMembershipListener(new
HazelcastMembershipListener());
+ }
+
+ @Override
+ protected void doStop() throws Exception {
+ if (listener != null) {
+ cluster.removeMembershipListener(listener);
+ listener = null;
+ }
- hazelcastInstance.getCluster().addMembershipListener(new
HazelcastMembershipListener());
+ super.doStop();
}
class HazelcastMembershipListener implements MembershipListener {
diff --git
a/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/queue/HazelcastQueueConsumer.java
b/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/queue/HazelcastQueueConsumer.java
index 5d98899c6fa8..93e790843885 100644
---
a/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/queue/HazelcastQueueConsumer.java
+++
b/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/queue/HazelcastQueueConsumer.java
@@ -16,6 +16,8 @@
*/
package org.apache.camel.component.hazelcast.queue;
+import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;
@@ -29,10 +31,16 @@ import
org.apache.camel.component.hazelcast.listener.CamelItemListener;
public class HazelcastQueueConsumer extends HazelcastDefaultConsumer {
+ // the minimum delay before polling again after a poll error, so a
pollingTimeout of 0 does not spin
+ private static final long MIN_POLL_ERROR_DELAY = 1000L;
+
private final Processor processor;
private ExecutorService executor;
- private QueueConsumerTask queueConsumerTask;
private HazelcastQueueConfiguration config;
+ private IQueue<Object> queue;
+ private UUID listener;
+ // counted down on stop, to end the delay after a poll error without
waiting for it
+ private CountDownLatch stopLatch;
public HazelcastQueueConsumer(HazelcastInstance hazelcastInstance,
Endpoint endpoint, Processor processor, String cacheName,
final HazelcastQueueConfiguration
configuration) {
@@ -44,17 +52,30 @@ public class HazelcastQueueConsumer extends
HazelcastDefaultConsumer {
@Override
protected void doStart() throws Exception {
super.doStart();
- executor = ((HazelcastQueueEndpoint)
getEndpoint()).createExecutor(this);
+ queue = hazelcastInstance.getQueue(cacheName);
- CamelItemListener camelItemListener = new CamelItemListener(this,
cacheName);
- queueConsumerTask = new QueueConsumerTask(camelItemListener);
- executor.submit(queueConsumerTask);
+ if (config.getQueueConsumerMode() ==
HazelcastQueueConsumerMode.LISTEN) {
+ // register the listener here, so that doStop can remove it
(CAMEL-15899)
+ listener = queue.addItemListener(new CamelItemListener(this,
cacheName), true);
+ } else if (config.getQueueConsumerMode() ==
HazelcastQueueConsumerMode.POLL) {
+ stopLatch = new CountDownLatch(1);
+ executor = ((HazelcastQueueEndpoint)
getEndpoint()).createExecutor(this);
+ executor.submit(new QueueConsumerTask(queue, stopLatch));
+ }
}
@Override
protected void doStop() throws Exception {
+ if (listener != null) {
+ queue.removeItemListener(listener);
+ listener = null;
+ }
+
super.doStop();
+ if (stopLatch != null) {
+ stopLatch.countDown();
+ }
if (executor != null) {
if (getEndpoint() != null && getEndpoint().getCamelContext() !=
null) {
getEndpoint().getCamelContext().getExecutorServiceManager().shutdownNow(executor);
@@ -67,41 +88,63 @@ public class HazelcastQueueConsumer extends
HazelcastDefaultConsumer {
class QueueConsumerTask implements Runnable {
- CamelItemListener camelItemListener;
+ private final IQueue<Object> queue;
+ private final CountDownLatch stopLatch;
- public QueueConsumerTask(CamelItemListener camelItemListener) {
- this.camelItemListener = camelItemListener;
+ QueueConsumerTask(IQueue<Object> queue, CountDownLatch stopLatch) {
+ this.queue = queue;
+ this.stopLatch = stopLatch;
}
@Override
public void run() {
- IQueue<Object> queue = hazelcastInstance.getQueue(cacheName);
- if (config.getQueueConsumerMode() ==
HazelcastQueueConsumerMode.LISTEN) {
- queue.addItemListener(camelItemListener, true);
- }
-
- if (config.getQueueConsumerMode() ==
HazelcastQueueConsumerMode.POLL) {
- while (isRunAllowed()) {
- try {
- final Object body =
queue.poll(config.getPollingTimeout(), TimeUnit.MILLISECONDS);
- // CAMEL-16035 - If the polling timeout is exceeded
with nothing to poll from the queue, the queue.poll() method return NULL
- if (body != null) {
- Exchange exchange = createExchange(false);
- exchange.getIn().setBody(body);
- try {
- processor.process(exchange);
- } catch (Exception e) {
- getExceptionHandler().handleException("Error
during processing", exchange, e);
- } finally {
- releaseExchange(exchange, false);
- }
+ while (isRunAllowed()) {
+ final Object body;
+ try {
+ body = queue.poll(config.getPollingTimeout(),
TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ // only doStop interrupts this thread
+ Thread.currentThread().interrupt();
+ return;
+ } catch (Exception e) {
+ // keep polling after an error (such as the client being
disconnected from the cluster)
+ if (isRunAllowed()) {
+ getExceptionHandler().handleException("Error polling
from the queue " + cacheName, e);
+ if (!waitBeforeNextPoll()) {
+ return;
}
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
+ }
+ continue;
+ }
+ // CAMEL-16035 - If the polling timeout is exceeded with
nothing to poll from the queue, the queue.poll() method return NULL
+ if (body != null) {
+ Exchange exchange = createExchange(false);
+ exchange.getIn().setBody(body);
+ try {
+ processor.process(exchange);
+ } catch (Exception e) {
+ getExceptionHandler().handleException("Error during
processing", exchange, e);
+ } finally {
+ releaseExchange(exchange, false);
}
}
}
}
+
+ /**
+ * Waits before the next poll after a poll error, and ends early when
the consumer stops.
+ *
+ * @return false if the thread was interrupted
+ */
+ private boolean waitBeforeNextPoll() {
+ try {
+ stopLatch.await(Math.max(config.getPollingTimeout(),
MIN_POLL_ERROR_DELAY), TimeUnit.MILLISECONDS);
+ return true;
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ return false;
+ }
+ }
}
}
diff --git
a/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastInstanceConsumerTest.java
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastInstanceConsumerTest.java
index 63dcc038df6d..068e1d0bb9d4 100644
---
a/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastInstanceConsumerTest.java
+++
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastInstanceConsumerTest.java
@@ -53,11 +53,13 @@ public class HazelcastInstanceConsumerTest extends
HazelcastCamelTestSupport {
private ArgumentCaptor<MembershipListener> argument;
+ private final UUID listenerId = UUID.randomUUID();
+
@Override
protected void trainHazelcastInstance(HazelcastInstance hazelcastInstance)
{
when(hazelcastInstance.getCluster()).thenReturn(cluster);
argument = ArgumentCaptor.forClass(MembershipListener.class);
-
when(cluster.addMembershipListener(any())).thenReturn(UUID.randomUUID());
+ when(cluster.addMembershipListener(any())).thenReturn(listenerId);
}
@Override
@@ -106,12 +108,20 @@ public class HazelcastInstanceConsumerTest extends
HazelcastCamelTestSupport {
this.checkHeaders(headers, HazelcastConstants.REMOVED);
}
+ @Test
+ public void testStopRemovesListener() throws Exception {
+ context.getRouteController().stopRoute("instance");
+
+ verify(cluster).removeMembershipListener(listenerId);
+ }
+
@Override
protected RouteBuilder createRouteBuilder() throws Exception {
return new RouteBuilder() {
@Override
public void configure() throws Exception {
- from(String.format("hazelcast-%sfoo",
HazelcastConstants.INSTANCE_PREFIX)).log("instance...").choice()
+ from(String.format("hazelcast-%sfoo",
HazelcastConstants.INSTANCE_PREFIX)).routeId("instance")
+ .log("instance...").choice()
.when(header(HazelcastConstants.LISTENER_ACTION).isEqualTo(HazelcastConstants.ADDED)).log("...added")
.to("mock:added").otherwise().log("...removed").to("mock:removed");
}
diff --git
a/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerPollErrorTest.java
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerPollErrorTest.java
new file mode 100644
index 000000000000..c7bc13ad4974
--- /dev/null
+++
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerPollErrorTest.java
@@ -0,0 +1,89 @@
+/*
+ * 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.hazelcast;
+
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import com.hazelcast.collection.IQueue;
+import com.hazelcast.core.HazelcastException;
+import com.hazelcast.core.HazelcastInstance;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mock;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * The queue consumer in poll mode must keep polling after a poll fails.
+ */
+public class HazelcastQueueConsumerPollErrorTest extends
HazelcastCamelTestSupport {
+
+ @Mock
+ private IQueue<String> queue;
+
+ private final LinkedBlockingQueue<String> items = new
LinkedBlockingQueue<>();
+ private final AtomicBoolean failed = new AtomicBoolean();
+
+ @Override
+ protected void trainHazelcastInstance(HazelcastInstance hazelcastInstance)
{
+ when(hazelcastInstance.<String> getQueue("foo")).thenReturn(queue);
+ try {
+ when(queue.poll(anyLong(),
any(TimeUnit.class))).thenAnswer(invocation -> {
+ // the first poll fails, as when the client is disconnected
from the cluster
+ if (failed.compareAndSet(false, true)) {
+ throw new HazelcastException("Simulated poll failure");
+ }
+ return items.poll(invocation.getArgument(0),
invocation.getArgument(1));
+ });
+ } catch (InterruptedException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ @Override
+ protected void verifyHazelcastInstance(HazelcastInstance
hazelcastInstance) {
+ verify(hazelcastInstance).getQueue("foo");
+ }
+
+ @Test
+ public void testKeepsPollingAfterError() throws Exception {
+ MockEndpoint out = getMockEndpoint("mock:result");
+ out.expectedBodiesReceived("bar");
+
+ items.add("bar");
+
+ MockEndpoint.assertIsSatisfied(context, 5, TimeUnit.SECONDS);
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+
from(String.format("hazelcast-%sfoo?queueConsumerMode=Poll&pollingTimeout=100",
+ HazelcastConstants.QUEUE_PREFIX))
+ .to("mock:result");
+ }
+ };
+ }
+}
diff --git
a/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerRestartTest.java
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerRestartTest.java
new file mode 100644
index 000000000000..b8c4d5821b18
--- /dev/null
+++
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerRestartTest.java
@@ -0,0 +1,112 @@
+/*
+ * 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.hazelcast;
+
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
+
+import com.hazelcast.collection.IQueue;
+import com.hazelcast.collection.ItemEvent;
+import com.hazelcast.collection.ItemListener;
+import com.hazelcast.core.HazelcastInstance;
+import com.hazelcast.core.ItemEventType;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mock;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.atLeastOnce;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * The queue consumer in listen mode must remove its item listener when it
stops, like the other Hazelcast consumers
+ * (CAMEL-15899).
+ */
+public class HazelcastQueueConsumerRestartTest extends
HazelcastCamelTestSupport {
+
+ @Mock
+ private IQueue<String> queue;
+
+ // the item listeners registered on the queue, as Hazelcast keeps them
+ private final Map<UUID, ItemListener<String>> listeners = new
ConcurrentHashMap<>();
+
+ @Override
+ @SuppressWarnings("unchecked")
+ protected void trainHazelcastInstance(HazelcastInstance hazelcastInstance)
{
+ when(hazelcastInstance.<String> getQueue("foo")).thenReturn(queue);
+ when(queue.addItemListener(any(ItemListener.class),
eq(true))).thenAnswer(invocation -> {
+ UUID id = UUID.randomUUID();
+ listeners.put(id, invocation.getArgument(0));
+ return id;
+ });
+ when(queue.removeItemListener(any(UUID.class))).thenAnswer(invocation
-> {
+ return listeners.remove(invocation.<UUID> getArgument(0)) != null;
+ });
+ }
+
+ @Override
+ protected void verifyHazelcastInstance(HazelcastInstance
hazelcastInstance) {
+ verify(hazelcastInstance, atLeastOnce()).getQueue("foo");
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testStopRemovesListener() throws Exception {
+ verify(queue, timeout(5000)).addItemListener(any(ItemListener.class),
eq(true));
+ assertEquals(1, listeners.size());
+
+ context.getRouteController().stopRoute("queue");
+
+ assertEquals(0, listeners.size());
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testItemReceivedOnceAfterRestart() throws Exception {
+ context.getRouteController().stopRoute("queue");
+ context.getRouteController().startRoute("queue");
+ verify(queue,
timeout(5000).times(2)).addItemListener(any(ItemListener.class), eq(true));
+
+ MockEndpoint added = getMockEndpoint("mock:added");
+ added.expectedMessageCount(1);
+ added.expectedHeaderReceived(HazelcastConstants.LISTENER_ACTION,
HazelcastConstants.ADDED);
+
+ // Hazelcast notifies every item listener registered on the queue
+ ItemEvent<String> event = new ItemEvent<>("foo", ItemEventType.ADDED,
"bar", null);
+ listeners.values().forEach(listener -> listener.itemAdded(event));
+
+ MockEndpoint.assertIsSatisfied(context, 5, TimeUnit.SECONDS);
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from(String.format("hazelcast-%sfoo",
HazelcastConstants.QUEUE_PREFIX)).routeId("queue")
+ .to("mock:added");
+ }
+ };
+ }
+}
diff --git
a/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerTest.java
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerTest.java
index 989fa3c5599a..2a150bb458dd 100644
---
a/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerTest.java
+++
b/components/camel-hazelcast/src/test/java/org/apache/camel/component/hazelcast/HazelcastQueueConsumerTest.java
@@ -19,7 +19,6 @@ package org.apache.camel.component.hazelcast;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
-import java.util.function.Consumer;
import com.hazelcast.collection.IQueue;
import com.hazelcast.collection.ItemEvent;
@@ -29,6 +28,7 @@ import com.hazelcast.core.ItemEventType;
import org.apache.camel.builder.RouteBuilder;
import org.apache.camel.component.mock.MockEndpoint;
import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -44,21 +44,11 @@ public class HazelcastQueueConsumerTest extends
HazelcastCamelTestSupport {
@Mock
private IQueue<String> queue;
- private volatile Consumer<ItemListener<String>> consumer;
-
@Override
@SuppressWarnings("unchecked")
protected void trainHazelcastInstance(HazelcastInstance hazelcastInstance)
{
when(hazelcastInstance.<String> getQueue("foo")).thenReturn(queue);
- when(queue.addItemListener(any(ItemListener.class),
eq(true))).thenAnswer(
- invocationOnMock -> {
- // Wait until the consumer is set
- while (consumer == null) {
- Thread.onSpinWait();
- }
- consumer.accept(invocationOnMock.getArgument(0,
ItemListener.class));
- return UUID.randomUUID();
- });
+ when(queue.addItemListener(any(ItemListener.class),
eq(true))).thenReturn(UUID.randomUUID());
}
@Override
@@ -70,25 +60,33 @@ public class HazelcastQueueConsumerTest extends
HazelcastCamelTestSupport {
@Test
public void add() throws InterruptedException {
- this.consumer = listener -> listener.itemAdded(new ItemEvent<>("foo",
ItemEventType.ADDED, "foo", null));
MockEndpoint out = getMockEndpoint("mock:added");
out.expectedMessageCount(1);
+ listener().itemAdded(new ItemEvent<>("foo", ItemEventType.ADDED,
"foo", null));
+
MockEndpoint.assertIsSatisfied(context, 2, TimeUnit.SECONDS);
this.checkHeaders(out.getExchanges().get(0).getIn().getHeaders(),
HazelcastConstants.ADDED);
}
@Test
public void remove() throws InterruptedException {
- this.consumer = listener -> listener.itemRemoved(new
ItemEvent<>("foo", ItemEventType.REMOVED, "foo", null));
-
MockEndpoint out = getMockEndpoint("mock:removed");
out.expectedMessageCount(1);
+ listener().itemRemoved(new ItemEvent<>("foo", ItemEventType.REMOVED,
"foo", null));
+
MockEndpoint.assertIsSatisfied(context, 2, TimeUnit.SECONDS);
this.checkHeaders(out.getExchanges().get(0).getIn().getHeaders(),
HazelcastConstants.REMOVED);
}
+ @SuppressWarnings("unchecked")
+ private ItemListener<String> listener() {
+ ArgumentCaptor<ItemListener<String>> captor =
ArgumentCaptor.forClass(ItemListener.class);
+ verify(queue).addItemListener(captor.capture(), eq(true));
+ return captor.getValue();
+ }
+
@Override
protected RouteBuilder createRouteBuilder() throws Exception {
return new RouteBuilder() {