From e0e8f4c5087c1c5dc787740d6bd862755bd8daf1 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Mon, 25 Aug 2014 15:15:31 +0000 Subject: Merging from trunk r1617822:1618206 in the Java tree git-svn-id: https://svn.apache.org/repos/asf/qpid/branches/0.30@1620339 13f79535-47bb-0310-9956-ffa450edef68 --- .../amqp_1_0/jms/impl/MessageConsumerImpl.java | 5 +- .../amqp_1_0/client/ConnectionErrorException.java | 2 +- .../org/apache/qpid/amqp_1_0/client/Receiver.java | 9 +- .../subjects/BDBHAVirtualHostNodeLogSubject.java | 32 ++ .../server/logging/subjects/GroupLogSubject.java | 32 ++ .../store/berkeleydb/AbstractBDBMessageStore.java | 2 +- .../berkeleydb/BDBHARemoteReplicationNodeImpl.java | 49 ++- .../berkeleydb/BDBHAVirtualHostNodeImpl.java | 152 ++++--- .../berkeleydb/BDBHARemoteReplicationNodeTest.java | 6 +- ...BDBHAVirtualHostNodeOperationalLoggingTest.java | 140 +----- .../store/berkeleydb/replication/GroupCreator.java | 6 +- .../berkeleydb/replication/JMXManagementTest.java | 14 - .../store/berkeleydb/replication/TwoNodeTest.java | 82 ++-- .../qpid/server/exchange/HeadersBinding.java | 6 + .../logging/messages/HighAvailabilityMessages.java | 137 +++--- .../HighAvailability_logmessages.properties | 55 ++- .../server/logging/subjects/LogSubjectFormat.java | 1 - .../subjects/VirtualHostNodeLogSubject.java | 33 -- .../server/model/AbstractConfiguredObject.java | 2 +- .../qpid/server/model/AbstractSystemConfig.java | 16 +- .../apache/qpid/server/queue/AbstractQueue.java | 2 +- .../server/store/AbstractJDBCMessageStore.java | 2 +- .../virtualhostnode/AbstractVirtualHostNode.java | 7 - .../qpid/server/protocol/v0_8/AMQChannel.java | 13 +- .../v0_8/UnacknowledgedMessageMapImpl.java | 2 +- .../v0_8/UnacknowledgedMessageMapTest.java | 84 ++++ .../protocol/v1_0/MessageConverter_from_1_0.java | 131 ++++-- .../v0_10_v1_0/MessageConverter_1_0_to_v0_10.java | 4 +- .../v0_8_v1_0/MessageConverter_1_0_to_v0_8.java | 2 +- .../plugin/servlet/DefinedFileServlet.java | 25 +- .../java/org/apache/qpid/client/AMQSession.java | 4 +- .../org/apache/qpid/client/AMQSession_0_8.java | 19 +- .../qpid/client/BasicMessageProducer_0_8.java | 27 +- .../client/UnsupportedAddressSyntaxException.java | 32 ++ .../java/org/apache/qpid/jms/BrokerDetails.java | 8 +- .../org/apache/qpid/filter/ConstantExpression.java | 24 +- .../management/common/sasl/PlainSaslClient.java | 8 +- .../common/sasl/UserPasswordCallbackHandler.java | 6 +- .../UsernameHashedPasswordCallbackHandler.java | 12 +- .../property/NumericGeneratedPropertySupport.java | 12 +- qpid/java/pom.xml | 1 - .../org/apache/qpid/test/utils/QpidTestCase.java | 16 +- .../qpid/test/utils/TestBrokerConfiguration.java | 8 +- .../org/apache/qpid/server/util/AveragedRun.java | 68 --- .../java/org/apache/qpid/server/util/RunStats.java | 57 --- .../java/org/apache/qpid/server/util/TimedRun.java | 52 --- .../qpid/test/utils/ConversationFactory.java | 484 --------------------- .../org/apache/qpid/util/ClasspathScanner.java | 239 ---------- 48 files changed, 713 insertions(+), 1417 deletions(-) create mode 100644 qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/BDBHAVirtualHostNodeLogSubject.java create mode 100644 qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/GroupLogSubject.java delete mode 100644 qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/VirtualHostNodeLogSubject.java create mode 100644 qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapTest.java create mode 100644 qpid/java/client/src/main/java/org/apache/qpid/client/UnsupportedAddressSyntaxException.java delete mode 100644 qpid/java/systests/src/test/java/org/apache/qpid/server/util/AveragedRun.java delete mode 100644 qpid/java/systests/src/test/java/org/apache/qpid/server/util/RunStats.java delete mode 100644 qpid/java/systests/src/test/java/org/apache/qpid/server/util/TimedRun.java delete mode 100644 qpid/java/systests/src/test/java/org/apache/qpid/test/utils/ConversationFactory.java delete mode 100644 qpid/java/systests/src/test/java/org/apache/qpid/util/ClasspathScanner.java (limited to 'qpid/java') diff --git a/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/MessageConsumerImpl.java b/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/MessageConsumerImpl.java index 1a72e129e7..508aaf7518 100644 --- a/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/MessageConsumerImpl.java +++ b/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/MessageConsumerImpl.java @@ -173,7 +173,10 @@ public class MessageConsumerImpl implements MessageConsumer, QueueReceiver, Topi } else { - throw new JMSException(e.getMessage(), error.getCondition().getValue().toString()); + JMSException jmsException = + new JMSException(e.getMessage(), error.getCondition().getValue().toString()); + jmsException.initCause(e); + throw jmsException; } } diff --git a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/ConnectionErrorException.java b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/ConnectionErrorException.java index 302060776a..82f29ea4b1 100644 --- a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/ConnectionErrorException.java +++ b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/ConnectionErrorException.java @@ -34,7 +34,7 @@ public class ConnectionErrorException extends ConnectionException public ConnectionErrorException(Error remoteError) { - super(remoteError.getDescription()); + super(remoteError.getDescription() == null ? remoteError.toString() : remoteError.getDescription()); _remoteError = remoteError; } diff --git a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Receiver.java b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Receiver.java index a2a15779d2..826a757850 100644 --- a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Receiver.java +++ b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Receiver.java @@ -194,9 +194,14 @@ public class Receiver implements DeliveryStateHandler } catch (InterruptedException e) { - throw new ConnectionErrorException(AmqpError.INTERNAL_ERROR,"Interrupted whil waiting for detach following failed attach"); + throw new ConnectionErrorException(AmqpError.INTERNAL_ERROR,"Interrupted while waiting for detach following failed attach"); } - throw new ConnectionErrorException(getError()); + throw new ConnectionErrorException(getError().getCondition(), + getError().getDescription() == null + ? "AMQP error: '" + getError().getCondition().toString() + + "' when attempting to create a receiver" + + (source != null ? " from: '" + source.getAddress() +"'" : "") + : getError().getDescription()); } else { diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/BDBHAVirtualHostNodeLogSubject.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/BDBHAVirtualHostNodeLogSubject.java new file mode 100644 index 0000000000..a209062993 --- /dev/null +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/BDBHAVirtualHostNodeLogSubject.java @@ -0,0 +1,32 @@ +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ +package org.apache.qpid.server.logging.subjects; + + +public class BDBHAVirtualHostNodeLogSubject extends AbstractLogSubject +{ + public static final String VIRTUAL_HOST_NODE_FORMAT = "grp(/{0})/vhn(/{1})"; + + public BDBHAVirtualHostNodeLogSubject(String groupName, String nodeName) + { + setLogStringWithFormat(VIRTUAL_HOST_NODE_FORMAT, groupName, nodeName); + } +} diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/GroupLogSubject.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/GroupLogSubject.java new file mode 100644 index 0000000000..51fd1fc2dc --- /dev/null +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/logging/subjects/GroupLogSubject.java @@ -0,0 +1,32 @@ +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ +package org.apache.qpid.server.logging.subjects; + + +public class GroupLogSubject extends AbstractLogSubject +{ + public static final String GROUP_FORMAT = "grp(/{0})"; + + public GroupLogSubject(String groupName) + { + setLogStringWithFormat(GROUP_FORMAT, groupName); + } +} diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/AbstractBDBMessageStore.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/AbstractBDBMessageStore.java index 835846a5ec..78cddc708e 100644 --- a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/AbstractBDBMessageStore.java +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/AbstractBDBMessageStore.java @@ -1333,7 +1333,7 @@ public abstract class AbstractBDBMessageStore implements MessageStore data = new byte[0]; } } - return ByteBuffer.wrap(data,offsetInMessage,size); + return ByteBuffer.wrap(data,offsetInMessage,Math.min(size,data.length-offsetInMessage)); } diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeImpl.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeImpl.java index 5263f5942f..06671998ec 100644 --- a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeImpl.java +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeImpl.java @@ -30,9 +30,14 @@ import java.util.Set; import java.util.concurrent.atomic.AtomicReference; import com.sleepycat.je.rep.MasterStateException; +import com.sleepycat.je.rep.ReplicatedEnvironment; import org.apache.log4j.Logger; import org.apache.qpid.server.configuration.IllegalConfigurationException; +import org.apache.qpid.server.logging.EventLogger; +import org.apache.qpid.server.logging.messages.HighAvailabilityMessages; +import org.apache.qpid.server.logging.subjects.BDBHAVirtualHostNodeLogSubject; +import org.apache.qpid.server.logging.subjects.GroupLogSubject; import org.apache.qpid.server.model.AbstractConfiguredObject; import org.apache.qpid.server.model.Broker; import org.apache.qpid.server.model.ConfiguredObject; @@ -40,6 +45,8 @@ import org.apache.qpid.server.model.IllegalStateTransitionException; import org.apache.qpid.server.model.ManagedAttributeField; import org.apache.qpid.server.model.State; import org.apache.qpid.server.model.StateTransition; +import org.apache.qpid.server.model.SystemConfig; +import org.apache.qpid.server.model.VirtualHostNode; import org.apache.qpid.server.security.access.Operation; import org.apache.qpid.server.store.berkeleydb.replication.ReplicatedEnvironmentFacade; @@ -53,13 +60,16 @@ public class BDBHARemoteReplicationNodeImpl extends AbstractConfiguredObject _state; private final boolean _isMonitor; private boolean _detached; + private BDBHAVirtualHostNodeLogSubject _virtualHostNodeLogSubject; + private GroupLogSubject _groupLogSubject; public BDBHARemoteReplicationNodeImpl(BDBHAVirtualHostNode virtualHostNode, Map attributes, ReplicatedEnvironmentFacade replicatedEnvironmentFacade) { @@ -92,7 +102,7 @@ public class BDBHARemoteReplicationNodeImpl extends AbstractConfiguredObject { public static final String VIRTUAL_HOST_NODE_TYPE = "BDB_HA"; + public static final String VIRTUAL_HOST_PRINCIPAL_NAME_FORMAT = "grp(/{0})/vhn(/{1})"; /** * Length of time we synchronously await the a JE mutation to complete. It is not considered an error if we exceed this timeout, although a @@ -87,6 +91,9 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode _environmentFacade = new AtomicReference<>(); private final AtomicReference _lastReplicatedEnvironmentState = new AtomicReference<>(ReplicatedEnvironment.State.UNKNOWN); + private BDBHAVirtualHostNodeLogSubject _virtualHostNodeLogSubject; + private GroupLogSubject _groupLogSubject; + private String _virtualHostNodePrincipalName; @ManagedAttributeField private String _storePath; @@ -267,7 +274,7 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode helpers = getRemoteNodeAddresses(); super.doDelete(); - getEventLogger().message(getVirtualHostNodeLogSubject(), HighAvailabilityMessages.DELETED(getName(), getGroupName())); + getEventLogger().message(getVirtualHostNodeLogSubject(), HighAvailabilityMessages.DELETED()); if (getState() == State.DELETED && !helpers.isEmpty()) { try @@ -413,10 +401,18 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode() + getTaskExecutor().submit(new VirtualHostNodeGroupTask() { @Override - public Void execute() + public void perform() { addRemoteReplicationNode(node); - return null; } }); } @@ -722,19 +726,18 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode() + getTaskExecutor().submit(new VirtualHostNodeGroupTask() { @Override - public Void execute() + public void perform() { recoverRemoteReplicationNode(node); - return null; } }); } @@ -745,19 +748,18 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode() + getTaskExecutor().submit(new VirtualHostNodeGroupTask() { @Override - public Void execute() + public void perform() { removeRemoteReplicationNode(node); - return null; } }); } @@ -768,12 +770,25 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode() + { + @Override + public Void run() + { + processNodeState(node, nodeState); + return null; + } + }); + } + + private void processNodeState(ReplicationNode node, NodeState nodeState) { BDBHARemoteReplicationNodeImpl remoteNode = getChildByName(BDBHARemoteReplicationNodeImpl.class, node.getName()); if (remoteNode != null) @@ -785,7 +800,7 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode() + { + @Override + public Void run() + { + processIntruderNode(node); + return null; + } + }); + } + + private void processIntruderNode(ReplicationNode node) { String hostAndPort = node.getHostName() + ":" + node.getPort(); - getEventLogger().message(getVirtualHostNodeLogSubject(), HighAvailabilityMessages.INTRUDER_DETECTED(node.getName(), hostAndPort, getGroupName())); + getEventLogger().message(getGroupLogSubject(), HighAvailabilityMessages.INTRUDER_DETECTED(node.getName(), hostAndPort)); boolean inManagementMode = getParent(Broker.class).isManagementMode(); if (inManagementMode) @@ -858,7 +885,7 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode nodeToAttributes(ReplicationNode replicationNode) @@ -872,4 +899,23 @@ public class BDBHAVirtualHostNodeImpl extends AbstractVirtualHostNode + { + @Override + public Void execute() + { + return Subject.doAs(SecurityManager.getSystemTaskSubject(_virtualHostNodePrincipalName), new PrivilegedAction() + { + @Override + public Void run() + { + perform(); + return null; + } + }); + } + + abstract void perform(); + } + } diff --git a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeTest.java b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeTest.java index 6259b49d61..0d64d87aef 100644 --- a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeTest.java +++ b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHARemoteReplicationNodeTest.java @@ -19,6 +19,7 @@ package org.apache.qpid.server.virtualhostnode.berkeleydb; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; @@ -33,6 +34,7 @@ import org.apache.qpid.server.configuration.updater.TaskExecutor; import org.apache.qpid.server.model.Broker; import org.apache.qpid.server.model.ConfiguredObjectFactory; import org.apache.qpid.server.model.VirtualHost; +import org.apache.qpid.server.model.VirtualHostNode; import org.apache.qpid.server.security.SecurityManager; import org.apache.qpid.server.security.access.Operation; import org.apache.qpid.server.store.DurableConfigurationStore; @@ -69,7 +71,7 @@ public class BDBHARemoteReplicationNodeTest extends QpidTestCase // Virtualhost needs the EventLogger from the SystemContext. when(_virtualHostNode.getParent(Broker.class)).thenReturn(_broker); - + doReturn(VirtualHostNode.class).when(_virtualHostNode).getCategoryClass(); ConfiguredObjectFactory objectFactory = _broker.getObjectFactory(); when(_virtualHostNode.getModel()).thenReturn(objectFactory.getModel()); when(_virtualHostNode.getTaskExecutor()).thenReturn(_taskExecutor); @@ -80,7 +82,7 @@ public class BDBHARemoteReplicationNodeTest extends QpidTestCase String remoteReplicationName = getName(); BDBHARemoteReplicationNode remoteReplicationNode = createRemoteReplicationNode(remoteReplicationName); - remoteReplicationNode.setAttribute(BDBHARemoteReplicationNode.ROLE, null, "MASTER"); + remoteReplicationNode.setAttribute(BDBHARemoteReplicationNode.ROLE, "UNKNOWN", "MASTER"); verify(_facade).transferMasterAsynchronously(remoteReplicationName); } diff --git a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHAVirtualHostNodeOperationalLoggingTest.java b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHAVirtualHostNodeOperationalLoggingTest.java index ef1021160c..ea7d74090d 100644 --- a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHAVirtualHostNodeOperationalLoggingTest.java +++ b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/virtualhostnode/berkeleydb/BDBHAVirtualHostNodeOperationalLoggingTest.java @@ -78,67 +78,18 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase _helper.assertNodeRole(node1, "MASTER"); - String expectedMessage = HighAvailabilityMessages.ADDED(node1.getName(), node1.getGroupName()).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.ADDED_LOG_HIERARCHY))); - - expectedMessage = HighAvailabilityMessages.ATTACHED(node1.getName(), node1.getGroupName(), "UNKNOWN").toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.ATTACHED_LOG_HIERARCHY))); + assertEquals("Unexpected VHN log subject", "[grp(/group)/vhn(/node1)] ", node1.getVirtualHostNodeLogSubject().getLogString()); + assertEquals("Unexpected group log subject", "[grp(/group)] ", node1.getGroupLogSubject().getLogString()); - - expectedMessage = HighAvailabilityMessages.STARTED(node1.getName(), node1.getGroupName()).toString(); + String expectedMessage = HighAvailabilityMessages.CREATED().toString(); verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.STARTED_LOG_HIERARCHY))); + argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.CREATED_LOG_HIERARCHY))); - expectedMessage = HighAvailabilityMessages.ROLE_CHANGED(node1.getName(), node1.getGroupName(), "UNKNOWN", "MASTER").toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), + expectedMessage = HighAvailabilityMessages.ROLE_CHANGED(node1.getName(), node1.getAddress(), "UNKNOWN", "MASTER").toString(); + verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getGroupLogSubject())), argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.ROLE_CHANGED_LOG_HIERARCHY))); } - public void testStop() throws Exception - { - int node1PortNumber = findFreePort(); - String helperAddress = "localhost:" + node1PortNumber; - String groupName = "group"; - String nodeName = "node1"; - - Map node1Attributes = _helper.createNodeAttributes(nodeName, groupName, helperAddress, helperAddress, nodeName, node1PortNumber); - BDBHAVirtualHostNodeImpl node1 = (BDBHAVirtualHostNodeImpl)_helper.createHaVHN(node1Attributes); - _helper.assertNodeRole(node1, "MASTER"); - reset(_eventLogger); - - node1.stop(); - - String expectedMessage = HighAvailabilityMessages.DETACHED(node1.getName(), node1.getGroupName()).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.DETACHED_LOG_HIERARCHY))); - - expectedMessage = HighAvailabilityMessages.STOPPED(node1.getName(), node1.getGroupName()).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.STOPPED_LOG_HIERARCHY))); - } - - public void testClose() throws Exception - { - int node1PortNumber = findFreePort(); - String helperAddress = "localhost:" + node1PortNumber; - String groupName = "group"; - String nodeName = "node1"; - - Map node1Attributes = _helper.createNodeAttributes(nodeName, groupName, helperAddress, helperAddress, nodeName, node1PortNumber); - BDBHAVirtualHostNodeImpl node1 = (BDBHAVirtualHostNodeImpl)_helper.createHaVHN(node1Attributes); - _helper.assertNodeRole(node1, "MASTER"); - - reset(_eventLogger); - - node1.close(); - - String expectedMessage = HighAvailabilityMessages.DETACHED(node1.getName(), node1.getGroupName()).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.DETACHED_LOG_HIERARCHY))); - } - public void testDelete() throws Exception { int node1PortNumber = findFreePort(); @@ -154,13 +105,10 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase node1.delete(); - String expectedMessage = HighAvailabilityMessages.DETACHED(node1.getName(), node1.getGroupName()).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.DETACHED_LOG_HIERARCHY))); - - expectedMessage = HighAvailabilityMessages.DELETED(node1.getName(), node1.getGroupName()).toString(); + String expectedMessage = HighAvailabilityMessages.DELETED().toString(); verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.DELETED_LOG_HIERARCHY))); + } public void testSetPriority() throws Exception @@ -181,7 +129,7 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase // make sure that task executor thread finishes all scheduled tasks node1.stop(); - String expectedMessage = HighAvailabilityMessages.PRIORITY_CHANGED(node1.getName(), node1.getGroupName(), "10").toString(); + String expectedMessage = HighAvailabilityMessages.PRIORITY_CHANGED("10").toString(); verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.PRIORITY_CHANGED_LOG_HIERARCHY))); } @@ -204,7 +152,7 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase // make sure that task executor thread finishes all scheduled tasks node1.stop(); - String expectedMessage = HighAvailabilityMessages.QUORUM_OVERRIDE_CHANGED(node1.getName(), node1.getGroupName(), "1").toString(); + String expectedMessage = HighAvailabilityMessages.QUORUM_OVERRIDE_CHANGED("1").toString(); verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.QUORUM_OVERRIDE_CHANGED_LOG_HIERARCHY))); } @@ -227,7 +175,7 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase // make sure that task executor thread finishes all scheduled tasks node1.stop(); - String expectedMessage = HighAvailabilityMessages.DESIGNATED_PRIMARY_CHANGED(node1.getName(), node1.getGroupName(), "true").toString(); + String expectedMessage = HighAvailabilityMessages.DESIGNATED_PRIMARY_CHANGED("true").toString(); verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.DESIGNATED_PRIMARY_CHANGED_LOG_HIERARCHY))); } @@ -254,14 +202,9 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase // make sure that task executor thread finishes all scheduled tasks node2.stop(); - // Verify ADDED message from node2 when its created - String expectedMessage = HighAvailabilityMessages.ADDED(node2.getName(), groupName).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node2.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.ADDED_LOG_HIERARCHY))); - // Verify ADDED message from node1 when it discovers node2 has been added - expectedMessage = HighAvailabilityMessages.ADDED(node2.getName(), groupName).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), + String expectedMessage = HighAvailabilityMessages.ADDED(node2.getName(), node2.getAddress()).toString(); + verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getGroupLogSubject())), argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.ADDED_LOG_HIERARCHY))); } @@ -292,9 +235,9 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase // make sure that task executor thread finishes all scheduled tasks node1.stop(); - String expectedMessage = HighAvailabilityMessages.DELETED(node2.getName(), groupName).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.DELETED_LOG_HIERARCHY))); + String expectedMessage = HighAvailabilityMessages.REMOVED(node2.getName(), node2.getAddress()).toString(); + verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getGroupLogSubject())), + argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.REMOVED_LOG_HIERARCHY))); } public void testRemoteNodeDetached() throws Exception @@ -324,9 +267,9 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase waitForNodeDetachedField(remoteNode, true); // verify that remaining node issues the DETACHED operational logging for remote node - String expectedMessage = HighAvailabilityMessages.DETACHED(node2.getName(), groupName).toString(); - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.DETACHED_LOG_HIERARCHY))); + String expectedMessage = HighAvailabilityMessages.LEFT(node2.getName(), node2.getAddress()).toString(); + verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getGroupLogSubject())), + argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.LEFT_LOG_HIERARCHY))); } @@ -361,48 +304,9 @@ public class BDBHAVirtualHostNodeOperationalLoggingTest extends QpidTestCase _helper.assertNodeRole(node2, "REPLICA", "MASTER"); waitForNodeDetachedField(remoteNode, false); - final String expectedMessage = HighAvailabilityMessages.ATTACHED(node2.getName(), groupName, "REPLICA").toString(); - final String expectedMessage2 = HighAvailabilityMessages.ATTACHED(node2.getName(), groupName, "UNKNOWN").toString(); - final String expectedMessage3 = HighAvailabilityMessages.ATTACHED(node2.getName(), groupName, "MASTER").toString(); - ArgumentMatcher matcher = new ArgumentMatcher() - { - private String _messageErrorDescription = null; - private String _hierarchyErrorDescription = null; - - @Override - public boolean matches(Object argument) - { - LogMessage logMessage = (LogMessage)argument; - String actualMessage = logMessage.toString(); - boolean expectedMessageMatches = expectedMessage.equals(actualMessage) - || expectedMessage2.equals(actualMessage) || expectedMessage3.equals(actualMessage); - if (!expectedMessageMatches) - { - _messageErrorDescription = "Actual message does not match any expected: " + actualMessage; - } - boolean expectedHierarchyMatches = HighAvailabilityMessages.ATTACHED_LOG_HIERARCHY.equals(logMessage.getLogHierarchy()); - if (!expectedHierarchyMatches) - { - _hierarchyErrorDescription = "Actual hierarchy does not match expected: " + logMessage.getLogHierarchy(); - } - return expectedMessageMatches && expectedHierarchyMatches; - } - - @Override - public void describeTo(Description description) - { - if (_messageErrorDescription != null) - { - description.appendText(_messageErrorDescription); - } - if (_hierarchyErrorDescription != null) - { - description.appendText(_hierarchyErrorDescription); - } - } - }; - verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getVirtualHostNodeLogSubject())), - argThat(matcher)); + final String expectedMessage = HighAvailabilityMessages.JOINED(node2.getName(), node2.getAddress()).toString(); + verify(_eventLogger).message(argThat(new LogSubjectMatcher(node1.getGroupLogSubject())), + argThat(new LogMessageMatcher(expectedMessage, HighAvailabilityMessages.JOINED_LOG_HIERARCHY))); } private void waitForNodeDetachedField(BDBHARemoteReplicationNodeImpl remoteNode, boolean expectedDetached) throws InterruptedException { diff --git a/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/GroupCreator.java b/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/GroupCreator.java index e78ef34759..f7dce4f3f5 100644 --- a/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/GroupCreator.java +++ b/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/GroupCreator.java @@ -75,9 +75,9 @@ public class GroupCreator private static final String MANY_BROKER_URL_FORMAT = "amqp://guest:guest@/%s?brokerlist='%s'&failover='roundrobin?cyclecount='%d''"; private static final String BROKER_PORTION_FORMAT = "tcp://localhost:%d?connectdelay='%d',retries='%d'"; - private static final int FAILOVER_CYCLECOUNT = 10; - private static final int FAILOVER_RETRIES = 1; - private static final int FAILOVER_CONNECTDELAY = 1000; + private static final int FAILOVER_CYCLECOUNT = 20; + private static final int FAILOVER_RETRIES = 0; + private static final int FAILOVER_CONNECTDELAY = 500; private static final String SINGLE_BROKER_URL_WITH_RETRY_FORMAT = "amqp://guest:guest@/%s?brokerlist='tcp://localhost:%d?connectdelay='%d',retries='%d''"; private static final String SINGLE_BROKER_URL_WITHOUT_RETRY_FORMAT = "amqp://guest:guest@/%s?brokerlist='tcp://localhost:%d'"; diff --git a/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/JMXManagementTest.java b/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/JMXManagementTest.java index c6f005c0e7..63de287be7 100644 --- a/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/JMXManagementTest.java +++ b/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/JMXManagementTest.java @@ -206,20 +206,6 @@ public class JMXManagementTest extends QpidBrokerTestCase } } - public void testSetDesignatedPrimary() throws Exception - { - int brokerPort = _clusterCreator.getBrokerPortNumbersForNodes().iterator().next(); - final ManagedBDBHAMessageStore storeBean = getStoreBeanForNodeAtBrokerPort(brokerPort); - assertFalse("Unexpected designated primary before change", storeBean.getDesignatedPrimary()); - storeBean.setDesignatedPrimary(true); - long limit = System.currentTimeMillis() + 5000; - while(!storeBean.getDesignatedPrimary() && System.currentTimeMillis() < limit) - { - Thread.sleep(100l); - } - assertTrue("Unexpected designated primary after change", storeBean.getDesignatedPrimary()); - } - public void testVirtualHostMbeanOnMasterTransfer() throws Exception { Connection connection = getConnection(_brokerFailoverUrl); diff --git a/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/TwoNodeTest.java b/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/TwoNodeTest.java index 0f8a1609de..248dbb4def 100644 --- a/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/TwoNodeTest.java +++ b/qpid/java/bdbstore/systests/src/test/java/org/apache/qpid/server/store/berkeleydb/replication/TwoNodeTest.java @@ -20,27 +20,29 @@ package org.apache.qpid.server.store.berkeleydb.replication; import java.io.File; +import java.util.Collections; +import java.util.Map; import javax.jms.Connection; import javax.jms.JMSException; -import javax.management.ObjectName; import org.apache.qpid.jms.ConnectionURL; -import org.apache.qpid.server.store.berkeleydb.jmx.ManagedBDBHAMessageStore; -import org.apache.qpid.test.utils.JMXTestUtils; +import org.apache.qpid.server.virtualhostnode.berkeleydb.BDBHAVirtualHostNode; import org.apache.qpid.test.utils.QpidBrokerTestCase; public class TwoNodeTest extends QpidBrokerTestCase { private static final String VIRTUAL_HOST = "test"; - private static final String MANAGED_OBJECT_QUERY = "org.apache.qpid:type=BDBHAMessageStore,name=" + ObjectName.quote(VIRTUAL_HOST); private static final int NUMBER_OF_NODES = 2; private final GroupCreator _groupCreator = new GroupCreator(this, VIRTUAL_HOST, NUMBER_OF_NODES); - private final JMXTestUtils _jmxUtils = new JMXTestUtils(this); - private ConnectionURL _brokerFailoverUrl; + /** Used when expectation is client will not (re)-connect */ + private ConnectionURL _positiveFailoverUrl; + + /** Used when expectation is client will not (re)-connect */ + private ConnectionURL _negativeFailoverUrl; @Override protected void setUp() throws Exception @@ -53,19 +55,6 @@ public class TwoNodeTest extends QpidBrokerTestCase super.setUp(); } - @Override - protected void tearDown() throws Exception - { - try - { - _jmxUtils.close(); - } - finally - { - super.tearDown(); - } - } - @Override public void startBroker() throws Exception { @@ -77,20 +66,21 @@ public class TwoNodeTest extends QpidBrokerTestCase setSystemProperty("java.util.logging.config.file", "etc" + File.separator + "log.properties"); _groupCreator.configureClusterNodes(); _groupCreator.setDesignatedPrimaryOnFirstBroker(designedPrimary); - _brokerFailoverUrl = _groupCreator.getConnectionUrlForAllClusterNodes(); + _positiveFailoverUrl = _groupCreator.getConnectionUrlForAllClusterNodes(); + _negativeFailoverUrl = _groupCreator.getConnectionUrlForAllClusterNodes(200, 0, 2); _groupCreator.startCluster(); } public void testMasterDesignatedPrimaryCanBeRestartedWithoutReplica() throws Exception { startCluster(true); - final Connection initialConnection = getConnection(_brokerFailoverUrl); + final Connection initialConnection = getConnection(_positiveFailoverUrl); int masterPort = _groupCreator.getBrokerPortNumberFromConnection(initialConnection); assertProducingConsuming(initialConnection); initialConnection.close(); _groupCreator.stopCluster(); _groupCreator.startNode(masterPort); - final Connection secondConnection = getConnection(_brokerFailoverUrl); + final Connection secondConnection = getConnection(_positiveFailoverUrl); assertProducingConsuming(secondConnection); secondConnection.close(); } @@ -98,12 +88,12 @@ public class TwoNodeTest extends QpidBrokerTestCase public void testClusterRestartWithoutDesignatedPrimary() throws Exception { startCluster(false); - final Connection initialConnection = getConnection(_brokerFailoverUrl); + final Connection initialConnection = getConnection(_positiveFailoverUrl); assertProducingConsuming(initialConnection); initialConnection.close(); _groupCreator.stopCluster(); _groupCreator.startClusterParallel(); - final Connection secondConnection = getConnection(_brokerFailoverUrl); + final Connection secondConnection = getConnection(_positiveFailoverUrl); assertProducingConsuming(secondConnection); secondConnection.close(); } @@ -112,7 +102,7 @@ public class TwoNodeTest extends QpidBrokerTestCase { startCluster(true); _groupCreator.stopNode(_groupCreator.getBrokerPortNumberOfSecondaryNode()); - final Connection connection = getConnection(_brokerFailoverUrl); + final Connection connection = getConnection(_positiveFailoverUrl); assertNotNull("Expected to get a valid connection to primary", connection); assertProducingConsuming(connection); } @@ -124,7 +114,7 @@ public class TwoNodeTest extends QpidBrokerTestCase try { - Connection connection = getConnection(_brokerFailoverUrl); + Connection connection = getConnection(_negativeFailoverUrl); assertProducingConsuming(connection); fail("Exception not thrown"); } @@ -143,7 +133,7 @@ public class TwoNodeTest extends QpidBrokerTestCase try { - getConnection(_brokerFailoverUrl); + getConnection(_negativeFailoverUrl); fail("Connection not expected"); } catch (JMSException e) @@ -155,41 +145,39 @@ public class TwoNodeTest extends QpidBrokerTestCase public void testInitialDesignatedPrimaryStateOfNodes() throws Exception { startCluster(true); - final ManagedBDBHAMessageStore primaryStoreBean = getStoreBeanForNodeAtBrokerPort(_groupCreator.getBrokerPortNumberOfPrimary()); - assertTrue("Expected primary node to be set as designated primary", primaryStoreBean.getDesignatedPrimary()); - final ManagedBDBHAMessageStore secondaryStoreBean = getStoreBeanForNodeAtBrokerPort(_groupCreator.getBrokerPortNumberOfSecondaryNode()); - assertFalse("Expected secondary node to NOT be set as designated primary", secondaryStoreBean.getDesignatedPrimary()); + Map primaryNodeAttributes = _groupCreator.getNodeAttributes(_groupCreator.getBrokerPortNumberOfPrimary()); + assertTrue("Expected primary node to be set as designated primary", + (Boolean) primaryNodeAttributes.get(BDBHAVirtualHostNode.DESIGNATED_PRIMARY)); + + Map secondaryNodeAttributes = _groupCreator.getNodeAttributes(_groupCreator.getBrokerPortNumberOfSecondaryNode()); + assertFalse("Expected secondary node to NOT be set as designated primary", + (Boolean) secondaryNodeAttributes.get(BDBHAVirtualHostNode.DESIGNATED_PRIMARY)); } public void testSecondaryDesignatedAsPrimaryAfterOriginalPrimaryStopped() throws Exception { startCluster(true); - final ManagedBDBHAMessageStore storeBean = getStoreBeanForNodeAtBrokerPort(_groupCreator.getBrokerPortNumberOfSecondaryNode()); + _groupCreator.stopNode(_groupCreator.getBrokerPortNumberOfPrimary()); - assertFalse("Expected node to NOT be set as designated primary", storeBean.getDesignatedPrimary()); - storeBean.setDesignatedPrimary(true); + Map secondaryNodeAttributes = _groupCreator.getNodeAttributes(_groupCreator.getBrokerPortNumberOfSecondaryNode()); + assertFalse("Expected node to NOT be set as designated primary", (Boolean) secondaryNodeAttributes.get(BDBHAVirtualHostNode.DESIGNATED_PRIMARY)); + + _groupCreator.setNodeAttributes(_groupCreator.getBrokerPortNumberOfSecondaryNode(), Collections.singletonMap(BDBHAVirtualHostNode.DESIGNATED_PRIMARY, true)); - long limit = System.currentTimeMillis() + 5000; - while( !storeBean.getDesignatedPrimary() && System.currentTimeMillis() < limit) + int timeout = 5000; + long limit = System.currentTimeMillis() + timeout; + while( !((Boolean)secondaryNodeAttributes.get(BDBHAVirtualHostNode.DESIGNATED_PRIMARY)) && System.currentTimeMillis() < limit) { Thread.sleep(100); + secondaryNodeAttributes = _groupCreator.getNodeAttributes(_groupCreator.getBrokerPortNumberOfSecondaryNode()); } - assertTrue("Expected node to now be set as designated primary", storeBean.getDesignatedPrimary()); + assertTrue("Expected secondary to transition to primary within " + timeout, (Boolean) secondaryNodeAttributes.get(BDBHAVirtualHostNode.DESIGNATED_PRIMARY)); - final Connection connection = getConnection(_brokerFailoverUrl); + final Connection connection = getConnection(_positiveFailoverUrl); assertNotNull("Expected to get a valid connection to new primary", connection); assertProducingConsuming(connection); } - private ManagedBDBHAMessageStore getStoreBeanForNodeAtBrokerPort( - final int activeBrokerPortNumber) throws Exception - { - _jmxUtils.open(activeBrokerPortNumber); - - ManagedBDBHAMessageStore storeBean = _jmxUtils.getManagedObject(ManagedBDBHAMessageStore.class, MANAGED_OBJECT_QUERY); - return storeBean; - } - } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/HeadersBinding.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/HeadersBinding.java index fa4e3f21dd..597fc44e4c 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/HeadersBinding.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/HeadersBinding.java @@ -281,4 +281,10 @@ class HeadersBinding return true; } + + @Override + public int hashCode() + { + return _binding == null ? 0 : _binding.hashCode(); + } } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailabilityMessages.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailabilityMessages.java index 9e497efcd2..b864a8c095 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailabilityMessages.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailabilityMessages.java @@ -44,15 +44,15 @@ public class HighAvailabilityMessages private static Locale _currentLocale = BrokerProperties.getLocale(); public static final String HIGHAVAILABILITY_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability"; - public static final String STOPPED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.stopped"; public static final String INTRUDER_DETECTED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.intruder_detected"; - public static final String STARTED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.started"; public static final String TRANSFER_MASTER_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.transfer_master"; public static final String QUORUM_OVERRIDE_CHANGED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.quorum_override_changed"; - public static final String DETACHED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.detached"; - public static final String MAJORITY_LOST_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.majority_lost"; + public static final String REMOVED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.removed"; + public static final String LEFT_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.left"; + public static final String JOINED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.joined"; + public static final String CREATED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.created"; + public static final String QUORUM_LOST_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.quorum_lost"; public static final String PRIORITY_CHANGED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.priority_changed"; - public static final String ATTACHED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.attached"; public static final String ADDED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.added"; public static final String DELETED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.deleted"; public static final String ROLE_CHANGED_LOG_HIERARCHY = DEFAULT_LOG_HIERARCHY_PREFIX + "highavailability.role_changed"; @@ -61,15 +61,15 @@ public class HighAvailabilityMessages static { Logger.getLogger(HIGHAVAILABILITY_LOG_HIERARCHY); - Logger.getLogger(STOPPED_LOG_HIERARCHY); Logger.getLogger(INTRUDER_DETECTED_LOG_HIERARCHY); - Logger.getLogger(STARTED_LOG_HIERARCHY); Logger.getLogger(TRANSFER_MASTER_LOG_HIERARCHY); Logger.getLogger(QUORUM_OVERRIDE_CHANGED_LOG_HIERARCHY); - Logger.getLogger(DETACHED_LOG_HIERARCHY); - Logger.getLogger(MAJORITY_LOST_LOG_HIERARCHY); + Logger.getLogger(REMOVED_LOG_HIERARCHY); + Logger.getLogger(LEFT_LOG_HIERARCHY); + Logger.getLogger(JOINED_LOG_HIERARCHY); + Logger.getLogger(CREATED_LOG_HIERARCHY); + Logger.getLogger(QUORUM_LOST_LOG_HIERARCHY); Logger.getLogger(PRIORITY_CHANGED_LOG_HIERARCHY); - Logger.getLogger(ATTACHED_LOG_HIERARCHY); Logger.getLogger(ADDED_LOG_HIERARCHY); Logger.getLogger(DELETED_LOG_HIERARCHY); Logger.getLogger(ROLE_CHANGED_LOG_HIERARCHY); @@ -80,14 +80,14 @@ public class HighAvailabilityMessages /** * Log a HighAvailability message of the Format: - *
HA-1012 : The node ''{0}'' from the replication group ''{1}'' is stopped.
+ *
HA-1008 : Intruder detected : Node ''{0}'' ({1})
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage STOPPED(String param1, String param2) + public static LogMessage INTRUDER_DETECTED(String param1, String param2) { - String rawMessage = _messages.getString("STOPPED"); + String rawMessage = _messages.getString("INTRUDER_DETECTED"); final Object[] messageArguments = {param1, param2}; // Create a new MessageFormat to ensure thread safety. @@ -105,23 +105,23 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return STOPPED_LOG_HIERARCHY; + return INTRUDER_DETECTED_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1007: Intruder node ''{0}'' from ''{1}'' is detected in replication group ''{2}''
+ *
HA-1007 : Master transfer requested : to ''{0}'' ({1})
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage INTRUDER_DETECTED(String param1, String param2, String param3) + public static LogMessage TRANSFER_MASTER(String param1, String param2) { - String rawMessage = _messages.getString("INTRUDER_DETECTED"); + String rawMessage = _messages.getString("TRANSFER_MASTER"); - final Object[] messageArguments = {param1, param2, param3}; + final Object[] messageArguments = {param1, param2}; // Create a new MessageFormat to ensure thread safety. // Sharing a MessageFormat and using applyPattern is not thread safe MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); @@ -137,23 +137,23 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return INTRUDER_DETECTED_LOG_HIERARCHY; + return TRANSFER_MASTER_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1013 : The node ''{0}'' from the replication group ''{1}'' is started.
+ *
HA-1011 : Minimum group  : {0}
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage STARTED(String param1, String param2) + public static LogMessage QUORUM_OVERRIDE_CHANGED(String param1) { - String rawMessage = _messages.getString("STARTED"); + String rawMessage = _messages.getString("QUORUM_OVERRIDE_CHANGED"); - final Object[] messageArguments = {param1, param2}; + final Object[] messageArguments = {param1}; // Create a new MessageFormat to ensure thread safety. // Sharing a MessageFormat and using applyPattern is not thread safe MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); @@ -169,23 +169,23 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return STARTED_LOG_HIERARCHY; + return QUORUM_OVERRIDE_CHANGED_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1014 : Transfer master to ''{0}'' is requested on node ''{1}'' from the replication group ''{2}''.
+ *
HA-1004 : Removed : Node : ''{0}'' ({1})
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage TRANSFER_MASTER(String param1, String param2, String param3) + public static LogMessage REMOVED(String param1, String param2) { - String rawMessage = _messages.getString("TRANSFER_MASTER"); + String rawMessage = _messages.getString("REMOVED"); - final Object[] messageArguments = {param1, param2, param3}; + final Object[] messageArguments = {param1, param2}; // Create a new MessageFormat to ensure thread safety. // Sharing a MessageFormat and using applyPattern is not thread safe MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); @@ -201,23 +201,23 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return TRANSFER_MASTER_LOG_HIERARCHY; + return REMOVED_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1009 : The quorum override attribute of node ''{0}'' from the replication group ''{1}'' is set to ''{2}''.
+ *
HA-1006 : Left : Node : ''{0}'' ({1})
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage QUORUM_OVERRIDE_CHANGED(String param1, String param2, String param3) + public static LogMessage LEFT(String param1, String param2) { - String rawMessage = _messages.getString("QUORUM_OVERRIDE_CHANGED"); + String rawMessage = _messages.getString("LEFT"); - final Object[] messageArguments = {param1, param2, param3}; + final Object[] messageArguments = {param1, param2}; // Create a new MessageFormat to ensure thread safety. // Sharing a MessageFormat and using applyPattern is not thread safe MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); @@ -233,21 +233,21 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return QUORUM_OVERRIDE_CHANGED_LOG_HIERARCHY; + return LEFT_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1003 : The node ''{0}'' detached from the replication group ''{1}''.
+ *
HA-1005 : Joined : Node : ''{0}'' ({1})
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage DETACHED(String param1, String param2) + public static LogMessage JOINED(String param1, String param2) { - String rawMessage = _messages.getString("DETACHED"); + String rawMessage = _messages.getString("JOINED"); final Object[] messageArguments = {param1, param2}; // Create a new MessageFormat to ensure thread safety. @@ -265,28 +265,23 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return DETACHED_LOG_HIERARCHY; + return JOINED_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1006 : A majority of nodes from replication group ''{0}'' is not available for node ''{1}''.
+ *
HA-1001 : Created
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage MAJORITY_LOST(String param1, String param2) + public static LogMessage CREATED() { - String rawMessage = _messages.getString("MAJORITY_LOST"); + String rawMessage = _messages.getString("CREATED"); - final Object[] messageArguments = {param1, param2}; - // Create a new MessageFormat to ensure thread safety. - // Sharing a MessageFormat and using applyPattern is not thread safe - MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); - - final String message = formatter.format(messageArguments); + final String message = rawMessage; return new LogMessage() { @@ -297,28 +292,23 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return MAJORITY_LOST_LOG_HIERARCHY; + return CREATED_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1008 : The priority attribute of node ''{0}'' from the replication group ''{1}'' is set to ''{2}''.
+ *
HA-1009 : Insufficient replicas contactable
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage PRIORITY_CHANGED(String param1, String param2, String param3) + public static LogMessage QUORUM_LOST() { - String rawMessage = _messages.getString("PRIORITY_CHANGED"); + String rawMessage = _messages.getString("QUORUM_LOST"); - final Object[] messageArguments = {param1, param2, param3}; - // Create a new MessageFormat to ensure thread safety. - // Sharing a MessageFormat and using applyPattern is not thread safe - MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); - - final String message = formatter.format(messageArguments); + final String message = rawMessage; return new LogMessage() { @@ -329,23 +319,23 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return PRIORITY_CHANGED_LOG_HIERARCHY; + return QUORUM_LOST_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1004 : The node ''{0}'' attached to the replication group ''{1}'' with role ''{2}''.
+ *
HA-1012 : Priority  : {0}
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage ATTACHED(String param1, String param2, String param3) + public static LogMessage PRIORITY_CHANGED(String param1) { - String rawMessage = _messages.getString("ATTACHED"); + String rawMessage = _messages.getString("PRIORITY_CHANGED"); - final Object[] messageArguments = {param1, param2, param3}; + final Object[] messageArguments = {param1}; // Create a new MessageFormat to ensure thread safety. // Sharing a MessageFormat and using applyPattern is not thread safe MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); @@ -361,14 +351,14 @@ public class HighAvailabilityMessages public String getLogHierarchy() { - return ATTACHED_LOG_HIERARCHY; + return PRIORITY_CHANGED_LOG_HIERARCHY; } }; } /** * Log a HighAvailability message of the Format: - *
HA-1001 : A new node ''{0}'' is added into a replication group ''{1}''.
+ *
HA-1003 : Added : Node : ''{0}'' ({1})
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * @@ -400,21 +390,16 @@ public class HighAvailabilityMessages /** * Log a HighAvailability message of the Format: - *
HA-1002 : An existing node ''{0}'' is removed from the replication group ''{1}''.
+ *
HA-1002 : Deleted
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage DELETED(String param1, String param2) + public static LogMessage DELETED() { String rawMessage = _messages.getString("DELETED"); - final Object[] messageArguments = {param1, param2}; - // Create a new MessageFormat to ensure thread safety. - // Sharing a MessageFormat and using applyPattern is not thread safe - MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); - - final String message = formatter.format(messageArguments); + final String message = rawMessage; return new LogMessage() { @@ -432,7 +417,7 @@ public class HighAvailabilityMessages /** * Log a HighAvailability message of the Format: - *
HA-1005 : The role of the node ''{0}'' from the replication group ''{1}'' has changed from ''{2}'' to ''{3}''.
+ *
HA-1010 : Role change reported: Node : ''{0}'' ({1}) : from ''{2}'' to ''{3}''
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * @@ -464,16 +449,16 @@ public class HighAvailabilityMessages /** * Log a HighAvailability message of the Format: - *
HA-1010 : The designated primary attribute of node ''{0}'' from the replication group ''{1}'' is set to ''{2}''.
+ *
HA-1013 : Designated primary : {0}
* Optional values are contained in [square brackets] and are numbered * sequentially in the method call. * */ - public static LogMessage DESIGNATED_PRIMARY_CHANGED(String param1, String param2, String param3) + public static LogMessage DESIGNATED_PRIMARY_CHANGED(String param1) { String rawMessage = _messages.getString("DESIGNATED_PRIMARY_CHANGED"); - final Object[] messageArguments = {param1, param2, param3}; + final Object[] messageArguments = {param1}; // Create a new MessageFormat to ensure thread safety. // Sharing a MessageFormat and using applyPattern is not thread safe MessageFormat formatter = new MessageFormat(rawMessage, _currentLocale); diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailability_logmessages.properties b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailability_logmessages.properties index 94df7cc38b..3c5b0d260f 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailability_logmessages.properties +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/messages/HighAvailability_logmessages.properties @@ -18,18 +18,47 @@ # # HA logging message i18n strings. -ADDED=HA-1001 : A new node ''{0}'' is added into a replication group ''{1}''. -DELETED=HA-1002 : An existing node ''{0}'' is removed from the replication group ''{1}''. -DETACHED=HA-1003 : The node ''{0}'' detached from the replication group ''{1}''. -ATTACHED=HA-1004 : The node ''{0}'' attached to the replication group ''{1}'' with role ''{2}''. -ROLE_CHANGED=HA-1005 : The role of the node ''{0}'' from the replication group ''{1}'' has changed from ''{2}'' to ''{3}''. -MAJORITY_LOST=HA-1006 : A majority of nodes from replication group ''{0}'' is not available for node ''{1}''. -INTRUDER_DETECTED=HA-1007: Intruder node ''{0}'' from ''{1}'' is detected in replication group ''{2}'' -PRIORITY_CHANGED=HA-1008 : The priority attribute of node ''{0}'' from the replication group ''{1}'' is set to ''{2}''. -QUORUM_OVERRIDE_CHANGED=HA-1009 : The quorum override attribute of node ''{0}'' from the replication group ''{1}'' is set to ''{2}''. -DESIGNATED_PRIMARY_CHANGED=HA-1010 : The designated primary attribute of node ''{0}'' from the replication group ''{1}'' is set to ''{2}''. -STOPPED=HA-1011 : The node ''{0}'' from the replication group ''{1}'' is stopped. -STARTED=HA-1012 : The node ''{0}'' from the replication group ''{1}'' is started. -TRANSFER_MASTER=HA-1013 : Transfer master to ''{0}'' is requested on node ''{1}'' from the replication group ''{2}''. +CREATED = HA-1001 : Created +DELETED = HA-1002 : Deleted + +# 0 - Node name +# 1 - Node address +ADDED = HA-1003 : Added : Node : ''{0}'' ({1}) + +# 0 - Node name +# 1 - Node address +REMOVED = HA-1004 : Removed : Node : ''{0}'' ({1}) + +# 0 - Node name +# 1 - Node address +JOINED = HA-1005 : Joined : Node : ''{0}'' ({1}) + +# 0 - Node name +# 1 - Node address +LEFT = HA-1006 : Left : Node : ''{0}'' ({1}) + +# 0 - Node name +# 1 - Node address +TRANSFER_MASTER = HA-1007 : Master transfer requested : to ''{0}'' ({1}) + +# 0 - Node name +# 1 - Node address +INTRUDER_DETECTED = HA-1008 : Intruder detected : Node ''{0}'' ({1}) +QUORUM_LOST = HA-1009 : Insufficient replicas contactable + +# 0 - Node name +# 1 - Node address +# 2 - Previous role value +# 3 - New role value +ROLE_CHANGED = HA-1010 : Role change reported: Node : ''{0}'' ({1}) : from ''{2}'' to ''{3}'' + +# 0 - new value +QUORUM_OVERRIDE_CHANGED = HA-1011 : Minimum group : {0} + +# 0 - new value +PRIORITY_CHANGED = HA-1012 : Priority : {0} + +# 0 - new value +DESIGNATED_PRIMARY_CHANGED = HA-1013 : Designated primary : {0} diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/LogSubjectFormat.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/LogSubjectFormat.java index edb78369ae..d59a09fce9 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/LogSubjectFormat.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/LogSubjectFormat.java @@ -116,5 +116,4 @@ public class LogSubjectFormat */ public static final String QUEUE_FORMAT = "vh(/{0})/qu({1})"; - public static final String VIRTUAL_HOST_NODE_FORMAT = "vhn(/{0}))"; } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/VirtualHostNodeLogSubject.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/VirtualHostNodeLogSubject.java deleted file mode 100644 index fad9a91841..0000000000 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/logging/subjects/VirtualHostNodeLogSubject.java +++ /dev/null @@ -1,33 +0,0 @@ -/* - * - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - * - */ -package org.apache.qpid.server.logging.subjects; - -import static org.apache.qpid.server.logging.subjects.LogSubjectFormat.VIRTUAL_HOST_NODE_FORMAT; - -import org.apache.qpid.server.model.VirtualHostNode; - -public class VirtualHostNodeLogSubject extends AbstractLogSubject -{ - public VirtualHostNodeLogSubject(String nodeName) - { - setLogStringWithFormat(VIRTUAL_HOST_NODE_FORMAT, nodeName); - } -} diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java index 6c8945582c..f944821c6f 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractConfiguredObject.java @@ -318,7 +318,7 @@ public abstract class AbstractConfiguredObject> im checkCandidate((Class) interfaceClass, candidates); } } - if(clazz.getSuperclass() != null & ConfiguredObject.class.isAssignableFrom(clazz.getSuperclass())) + if(clazz.getSuperclass() != null && ConfiguredObject.class.isAssignableFrom(clazz.getSuperclass())) { findBestFitInterface((Class) clazz.getSuperclass(), candidates); } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractSystemConfig.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractSystemConfig.java index 0f4ecb09dc..b0dda69ee6 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractSystemConfig.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/AbstractSystemConfig.java @@ -186,6 +186,7 @@ public abstract class AbstractSystemConfig> ConfiguredObjectRecordConverter converter = new ConfiguredObjectRecordConverter(BrokerModel.getInstance()); Reader reader; + try { URL url = new URL(initialConfigurationLocation); @@ -196,9 +197,18 @@ public abstract class AbstractSystemConfig> reader = new FileReader(initialConfigurationLocation); } - Collection records = converter.readFromJson(org.apache.qpid.server.model.Broker.class, - systemConfig, reader); - return records.toArray(new ConfiguredObjectRecord[records.size()]); + try + { + Collection records = + converter.readFromJson(org.apache.qpid.server.model.Broker.class, + systemConfig, reader); + return records.toArray(new ConfiguredObjectRecord[records.size()]); + } + finally + { + reader.close(); + } + } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java index 545a1d941d..c49c2790df 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java @@ -2925,7 +2925,7 @@ public abstract class AbstractQueue> if(existingPolicy != _exclusive) { ExclusivityPolicy newPolicy = _exclusive; - _exclusive = newPolicy; + _exclusive = existingPolicy; updateExclusivityPolicy(newPolicy); } return true; diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java index 57142e6e1f..d0acf1a46f 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java @@ -1564,7 +1564,7 @@ public abstract class AbstractJDBCMessageStore implements MessageStore data = new byte[0]; } } - return ByteBuffer.wrap(data,offsetInMessage,size); + return ByteBuffer.wrap(data,offsetInMessage,Math.min(size,data.length-offsetInMessage)); } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhostnode/AbstractVirtualHostNode.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhostnode/AbstractVirtualHostNode.java index 38101525cd..ad9df793c8 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhostnode/AbstractVirtualHostNode.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhostnode/AbstractVirtualHostNode.java @@ -24,7 +24,6 @@ import org.apache.log4j.Logger; import org.apache.qpid.server.logging.EventLogger; import org.apache.qpid.server.logging.messages.ConfigStoreMessages; import org.apache.qpid.server.logging.subjects.MessageStoreLogSubject; -import org.apache.qpid.server.logging.subjects.VirtualHostNodeLogSubject; import org.apache.qpid.server.model.AbstractConfiguredObject; import org.apache.qpid.server.model.Broker; import org.apache.qpid.server.model.ConfiguredObject; @@ -54,7 +53,6 @@ public abstract class AbstractVirtualHostNode _broker; private final AtomicReference _state = new AtomicReference(State.UNINITIALIZED); private final EventLogger _eventLogger; - private final VirtualHostNodeLogSubject _virtualHostNodeLogSubject; private DurableConfigurationStore _durableConfigurationStore; @@ -67,7 +65,6 @@ public abstract class AbstractVirtualHostNode systemConfig = _broker.getParent(SystemConfig.class); _eventLogger = systemConfig.getEventLogger(); - _virtualHostNodeLogSubject = new VirtualHostNodeLogSubject(getName()); } @@ -248,8 +245,4 @@ public abstract class AbstractVirtualHostNode> } else { - getVirtualHost().getEventLogger().message(ExchangeMessages.DISCARDMSG(_currentMessage.getExchangeName().asString(), - _currentMessage.getMessagePublishInfo().getRoutingKey() - == null - ? null - : _currentMessage.getMessagePublishInfo() - .getRoutingKey() - .toString())); + AMQShortString exchangeName = _currentMessage.getExchangeName(); + AMQShortString routingKey = _currentMessage.getMessagePublishInfo().getRoutingKey(); + + getVirtualHost().getEventLogger().message( + ExchangeMessages.DISCARDMSG(exchangeName == null ? null : exchangeName.asString(), + routingKey == null ? null : routingKey.asString())); } } diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapImpl.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapImpl.java index 1bd9ab079e..c33af48d8e 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapImpl.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapImpl.java @@ -158,7 +158,7 @@ public class UnacknowledgedMessageMapImpl implements UnacknowledgedMessageMap acknowledged.add(instance); } } - return ackedMessageMap.values(); + return acknowledged; } private void collect(long key, Map msgs) diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapTest.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapTest.java new file mode 100644 index 0000000000..ca52173e66 --- /dev/null +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/UnacknowledgedMessageMapTest.java @@ -0,0 +1,84 @@ +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ +package org.apache.qpid.server.protocol.v0_8; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.util.Collection; + +import junit.framework.TestCase; + +import org.apache.qpid.server.message.MessageInstance; + +public class UnacknowledgedMessageMapTest extends TestCase +{ + public void testDeletedMessagesCantBeAcknowledged() + { + UnacknowledgedMessageMap map = new UnacknowledgedMessageMapImpl(100); + final int expectedSize = 5; + MessageInstance[] msgs = populateMap(map,expectedSize); + assertEquals(expectedSize,map.size()); + Collection acknowledged = map.acknowledge(100, true); + assertEquals(expectedSize, acknowledged.size()); + assertEquals(0,map.size()); + for(int i = 0; i < expectedSize; i++) + { + assertTrue("Message " + i + " is missing", acknowledged.contains(msgs[i])); + } + + map = new UnacknowledgedMessageMapImpl(100); + msgs = populateMap(map,expectedSize); + // simulate some messages being ttl expired + when(msgs[2].lockAcquisition()).thenReturn(Boolean.FALSE); + when(msgs[4].lockAcquisition()).thenReturn(Boolean.FALSE); + + assertEquals(expectedSize,map.size()); + + + acknowledged = map.acknowledge(100, true); + assertEquals(expectedSize-2, acknowledged.size()); + assertEquals(0,map.size()); + for(int i = 0; i < expectedSize; i++) + { + assertEquals(i != 2 && i != 4, acknowledged.contains(msgs[i])); + } + + } + + public MessageInstance[] populateMap(final UnacknowledgedMessageMap map, int size) + { + MessageInstance[] msgs = new MessageInstance[size]; + for(int i = 0; i < size; i++) + { + msgs[i] = createMessageInstance(i); + map.add((long)i,msgs[i]); + } + return msgs; + } + + private MessageInstance createMessageInstance(final int id) + { + MessageInstance instance = mock(MessageInstance.class); + when(instance.lockAcquisition()).thenReturn(Boolean.TRUE); + return instance; + } +} diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageConverter_from_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageConverter_from_1_0.java index 3974ab0af6..266f3b6868 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageConverter_from_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageConverter_from_1_0.java @@ -20,10 +20,32 @@ */ package org.apache.qpid.server.protocol.v1_0; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.ObjectOutputStream; +import java.nio.ByteBuffer; +import java.nio.charset.Charset; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Date; +import java.util.HashSet; +import java.util.Iterator; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.ListIterator; +import java.util.Map; +import java.util.Set; +import java.util.UUID; + import org.apache.qpid.amqp_1_0.messaging.SectionDecoderImpl; import org.apache.qpid.amqp_1_0.type.AmqpErrorException; import org.apache.qpid.amqp_1_0.type.Binary; import org.apache.qpid.amqp_1_0.type.Section; +import org.apache.qpid.amqp_1_0.type.Symbol; +import org.apache.qpid.amqp_1_0.type.UnsignedByte; +import org.apache.qpid.amqp_1_0.type.UnsignedInteger; +import org.apache.qpid.amqp_1_0.type.UnsignedLong; +import org.apache.qpid.amqp_1_0.type.UnsignedShort; import org.apache.qpid.amqp_1_0.type.messaging.AmqpSequence; import org.apache.qpid.amqp_1_0.type.messaging.AmqpValue; import org.apache.qpid.amqp_1_0.type.messaging.Data; @@ -32,17 +54,6 @@ import org.apache.qpid.transport.codec.BBEncoder; import org.apache.qpid.typedmessage.TypedBytesContentWriter; import org.apache.qpid.typedmessage.TypedBytesFormatException; -import java.io.ByteArrayOutputStream; -import java.io.IOException; -import java.io.ObjectOutputStream; -import java.nio.ByteBuffer; -import java.nio.charset.Charset; -import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.ListIterator; -import java.util.Map; - public class MessageConverter_from_1_0 { private static final Charset UTF_8 = Charset.forName("UTF-8"); @@ -91,7 +102,7 @@ public class MessageConverter_from_1_0 Section firstBodySection = sections.get(0); if(firstBodySection instanceof AmqpValue) { - bodyObject = fixObject(((AmqpValue)firstBodySection).getValue()); + bodyObject = convertValue(((AmqpValue)firstBodySection).getValue()); } else if(firstBodySection instanceof Data) { @@ -115,7 +126,7 @@ public class MessageConverter_from_1_0 { totalSequence.addAll(((AmqpSequence)section).getValue()); } - bodyObject = fixObject(totalSequence); + bodyObject = convertValue(totalSequence); } } @@ -127,40 +138,94 @@ public class MessageConverter_from_1_0 return bodyObject; } - private static Object fixObject(final Object value) + private static final Set STANDARD_TYPES = new HashSet<>(Arrays.asList(Boolean.class, + Byte.class, + Short.class, + Integer.class, + Long.class, + Float.class, + Double.class, + Character.class, + String.class, + byte[].class, + UUID.class)); + + private static Map convertMap(final Map map) { - if(value instanceof Binary) + Map resultMap = new LinkedHashMap(); + Iterator iterator = map.entrySet().iterator(); + while(iterator.hasNext()) { - final Binary binaryValue = (Binary) value; - byte[] data = new byte[binaryValue.getLength()]; - binaryValue.asByteBuffer().get(data); - return data; + Map.Entry entry = iterator.next(); + resultMap.put(convertValue(entry.getKey()), convertValue(entry.getValue())); + } - else if(value instanceof List) + return resultMap; + } + + public static Object convertValue(final Object value) + { + if(value != null && !STANDARD_TYPES.contains(value)) { - List listValue = (List) value; - List fixedValue = new ArrayList(listValue.size()); - for(Object o : listValue) + if(value instanceof Map) { - fixedValue.add(fixObject(o)); + return convertMap((Map)value); } - return fixedValue; - } - else if(value instanceof Map) - { - Map mapValue = (Map) value; - Map fixedValue = new LinkedHashMap(mapValue.size()); - for(Map.Entry entry : mapValue.entrySet()) + else if(value instanceof List) + { + return convertList((List)value); + } + else if(value instanceof UnsignedByte) + { + return ((UnsignedByte)value).shortValue(); + } + else if(value instanceof UnsignedShort) + { + return ((UnsignedShort)value).intValue(); + } + else if(value instanceof UnsignedInteger) + { + return ((UnsignedInteger)value).longValue(); + } + else if(value instanceof UnsignedLong) + { + return ((UnsignedLong)value).longValue(); + } + else if(value instanceof Symbol) + { + return value.toString(); + } + else if(value instanceof Date) { - fixedValue.put(fixObject(entry.getKey()),fixObject(entry.getValue())); + return ((Date)value).getTime(); + } + else if(value instanceof Binary) + { + Binary binary = (Binary)value; + byte[] data = new byte[binary.getLength()]; + binary.asByteBuffer().get(data); + return data; + } + else + { + // Throw exception instead? + return value.toString(); } - return fixedValue; } else { return value; } + } + private static List convertList(final List list) + { + List result = new ArrayList(list.size()); + for(Object entry : list) + { + result.add(convertValue(entry)); + } + return result; } public static byte[] convertToBody(Object object) diff --git a/qpid/java/broker-plugins/amqp-msg-conv-0-10-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_10_v1_0/MessageConverter_1_0_to_v0_10.java b/qpid/java/broker-plugins/amqp-msg-conv-0-10-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_10_v1_0/MessageConverter_1_0_to_v0_10.java index 8d77a8cfaf..54d4638bb8 100644 --- a/qpid/java/broker-plugins/amqp-msg-conv-0-10-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_10_v1_0/MessageConverter_1_0_to_v0_10.java +++ b/qpid/java/broker-plugins/amqp-msg-conv-0-10-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_10_v1_0/MessageConverter_1_0_to_v0_10.java @@ -21,6 +21,7 @@ package org.apache.qpid.server.protocol.converter.v0_10_v1_0; import java.nio.ByteBuffer; +import java.util.Map; import org.apache.qpid.server.message.AMQMessageHeader; import org.apache.qpid.server.plugin.MessageConverter; @@ -176,7 +177,8 @@ public class MessageConverter_1_0_to_v0_10 implements MessageConverter) MessageConverter_from_1_0.convertValue(serverMsg.getMessageHeader() + .getHeadersAsMap())); Header header = new Header(deliveryProps, messageProps, null); return new MessageMetaData_0_10(header, size, serverMsg.getArrivalTime()); diff --git a/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v1_0/MessageConverter_1_0_to_v0_8.java b/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v1_0/MessageConverter_1_0_to_v0_8.java index 5b1c25e879..2de21e1a9f 100644 --- a/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v1_0/MessageConverter_1_0_to_v0_8.java +++ b/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-1-0/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v1_0/MessageConverter_1_0_to_v0_8.java @@ -194,7 +194,7 @@ public class MessageConverter_1_0_to_v0_8 implements MessageConverter 0) + while ((read = fileInput.read(buffer)) > 0) + { + output.write(buffer, 0, read); + } + } + else { - output.write(buffer, 0, read); + response.sendError(HttpServletResponse.SC_NOT_FOUND, "unknown file: " + _filename); } } - else - { - response.sendError(HttpServletResponse.SC_NOT_FOUND, "unknown file: " + _filename); - } } } } diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java index 945645ccb1..35252204ac 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java @@ -2101,7 +2101,7 @@ public abstract class AMQSession { /** Used for debugging. */ @@ -736,10 +737,10 @@ public class AMQSession_0_8 extends AMQSession:///[]/[][?