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 d1f0dd80bda7 CAMEL-24401: Fix permanent temporary replyTo strand in
camel-jms (#25531)
d1f0dd80bda7 is described below
commit d1f0dd80bda7d7204f0a7bf638660df4b37d4fb5
Author: Omar Atie <[email protected]>
AuthorDate: Mon Aug 24 13:07:38 2026 -0700
CAMEL-24401: Fix permanent temporary replyTo strand in camel-jms (#25531)
camel-jms InOut producers using temporary reply queues could enter a
permanent "Failed to resolve replyTo destination" loop after a JMS
connection fault. TemporaryQueueReplyManager tracked a one-shot
refreshWanted flag that was cleared before the replacement temporary
queue was published, and could stay armed with no consumer left to
invoke the destination resolver.
Replaces the flag with a monotonic refresh-generation counter that is
only acknowledged once a replacement queue is created and published to
replyTo for that generation. When a refresh is scheduled, listener-
container recovery is deferred until after the endpoint recovery
interval and skipped while Spring is already recovering, so a pending
refresh can no longer outlive the last resolver call. Adds a guarded
DefaultJmsMessageListenerContainer.recoverReplyDestinationAfterRefresh()
hook into Spring's recovery path to cover the case where no consumer
remains to trigger a resolve.
Adds TemporaryQueueReplyManagerRefreshTest (failed create retries,
latest-generation publish, permanent failure, generation change during
resolve) and a JmsTemporaryReplyToRequestReplyIT consecutive
request-reply check.
Co-authored-by: Cursor <[email protected]>
---
.../jms/DefaultJmsMessageListenerContainer.java | 11 +
.../jms/reply/TemporaryQueueReplyManager.java | 162 +++++++++--
.../reply/JmsTemporaryReplyToRequestReplyIT.java | 68 +++++
.../TemporaryQueueReplyManagerRefreshTest.java | 322 +++++++++++++++++++++
4 files changed, 541 insertions(+), 22 deletions(-)
diff --git
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/DefaultJmsMessageListenerContainer.java
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/DefaultJmsMessageListenerContainer.java
index a98c1a60c186..fe48be13475e 100644
---
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/DefaultJmsMessageListenerContainer.java
+++
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/DefaultJmsMessageListenerContainer.java
@@ -161,4 +161,15 @@ public class DefaultJmsMessageListenerContainer extends
DefaultMessageListenerCo
}
super.stopSharedConnection();
}
+
+ /**
+ * Forces the listener container to refresh its connection and recreate
consumers, which will re-resolve the
+ * temporary reply destination. Used when a pending reply-destination
refresh must be consumed but cached consumers
+ * would otherwise never call the destination resolver again.
+ */
+ public void recoverReplyDestinationAfterRefresh() {
+ if (isRunning() && !isRecovering()) {
+ recoverAfterListenerSetupFailure();
+ }
+ }
}
diff --git
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/TemporaryQueueReplyManager.java
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/TemporaryQueueReplyManager.java
index 6f5b4ef9275f..9d4301131107 100644
---
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/TemporaryQueueReplyManager.java
+++
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/TemporaryQueueReplyManager.java
@@ -16,7 +16,10 @@
*/
package org.apache.camel.component.jms.reply;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@@ -40,6 +43,7 @@ import org.apache.camel.support.service.ServiceHelper;
import org.apache.camel.support.service.ServiceSupport;
import org.springframework.jms.listener.AbstractMessageListenerContainer;
import org.springframework.jms.listener.DefaultMessageListenerContainer;
+import org.springframework.jms.listener.SimpleMessageListenerContainer;
import org.springframework.jms.support.destination.DestinationResolver;
/**
@@ -48,6 +52,8 @@ import
org.springframework.jms.support.destination.DestinationResolver;
public class TemporaryQueueReplyManager extends ReplyManagerSupport {
final TemporaryReplyQueueDestinationResolver destinationResolver;
+ private ExecutorService refreshRecoveryExecutor;
+ private final AtomicBoolean refreshRecoveryScheduled = new AtomicBoolean();
public TemporaryQueueReplyManager(CamelContext camelContext,
TemporaryQueueResolver resolver) {
super(camelContext);
@@ -57,9 +63,95 @@ public class TemporaryQueueReplyManager extends
ReplyManagerSupport {
@Override
protected void doStop() throws Exception {
super.doStop();
+ if (refreshRecoveryExecutor != null) {
+
camelContext.getExecutorServiceManager().shutdownNow(refreshRecoveryExecutor);
+ refreshRecoveryExecutor = null;
+ }
ServiceHelper.stopService(destinationResolver);
}
+ private void triggerReplyDestinationRecovery() {
+ if (listenerContainer == null || isStopping() || isStopped()) {
+ return;
+ }
+ if (!destinationResolver.isRefreshPending()) {
+ return;
+ }
+ if (!refreshRecoveryScheduled.compareAndSet(false, true)) {
+ return;
+ }
+ try {
+
getRefreshRecoveryExecutor().execute(this::runReplyDestinationRecovery);
+ } catch (RejectedExecutionException e) {
+ refreshRecoveryScheduled.set(false);
+ }
+ }
+
+ private void runReplyDestinationRecovery() {
+ try {
+ long delay = endpoint.getRecoveryInterval() >= 0 ?
endpoint.getRecoveryInterval() : 5000L;
+ if (!sleepQuietly(delay)) {
+ return;
+ }
+ int attempts = 0;
+ while (destinationResolver.isRefreshPending() && !isStopping() &&
!isStopped()
+ && listenerContainer != null &&
listenerContainer.isRunning() && attempts < 20) {
+ try {
+ if (listenerContainer instanceof
DefaultJmsMessageListenerContainer dmlc) {
+ if (dmlc.isRecovering()) {
+ if (!sleepQuietly(delay)) {
+ return;
+ }
+ attempts++;
+ continue;
+ }
+ dmlc.recoverReplyDestinationAfterRefresh();
+ } else if (listenerContainer instanceof
SimpleMessageListenerContainer smlc) {
+ smlc.stop();
+ smlc.start();
+ } else {
+ listenerContainer.stop();
+ listenerContainer.start();
+ }
+ } catch (Exception e) {
+ log.warn("Failed to trigger recovery of temporary reply
destination on endpoint: {}",
+ endpoint.getEndpointUri(), e);
+ }
+ if (!destinationResolver.isRefreshPending()) {
+ break;
+ }
+ if (!sleepQuietly(delay)) {
+ return;
+ }
+ attempts++;
+ }
+ } finally {
+ refreshRecoveryScheduled.set(false);
+ if (destinationResolver.isRefreshPending() && !isStopping() &&
!isStopped()
+ && listenerContainer != null &&
listenerContainer.isRunning()) {
+ triggerReplyDestinationRecovery();
+ }
+ }
+ }
+
+ private boolean sleepQuietly(long millis) {
+ try {
+ Thread.sleep(millis);
+ return !isStopping() && !isStopped();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ return false;
+ }
+ }
+
+ private ExecutorService getRefreshRecoveryExecutor() {
+ if (refreshRecoveryExecutor == null) {
+ String name = "JmsTemporaryReplyToRefresh[" +
endpoint.getDestinationName() + "]";
+ refreshRecoveryExecutor =
camelContext.getExecutorServiceManager().newSingleThreadExecutor(this, name);
+ }
+ return refreshRecoveryExecutor;
+ }
+
@Override
protected ReplyHandler createReplyHandler(
ReplyManager replyManager, Exchange exchange, AsyncCallback
callback,
@@ -283,51 +375,76 @@ public class TemporaryQueueReplyManager extends
ReplyManagerSupport {
// the destination, it would deadlock trying to acquire
BaseService.lock.
private final Lock destinationLock = new ReentrantLock();
private volatile TemporaryQueue queue;
- private final AtomicBoolean refreshWanted = new AtomicBoolean();
+ private final AtomicLong refreshGeneration = new AtomicLong();
+ private volatile long publishedGeneration;
private final TemporaryQueueResolver custom;
public TemporaryReplyQueueDestinationResolver(TemporaryQueueResolver
custom) {
this.custom = custom;
}
+ boolean isRefreshPending() {
+ return refreshGeneration.get() != publishedGeneration;
+ }
+
@Override
public Destination resolveDestinationName(Session session, String
destinationName, boolean pubSubDomain)
throws JMSException {
- // fast path: queue already resolved and no refresh needed
- TemporaryQueue answer = queue;
- if (answer != null && !refreshWanted.get()) {
- return answer;
- }
destinationLock.lock();
try {
- if (queue == null || refreshWanted.get()) {
- refreshWanted.set(false);
- if (custom != null) {
- if (queue != null) {
- try {
- custom.delete(queue);
- } catch (Exception e) {
- // ignore
- }
+ TemporaryQueue answer = queue;
+ if (answer != null && !isRefreshPending()) {
+ return answer;
+ }
+ long generationToHandle = refreshGeneration.get();
+ TemporaryQueue previousQueue = queue;
+ if (previousQueue != null) {
+ try {
+ if (custom != null) {
+ custom.delete(previousQueue);
+ } else {
+ previousQueue.delete();
}
- queue = custom.createTemporaryQueue(session);
- } else {
- queue = session.createTemporaryQueue();
+ } catch (Exception e) {
+ // ignore
}
- setReplyTo(queue);
+ queue = null;
+ }
+ TemporaryQueue refreshedQueue;
+ if (custom != null) {
+ refreshedQueue = custom.createTemporaryQueue(session);
+ } else {
+ refreshedQueue = session.createTemporaryQueue();
+ }
+ if (refreshGeneration.get() == generationToHandle) {
+ queue = refreshedQueue;
+ setReplyTo(refreshedQueue);
+ publishedGeneration = generationToHandle;
if (log.isDebugEnabled()) {
- log.debug("Refreshed Temporary ReplyTo Queue. New
queue: {}", queue.getQueueName());
+ log.debug("Refreshed Temporary ReplyTo Queue. New
queue: {}", refreshedQueue.getQueueName());
+ }
+ return refreshedQueue;
+ }
+ // a newer refresh was requested while creating the queue;
discard this attempt
+ try {
+ if (custom != null) {
+ custom.delete(refreshedQueue);
+ } else {
+ refreshedQueue.delete();
}
+ } catch (Exception e) {
+ // ignore
}
+ return null;
} finally {
destinationLock.unlock();
}
- return queue;
}
public void scheduleRefresh() {
- refreshWanted.set(true);
+ refreshGeneration.incrementAndGet();
replyTo = null;
+ triggerReplyDestinationRecovery();
}
@Override
@@ -346,6 +463,7 @@ public class TemporaryQueueReplyManager extends
ReplyManagerSupport {
}
queue = null;
}
+ publishedGeneration = refreshGeneration.get();
} finally {
destinationLock.unlock();
}
diff --git
a/components/camel-jms/src/test/java/org/apache/camel/component/jms/reply/JmsTemporaryReplyToRequestReplyIT.java
b/components/camel-jms/src/test/java/org/apache/camel/component/jms/reply/JmsTemporaryReplyToRequestReplyIT.java
new file mode 100644
index 000000000000..880a03c4a394
--- /dev/null
+++
b/components/camel-jms/src/test/java/org/apache/camel/component/jms/reply/JmsTemporaryReplyToRequestReplyIT.java
@@ -0,0 +1,68 @@
+/*
+ * 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.jms.reply;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ProducerTemplate;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.jms.JmsComponent;
+import org.apache.camel.test.infra.artemis.common.ConnectionFactoryHelper;
+import org.apache.camel.test.infra.artemis.services.ArtemisService;
+import org.apache.camel.test.infra.artemis.services.ArtemisServiceFactory;
+import org.apache.camel.test.infra.core.DefaultCamelContextExtension;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static
org.apache.camel.component.jms.JmsComponent.jmsComponentAutoAcknowledge;
+import static org.assertj.core.api.Assertions.assertThat;
+
+class JmsTemporaryReplyToRequestReplyIT {
+
+ private static final String REQUEST_QUEUE =
"JmsTemporaryReplyToRequestReplyIT.request";
+
+ @RegisterExtension
+ static ArtemisService service = ArtemisServiceFactory.createVMService();
+
+ @RegisterExtension
+ static DefaultCamelContextExtension contextExtension = new
DefaultCamelContextExtension();
+
+ private ProducerTemplate template;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ CamelContext context = contextExtension.getContext();
+ JmsComponent component =
jmsComponentAutoAcknowledge(ConnectionFactoryHelper.createConnectionFactory(service));
+ context.addComponent("jms", component);
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("jms:queue:" + REQUEST_QUEUE).routeId("server")
+ .transform(simple("echo:${body}"));
+ }
+ });
+ template = contextExtension.getProducerTemplate();
+ }
+
+ @Test
+ void shouldSupportConsecutiveTemporaryReplyRequests() {
+ assertThat(template.requestBody("jms:queue:" + REQUEST_QUEUE, "first",
String.class))
+ .isEqualTo("echo:first");
+ assertThat(template.requestBody("jms:queue:" + REQUEST_QUEUE,
"second", String.class))
+ .isEqualTo("echo:second");
+ }
+}
diff --git
a/components/camel-jms/src/test/java/org/apache/camel/component/jms/reply/TemporaryQueueReplyManagerRefreshTest.java
b/components/camel-jms/src/test/java/org/apache/camel/component/jms/reply/TemporaryQueueReplyManagerRefreshTest.java
new file mode 100644
index 000000000000..5daff22434d3
--- /dev/null
+++
b/components/camel-jms/src/test/java/org/apache/camel/component/jms/reply/TemporaryQueueReplyManagerRefreshTest.java
@@ -0,0 +1,322 @@
+/*
+ * 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.jms.reply;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import jakarta.jms.Connection;
+import jakarta.jms.Destination;
+import jakarta.jms.JMSException;
+import jakarta.jms.Session;
+import jakarta.jms.TemporaryQueue;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.component.jms.DefaultJmsMessageListenerContainer;
+import org.apache.camel.component.jms.JmsComponent;
+import org.apache.camel.component.jms.JmsConfiguration;
+import org.apache.camel.component.jms.JmsEndpoint;
+import org.apache.camel.component.jms.TemporaryQueueResolver;
+import org.apache.camel.support.service.ServiceHelper;
+import org.apache.camel.test.infra.artemis.common.ConnectionFactoryHelper;
+import org.apache.camel.test.infra.artemis.services.ArtemisService;
+import org.apache.camel.test.infra.artemis.services.ArtemisServiceFactory;
+import org.apache.camel.test.infra.core.DefaultCamelContextExtension;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static
org.apache.camel.component.jms.JmsComponent.jmsComponentAutoAcknowledge;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.awaitility.Awaitility.await;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+class TemporaryQueueReplyManagerRefreshTest {
+
+ @RegisterExtension
+ static DefaultCamelContextExtension contextExtension = new
DefaultCamelContextExtension();
+
+ @RegisterExtension
+ static ArtemisService service = ArtemisServiceFactory.createVMService();
+
+ private TemporaryQueueReplyManager replyManager;
+ private JmsEndpoint startedEndpoint;
+
+ @BeforeEach
+ void setUp() {
+ CamelContext context = contextExtension.getContext();
+ JmsEndpoint endpoint = mock(JmsEndpoint.class);
+ when(endpoint.getCamelContext()).thenReturn(context);
+ when(endpoint.getConfiguration()).thenReturn(new JmsConfiguration());
+
when(endpoint.getDestinationName()).thenReturn("TemporaryQueueReplyManagerRefreshTest");
+
when(endpoint.getEndpointUri()).thenReturn("jms:queue:TemporaryQueueReplyManagerRefreshTest");
+
+ replyManager = new TemporaryQueueReplyManager(context, null);
+ replyManager.setEndpoint(endpoint);
+ }
+
+ @AfterEach
+ void tearDown() throws Exception {
+ if (replyManager != null && replyManager.isStarted()) {
+ replyManager.stop();
+ }
+ if (startedEndpoint != null && startedEndpoint.isStarted()) {
+ startedEndpoint.stop();
+ }
+ }
+
+ @Test
+ void
shouldRecoverReplyDestinationAfterRefreshWithRunningListenerContainer() throws
Exception {
+ CamelContext context = contextExtension.getContext();
+ JmsComponent component =
jmsComponentAutoAcknowledge(ConnectionFactoryHelper.createConnectionFactory(service));
+ context.addComponent("jms", component);
+
+ startedEndpoint = (JmsEndpoint) component.createEndpoint(
+
"jms:queue:TemporaryQueueReplyManagerRefreshTest?recoveryInterval=100");
+ ServiceHelper.startService(startedEndpoint);
+
+ replyManager = new TemporaryQueueReplyManager(context, null);
+ replyManager.setEndpoint(startedEndpoint);
+
+ ScheduledExecutorService scheduledExecutorService
+ =
context.getExecutorServiceManager().newSingleThreadScheduledExecutor(replyManager,
"test-timeout-checker");
+ ExecutorService onTimeoutExecutorService
+ =
context.getExecutorServiceManager().newThreadPool(replyManager,
"test-on-timeout", 0, 1);
+ replyManager.setScheduledExecutorService(scheduledExecutorService);
+ replyManager.setOnTimeoutExecutorService(onTimeoutExecutorService);
+ replyManager.start();
+
+
assertThat(replyManager.listenerContainer).isInstanceOf(DefaultJmsMessageListenerContainer.class);
+ DefaultJmsMessageListenerContainer listenerContainer
+ = (DefaultJmsMessageListenerContainer)
replyManager.listenerContainer;
+ assertThat(listenerContainer.isRunning()).isTrue();
+
+ TemporaryQueueReplyManager.TemporaryReplyQueueDestinationResolver
destinationResolver
+ = replyManager.destinationResolver;
+
+ Destination firstReplyTo = replyManager.getReplyTo();
+ assertThat(firstReplyTo).isNotNull();
+ assertThat(destinationResolver.isRefreshPending()).isFalse();
+
+ destinationResolver.scheduleRefresh();
+ assertThat(destinationResolver.isRefreshPending()).isTrue();
+
+ await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> {
+ assertThat(destinationResolver.isRefreshPending()).isFalse();
+ assertThat(replyManager.getReplyTo()).isNotNull();
+ assertThat(listenerContainer.isRunning()).isTrue();
+ });
+
+ assertThat(replyManager.getReplyTo()).isNotEqualTo(firstReplyTo);
+ }
+
+ @Test
+ void shouldRetryAfterFailedTemporaryQueueCreation() throws Exception {
+ FailOnSecondCreateResolver resolver = new FailOnSecondCreateResolver();
+ replyManager = new
TemporaryQueueReplyManager(contextExtension.getContext(), resolver);
+ replyManager.setEndpoint(createEndpoint());
+
+ try (Connection connection =
ConnectionFactoryHelper.createConnectionFactory(service).createConnection()) {
+ connection.start();
+ Session session = connection.createSession(false,
Session.AUTO_ACKNOWLEDGE);
+
+ TemporaryQueueReplyManager.TemporaryReplyQueueDestinationResolver
destinationResolver
+ = replyManager.destinationResolver;
+
+ destinationResolver.resolveDestinationName(session, "temporary",
false);
+ assertThat(replyManager.getReplyTo()).isNotNull();
+ assertThat(resolver.createAttempts.get()).isEqualTo(1);
+
+ destinationResolver.scheduleRefresh();
+ assertThat(destinationResolver.isRefreshPending()).isTrue();
+ assertThat(replyManager.getReplyTo()).isNull();
+
+ assertThatThrownBy(() ->
destinationResolver.resolveDestinationName(session, "temporary", false))
+ .isInstanceOf(JMSException.class);
+ assertThat(destinationResolver.isRefreshPending()).isTrue();
+ assertThat(replyManager.getReplyTo()).isNull();
+ assertThat(resolver.createAttempts.get()).isEqualTo(2);
+
+ destinationResolver.resolveDestinationName(session, "temporary",
false);
+ assertThat(destinationResolver.isRefreshPending()).isFalse();
+ assertThat(replyManager.getReplyTo()).isNotNull();
+ assertThat(resolver.createAttempts.get()).isEqualTo(3);
+ }
+ }
+
+ @Test
+ void shouldPublishLatestRefreshGeneration() throws Exception {
+ CountingTemporaryQueueResolver resolver = new
CountingTemporaryQueueResolver();
+ replyManager = new
TemporaryQueueReplyManager(contextExtension.getContext(), resolver);
+ replyManager.setEndpoint(createEndpoint());
+
+ try (Connection connection =
ConnectionFactoryHelper.createConnectionFactory(service).createConnection()) {
+ connection.start();
+ Session session = connection.createSession(false,
Session.AUTO_ACKNOWLEDGE);
+
+ TemporaryQueueReplyManager.TemporaryReplyQueueDestinationResolver
destinationResolver
+ = replyManager.destinationResolver;
+
+ TemporaryQueue firstQueue
+ = (TemporaryQueue)
destinationResolver.resolveDestinationName(session, "temporary", false);
+ assertThat(replyManager.getReplyTo()).isEqualTo(firstQueue);
+
+ destinationResolver.scheduleRefresh();
+ destinationResolver.scheduleRefresh();
+
+ TemporaryQueue resolvedQueue
+ = (TemporaryQueue)
destinationResolver.resolveDestinationName(session, "temporary", false);
+
+ assertThat(destinationResolver.isRefreshPending()).isFalse();
+ assertThat(replyManager.getReplyTo()).isEqualTo(resolvedQueue);
+ assertThat(resolvedQueue).isNotEqualTo(firstQueue);
+ assertThat(resolver.createAttempts.get()).isEqualTo(2);
+ }
+ }
+
+ @Test
+ void shouldKeepRefreshPendingWhenCreationKeepsFailing() throws Exception {
+ AlwaysFailingTemporaryQueueResolver resolver = new
AlwaysFailingTemporaryQueueResolver();
+ replyManager = new
TemporaryQueueReplyManager(contextExtension.getContext(), resolver);
+ replyManager.setEndpoint(createEndpoint());
+
+ try (Connection connection =
ConnectionFactoryHelper.createConnectionFactory(service).createConnection()) {
+ connection.start();
+ Session session = connection.createSession(false,
Session.AUTO_ACKNOWLEDGE);
+
+ TemporaryQueueReplyManager.TemporaryReplyQueueDestinationResolver
destinationResolver
+ = replyManager.destinationResolver;
+
+ destinationResolver.scheduleRefresh();
+ assertThatThrownBy(() ->
destinationResolver.resolveDestinationName(session, "temporary", false))
+ .isInstanceOf(JMSException.class);
+
+ assertThat(destinationResolver.isRefreshPending()).isTrue();
+ assertThat(replyManager.getReplyTo()).isNull();
+ }
+ }
+
+ @Test
+ void shouldDiscardPublishWhenRefreshGenerationChangesDuringResolve()
throws Exception {
+ InterleavingRefreshResolver resolver = new
InterleavingRefreshResolver();
+ replyManager = new
TemporaryQueueReplyManager(contextExtension.getContext(), resolver);
+ replyManager.setEndpoint(createEndpoint());
+ resolver.attach(replyManager.destinationResolver);
+
+ try (Connection connection =
ConnectionFactoryHelper.createConnectionFactory(service).createConnection()) {
+ connection.start();
+ Session session = connection.createSession(false,
Session.AUTO_ACKNOWLEDGE);
+
+ TemporaryQueueReplyManager.TemporaryReplyQueueDestinationResolver
destinationResolver
+ = replyManager.destinationResolver;
+
+ destinationResolver.resolveDestinationName(session, "temporary",
false);
+ assertThat(replyManager.getReplyTo()).isNotNull();
+
+ destinationResolver.scheduleRefresh();
+ assertThat(destinationResolver.resolveDestinationName(session,
"temporary", false)).isNull();
+ assertThat(destinationResolver.isRefreshPending()).isTrue();
+ assertThat(replyManager.getReplyTo()).isNull();
+
+ TemporaryQueue resolvedQueue
+ = (TemporaryQueue)
destinationResolver.resolveDestinationName(session, "temporary", false);
+ assertThat(resolvedQueue).isNotNull();
+ assertThat(destinationResolver.isRefreshPending()).isFalse();
+ assertThat(replyManager.getReplyTo()).isEqualTo(resolvedQueue);
+ }
+ }
+
+ private JmsEndpoint createEndpoint() {
+ JmsEndpoint endpoint = mock(JmsEndpoint.class);
+
when(endpoint.getCamelContext()).thenReturn(contextExtension.getContext());
+ when(endpoint.getConfiguration()).thenReturn(new JmsConfiguration());
+
when(endpoint.getDestinationName()).thenReturn("TemporaryQueueReplyManagerRefreshTest");
+
when(endpoint.getEndpointUri()).thenReturn("jms:queue:TemporaryQueueReplyManagerRefreshTest");
+ return endpoint;
+ }
+
+ private static final class FailOnSecondCreateResolver implements
TemporaryQueueResolver {
+ private final AtomicInteger createAttempts = new AtomicInteger();
+
+ @Override
+ public TemporaryQueue createTemporaryQueue(Session session) throws
JMSException {
+ if (createAttempts.incrementAndGet() == 2) {
+ throw new JMSException("simulated broker failure");
+ }
+ return session.createTemporaryQueue();
+ }
+
+ @Override
+ public void delete(TemporaryQueue queue) {
+ // noop
+ }
+ }
+
+ private static final class CountingTemporaryQueueResolver implements
TemporaryQueueResolver {
+ private final AtomicInteger createAttempts = new AtomicInteger();
+
+ @Override
+ public TemporaryQueue createTemporaryQueue(Session session) throws
JMSException {
+ createAttempts.incrementAndGet();
+ return session.createTemporaryQueue();
+ }
+
+ @Override
+ public void delete(TemporaryQueue queue) {
+ // noop
+ }
+ }
+
+ private static final class AlwaysFailingTemporaryQueueResolver implements
TemporaryQueueResolver {
+ @Override
+ public TemporaryQueue createTemporaryQueue(Session session) throws
JMSException {
+ throw new JMSException("simulated broker failure");
+ }
+
+ @Override
+ public void delete(TemporaryQueue queue) {
+ // noop
+ }
+ }
+
+ private static final class InterleavingRefreshResolver implements
TemporaryQueueResolver {
+ private
TemporaryQueueReplyManager.TemporaryReplyQueueDestinationResolver
destinationResolver;
+ private final AtomicInteger createAttempts = new AtomicInteger();
+
+ void
attach(TemporaryQueueReplyManager.TemporaryReplyQueueDestinationResolver
destinationResolver) {
+ this.destinationResolver = destinationResolver;
+ }
+
+ @Override
+ public TemporaryQueue createTemporaryQueue(Session session) throws
JMSException {
+ if (createAttempts.incrementAndGet() == 2) {
+ destinationResolver.scheduleRefresh();
+ }
+ return session.createTemporaryQueue();
+ }
+
+ @Override
+ public void delete(TemporaryQueue queue) {
+ // noop
+ }
+ }
+}