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() {

Reply via email to