Github user mtaylor commented on a diff in the pull request:

    https://github.com/apache/activemq-artemis/pull/541#discussion_r64621914
  
    --- Diff: 
tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt/imported/MQTTTest.java
 ---
    @@ -1001,91 +1005,93 @@ public void testCleanSession() throws Exception {
        outstanding task to add cross protocol support.  This task should 
rework these tests.  The tests are included here
        and commented out to ensure ActiveMQ and Artemis tests are in sync. */
     
    -   //   @Test(timeout = 60 * 1000)
    -   //   public void testSendMQTTReceiveJMS() throws Exception {
    -   //      doTestSendMQTTReceiveJMS("foo.*");
    -   //   }
    +   @Test(timeout = 60 * 1000)
    +   public void testSendMQTTReceiveJMS() throws Exception {
    +      doTestSendMQTTReceiveJMS("foo.*", "foo/bar");
    +   }
     
    -   //   public void doTestSendMQTTReceiveJMS(String destinationName) 
throws Exception {
    -   //      final MQTTClientProvider provider = getMQTTClientProvider();
    -   //      initializeConnection(provider);
    -   //
    -   //      // send retained message
    -   //      final String RETAINED = "RETAINED";
    -   //      provider.publish("foo/bah", RETAINED.getBytes(), AT_LEAST_ONCE, 
true);
    -   //
    -   //      ActiveMQConnection activeMQConnection = (ActiveMQConnection) 
cf.createConnection();
    -   //      // MUST set to true to receive retained messages
    -   //      activeMQConnection.setUseRetroactiveConsumer(true);
    -   //      activeMQConnection.start();
    -   //      Session s = activeMQConnection.createSession(false, 
Session.AUTO_ACKNOWLEDGE);
    -   //      javax.jms.Topic jmsTopic = s.createTopic(destinationName);
    -   //      MessageConsumer consumer = s.createConsumer(jmsTopic);
    -   //
    -   //      // check whether we received retained message on JMS subscribe
    -   //      ActiveMQMessage message = (ActiveMQMessage) 
consumer.receive(5000);
    -   //      assertNotNull("Should get retained message", message);
    -   //      ByteSequence bs = message.getContent();
    -   //      assertEquals(RETAINED, new String(bs.data, bs.offset, 
bs.length));
    -   //      
assertTrue(message.getBooleanProperty(RetainedMessageSubscriptionRecoveryPolicy.RETAINED_PROPERTY));
    -   //
    -   //      for (int i = 0; i < NUM_MESSAGES; i++) {
    -   //         String payload = "Test Message: " + i;
    -   //         provider.publish("foo/bah", payload.getBytes(), 
AT_LEAST_ONCE);
    -   //         message = (ActiveMQMessage) consumer.receive(5000);
    -   //         assertNotNull("Should get a message", message);
    -   //         bs = message.getContent();
    -   //         assertEquals(payload, new String(bs.data, bs.offset, 
