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);