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 ddd0fee0d23e CAMEL-25288: camel-dapr - consumers must only close the 
clients they created, and the configuration consumer must take the client of 
the started endpoint (#27320)
ddd0fee0d23e is described below

commit ddd0fee0d23e2355db16b4e681f97bf6f090c2e5
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:34:00 2026 +0530

    CAMEL-25288: camel-dapr - consumers must only close the clients they 
created, and the configuration consumer must take the client of the started 
endpoint (#27320)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../apache/camel/component/dapr/DaprEndpoint.java  |  36 ++++++
 .../dapr/consumer/DaprConfigurationConsumer.java   |  39 ++++--
 .../dapr/consumer/DaprPubSubConsumer.java          |  10 +-
 .../dapr/DaprEndpointClientOwnershipTest.java      |  91 +++++++++++++
 .../consumer/DaprConfigurationConsumerTest.java    |  11 +-
 .../consumer/DaprConsumerSharedClientTest.java     | 142 +++++++++++++++++++++
 .../dapr/consumer/DaprPubSubConsumerTest.java      |   4 +-
 7 files changed, 319 insertions(+), 14 deletions(-)

diff --git 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/DaprEndpoint.java
 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/DaprEndpoint.java
index 2b658452f32e..8a36b43f6842 100644
--- 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/DaprEndpoint.java
+++ 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/DaprEndpoint.java
@@ -31,6 +31,8 @@ import org.apache.camel.spi.HeaderFilterStrategyAware;
 import org.apache.camel.spi.UriEndpoint;
 import org.apache.camel.spi.UriParam;
 import org.apache.camel.support.DefaultEndpoint;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 /**
  * Dapr component which interfaces with Dapr Building Blocks.
@@ -39,6 +41,8 @@ import org.apache.camel.support.DefaultEndpoint;
         Category.CLOUD, Category.SAAS }, headersClass = DaprConstants.class)
 public class DaprEndpoint extends DefaultEndpoint implements 
HeaderFilterStrategyAware {
 
+    private static final Logger LOG = 
LoggerFactory.getLogger(DaprEndpoint.class);
+
     @UriParam
     private DaprConfiguration configuration;
     @UriParam(label = "advanced",
@@ -83,6 +87,38 @@ public class DaprEndpoint extends DefaultEndpoint implements 
HeaderFilterStrateg
                 = configuration.getWorkflowClient() != null ? 
configuration.getWorkflowClient() : new DaprWorkflowClient();
     }
 
+    @Override
+    public void doStop() throws Exception {
+        // close only the clients that this endpoint created: a configured (or 
autowired) client can be shared with
+        // other endpoints and is closed by its owner. The clients are created 
again when the endpoint is started again
+        if (client != configuration.getClient()) {
+            close(client);
+        }
+        if (previewClient != configuration.getPreviewClient()) {
+            close(previewClient);
+        }
+        if (workflowClient != configuration.getWorkflowClient()) {
+            close(workflowClient);
+        }
+        client = null;
+        previewClient = null;
+        workflowClient = null;
+        super.doStop();
+    }
+
+    private static void close(AutoCloseable daprClient) {
+        if (daprClient == null) {
+            return;
+        }
+        try {
+            daprClient.close();
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+        } catch (Exception e) {
+            LOG.debug("Failed to close {} due to: {}", daprClient, 
e.getMessage(), e);
+        }
+    }
+
     /**
      * The endpoint configurations
      */
diff --git 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumer.java
 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumer.java
index 237316e6cc77..02da12ece5bb 100644
--- 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumer.java
+++ 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumer.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.dapr.consumer;
 
+import java.time.Duration;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -33,21 +34,23 @@ import org.apache.camel.support.DefaultConsumer;
 import org.apache.camel.util.ObjectHelper;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
+import reactor.core.Disposable;
 import reactor.core.publisher.Flux;
 
 public class DaprConfigurationConsumer extends DefaultConsumer {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(DaprConfigurationConsumer.class);
+    private static final Duration UNSUBSCRIBE_TIMEOUT = Duration.ofSeconds(10);
     private final String configStore;
     private final List<String> configKeys;
-    private final DaprClient client;
-    private String subscriptionId;
+    private DaprClient client;
+    private volatile String subscriptionId;
+    private Disposable subscription;
 
     public DaprConfigurationConsumer(final DaprEndpoint endpoint, final 
Processor processor) {
         super(endpoint, processor);
         configStore = endpoint.getConfiguration().getConfigStore();
         configKeys = endpoint.getConfiguration().getConfigKeysAsList();
-        client = endpoint.getClient();
     }
 
     @Override
@@ -66,13 +69,17 @@ public class DaprConfigurationConsumer extends 
DefaultConsumer {
 
         LOG.debug("Creating connection to Dapr Configuration");
 
+        // the endpoint creates its client (or takes the configured one) when 
it starts, which can be after this
+        // consumer was created
+        client = getEndpoint().getClient();
+
         SubscribeConfigurationRequest configRequest = new 
SubscribeConfigurationRequest(configStore, configKeys);
-        Flux<SubscribeConfigurationResponse> subscription = 
client.subscribeConfiguration(configRequest);
+        Flux<SubscribeConfigurationResponse> responses = 
client.subscribeConfiguration(configRequest);
 
-        subscription.subscribe((response) -> {
-            // first ever response contains the subscription id
+        subscription = responses.subscribe((response) -> {
+            // every response contains the subscription id, the first ever 
response has no items
+            subscriptionId = response.getSubscriptionId();
             if (response.getItems() == null || response.getItems().isEmpty()) {
-                subscriptionId = response.getSubscriptionId();
                 LOG.debug("App subscribed to config changes with subscription 
id: {}", subscriptionId);
             } else {
                 final Exchange exchange = createServiceBusExchange(response);
@@ -86,9 +93,21 @@ public class DaprConfigurationConsumer extends 
DefaultConsumer {
 
     @Override
     protected void doStop() throws Exception {
-        if (client != null) {
-            client.unsubscribeConfiguration(subscriptionId, configStore);
-            client.close();
+        // the client belongs to the endpoint (or is configured or autowired, 
and shared with other endpoints): end the
+        // subscription but do not close the client, it is used again when the 
route is started again
+        if (subscription != null) {
+            subscription.dispose();
+            subscription = null;
+        }
+        String id = subscriptionId;
+        if (id != null) {
+            subscriptionId = null;
+            try {
+                client.unsubscribeConfiguration(id, 
configStore).block(UNSUBSCRIBE_TIMEOUT);
+            } catch (Exception e) {
+                LOG.warn("Failed to unsubscribe from config changes with 
subscription id: {} due to: {}", id,
+                        e.getMessage(), e);
+            }
         }
 
         // shutdown camel consumer
diff --git 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
index 5448b99b3631..facdf5d39f37 100644
--- 
a/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
+++ 
b/components/camel-dapr/src/main/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumer.java
@@ -42,6 +42,7 @@ public class DaprPubSubConsumer extends DefaultConsumer {
     private final String pubSubName;
     private final String topic;
     private DaprPreviewClient client;
+    private boolean closeClient;
     private Closeable subscription;
 
     public DaprPubSubConsumer(final DaprEndpoint endpoint, final Processor 
processor) {
@@ -68,6 +69,8 @@ public class DaprPubSubConsumer extends DefaultConsumer {
 
         if (client == null) {
             client = new DaprClientBuilder().buildPreviewClient();
+            // created by this consumer: closed when it stops
+            closeClient = true;
         }
         subscription = client.subscribeToEvents(pubSubName, topic, new 
DaprSubscriptionListener(), TypeRef.get(byte[].class));
     }
@@ -76,9 +79,14 @@ public class DaprPubSubConsumer extends DefaultConsumer {
     protected void doStop() throws Exception {
         if (subscription != null) {
             subscription.close();
+            subscription = null;
         }
-        if (client != null) {
+        // only close a client that this consumer created: a configured or 
autowired client can be shared with other
+        // endpoints, and is used again when the route is started again
+        if (closeClient && client != null) {
             client.close();
+            client = null;
+            closeClient = false;
         }
 
         // shutdown camel consumer
diff --git 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/DaprEndpointClientOwnershipTest.java
 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/DaprEndpointClientOwnershipTest.java
new file mode 100644
index 000000000000..89a5c63dc1c8
--- /dev/null
+++ 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/DaprEndpointClientOwnershipTest.java
@@ -0,0 +1,91 @@
+/*
+ * 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.dapr;
+
+import io.dapr.client.DaprClient;
+import io.dapr.client.DaprClientBuilder;
+import io.dapr.client.DaprPreviewClient;
+import io.dapr.workflows.client.DaprWorkflowClient;
+import org.apache.camel.test.junit6.CamelTestSupport;
+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.assertSame;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * An endpoint closes the Dapr clients that it created when it stops, and 
never the configured ones, which can be shared
+ * with other endpoints.
+ */
+class DaprEndpointClientOwnershipTest extends CamelTestSupport {
+
+    private static final String URI = 
"dapr:invokeService?serviceToInvoke=myService&methodToInvoke=myMethod";
+
+    private final DaprClient client = mock(DaprClient.class);
+    private final DaprPreviewClient previewClient = 
mock(DaprPreviewClient.class);
+    private final DaprWorkflowClient workflowClient = 
mock(DaprWorkflowClient.class);
+
+    @Test
+    void closesTheClientsItCreated() throws Exception {
+        try (MockedConstruction<DaprClientBuilder> builders = 
mockConstruction(DaprClientBuilder.class, (builder, ctx) -> {
+            when(builder.build()).thenReturn(client);
+            when(builder.buildPreviewClient()).thenReturn(previewClient);
+        });
+             MockedConstruction<DaprWorkflowClient> workflowClients = 
mockConstruction(DaprWorkflowClient.class)) {
+
+            DaprEndpoint endpoint = context.getEndpoint(URI, 
DaprEndpoint.class);
+            assertSame(client, endpoint.getClient());
+            assertSame(previewClient, endpoint.getPreviewClient());
+            assertEquals(1, workflowClients.constructed().size());
+
+            endpoint.stop();
+
+            verify(client).close();
+            verify(previewClient).close();
+            verify(workflowClients.constructed().get(0)).close();
+        }
+    }
+
+    @Test
+    void doesNotCloseTheConfiguredClients() throws Exception {
+        context.getRegistry().bind("myClient", client);
+        context.getRegistry().bind("myPreviewClient", previewClient);
+        context.getRegistry().bind("myWorkflowClient", workflowClient);
+
+        DaprEndpoint endpoint = context.getEndpoint(
+                URI + 
"&client=#myClient&previewClient=#myPreviewClient&workflowClient=#myWorkflowClient",
+                DaprEndpoint.class);
+        assertSame(client, endpoint.getClient());
+
+        endpoint.stop();
+
+        verify(client, never()).close();
+        verify(previewClient, never()).close();
+        verify(workflowClient, never()).close();
+
+        // and uses them again when it is started again
+        endpoint.start();
+        assertSame(client, endpoint.getClient());
+        assertSame(previewClient, endpoint.getPreviewClient());
+        assertSame(workflowClient, endpoint.getWorkflowClient());
+    }
+}
diff --git 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumerTest.java
 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumerTest.java
index 7fe2ac23d2b8..5070c41765d7 100644
--- 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumerTest.java
+++ 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprConfigurationConsumerTest.java
@@ -23,6 +23,7 @@ import io.dapr.client.DaprClient;
 import io.dapr.client.domain.ConfigurationItem;
 import io.dapr.client.domain.SubscribeConfigurationRequest;
 import io.dapr.client.domain.SubscribeConfigurationResponse;
+import io.dapr.client.domain.UnsubscribeConfigurationResponse;
 import org.apache.camel.AsyncCallback;
 import org.apache.camel.AsyncProcessor;
 import org.apache.camel.CamelContext;
@@ -39,12 +40,15 @@ import org.junit.jupiter.api.Test;
 import org.mockito.ArgumentCaptor;
 import reactor.core.publisher.Flux;
 import reactor.core.publisher.FluxSink;
+import reactor.core.publisher.Mono;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -85,6 +89,8 @@ public class DaprConfigurationConsumerTest extends 
CamelTestSupport {
         final FluxSink<SubscribeConfigurationResponse>[] sinkHolder = new 
FluxSink[1];
         Flux<SubscribeConfigurationResponse> flux = Flux.create(sink -> 
sinkHolder[0] = sink);
         when(mockClient.subscribeConfiguration(any())).thenReturn(flux);
+        when(mockClient.unsubscribeConfiguration(anyString(), anyString()))
+                .thenReturn(Mono.just(new 
UnsubscribeConfigurationResponse(true, "")));
 
         consumer.doStart();
 
@@ -108,7 +114,8 @@ public class DaprConfigurationConsumerTest extends 
CamelTestSupport {
         assertEquals(mockBody, exchange.getIn().getBody());
 
         consumer.doStop();
-        verify(mockClient).unsubscribeConfiguration(any(), any());
-        verify(mockClient).close();
+        verify(mockClient).unsubscribeConfiguration("mySubId", "myStore");
+        // the client belongs to the endpoint: it is not closed
+        verify(mockClient, never()).close();
     }
 }
diff --git 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprConsumerSharedClientTest.java
 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprConsumerSharedClientTest.java
new file mode 100644
index 000000000000..c399761266b5
--- /dev/null
+++ 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprConsumerSharedClientTest.java
@@ -0,0 +1,142 @@
+/*
+ * 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.dapr.consumer;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import io.dapr.client.DaprClient;
+import io.dapr.client.DaprPreviewClient;
+import io.dapr.client.Subscription;
+import io.dapr.client.SubscriptionListener;
+import io.dapr.client.domain.ConfigurationItem;
+import io.dapr.client.domain.SubscribeConfigurationRequest;
+import io.dapr.client.domain.SubscribeConfigurationResponse;
+import io.dapr.client.domain.UnsubscribeConfigurationResponse;
+import io.dapr.utils.TypeRef;
+import io.dapr.workflows.client.DaprWorkflowClient;
+import org.apache.camel.BindToRegistry;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.FluxSink;
+import reactor.core.publisher.Mono;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+
+/**
+ * The Dapr clients in the registry are autowired and shared by all the dapr 
endpoints: a consumer must not close them
+ * when it stops, and must be able to consume again when its route is started 
again.
+ */
+class DaprConsumerSharedClientTest extends CamelTestSupport {
+
+    private final DaprClient client = mock(DaprClient.class);
+    private final DaprPreviewClient previewClient = 
mock(DaprPreviewClient.class);
+    private final DaprWorkflowClient workflowClient = 
mock(DaprWorkflowClient.class);
+
+    private final List<FluxSink<SubscribeConfigurationResponse>> 
configSubscriptions = new CopyOnWriteArrayList<>();
+    private final AtomicBoolean configSubscriptionCancelled = new 
AtomicBoolean();
+    private final List<String> unsubscribed = new CopyOnWriteArrayList<>();
+
+    @BindToRegistry("daprClient")
+    public DaprClient daprClient() {
+        doAnswer(inv -> Flux.<SubscribeConfigurationResponse> create(sink -> {
+            sink.onCancel(() -> configSubscriptionCancelled.set(true));
+            configSubscriptions.add(sink);
+        
})).when(client).subscribeConfiguration(any(SubscribeConfigurationRequest.class));
+        doAnswer(inv -> Mono.fromCallable(() -> {
+            unsubscribed.add(inv.getArgument(0));
+            return new UnsubscribeConfigurationResponse(true, "");
+        })).when(client).unsubscribeConfiguration(anyString(), anyString());
+        return client;
+    }
+
+    @BindToRegistry("daprPreviewClient")
+    public DaprPreviewClient daprPreviewClient() {
+        doAnswer(inv -> mock(Subscription.class)).when(previewClient)
+                .subscribeToEvents(anyString(), anyString(), 
any(SubscriptionListener.class), any(TypeRef.class));
+        return previewClient;
+    }
+
+    @BindToRegistry("daprWorkflowClient")
+    public DaprWorkflowClient daprWorkflowClient() {
+        return workflowClient;
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("dapr:pubSub?pubSubName=myPubSub&topic=myTopic").routeId("pubSub")
+                        .to("mock:pubSub");
+                
from("dapr:configuration?configStore=myStore&configKeys=myKey").routeId("configuration").autoStartup(false)
+                        .to("mock:configuration");
+            }
+        };
+    }
+
+    @Test
+    void pubSubConsumerDoesNotCloseTheSharedClient() throws Exception {
+        context.getRouteController().stopRoute("pubSub");
+
+        verify(previewClient, never()).close();
+
+        context.getRouteController().startRoute("pubSub");
+
+        verify(previewClient, times(2)).subscribeToEvents(anyString(), 
anyString(), any(SubscriptionListener.class),
+                any(TypeRef.class));
+    }
+
+    @Test
+    void configurationConsumerUnsubscribesAndDoesNotCloseTheSharedClient() 
throws Exception {
+        // the consumer is created before its endpoint is started (and has 
created or taken its client)
+        context.getRouteController().startRoute("configuration");
+        assertEquals(1, configSubscriptions.size());
+        configSubscriptions.get(0).next(new 
SubscribeConfigurationResponse("sub-1", Map.of()));
+
+        context.getRouteController().stopRoute("configuration");
+
+        // the subscription is ended, the client stays open for the other 
endpoints
+        assertTrue(configSubscriptionCancelled.get());
+        assertEquals(List.of("sub-1"), unsubscribed);
+        verify(client, never()).close();
+
+        // and the route consumes again when it is started again
+        context.getRouteController().startRoute("configuration");
+        assertEquals(2, configSubscriptions.size());
+
+        MockEndpoint mock = getMockEndpoint("mock:configuration");
+        mock.expectedMessageCount(1);
+        configSubscriptions.get(1).next(new SubscribeConfigurationResponse(
+                "sub-2", Map.of("myKey", new ConfigurationItem("myKey", 
"myValue", "1"))));
+        mock.assertIsSatisfied();
+        assertEquals(Map.of("myKey", "myValue"), 
mock.getReceivedExchanges().get(0).getIn().getBody());
+    }
+}
diff --git 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
index 9461469ec887..13d308ef4993 100644
--- 
a/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
+++ 
b/components/camel-dapr/src/test/java/org/apache/camel/component/dapr/consumer/DaprPubSubConsumerTest.java
@@ -54,6 +54,7 @@ import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -155,7 +156,8 @@ public class DaprPubSubConsumerTest extends 
CamelTestSupport {
 
         consumer.doStop();
         verify(mockSubscription).close();
-        verify(mockClient).close();
+        // the client is configured, not created by the consumer: it is not 
closed
+        verify(mockClient, never()).close();
     }
 
     @Test

Reply via email to