This is an automated email from the ASF dual-hosted git repository.

cshannon pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/activemq.git


The following commit(s) were added to refs/heads/main by this push:
     new c140d73fe AMQ-9157 - Include consumer id as part of Dispatched advisory
c140d73fe is described below

commit c140d73feca4a18044701ebf1df2859ec8c33207
Author: Christopher L. Shannon (cshannon) <[email protected]>
AuthorDate: Fri Nov 11 13:53:47 2022 -0500

    AMQ-9157 - Include consumer id as part of Dispatched advisory
---
 .../src/main/java/org/apache/activemq/advisory/AdvisoryBroker.java | 7 +++++--
 .../src/main/java/org/apache/activemq/broker/Broker.java           | 3 ++-
 .../src/main/java/org/apache/activemq/broker/BrokerFilter.java     | 4 ++--
 .../src/main/java/org/apache/activemq/broker/EmptyBroker.java      | 2 +-
 .../src/main/java/org/apache/activemq/broker/ErrorBroker.java      | 2 +-
 .../java/org/apache/activemq/broker/region/BaseDestination.java    | 4 ++--
 .../main/java/org/apache/activemq/broker/region/Destination.java   | 3 ++-
 .../java/org/apache/activemq/broker/region/DestinationFilter.java  | 4 ++--
 .../org/apache/activemq/broker/region/PrefetchSubscription.java    | 2 +-
 .../java/org/apache/activemq/broker/region/TopicSubscription.java  | 4 ++--
 .../java/org/apache/activemq/broker/util/LoggingBrokerPlugin.java  | 4 ++--
 .../src/test/java/org/apache/activemq/advisory/AdvisoryTests.java  | 5 +++++
 12 files changed, 27 insertions(+), 17 deletions(-)

diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/advisory/AdvisoryBroker.java
 
b/activemq-broker/src/main/java/org/apache/activemq/advisory/AdvisoryBroker.java
index e1ad10e33..14c856413 100644
--- 
a/activemq-broker/src/main/java/org/apache/activemq/advisory/AdvisoryBroker.java
+++ 
b/activemq-broker/src/main/java/org/apache/activemq/advisory/AdvisoryBroker.java
@@ -489,8 +489,8 @@ public class AdvisoryBroker extends BrokerFilter {
     }
 
     @Override
