diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2014-10-17 15:53:42 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2014-10-17 15:53:42 +0000 |
| commit | 9bd52fa485d73b3eb5c68d698e63243052a1db9c (patch) | |
| tree | fea4a994556644221dbcf41e74f0c43a79cfd752 /qpid/java/broker-core/src | |
| parent | 95fc93485ab66966713611a4e1429d917dabde64 (diff) | |
| download | qpid-python-9bd52fa485d73b3eb5c68d698e63243052a1db9c.tar.gz | |
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
Diffstat (limited to 'qpid/java/broker-core/src')
5 files changed, 49 insertions, 11 deletions
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); @@ -356,6 +359,33 @@ public class ChannelMessages /** * Log a Channel message of the Format: + * <pre>CHN-1012 : Flow Control Ignored. Channel will be closed.</pre> + * 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: * <pre>CHN-1010 : Discarded message : {0,number} as no binding on alternate exchange : {1}</pre> * Optional values are contained in [square brackets] and are numbered * sequentially in the method call. 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<X extends Broker<X>> extends ConfiguredObject<X>, 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<X extends Broker<X>> extends ConfiguredObject<X>, 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<T extends AMQSessionModel<T,C>, C extends AMQConnectionModel<C,T>> extends Comparable<T>, Deletable<T> +public interface AMQSessionModel<T extends AMQSessionModel<T,C>, C extends AMQConnectionModel<C,T>> extends Comparable<AMQSessionModel>, Deletable<T> { 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); |
