[ 
https://issues.apache.org/jira/browse/ARTEMIS-4523?focusedWorklogId=893703&page=com.atlassian.jira.plugin.system.issuetabpanels:worklog-tabpanel#worklog-893703
 ]

ASF GitHub Bot logged work on ARTEMIS-4523:
-------------------------------------------

                Author: ASF GitHub Bot
            Created on: 04/Dec/23 10:22
            Start Date: 04/Dec/23 10:22
    Worklog Time Spent: 10m 
      Work Description: gtully commented on code in PR #4700:
URL: https://github.com/apache/activemq-artemis/pull/4700#discussion_r1413655996


##########
tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/openwire/FastReconnectOpenWireTest.java:
##########
@@ -0,0 +1,281 @@
+/*
+ * 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.activemq.artemis.tests.integration.openwire;
+
+import javax.jms.Connection;
+import javax.jms.MessageConsumer;
+import javax.jms.MessageProducer;
+import javax.jms.Queue;
+import javax.jms.Session;
+import javax.jms.TextMessage;
+import java.lang.invoke.MethodHandles;
+import java.util.ArrayList;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.activemq.ActiveMQConnection;
+import org.apache.activemq.ActiveMQConnectionFactory;
+import org.apache.activemq.RedeliveryPolicy;
+import org.apache.activemq.artemis.api.core.QueueConfiguration;
+import org.apache.activemq.artemis.api.core.RoutingType;
+import org.apache.activemq.artemis.api.core.SimpleString;
+import org.apache.activemq.artemis.core.settings.impl.AddressSettings;
+import org.apache.activemq.artemis.tests.util.Wait;
+import org.apache.activemq.command.ActiveMQQueue;
+import org.apache.activemq.transport.tcp.TcpTransport;
+import org.junit.Test;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class FastReconnectOpenWireTest extends OpenWireTestBase {
+
+   // change this number to give the test a bit more of spinning
+   private static final int NUM_ITERATIONS = 100;
+
+   private static final Logger logger = 
LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
+
+   @Override
+   protected void configureAddressSettings(Map<String, AddressSettings> 
addressSettingsMap) {
+      super.configureAddressSettings(addressSettingsMap);
+      // force send to dlq early
+      addressSettingsMap.put("exampleQueue", new 
AddressSettings().setAutoCreateQueues(false).setAutoCreateAddresses(false).setDeadLetterAddress(new
 
SimpleString("ActiveMQ.DLQ")).setAutoCreateAddresses(true).setMaxDeliveryAttempts(2));
+      // force send to dlq late
+      addressSettingsMap.put("exampleQueueTwo", new 
AddressSettings().setAutoCreateQueues(false).setAutoCreateAddresses(false).setDeadLetterAddress(new
 
SimpleString("ActiveMQ.DLQ")).setAutoCreateAddresses(true).setMaxDeliveryAttempts(-1));
+   }
+
+
+   @Test(timeout = 80_000)
+   public void testFastReconnectCreateConsumerNoErrors() throws Exception {
+
+      final ArrayList<Throwable> errors = new ArrayList<>();
+      SimpleString durableQueue = new SimpleString("exampleQueueTwo");
+      this.server.createQueue(new 
QueueConfiguration(durableQueue).setRoutingType(RoutingType.ANYCAST));
+
+      Queue queue = new ActiveMQQueue(durableQueue.toString());
+
+      final ActiveMQConnectionFactory exFact = new 
ActiveMQConnectionFactory("failover:(tcp://localhost:61616?closeAsync=false)?startupMaxReconnectAttempts=-1&maxReconnectAttempts=-1&timeout=5000");
+      exFact.setWatchTopicAdvisories(false);
+      exFact.setConnectResponseTimeout(10000);
+      exFact.setClientID("myID");
+
+      RedeliveryPolicy redeliveryPolicy = new RedeliveryPolicy();
+      redeliveryPolicy.setRedeliveryDelay(0);
+      redeliveryPolicy.setMaximumRedeliveries(-1);
+      exFact.setRedeliveryPolicy(redeliveryPolicy);
+
+      publish(1000, durableQueue.toString());
+
+      final AtomicInteger numIterations = new AtomicInteger(NUM_ITERATIONS);
+      ExecutorService executor = Executors.newCachedThreadPool();
+      runAfter(executor::shutdownNow);
+
+      final int concurrent = 2;
+      for (int i = 0; i < concurrent; i++) {
+         executor.execute(() -> {
+            while (numIterations.decrementAndGet() > 0) {
+               try (
+                  Connection conn = exFact.createConnection();
+                  Session consumerConnectionSession = 
conn.createSession(Session.SESSION_TRANSACTED);
+                  MessageConsumer messageConsumer = 
consumerConnectionSession.createConsumer(queue)) {
+
+                  messageConsumer.receiveNoWait();
+
+                  if (numIterations.get() % 2 == 0) {
+                     TimeUnit.MILLISECONDS.sleep(30);
+                  } else {
+                     TimeUnit.MILLISECONDS.sleep(50);
+                  }
+
+                  try {
+                     // force a local socket close such that the broker sees 
an exception on the connection and fails the consumer via serverConsumer close
+                     ((ActiveMQConnection) 
conn).getTransport().narrow(TcpTransport.class).stop();
+                  } catch (Throwable expected) {
+                  }
+
+               } catch (javax.jms.InvalidClientIDException expected) {
+                  // deliberate clash across concurrent connections
+               } catch (Throwable unexpected) {
+                  unexpected.printStackTrace();
+                  errors.add(unexpected);
+                  numIterations.set(0);
+               }
+            }
+         });
+      }
+
+      executor.shutdown();
+      assertTrue(executor.awaitTermination(60, TimeUnit.SECONDS));
+
+      Wait.assertEquals(0, () -> 
server.locateQueue(durableQueue).getConsumerCount());
+      assertTrue(errors.isEmpty());
+
+   }
+
+   @Test(timeout = 60_000)
+   public void testFastReconnectCreateConsumerNoErrorsNoClientIdA() throws 
Exception {

Review Comment:
   this test is now a duplicate of the one that comes after, it can safely be 
deleted.





Issue Time Tracking
-------------------

    Worklog Id:     (was: 893703)
    Time Spent: 0.5h  (was: 20m)

> Failover leaving consumers and producers behind
> -----------------------------------------------
>
>                 Key: ARTEMIS-4523
>                 URL: https://issues.apache.org/jira/browse/ARTEMIS-4523
>             Project: ActiveMQ Artemis
>          Issue Type: Bug
>    Affects Versions: 2.31.2
>            Reporter: Clebert Suconic
>            Assignee: Clebert Suconic
>            Priority: Major
>             Fix For: 2.32.0
>
>          Time Spent: 0.5h
>  Remaining Estimate: 0h
>
> Failover in OpenWire could leave consumers, producers and internal sessions 
> hanging behind.
> If the failover happened in a race condition where the new connection 
> (reconnection) happened before the client failure.. you could end with 
> consumers without connections.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to