summaryrefslogtreecommitdiff
path: root/qpid/java/broker-core/src
diff options
context:
space:
mode:
authorRobert Godfrey <rgodfrey@apache.org>2014-10-17 15:53:42 +0000
committerRobert Godfrey <rgodfrey@apache.org>2014-10-17 15:53:42 +0000
commit9bd52fa485d73b3eb5c68d698e63243052a1db9c (patch)
treefea4a994556644221dbcf41e74f0c43a79cfd752 /qpid/java/broker-core/src
parent95fc93485ab66966713611a4e1429d917dabde64 (diff)
downloadqpid-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')
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/ChannelMessages.java38
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/Channel_logmessages.properties2
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/Broker.java17
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/protocol/AMQSessionModel.java2
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/util/BrokerTestHelper.java1
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);