-    public void messageDispatched(ConnectionContext context, MessageReference 
messageReference) {
-        super.messageDispatched(context, messageReference);
+    public void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference) {
+        super.messageDispatched(context, sub, messageReference);
         try {
             if (!messageReference.isAdvisory()) {
                 BaseDestination baseDestination = (BaseDestination) 
messageReference.getMessage().getRegionDestination();
@@ -502,6 +502,9 @@ public class AdvisoryBroker extends BrokerFilter {
                 ActiveMQMessage advisoryMessage = new ActiveMQMessage();
                 
advisoryMessage.setStringProperty(AdvisorySupport.MSG_PROPERTY_MESSAGE_ID, 
payload.getMessageId().toString());
                 
advisoryMessage.setStringProperty(AdvisorySupport.MSG_PROPERTY_DESTINATION, 
baseDestination.getActiveMQDestination().getQualifiedName());
+                if (sub.getConsumerInfo() != null) {
+                    
advisoryMessage.setStringProperty(AdvisorySupport.MSG_PROPERTY_CONSUMER_ID, 
sub.getConsumerInfo().getConsumerId().toString());
+                }
                 fireAdvisory(context, topic, payload, null, advisoryMessage);
             }
         } catch (Exception e) {
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/Broker.java 
b/activemq-broker/src/main/java/org/apache/activemq/broker/Broker.java
index 0c884348e..36e773a6f 100644
--- a/activemq-broker/src/main/java/org/apache/activemq/broker/Broker.java
+++ b/activemq-broker/src/main/java/org/apache/activemq/broker/Broker.java
@@ -356,9 +356,10 @@ public interface Broker extends Region, Service {
     /**
      * Called when message is dispatched to a consumer
      * @param context
+     * @param sub
      * @param messageReference
      */
-    void messageDispatched(ConnectionContext context, MessageReference 
messageReference);
+    void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference);
 
     /**
      * Called when a message is discarded - e.g. running low on memory
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/BrokerFilter.java 
b/activemq-broker/src/main/java/org/apache/activemq/broker/BrokerFilter.java
index 9f27bf918..b9374e352 100644
--- a/activemq-broker/src/main/java/org/apache/activemq/broker/BrokerFilter.java
+++ b/activemq-broker/src/main/java/org/apache/activemq/broker/BrokerFilter.java
@@ -352,8 +352,8 @@ public class BrokerFilter implements Broker {
     }
 
     @Override
-    public void messageDispatched(ConnectionContext context, MessageReference 
messageReference) {
-        getNext().messageDispatched(context, messageReference);
+    public void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference) {
+        getNext().messageDispatched(context, sub, messageReference);
     }
 
     @Override
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/EmptyBroker.java 
b/activemq-broker/src/main/java/org/apache/activemq/broker/EmptyBroker.java
index 1e59f5bfc..4872a5a0f 100644
--- a/activemq-broker/src/main/java/org/apache/activemq/broker/EmptyBroker.java
+++ b/activemq-broker/src/main/java/org/apache/activemq/broker/EmptyBroker.java
@@ -309,7 +309,7 @@ public class EmptyBroker implements Broker {
     }
 
     @Override
-    public void messageDispatched(ConnectionContext context, MessageReference 
messageReference) {
+    public void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference) {
 
     }
 
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/ErrorBroker.java 
b/activemq-broker/src/main/java/org/apache/activemq/broker/ErrorBroker.java
index 48cf76bd2..8c138e394 100644
--- a/activemq-broker/src/main/java/org/apache/activemq/broker/ErrorBroker.java
+++ b/activemq-broker/src/main/java/org/apache/activemq/broker/ErrorBroker.java
@@ -349,7 +349,7 @@ public class ErrorBroker implements Broker {
     }
 
     @Override
-    public void messageDispatched(ConnectionContext context,MessageReference 
messageReference) {
+    public void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference) {
         throw new BrokerStoppedException(this.message);
     }
 
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/BaseDestination.java
 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/BaseDestination.java
index 3b2c3ce8c..06088da57 100644
--- 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/BaseDestination.java
+++ 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/BaseDestination.java
@@ -558,9 +558,9 @@ public abstract class BaseDestination implements 
Destination {
     }
 
     @Override
-    public void messageDispatched(ConnectionContext context, MessageReference 
messageReference) {
+    public void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference) {
         if (advisoryForDispatched) {
-            broker.messageDispatched(context, messageReference);
+            broker.messageDispatched(context, sub, messageReference);
         }
     }
 
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/Destination.java
 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/Destination.java
index 7e13be7ba..70f807be8 100644
--- 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/Destination.java
+++ 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/Destination.java
@@ -194,9 +194,10 @@ public interface Destination extends Service, Task, 
Message.MessageDestination {
      * Called when message is dispatched to a consumer
      *
      * @param context
+     * @param sub
      * @param messageReference
      */
-    void messageDispatched(ConnectionContext context, MessageReference 
messageReference);
+    void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference);
 
     /**
      * Called when a message is discarded - e.g. running low on memory This 
will
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/DestinationFilter.java
 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/DestinationFilter.java
index 0a9443640..6b288a234 100644
--- 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/DestinationFilter.java
+++ 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/DestinationFilter.java
@@ -325,8 +325,8 @@ public class DestinationFilter implements Destination {
     }
 
     @Override
-    public void messageDispatched(ConnectionContext context, MessageReference 
messageReference) {
-        next.messageDispatched(context, messageReference);
+    public void messageDispatched(ConnectionContext context, Subscription sub, 
MessageReference messageReference) {
+        next.messageDispatched(context, sub, messageReference);
     }
 
     @Override
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/PrefetchSubscription.java
 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/PrefetchSubscription.java
index 5c5e66777..93b3b2ae5 100644
--- 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/PrefetchSubscription.java
+++ 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/PrefetchSubscription.java
@@ -760,7 +760,7 @@ public abstract class PrefetchSubscription extends 
AbstractSubscription {
             if (node != QueueMessageReference.NULL_MESSAGE) {
                 
nodeDest.getDestinationStatistics().getDispatched().increment();
                 incrementPrefetchCounter(node);
-                nodeDest.messageDispatched(context, node);
+                nodeDest.messageDispatched(context, this, node);
                 LOG.trace("{} dispatched: {} - {}, dispatched: {}, inflight: 
{}",
                         info.getConsumerId(), message.getMessageId(), 
message.getDestination(),
                         
getSubscriptionStatistics().getDispatched().getCount(), dispatched.size());
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/TopicSubscription.java
 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/TopicSubscription.java
index 73b6017c4..faa29edef 100644
--- 
a/activemq-broker/src/main/java/org/apache/activemq/broker/region/TopicSubscription.java
+++ 
b/activemq-broker/src/main/java/org/apache/activemq/broker/region/TopicSubscription.java
@@ -702,7 +702,7 @@ public class TopicSubscription extends AbstractSubscription 
{
                         Destination regionDestination = (Destination) 
node.getRegionDestination();
                         
regionDestination.getDestinationStatistics().getDispatched().increment();
                         
regionDestination.getDestinationStatistics().getInflight().increment();
-                        regionDestination.messageDispatched(context, node);
+                        regionDestination.messageDispatched(context, 
TopicSubscription.this, node);
                         node.decrementReferenceCount();
                     }
 
@@ -724,7 +724,7 @@ public class TopicSubscription extends AbstractSubscription 
{
                 Destination regionDestination = (Destination) 
node.getRegionDestination();
                 
regionDestination.getDestinationStatistics().getDispatched().increment();
                 
regionDestination.getDestinationStatistics().getInflight().increment();
-                regionDestination.messageDispatched(context, node);
+                regionDestination.messageDispatched(context, this, node);
                 node.decrementReferenceCount();
             }
         }
diff --git 
a/activemq-broker/src/main/java/org/apache/activemq/broker/util/LoggingBrokerPlugin.java
 
b/activemq-broker/src/main/java/org/apache/activemq/broker/util/LoggingBrokerPlugin.java
index e40573bd9..37fdc8562 100644
--- 
a/activemq-broker/src/main/java/org/apache/activemq/broker/util/LoggingBrokerPlugin.java
+++ 
b/activemq-broker/src/main/java/org/apache/activemq/broker/util/LoggingBrokerPlugin.java
@@ -542,12 +542,12 @@ public class LoggingBrokerPlugin extends 
BrokerPluginSupport {
     }
 
     @Override
-    public void messageDispatched(ConnectionContext context, MessageReference 
messageReference) {
+    public void messageDispatched(ConnectionContext context,  Subscription 
sub, MessageReference messageReference) {
         if (isLogAll() || isLogConsumerEvents() || isLogInternalEvents()) {
             String msg = messageReference.getMessage().toString();
             LOG.info("Message dispatched: {}", msg);
         }
-        super.messageDispatched(context, messageReference);
+        super.messageDispatched(context, sub, messageReference);
     }
 
     @Override
diff --git 
a/activemq-unit-tests/src/test/java/org/apache/activemq/advisory/AdvisoryTests.java
 
b/activemq-unit-tests/src/test/java/org/apache/activemq/advisory/AdvisoryTests.java
index fd42cf3bf..c5faafe3e 100644
--- 
a/activemq-unit-tests/src/test/java/org/apache/activemq/advisory/AdvisoryTests.java
+++ 
b/activemq-unit-tests/src/test/java/org/apache/activemq/advisory/AdvisoryTests.java
@@ -271,6 +271,11 @@ public class AdvisoryTests {
         
assertTrue(((String)message.getProperty(AdvisorySupport.MSG_PROPERTY_ORIGIN_BROKER_URL)).startsWith("tcp://"));
         
assertEquals(message.getProperty(AdvisorySupport.MSG_PROPERTY_DESTINATION), 
dest.getQualifiedName());
 
+        //Make sure consumer id exists if dispatched advisory
+        if (AdvisorySupport.isMessageDispatchedAdvisoryTopic(advisoryTopic)) {
+            
assertNotNull(message.getStringProperty(AdvisorySupport.MSG_PROPERTY_CONSUMER_ID));
+        }
+
         //Add assertion to make sure body is included for advisory topics
         //when includeBodyForAdvisory is true
         assertIncludeBodyForAdvisory(payload);

Reply via email to