summaryrefslogtreecommitdiff
path: root/qpid/java/broker-plugins/amqp-0-8-protocol
diff options
context:
space:
mode:
authorKeith Wall <kwall@apache.org>2015-02-10 16:15:08 +0000
committerKeith Wall <kwall@apache.org>2015-02-10 16:15:08 +0000
commit085486ebe5ff21133b9caf1c31625ac6ea356568 (patch)
tree7acbe9ca99a345dca71f9f80cd3e29ea4e3710f0 /qpid/java/broker-plugins/amqp-0-8-protocol
parent60c62c03ca404e98e4fbd1abf4a5ebf50763d604 (diff)
parente2e6d542b8cde9e702d1c3b63376e9d8380ba1c7 (diff)
downloadqpid-python-085486ebe5ff21133b9caf1c31625ac6ea356568.tar.gz
merge from trunk
git-svn-id: https://svn.apache.org/repos/asf/qpid/branches/QPID-6262-JavaBrokerNIO@1658748 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker-plugins/amqp-0-8-protocol')
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java11
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQProtocolEngine.java58
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java6
3 files changed, 48 insertions, 27 deletions
diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
index f7f65e29c2..a149214455 100644
--- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
+++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
@@ -3256,17 +3256,6 @@ public class AMQChannel
+ autoDelete
+ ")");
}
- else if (queue.isDurable() != durable)
- {
- closeChannel(AMQConstant.ALREADY_EXISTS,
- "Cannot re-declare queue '"
- + queue.getName()
- + "' with different durability (was: "
- + queue.isDurable()
- + " requested "
- + durable
- + ")");
- }
else
{
setDefaultQueue(queue);
diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQProtocolEngine.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQProtocolEngine.java
index 1aa4ef0b3f..233f68aeb6 100644
--- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQProtocolEngine.java
+++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQProtocolEngine.java
@@ -96,6 +96,16 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
AMQConnectionModel<AMQProtocolEngine, AMQChannel>,
ServerMethodProcessor<ServerChannelMethodProcessor>
{
+ enum ConnectionState
+ {
+ INIT,
+ AWAIT_START_OK,
+ AWAIT_SECURE_OK,
+ AWAIT_TUNE_OK,
+ AWAIT_OPEN,
+ OPEN
+ }
+
private static final Logger _logger = Logger.getLogger(AMQProtocolEngine.class);
// to save boxing the channelId and looking up in a map... cache in an array the low numbered
@@ -123,6 +133,8 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
private final AMQChannel[] _cachedChannels = new AMQChannel[CHANNEL_CACHE_SIZE + 1];
+ private ConnectionState _state = ConnectionState.INIT;
+
/**
* The channels that the latest call to {@link #received(ByteBuffer)} applied to.
* Used so we know which channels we need to call {@link AMQChannel#receivedComplete()}
@@ -486,14 +498,9 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
serverProperties,
mechanisms.getBytes(),
locales.getBytes());
- try
- {
- responseBody.generateFrame(0).writePayload(_sender);
- }
- catch (IOException e)
- {
- throw new ServerScopedRuntimeException(e);
- }
+ writeFrame(responseBody.generateFrame(0));
+ _state = ConnectionState.AWAIT_START_OK;
+
_sender.flush();
}
@@ -501,14 +508,7 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
{
_logger.info("Received unsupported protocol initiation for protocol version: " + getProtocolVersion());
- try
- {
- new ProtocolInitiation(ProtocolVersion.getLatestSupportedVersion()).writePayload(_sender);
- }
- catch (IOException ioex)
- {
- throw new ServerScopedRuntimeException(ioex);
- }
+ writeFrame(new ProtocolInitiation(ProtocolVersion.getLatestSupportedVersion()));
_sender.flush();
}
}
@@ -1498,6 +1498,7 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
{
_logger.debug("RECV[" + channelId + "] ChannelOpen");
}
+ assertState(ConnectionState.OPEN);
// Protect the broker against out of order frame request.
if (_virtualHost == null)
@@ -1534,6 +1535,15 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
}
}
+ void assertState(final ConnectionState requiredState)
+ {
+ if(_state != requiredState)
+ {
+ closeConnection(AMQConstant.COMMAND_INVALID, "Command Invalid", 0);
+
+ }
+ }
+
@Override
public void receiveConnectionOpen(AMQShortString virtualHostName,
AMQShortString capabilities,
@@ -1586,6 +1596,7 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
AMQMethodBody responseBody = methodRegistry.createConnectionOpenOkBody(virtualHostName);
writeFrame(responseBody.generateFrame(0));
+ _state = ConnectionState.OPEN;
}
catch (AccessControlException e)
{
@@ -1656,6 +1667,8 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
_logger.debug("RECV ConnectionSecureOk[ response: ******** ] ");
}
+ assertState(ConnectionState.AWAIT_SECURE_OK);
+
Broker<?> broker = getBroker();
SubjectCreator subjectCreator = getSubjectCreator();
@@ -1696,6 +1709,7 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
frameMax,
broker.getConnection_heartBeatDelay());
writeFrame(tuneBody.generateFrame(0));
+ _state = ConnectionState.AWAIT_TUNE_OK;
setAuthorizedSubject(authResult.getSubject());
disposeSaslServer();
break;
@@ -1744,6 +1758,8 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
+ " ]");
}
+ assertState(ConnectionState.AWAIT_START_OK);
+
Broker<?> broker = getBroker();
_logger.info("SASL Mechanism selected: " + mechanism);
@@ -1805,11 +1821,14 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
frameMax,
broker.getConnection_heartBeatDelay());
writeFrame(tuneBody.generateFrame(0));
+ _state = ConnectionState.AWAIT_TUNE_OK;
break;
case CONTINUE:
ConnectionSecureBody
secureBody = methodRegistry.createConnectionSecureBody(authResult.getChallenge());
writeFrame(secureBody.generateFrame(0));
+
+ _state = ConnectionState.AWAIT_SECURE_OK;
}
}
}
@@ -1828,6 +1847,8 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
_logger.debug("RECV ConnectionTuneOk[" +" channelMax: " + channelMax + " frameMax: " + frameMax + " heartbeat: " + heartbeat + " ]");
}
+ assertState(ConnectionState.AWAIT_TUNE_OK);
+
initHeartbeats(heartbeat);
int brokerFrameMax = getBroker().getContextValue(Integer.class, Broker.BROKER_FRAME_SIZE);
@@ -1859,7 +1880,10 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
setMaximumNumberOfChannels( ((channelMax == 0l) || (channelMax > 0xFFFFL))
? 0xFFFFL
: channelMax);
+
}
+ _state = ConnectionState.AWAIT_OPEN;
+
}
public int getBinaryDataLimit()
@@ -1959,6 +1983,8 @@ public class AMQProtocolEngine implements ServerProtocolEngine,
@Override
public ServerChannelMethodProcessor getChannelMethodProcessor(final int channelId)
{
+ assertState(ConnectionState.OPEN);
+
ServerChannelMethodProcessor channelMethodProcessor = getChannel(channelId);
if(channelMethodProcessor == null)
{
diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java
index 0f198a8d46..f8098eb2ec 100644
--- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java
+++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java
@@ -276,6 +276,12 @@ public class InternalTestProtocolSession extends AMQProtocolEngine implements Pr
}
}
+ void assertState(final ConnectionState requiredState)
+ {
+ // no-op
+ }
+
+
private static final AtomicInteger portNumber = new AtomicInteger(0);
private static class TestNetworkConnection implements NetworkConnection