From 9bd52fa485d73b3eb5c68d698e63243052a1db9c Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Fri, 17 Oct 2014 15:53:42 +0000 Subject: QPID-6163 : [Java Broker] Disconnect clients which do not obey flow control git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1632618 13f79535-47bb-0310-9956-ffa450edef68 --- .../server/logging/messages/ChannelMessages.java | 38 +++++++++++++++++++--- .../messages/Channel_logmessages.properties | 2 ++ .../java/org/apache/qpid/server/model/Broker.java | 17 ++++++---- .../qpid/server/protocol/AMQSessionModel.java | 2 +- .../apache/qpid/server/util/BrokerTestHelper.java | 1 + 5 files changed, 49 insertions(+), 11 deletions(-) (limited to 'qpid/java/broker-core/src') diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/ChannelMessages.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/ChannelMessages.java index 0cd0828623..6ae1ac4f02 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/ChannelMessages.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/ChannelMessages.java @@ -22,14 +22,15 @@ package org.apache.qpid.server.logging.messages; import static org.apache.qpid.server.logging.AbstractMessageLogger.DEFAULT_LOG_HIERARCHY_PREFIX; -import org.apache.log4j.Logger; -import org.apache.qpid.server.configuration.BrokerProperties; -import org.apache.qpid.server.logging.LogMessage; - import java.text.MessageFormat; import java.util.Locale; import java.util.ResourceBundle; +import org.apache.log4j.Logger; + +import org.apache.qpid.server.configuration.BrokerProperties; +import org.apache.qpid.server.logging.LogMessage; + /** * DO NOT EDIT DIRECTLY, THIS FILE WAS GENERATED. * @@ -53,6 +54,7 @@ public class ChannelMessages public static final String DEADLETTERMSG_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "channel.deadlettermsg"; public static final String DISCARDMSG_NOALTEXCH_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "channel.discardmsg_noaltexch"; public static final String IDLE_TXN_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "channel.idle_txn"; + public static final String FLOW_CONTROL_IGNORED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "channel.flow_control_ignored"; public static final String DISCARDMSG_NOROUTE_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "channel.discardmsg_noroute"; public static final String OPEN_TXN_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "channel.open_txn"; public static final String FLOW_REMOVED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "channel.flow_removed"; @@ -69,6 +71,7 @@ public class ChannelMessages Logger.getLogger(DEADLETTERMSG_LOG_HIERARCHY); Logger.getLogger(DISCARDMSG_NOALTEXCH_LOG_HIERARCHY); Logger.getLogger(IDLE_TXN_LOG_HIERARCHY); + Logger.getLogger(FLOW_CONTROL_IGNORED_LOG_HIERARCHY); Logger.getLogger(DISCARDMSG_NOROUTE_LOG_HIERARCHY); Logger.getLogger(OPEN_TXN_LOG_HIERARCHY); Logger.getLogger(FLOW_REMOVED_LOG_HIERARCHY); @@ -354,6 +357,33 @@ public class ChannelMessages }; } + /** + * Log a Channel message of the Format: + *
CHN-1012 : Flow Control Ignored. Channel will be closed.
+ * Optional values are contained in [square brackets] and are numbered + * sequentially in the method call. + * + */ + public static LogMessage FLOW_CONTROL_IGNORED() + { + String rawMessage = _messages.getString("FLOW_CONTROL_IGNORED"); + + final String message = rawMessage; + + return new LogMessage() + { + public String toString() + { + return message; + } + + public String getLogHierarchy() + { + return FLOW_CONTROL_IGNORED_LOG_HIERARCHY; + } + }; + } + /** * Log a Channel message of the Format: *
CHN-1010 : Discarded message : {0,number} as no binding on alternate exchange : {1}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/Channel_logmessages.properties b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/Channel_logmessages.properties index 397c12d73c..5c6e066541 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/Channel_logmessages.properties +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/Channel_logmessages.properties @@ -38,3 +38,5 @@ IDLE_TXN = CHN-1008 : Idle Transaction : {0,number} ms DISCARDMSG_NOALTEXCH = CHN-1009 : Discarded message : {0,number} as no alternate exchange configured for queue : {1} routing key : {2} DISCARDMSG_NOROUTE = CHN-1010 : Discarded message : {0,number} as no binding on alternate exchange : {1} DEADLETTERMSG = CHN-1011 : Message : {0,number} moved to dead letter queue : {1} + +FLOW_CONTROL_IGNORED = CHN-1012 : Flow Control Ignored. Channel will be closed. diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Broker.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Broker.java index 4c5293ff65..4e4acb3e21 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Broker.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Broker.java @@ -51,6 +51,8 @@ public interface Broker> extends ConfiguredObject, EventL String MODEL_VERSION = "modelVersion"; String CONFIDENTIAL_CONFIGURATION_ENCRYPTION_PROVIDER = "confidentialConfigurationEncryptionProvider"; + String CHANNEL_FLOW_CONTROL_ENFORCEMENT_TIMEOUT = "channel.flowControlEnforcementTimeout"; + String CONNECTION_SESSION_COUNT_LIMIT = "connection.sessionCountLimit"; String CONNECTION_HEART_BEAT_DELAY = "connection.heartBeatDelay"; String CONNECTION_CLOSE_WHEN_NO_ROUTE = "connection.closeWhenNoRoute"; @@ -63,19 +65,22 @@ public interface Broker> extends ConfiguredObject, EventL String QPID_JMX_PORT = "qpid.jmx_port"; @ManagedContextDefault(name = "broker.name") - static final String DEFAULT_BROKER_NAME = "Broker"; + String DEFAULT_BROKER_NAME = "Broker"; @ManagedContextDefault(name = QPID_AMQP_PORT) - public static final String DEFAULT_AMQP_PORT_NUMBER = "5672"; + String DEFAULT_AMQP_PORT_NUMBER = "5672"; @ManagedContextDefault(name = QPID_HTTP_PORT) - public static final String DEFAULT_HTTP_PORT_NUMBER = "8080"; + String DEFAULT_HTTP_PORT_NUMBER = "8080"; @ManagedContextDefault(name = QPID_RMI_PORT) - public static final String DEFAULT_RMI_PORT_NUMBER = "8999"; + String DEFAULT_RMI_PORT_NUMBER = "8999"; @ManagedContextDefault(name = QPID_JMX_PORT) - public static final String DEFAULT_JMX_PORT_NUMBER = "9099"; + String DEFAULT_JMX_PORT_NUMBER = "9099"; @ManagedContextDefault(name = BROKER_FLOW_TO_DISK_THRESHOLD) - public static final long DEFAULT_FLOW_TO_DISK_THRESHOLD = (long)(0.4 * (double)Runtime.getRuntime().maxMemory()); + long DEFAULT_FLOW_TO_DISK_THRESHOLD = (long)(0.4 * (double)Runtime.getRuntime().maxMemory()); + + @ManagedContextDefault(name = CHANNEL_FLOW_CONTROL_ENFORCEMENT_TIMEOUT) + long DEFAULT_CHANNEL_FLOW_CONTROL_ENFORCEMENT_TIMEOUT = 5000l; String BROKER_FRAME_SIZE = "qpid.broker_frame_size"; @ManagedContextDefault(name = BROKER_FRAME_SIZE) diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/protocol/AMQSessionModel.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/protocol/AMQSessionModel.java index a9cd32f8f9..f13af479ad 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/protocol/AMQSessionModel.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/protocol/AMQSessionModel.java @@ -36,7 +36,7 @@ import org.apache.qpid.server.util.Deletable; * Extends {@link Comparable} to allow objects to be inserted into a {@link ConcurrentSkipListSet} * when monitoring the blocking and blocking of queues/sessions in {@link AMQQueue}. */ -public interface AMQSessionModel, C extends AMQConnectionModel> extends Comparable, Deletable +public interface AMQSessionModel, C extends AMQConnectionModel> extends Comparable, Deletable { public UUID getId(); diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java index 0bee92a2e9..2f44218cf1 100644 --- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java +++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java @@ -97,6 +97,7 @@ public class BrokerTestHelper when(broker.getEventLogger()).thenReturn(eventLogger); when(broker.getCategoryClass()).thenReturn(Broker.class); when(broker.getParent(SystemConfig.class)).thenReturn(systemConfig); + when(broker.getContextValue(eq(Long.class), eq(Broker.CHANNEL_FLOW_CONTROL_ENFORCEMENT_TIMEOUT))).thenReturn(0l); when(broker.getTaskExecutor()).thenReturn(TASK_EXECUTOR); when(systemConfig.getTaskExecutor()).thenReturn(TASK_EXECUTOR); -- cgit v1.2.1