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

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

                Author: ASF GitHub Bot
            Created on: 06/Dec/24 16:10
            Start Date: 06/Dec/24 16:10
    Worklog Time Spent: 10m 
      Work Description: lavocatt commented on code in PR #5220:
URL: https://github.com/apache/activemq-artemis/pull/5220#discussion_r1873608014


##########
tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/amqp/connect/AMQPNoForwardMirroringTest.java:
##########
@@ -0,0 +1,278 @@
+/*
+ * 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.amqp.connect;
+
+import javax.jms.Connection;
+import javax.jms.ConnectionFactory;
+import javax.jms.MessageConsumer;
+import javax.jms.MessageProducer;
+import javax.jms.Session;
+import javax.jms.TextMessage;
+
+import 
org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPBrokerConnectConfiguration;
+import 
org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPMirrorBrokerConnectionElement;
+import org.apache.activemq.artemis.core.server.ActiveMQServer;
+import org.apache.activemq.artemis.core.server.Queue;
+import 
org.apache.activemq.artemis.tests.integration.amqp.AmqpClientTestSupport;
+import org.apache.activemq.artemis.tests.util.CFUtil;
+import org.apache.activemq.artemis.utils.Wait;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+public class AMQPNoForwardMirroringTest extends AmqpClientTestSupport {
+
+   protected static final int AMQP_PORT_2 = 5673;
+   protected static final int AMQP_PORT_3 = 5674;
+
+   ActiveMQServer server_2;
+   ActiveMQServer server_3;
+
+   @Override
+   protected ActiveMQServer createServer() throws Exception {
+      return createServer(AMQP_PORT, false);
+   }
+
+   @Override
+   protected String getConfiguredProtocols() {
+      return "AMQP,CORE,OPENWIRE";
+   }
+
+   @Test
+   public void testNoForward() throws Exception {
+      server_2 = createServer(AMQP_PORT_2, false);
+      server_3 = createServer(AMQP_PORT_3, false);
+
+      // name the servers, for convenience during debugging
+      server.getConfiguration().setName("1");
+      server_2.getConfiguration().setName("2");
+      server_3.getConfiguration().setName("3");
+
+      /**
+       *
+       * the mirroring topology:
+       *
+       * v---------|
+       * 1 ------> 2       The link between 1 and 2 is noForward=true
+       *           ^
+       *           v
+       *           3
+       */
+
+      {
+         AMQPBrokerConnectConfiguration amqpConnection = new 
AMQPBrokerConnectConfiguration(getTestMethodName() + "1to2", "tcp://localhost:" 
+ AMQP_PORT_2).setRetryInterval(100);
+         amqpConnection.addElement(new 
AMQPMirrorBrokerConnectionElement().setAddressFilter(getQueueName()).setNoForward(true));
+         server.getConfiguration().addAMQPConnection(amqpConnection);
+      }
+
+      {
+         AMQPBrokerConnectConfiguration amqpConnection = new 
AMQPBrokerConnectConfiguration(getTestMethodName() + "2to1", "tcp://localhost:" 
+ AMQP_PORT).setRetryInterval(100);
+         amqpConnection.addElement(new 
AMQPMirrorBrokerConnectionElement().setAddressFilter(getQueueName()));
+         server_2.getConfiguration().addAMQPConnection(amqpConnection);
+
+         amqpConnection = new 
AMQPBrokerConnectConfiguration(getTestMethodName() + "2to3", "tcp://localhost:" 
+ AMQP_PORT_3).setRetryInterval(100);
+         amqpConnection.addElement(new 
AMQPMirrorBrokerConnectionElement().setAddressFilter(getQueueName()));
+         server_2.getConfiguration().addAMQPConnection(amqpConnection);
+      }
+
+      server.start();
+      server_2.start();
+      server_3.start();
+
+      createAddressAndQueues(server);
+      Wait.assertTrue(() -> server.locateQueue(getQueueName()) != null);
+      Wait.assertTrue(() -> server_2.locateQueue(getQueueName()) != null);
+      // queue creation doesn't reach 3 b/c of the noForward link between 1 
and 2.
+      Wait.assertTrue(() -> server_3.locateQueue(getQueueName()) == null);
+
+      Queue q1 = server.locateQueue(getQueueName());
+      assertNotNull(q1);
+
+      Queue q2 = server_2.locateQueue(getQueueName());
+      assertNotNull(q2);
+
+      ConnectionFactory factory = 
CFUtil.createConnectionFactory(randomProtocol(), "tcp://localhost:" + 
AMQP_PORT);
+      ConnectionFactory factory2 = 
CFUtil.createConnectionFactory(randomProtocol(), "tcp://localhost:" + 
AMQP_PORT_2);
+      ConnectionFactory factory3 = 
CFUtil.createConnectionFactory(randomProtocol(), "tcp://localhost:" + 
AMQP_PORT_3);
+
+      // send from 1, 2 receives, 3 don't.
+      try (Connection conn = factory.createConnection()) {
+         Session session = conn.createSession();
+         MessageProducer producer = 
session.createProducer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            producer.send(session.createTextMessage("message " + i));
+         }
+      }
+      Wait.assertEquals(10L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(10L, q2::getMessageCount, 1000, 100);
+
+      // consume from 2, 1 and 2 counters go back to 0
+      try (Connection conn = factory2.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            TextMessage message = (TextMessage) consumer.receive(1000);
+            assertNotNull(message);
+            assertEquals("message " + i, message.getText());
+         }
+      }
+
+      Wait.assertEquals(0L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q2::getMessageCount, 1000, 100);
+
+      Thread.sleep(100); // some time to allow eventual loops
+
+      // queue creation was originated on server, with noForward in place,
+      // the messages never reached server_3, for the rest of the test suite,
+      // we need server_3 to have access to the queue
+      createAddressAndQueues(server_3);
+      Wait.assertTrue(() -> server_3.locateQueue(getQueueName()) != null);
+      Queue q3 = server_3.locateQueue(getQueueName());
+      assertNotNull(q3);
+
+      // produce on 2. 1, 2 and 3 receive messages.
+      try (Connection conn = factory2.createConnection()) {
+         Session session = conn.createSession();
+         MessageProducer producer = 
session.createProducer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            producer.send(session.createTextMessage("message " + i));
+         }
+      }
+
+      Wait.assertEquals(10L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(10L, q2::getMessageCount, 1000, 100);
+      Wait.assertEquals(10L, q3::getMessageCount, 1000, 100);
+
+      // consume on 1. 1, 2, and 3 counters are back to 0
+      try (Connection conn = factory.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            TextMessage message = (TextMessage) consumer.receive(1000);
+            assertNotNull(message);
+            assertEquals("message " + i, message.getText());
+         }
+         consumer.close();
+      }
+      Wait.assertEquals(0L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q2::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q3::getMessageCount, 1000, 100);
+
+      // produce on 2. 1, 2 and 3 receive messages.
+      try (Connection conn = factory2.createConnection()) {
+         Session session = conn.createSession();
+         MessageProducer producer = 
session.createProducer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            producer.send(session.createTextMessage("message " + i));
+         }
+      }
+
+      Wait.assertEquals(10L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(10L, q2::getMessageCount, 1000, 100);
+      Wait.assertEquals(10L, q3::getMessageCount, 1000, 100);
+
+      // consume on 3. 1, 2 counters are still at 10 and 3 is at 0.
+      try (Connection conn = factory3.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            TextMessage message = (TextMessage) consumer.receive(1000);
+            assertNotNull(message);
+            assertEquals("message " + i, message.getText());
+         }
+         consumer.close();
+      }
+      Wait.assertEquals(10L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(10L, q2::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q3::getMessageCount, 1000, 100);
+
+      // consume on 2. 1, 2 and 3 counters are back to 0
+      try (Connection conn = factory2.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            TextMessage message = (TextMessage) consumer.receive(1000);
+            assertNotNull(message);
+            assertEquals("message " + i, message.getText());
+         }
+         consumer.close();
+      }
+      Wait.assertEquals(0L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q2::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q3::getMessageCount, 1000, 100);
+
+      // produce on 3. only 3 has messages.
+      try (Connection conn = factory3.createConnection()) {
+         Session session = conn.createSession();
+         MessageProducer producer = 
session.createProducer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            producer.send(session.createTextMessage("message " + i));
+         }
+      }
+
+      Wait.assertEquals(0L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q2::getMessageCount, 1000, 100);
+      Wait.assertEquals(10L, q3::getMessageCount, 1000, 100);
+
+      // consume on 3. 1, 2, and 3 counters are back to 0
+      try (Connection conn = factory3.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         for (int i = 0; i < 10; i++) {
+            TextMessage message = (TextMessage) consumer.receive(1000);
+            assertNotNull(message);
+            assertEquals("message " + i, message.getText());
+         }
+         consumer.close();
+      }
+      Wait.assertEquals(0L, q1::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q2::getMessageCount, 1000, 100);
+      Wait.assertEquals(0L, q3::getMessageCount, 1000, 100);
+
+      try (Connection conn = factory.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         assertNull(consumer.receiveNoWait());
+         consumer.close();
+      }
+
+      try (Connection conn = factory2.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         assertNull(consumer.receiveNoWait());
+         consumer.close();
+      }
+
+      try (Connection conn = factory3.createConnection()) {
+         Session session = conn.createSession();
+         conn.start();
+         MessageConsumer consumer = 
session.createConsumer(session.createQueue(getQueueName()));
+         assertNull(consumer.receiveNoWait());
+         consumer.close();
+      }

Review Comment:
   I'm trying to do that, but that makes the test fails. I'll show you next 
week.





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

    Worklog Id:     (was: 947070)
    Time Spent: 7h 20m  (was: 7h 10m)

> AMQ Broker Mirroring: One to Many - avoid the infinite loop of the messages
> ---------------------------------------------------------------------------
>
>                 Key: ARTEMIS-5037
>                 URL: https://issues.apache.org/jira/browse/ARTEMIS-5037
>             Project: ActiveMQ Artemis
>          Issue Type: Bug
>            Reporter: Thomas Lavocat
>            Assignee: Thomas Lavocat
>            Priority: Major
>              Labels: pull-request-available
>          Time Spent: 7h 20m
>  Remaining Estimate: 0h
>
> AMQ Broker Mirroring: One to Many - avoid the infinite loop of the messages:
> if we have a->b-c->a..
> a message will circulate forever in the mirrors



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

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
For further information, visit: https://activemq.apache.org/contact


Reply via email to