bs.length));
    -   //      }
    -   //
    -   //      activeMQConnection.close();
    -   //      provider.disconnect();
    -   //   }
    +   public void doTestSendMQTTReceiveJMS(String jmsTopicAddress, String 
mqttAddress) throws Exception {
    +      final MQTTClientProvider provider = getMQTTClientProvider();
    +      initializeConnection(provider);
     
    -   // TODO As with other tests, this should be enabled as part of the 
cross protocol support with MQTT.
    -   //   @Test(timeout = 2 * 60 * 1000)
    -   //   public void testSendJMSReceiveMQTT() throws Exception {
    -   //      doTestSendJMSReceiveMQTT("foo.far");
    -   //   }
    +      // send retained message
    +      final String address = "jms/queue/" + mqttAddress;
    +      final String RETAINED = "RETAINED";
     
    -   // TODO As with other tests, this should be enabled as part of the 
cross protocol support with MQTT.
    -   //   public void doTestSendJMSReceiveMQTT(String destinationName) 
throws Exception {
    -   //      final MQTTClientProvider provider = getMQTTClientProvider();
    -   //      initializeConnection(provider);
    -   //
    -   //      ActiveMQConnection activeMQConnection = (ActiveMQConnection) 
cf.createConnection();
    -   //      activeMQConnection.setUseRetroactiveConsumer(true);
    -   //      activeMQConnection.start();
    -   //      Session s = activeMQConnection.createSession(false, 
Session.AUTO_ACKNOWLEDGE);
    -   //      javax.jms.Topic jmsTopic = s.createTopic(destinationName);
    -   //      MessageProducer producer = s.createProducer(jmsTopic);
    -   //
    -   //      // send retained message from JMS
    -   //      final String RETAINED = "RETAINED";
    -   //      TextMessage sendMessage = s.createTextMessage(RETAINED);
    -   //      // mark the message to be retained
    -   //      
sendMessage.setBooleanProperty(RetainedMessageSubscriptionRecoveryPolicy.RETAIN_PROPERTY,
 true);
    -   //      // MQTT QoS can be set using 
MQTTProtocolConverter.QOS_PROPERTY_NAME property
    -   //      
sendMessage.setIntProperty(MQTTProtocolConverter.QOS_PROPERTY_NAME, 0);
    -   //      producer.send(sendMessage);
    -   //
    -   //      provider.subscribe("foo/+", AT_MOST_ONCE);
    -   //      byte[] message = provider.receive(10000);
    -   //      assertNotNull("Should get retained message", message);
    -   //      assertEquals(RETAINED, new String(message));
    -   //
    -   //      for (int i = 0; i < NUM_MESSAGES; i++) {
    -   //         String payload = "This is Test Message: " + i;
    -   //         sendMessage = s.createTextMessage(payload);
    -   //         producer.send(sendMessage);
    -   //         message = provider.receive(5000);
    -   //         assertNotNull("Should get a message", message);
    -   //
    -   //         assertEquals(payload, new String(message));
    -   //      }
    -   //      provider.disconnect();
    -   //      activeMQConnection.close();
    -   //   }
    +      final byte[] payload = RETAINED.getBytes();
    +
    +      Connection connection = cf.createConnection();
    +      // MUST set to true to receive retained messages
    +      connection.start();
    +
    +      Session s = connection.createSession(false, 
Session.AUTO_ACKNOWLEDGE);
    +      javax.jms.Queue jmsQueue = s.createQueue(jmsTopicAddress);
    +      MessageConsumer consumer = s.createConsumer(jmsQueue);
    +
    +      provider.publish(address, RETAINED.getBytes(), AT_LEAST_ONCE, true);
    +
    +      // check whether we received retained message on JMS subscribe
    +      BytesMessage message = (BytesMessage) consumer.receive(5000);
    +      assertNotNull("Should get retained message", message);
    +
    +      byte[] b = new byte[8];
    +      message.readBytes(b);
    +      assertArrayEquals(payload, b);
    +
    +      for (int i = 0; i < NUM_MESSAGES; i++) {
    +         String p = "Test Message: " + i;
    +         provider.publish(address, p.getBytes(), AT_LEAST_ONCE);
    +         message = (BytesMessage) consumer.receive(5000);
    +         assertNotNull("Should get a message", message);
    +
    +         byte[] bytePayload = new byte[p.getBytes().length];
    +         message.readBytes(bytePayload);
    +         assertArrayEquals(payload, b);
    +      }
    +
    +      connection.close();
    +      provider.disconnect();
    +   }
    +
    +//   @Test(timeout = 2 * 60 * 1000)
    +//   public void testSendJMSReceiveMQTT() throws Exception {
    --- End diff --
    
    I uncomment it in the next commit, this was a commit I did a while back and 
left these comments in to remind me.  I can rebase if you want?


---
If your project is set up for it, you can reply to this email and have your
reply appear on GitHub as well. If your project does not have this feature
enabled and wishes so, or if the feature is enabled but not working, please
contact infrastructure at [email protected] or file a JIRA ticket
with INFRA.
---

Reply via email to