diff options
| author | Keith Wall <kwall@apache.org> | 2013-12-13 18:50:15 +0000 |
|---|---|---|
| committer | Keith Wall <kwall@apache.org> | 2013-12-13 18:50:15 +0000 |
| commit | 5f205440ea35de71c5d2a58bcb68afa1a5f32448 (patch) | |
| tree | 66eb7ca7378e8303f171854d23893cfd2cb971f5 | |
| parent | fda8441f2bc5c8cfcc1cd060f6a918d3ceab84bc (diff) | |
| download | qpid-python-5f205440ea35de71c5d2a58bcb68afa1a5f32448.tar.gz | |
QPID-5411: Initial code introducing replication node into broker model
git-svn-id: https://svn.apache.org/repos/asf/qpid/branches/java-broker-bdb-ha@1550805 13f79535-47bb-0310-9956-ffa450edef68
13 files changed, 577 insertions, 35 deletions
diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHost.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHost.java index aa51b7bb6e..2f85c91a02 100644 --- a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHost.java +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHost.java @@ -32,6 +32,7 @@ import org.apache.qpid.server.logging.actors.AbstractActor; import org.apache.qpid.server.logging.actors.CurrentActor; import org.apache.qpid.server.logging.subjects.MessageStoreLogSubject; import org.apache.qpid.server.model.VirtualHost; +import org.apache.qpid.server.replication.ReplicationGroupListener; import org.apache.qpid.server.stats.StatisticsGatherer; import org.apache.qpid.server.store.DurableConfigurationRecoverer; import org.apache.qpid.server.store.DurableConfigurationStore; @@ -80,8 +81,6 @@ public class BDBHAVirtualHost extends AbstractVirtualHost _messageStore.addEventListener(new AfterActivationListener(), Event.AFTER_ACTIVATE); _messageStore.addEventListener(new BeforeCloseListener(), Event.BEFORE_CLOSE); - - _messageStore.addEventListener(new AfterInitialisationListener(), Event.AFTER_INIT); _messageStore.addEventListener(new BeforePassivationListener(), Event.BEFORE_PASSIVATE); @@ -99,7 +98,10 @@ public class BDBHAVirtualHost extends AbstractVirtualHost recoveryHandler ); - ((ReplicatedEnvironmentFacade)_messageStore.getEnvironmentFacade()).setStateChangeListener(new BDBHAMessageStoreStateChangeListener()); + // Make the virtualhost model object a replication group listener + ReplicatedEnvironmentFacade environmentFacade = (ReplicatedEnvironmentFacade) _messageStore.getEnvironmentFacade(); + environmentFacade.setReplicationGroupListener((ReplicationGroupListener) virtualHost); + environmentFacade.setStateChangeListener(new BDBHAMessageStoreStateChangeListener()); } diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacade.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacade.java index 8540c00f05..3677a0bcf4 100644 --- a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacade.java +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacade.java @@ -37,7 +37,10 @@ import java.util.concurrent.atomic.AtomicReference; import org.apache.log4j.Logger; import org.apache.qpid.AMQStoreException; +import org.apache.qpid.server.replication.ReplicationGroupListener; import org.apache.qpid.server.store.StoreFuture; +import org.apache.qpid.server.store.berkeleydb.replication.RemoteReplicationNode; +import org.apache.qpid.server.store.berkeleydb.replication.RemoteReplicationNodeFactory; import com.sleepycat.je.Database; import com.sleepycat.je.DatabaseConfig; @@ -56,6 +59,7 @@ import com.sleepycat.je.rep.NetworkRestore; import com.sleepycat.je.rep.NetworkRestoreConfig; import com.sleepycat.je.rep.ReplicatedEnvironment; import com.sleepycat.je.rep.ReplicationConfig; +import com.sleepycat.je.rep.ReplicationGroup; import com.sleepycat.je.rep.ReplicationMutableConfig; import com.sleepycat.je.rep.ReplicationNode; import com.sleepycat.je.rep.StateChangeEvent; @@ -64,6 +68,8 @@ import com.sleepycat.je.rep.util.ReplicationGroupAdmin; public class ReplicatedEnvironmentFacade implements EnvironmentFacade, StateChangeListener { + private static final Logger LOGGER = Logger.getLogger(ReplicatedEnvironmentFacade.class); + @SuppressWarnings("serial") private static final Map<String, String> REPCONFIG_DEFAULTS = Collections.unmodifiableMap(new HashMap<String, String>() {{ @@ -101,7 +107,6 @@ public class ReplicatedEnvironmentFacade implements EnvironmentFacade, StateChan put(ReplicationConfig.ENV_UNKNOWN_STATE_TIMEOUT, "5 s"); }}); - private static final Logger LOGGER = Logger.getLogger(ReplicatedEnvironmentFacade.class); public static final String TYPE = "BDB-HA"; // TODO: get rid of these names @@ -111,28 +116,34 @@ public class ReplicatedEnvironmentFacade implements EnvironmentFacade, StateChan private volatile ReplicatedEnvironment _environment; private CommitThreadWrapper _commitThreadWrapper; - private String _groupName; - private String _nodeName; - private String _nodeHostPort; - private String _helperHostPort; - private Durability _durability; - private boolean _designatedPrimary; - private boolean _coalescingSync; + private final String _groupName; + private final String _nodeName; + private final String _nodeHostPort; + private final String _helperHostPort; + private final Durability _durability; + private final boolean _designatedPrimary; + private final boolean _coalescingSync; private volatile StateChangeListener _stateChangeListener; - private String _environmentPath; - private Map<String, String> _environmentParameters; - private Map<String, String> _replicationEnvironmentParameters; - private String _name; - private ExecutorService _executor = Executors.newFixedThreadPool(1); - private AtomicReference<State> _state = new AtomicReference<State>(State.INITIAL); + private final String _environmentPath; + private final Map<String, String> _environmentParameters; + private final Map<String, String> _replicationEnvironmentParameters; + private final String _name; + private final ExecutorService _executor = Executors.newFixedThreadPool(1); + private final AtomicReference<State> _state = new AtomicReference<State>(State.INITIAL); private final ConcurrentMap<String, Database> _databases = new ConcurrentHashMap<String, Database>(); - public ReplicatedEnvironmentFacade(String name, String environmentPath, String groupName, String nodeName, String nodeHostPort, - String helperHostPort, Durability durability, boolean designatedPrimary, boolean coalescingSync, - Map<String, String> environmentParameters, Map<String, String> replicationEnvironmentParameters) - { + private ReplicationGroupListener _replicationGroupListener; + private final RemoteReplicationNodeFactory _remoteReplicationNodeFactory; + + public ReplicatedEnvironmentFacade(String name, String environmentPath, + String groupName, String nodeName, String nodeHostPort, + String helperHostPort, Durability durability, + Boolean designatedPrimary, Boolean coalescingSync, + Map<String, String> envConfigMap, + Map<String, String> replicationConfig, RemoteReplicationNodeFactory remoteReplicationNodeFactory) + { _name = name; _environmentPath = environmentPath; _groupName = groupName; @@ -142,15 +153,17 @@ public class ReplicatedEnvironmentFacade implements EnvironmentFacade, StateChan _durability = durability; _designatedPrimary = designatedPrimary; _coalescingSync = coalescingSync; - _environmentParameters = environmentParameters; - _replicationEnvironmentParameters = replicationEnvironmentParameters; + _environmentParameters = envConfigMap; + _replicationEnvironmentParameters = replicationConfig; + _remoteReplicationNodeFactory = remoteReplicationNodeFactory; _state.set(State.OPENING); _environment = createEnvironment(environmentPath, groupName, nodeName, nodeHostPort, helperHostPort, durability, - designatedPrimary, environmentParameters, replicationEnvironmentParameters); - startCommitThread(name, _environment); + designatedPrimary, _environmentParameters, _replicationEnvironmentParameters); + startCommitThread(_name, _environment); } + @Override public StoreFuture commit(final Transaction tx, final boolean syncCommit) throws AMQStoreException { @@ -430,6 +443,42 @@ public class ReplicatedEnvironmentFacade implements EnvironmentFacade, StateChan return _state.get(); } + /** + * Sets the replication group listener. Whenever a new listener is set, the listener + * will hear {@link ReplicationGroupListener#onReplicationNodeRecovered(org.apache.qpid.server.model.ReplicationNode) + * for every existing remote node. + * + * @param replicationGroupListener listener + */ + public void setReplicationGroupListener(ReplicationGroupListener replicationGroupListener) + { + _replicationGroupListener = replicationGroupListener; + if (_replicationGroupListener != null) + { + notifyExistingRemoteReplicationNodes(_replicationGroupListener); + } + } + + private void notifyExistingRemoteReplicationNodes(ReplicationGroupListener listener) + { + ReplicationGroup group = _environment.getGroup(); + Set<ReplicationNode> nodes = new HashSet<ReplicationNode>(group.getElectableNodes()); + String localNodeName = getNodeName(); + + for (ReplicationNode replicationNode : nodes) + { + String discoveredNodeName = replicationNode.getName(); + if (!discoveredNodeName.equals(localNodeName)) + { + // TODO remote replication nodes should be cached + RemoteReplicationNode remoteNode = _remoteReplicationNodeFactory.create(group.getName(), + replicationNode.getName(), + replicationNode.getHostName(), replicationNode.getPort()); + listener.onReplicationNodeRecovered(remoteNode); + } + } + } + private ReplicationGroupAdmin createReplicationGroupAdmin() { final Set<InetSocketAddress> helpers = new HashSet<InetSocketAddress>(); @@ -652,5 +701,4 @@ public class ReplicatedEnvironmentFacade implements EnvironmentFacade, StateChan CLOSING, CLOSED } - } diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeFactory.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeFactory.java index 5ea0100c64..adabbc8486 100644 --- a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeFactory.java +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeFactory.java @@ -25,6 +25,8 @@ import java.util.Map; import org.apache.qpid.server.configuration.IllegalConfigurationException; import org.apache.qpid.server.model.VirtualHost; +import org.apache.qpid.server.store.berkeleydb.replication.RemoteReplicationNode; +import org.apache.qpid.server.store.berkeleydb.replication.RemoteReplicationNodeFactory; import com.sleepycat.je.Durability; import com.sleepycat.je.Durability.ReplicaAckPolicy; @@ -81,7 +83,7 @@ public class ReplicatedEnvironmentFacadeFactory implements EnvironmentFacadeFact } return new ReplicatedEnvironmentFacade(name, storeLocation, groupName, nodeName, nodeHostPort, helperHostPort, durability, - designatedPrimary, coalescingSync, envConfigMap, replicationConfig); + designatedPrimary, coalescingSync, envConfigMap, replicationConfig, new RemoteReplicationNodeFactoryImpl(virtualHost)); } private String getValidatedStringAttribute(org.apache.qpid.server.model.VirtualHost virtualHost, String attributeName) @@ -126,4 +128,19 @@ public class ReplicatedEnvironmentFacadeFactory implements EnvironmentFacadeFact return defaultVal; } + private static class RemoteReplicationNodeFactoryImpl implements RemoteReplicationNodeFactory + { + private VirtualHost _virtualHost; + + public RemoteReplicationNodeFactoryImpl(VirtualHost virtualHost) + { + _virtualHost = virtualHost; + } + + @Override + public RemoteReplicationNode create(String groupName, String nodeName, String host, int port) + { + return new RemoteReplicationNode(groupName, nodeName, host, port, _virtualHost); + } + } } diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNode.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNode.java new file mode 100644 index 0000000000..5649c3302d --- /dev/null +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNode.java @@ -0,0 +1,220 @@ +package org.apache.qpid.server.store.berkeleydb.replication; + +import java.security.AccessControlException; +import java.util.Collection; +import java.util.Map; +import java.util.UUID; + +import org.apache.qpid.server.model.ConfigurationChangeListener; +import org.apache.qpid.server.model.ConfiguredObject; +import org.apache.qpid.server.model.IllegalStateTransitionException; +import org.apache.qpid.server.model.LifetimePolicy; +import org.apache.qpid.server.model.ReplicationNode; +import org.apache.qpid.server.model.State; +import org.apache.qpid.server.model.Statistics; +import org.apache.qpid.server.model.UUIDGenerator; +import org.apache.qpid.server.model.VirtualHost; + +/** + * Represents a remote replication node in a BDB group. + */ +public class RemoteReplicationNode implements ReplicationNode +{ + + private final UUID _id; + private final String _groupName; + private final String _nodeName; + private final String _host; + private final int _port; + private final VirtualHost _virtualHost; + + public RemoteReplicationNode(String groupName, String nodeName, String host, int port, VirtualHost virtualHost) + { + super(); + _id = UUIDGenerator.generateReplicationNodeId(groupName, nodeName); + _groupName = groupName; + _nodeName = nodeName; + _host = host; + _port = port; + _virtualHost = virtualHost; + } + + @Override + public UUID getId() + { + return _id; + } + + @Override + public String getName() + { + return _nodeName; + } + + @Override + public String setName(String currentName, String desiredName) + throws IllegalStateException, AccessControlException + { + throw new UnsupportedOperationException(); + } + + @Override + public State getDesiredState() + { + throw new UnsupportedOperationException(); + } + + @Override + public State setDesiredState(State currentState, State desiredState) + throws IllegalStateTransitionException, AccessControlException + { + throw new UnsupportedOperationException(); + } + + @Override + public State getActualState() + { + throw new UnsupportedOperationException(); + } + + @Override + public void addChangeListener(ConfigurationChangeListener listener) + { + throw new UnsupportedOperationException(); + } + + @Override + public boolean removeChangeListener(ConfigurationChangeListener listener) + { + throw new UnsupportedOperationException(); + } + + @SuppressWarnings("unchecked") + @Override + public <T extends ConfiguredObject> T getParent(Class<T> clazz) + { + if (clazz == VirtualHost.class) + { + return (T) _virtualHost; + } + throw new IllegalArgumentException(); + } + + @Override + public boolean isDurable() + { + return true; + } + + @Override + public void setDurable(boolean durable) throws IllegalStateException, + AccessControlException, IllegalArgumentException + { + throw new UnsupportedOperationException(); + } + + @Override + public LifetimePolicy getLifetimePolicy() + { + return LifetimePolicy.PERMANENT; + } + + @Override + public LifetimePolicy setLifetimePolicy(LifetimePolicy expected, + LifetimePolicy desired) throws IllegalStateException, + AccessControlException, IllegalArgumentException + { + throw new UnsupportedOperationException(); + } + + @Override + public long getTimeToLive() + { + return 0; + } + + @Override + public long setTimeToLive(long expected, long desired) + throws IllegalStateException, AccessControlException, + IllegalArgumentException + { + throw new UnsupportedOperationException(); + } + + @Override + public Collection<String> getAttributeNames() + { + return ReplicationNode.AVAILABLE_ATTRIBUTES; + } + + @Override + public Object getAttribute(String name) + { + if (ReplicationNode.ID.equals(name)) + { + return getId(); + } + else if (ReplicationNode.NAME.equals(name)) + { + return getName(); + } + else if (ReplicationNode.LIFETIME_POLICY.equals(name)) + { + return getLifetimePolicy(); + } + else if (ReplicationNode.DURABLE.equals(name)) + { + return isDurable(); + } + else if (ReplicationNode.HOST_PORT.equals(name)) + { + return _host + ":" + _port; + } + else if (ReplicationNode.GROUP_NAME.equals(name)) + { + return _groupName; + } + throw new UnsupportedOperationException(); + } + + @Override + public Map<String, Object> getActualAttributes() + { + throw new UnsupportedOperationException(); + } + + @Override + public Object setAttribute(String name, Object expected, Object desired) + throws IllegalStateException, AccessControlException, + IllegalArgumentException + { + throw new UnsupportedOperationException(); + } + + @Override + public Statistics getStatistics() + { + throw new UnsupportedOperationException(); + } + + @Override + public <C extends ConfiguredObject> Collection<C> getChildren(Class<C> clazz) + { + throw new UnsupportedOperationException(); + } + + @Override + public <C extends ConfiguredObject> C createChild(Class<C> childClass, + Map<String, Object> attributes, ConfiguredObject... otherParents) + { + throw new UnsupportedOperationException(); + } + + @Override + public void setAttributes(Map<String, Object> attributes) + throws IllegalStateException, AccessControlException, + IllegalArgumentException + { + throw new UnsupportedOperationException(); + } +} diff --git a/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNodeFactory.java b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNodeFactory.java new file mode 100644 index 0000000000..78d1034230 --- /dev/null +++ b/qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNodeFactory.java @@ -0,0 +1,26 @@ +/* + * + * 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.store.berkeleydb.replication; + +public interface RemoteReplicationNodeFactory +{ + RemoteReplicationNode create(String groupName, String nodeName, String host, int port); +} diff --git a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAMessageStoreTest.java b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAMessageStoreTest.java index b13814f6cb..436b342560 100644 --- a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAMessageStoreTest.java +++ b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAMessageStoreTest.java @@ -22,11 +22,12 @@ package org.apache.qpid.server.store.berkeleydb; import java.io.File; import java.net.InetAddress; - import java.util.HashMap; import java.util.Map; + import org.apache.commons.configuration.XMLConfiguration; import org.apache.qpid.server.configuration.VirtualHostConfiguration; +import org.apache.qpid.server.replication.ReplicationGroupListener; import org.apache.qpid.server.util.BrokerTestHelper; import org.apache.qpid.server.virtualhost.VirtualHost; import org.apache.qpid.test.utils.QpidTestCase; @@ -65,7 +66,7 @@ public class BDBHAMessageStoreTest extends QpidTestCase FileUtils.delete(new File(_workDir), true); _configXml = new XMLConfiguration(); - _modelVhost = mock(org.apache.qpid.server.model.VirtualHost.class); + _modelVhost = mock(TestVirtualHost.class); BrokerTestHelper.setUp(); @@ -95,7 +96,7 @@ public class BDBHAMessageStoreTest extends QpidTestCase String vhostName = "test" + _masterPort; VirtualHostConfiguration configuration = new VirtualHostConfiguration(vhostName, _configXml.subset("virtualhosts.virtualhost." + vhostName), BrokerTestHelper.createBrokerMock()); - _virtualHost = BrokerTestHelper.createVirtualHost(configuration,null,_modelVhost); + _virtualHost = BrokerTestHelper.createVirtualHost(configuration, null, _modelVhost); BDBMessageStore store = (BDBMessageStore) _virtualHost.getMessageStore(); // test whether JVM system settings were applied @@ -125,7 +126,7 @@ public class BDBHAMessageStoreTest extends QpidTestCase _configXml.addProperty("virtualhosts.virtualhost.name", vhostName); _configXml.addProperty(vhostPrefix + ".type", BDBHAVirtualHostFactory.TYPE); - when(_modelVhost.getAttribute(eq(_modelVhost.STORE_PATH))).thenReturn(_workDir + File.separator + when(_modelVhost.getAttribute(eq(org.apache.qpid.server.model.VirtualHost.STORE_PATH))).thenReturn(_workDir + File.separator + port); when(_modelVhost.getAttribute(eq("haGroupName"))).thenReturn(_groupName); when(_modelVhost.getAttribute(eq("haNodeName"))).thenReturn(nodeName); @@ -163,4 +164,9 @@ public class BDBHAMessageStoreTest extends QpidTestCase } return _host + ":" + _masterPort; } + + // TODO: a temporary work around against issue with casting VirtualHost model object to ReplicationGroupListener in the BDBHAVH + public static interface TestVirtualHost extends org.apache.qpid.server.model.VirtualHost, ReplicationGroupListener + { + } } diff --git a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeTest.java b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeTest.java index c62e841ec3..ce9c35355f 100644 --- a/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeTest.java +++ b/qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeTest.java @@ -20,6 +20,10 @@ */ package org.apache.qpid.server.store.berkeleydb; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; + import java.io.File; import java.util.ArrayList; import java.util.Collections; @@ -36,6 +40,9 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import org.apache.qpid.AMQStoreException; +import org.apache.qpid.server.replication.ReplicationGroupListener; +import org.apache.qpid.server.store.berkeleydb.replication.RemoteReplicationNode; +import org.apache.qpid.server.store.berkeleydb.replication.RemoteReplicationNodeFactory; import org.apache.qpid.test.utils.QpidTestCase; import com.sleepycat.bind.tuple.IntegerBinding; @@ -65,6 +72,7 @@ public class ReplicatedEnvironmentFacadeTest extends EnvironmentFacadeTestCase private static final boolean TEST_DESIGNATED_PRIMARY = true; private static final boolean TEST_COALESCING_SYNC = true; private final Map<String, ReplicatedEnvironmentFacade> _nodes = new HashMap<String, ReplicatedEnvironmentFacade>(); + private RemoteReplicationNodeFactory _remoteReplicationNodeFactory = mock(RemoteReplicationNodeFactory.class);; public void tearDown() throws Exception { @@ -136,6 +144,24 @@ public class ReplicatedEnvironmentFacadeTest extends EnvironmentFacadeTestCase assertEquals("Unexpected group members", expectedGroupMembers, new HashSet<Map<String, String>>(groupMembers)); } + public void testReplicationGroupListenerHearsAboutExistingRemoteReplicationNodes() throws Exception + { + ReplicatedEnvironmentFacade replicatedEnvironmentFacade = getEnvironmentFacade(); + String nodeName2 = TEST_NODE_NAME + "_2"; + String host = "localhost"; + int port = getNextAvailable(TEST_NODE_PORT + 1); + String node2NodeHostPort = host + ":" + port; + joinReplica(nodeName2, node2NodeHostPort); + + List<Map<String, String>> groupMembers = replicatedEnvironmentFacade.getGroupMembers(); + assertEquals("Unexpected number of nodes at start of test", 2, groupMembers.size()); + + ReplicationGroupListener listener = mock(ReplicationGroupListener.class); + replicatedEnvironmentFacade.setReplicationGroupListener(listener); + verify(listener).onReplicationNodeRecovered(any(RemoteReplicationNode.class)); + verify(_remoteReplicationNodeFactory).create(TEST_GROUP_NAME, nodeName2, host, port); + } + public void testRemoveNodeFromGroup() throws Exception { ReplicatedEnvironmentFacade environmentFacade = getEnvironmentFacade(); @@ -418,7 +444,7 @@ public class ReplicatedEnvironmentFacadeTest extends EnvironmentFacadeTestCase ReplicatedEnvironmentFacade ref = new ReplicatedEnvironmentFacade(getName(), nodePath, TEST_GROUP_NAME, nodeName, nodeHostPort, TEST_NODE_HELPER_HOST_PORT, TEST_DURABILITY, designatedPrimary, TEST_COALESCING_SYNC, - Collections.<String, String> emptyMap(), repConfig); + Collections.<String, String> emptyMap(), repConfig, _remoteReplicationNodeFactory); return ref; } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/ReplicationNode.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/ReplicationNode.java new file mode 100644 index 0000000000..57da40ff53 --- /dev/null +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/ReplicationNode.java @@ -0,0 +1,114 @@ +/* + * + * 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.model; + +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; + +public interface ReplicationNode extends ConfiguredObject +{ + String ID = "id"; + String NAME = "name"; + String STATE = "state"; + String CREATED = "created"; + String DURABLE = "durable"; + String LIFETIME_POLICY = "lifetimePolicy"; + String TIME_TO_LIVE = "timeToLive"; + String TYPE = "type"; + String UPDATED = "updated"; + + /** Name of the group to which this replication node belongs */ + String GROUP_NAME = "groupName"; + + /** Node host name/IP and port separated by semicolon*/ + String HOST_PORT = "hostPort"; + + /** Node helper host name/IP and port separated by semicolon*/ + String HELPER_HOST_PORT = "helperHostPort"; + + /** Durability settings*/ + String DURABILITY = "durability"; + + /** Sync multiple transactions on disc at the same time*/ + String COALESCING_SYNC = "coalescingSync"; + + /** A designated primary setting for 2-nodes group*/ + String DESIGNATED_PRIMARY = "designatedPrimary"; + + /** Node priority*/ + String PRIORITY = "priority"; + + /** The overridden minimum number of group nodes required to commit transaction on this node instead of simple majority*/ + String QUORUM_OVERRIDE = "quorumOverride"; + + /** Node role: MASTER,REPLICA,UNKNOWN,DETACHED*/ + String ROLE = "role"; + + /** Time when node joined the group */ + String JOIN_TIME = "joinTime"; + + /** Last known replication transaction id */ + String LAST_KNOWN_REPLICATION_TRANSACTION_ID= "lastKnownReplicationTransactionId"; + + /** Map with additional implementation specific node settings */ + String PARAMETERS = "parameters"; + + /** Map with additional implementation specific replication parameters*/ + String REPLICATION_PARAMETERS = "replicationParameters"; + + /** Name of the node which is a master at the moment*/ + String MASTER_NODE = "masterNode"; + + /** Store path */ + String STORE_PATH = "storePath"; + + // Attributes + public static final Collection<String> AVAILABLE_ATTRIBUTES = + Collections.unmodifiableList( + Arrays.asList( + ID, + NAME, + TYPE, + STATE, + DURABLE, + LIFETIME_POLICY, + TIME_TO_LIVE, + CREATED, + UPDATED, + GROUP_NAME, + HOST_PORT, + HELPER_HOST_PORT, + DURABILITY, + COALESCING_SYNC, + DESIGNATED_PRIMARY, + PRIORITY, + QUORUM_OVERRIDE, + ROLE, + JOIN_TIME, + LAST_KNOWN_REPLICATION_TRANSACTION_ID, + PARAMETERS, + REPLICATION_PARAMETERS, + MASTER_NODE, + STORE_PATH + )); + +} diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/UUIDGenerator.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/UUIDGenerator.java index 7def89025d..2af644518a 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/UUIDGenerator.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/UUIDGenerator.java @@ -97,4 +97,9 @@ public class UUIDGenerator { return createUUID(PreferencesProvider.class.getName(), authenticationProviderName, preferencesProviderName); } + + public static UUID generateReplicationNodeId(String groupName, String nodeName) + { + return createUUID(ReplicationNode.class.getName(), groupName, nodeName); + } } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/VirtualHost.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/VirtualHost.java index ae07005679..b3be6a0b88 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/VirtualHost.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/VirtualHost.java @@ -131,6 +131,7 @@ public interface VirtualHost extends ConfiguredObject Collection<Connection> getConnections(); Collection<Queue> getQueues(); Collection<Exchange> getExchanges(); + Collection<ReplicationNode> getReplicationNodes(); Exchange createExchange(String name, State initialState, boolean durable, LifetimePolicy lifetime, long ttl, String type, Map<String, Object> attributes) diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java index 16151dbb63..e68f772e62 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java @@ -55,6 +55,7 @@ import org.apache.qpid.server.model.Port; import org.apache.qpid.server.model.Protocol; import org.apache.qpid.server.model.Queue; import org.apache.qpid.server.model.QueueType; +import org.apache.qpid.server.model.ReplicationNode; import org.apache.qpid.server.model.State; import org.apache.qpid.server.model.Statistics; import org.apache.qpid.server.model.UUIDGenerator; @@ -66,6 +67,7 @@ import org.apache.qpid.server.protocol.AMQConnectionModel; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.queue.AMQQueueFactory; import org.apache.qpid.server.queue.QueueEntry; +import org.apache.qpid.server.replication.ReplicationGroupListener; import org.apache.qpid.server.security.SecurityManager; import org.apache.qpid.server.security.access.Operation; import org.apache.qpid.server.security.auth.AuthenticatedPrincipal; @@ -82,7 +84,7 @@ import org.apache.qpid.server.virtualhost.VirtualHostListener; import org.apache.qpid.server.virtualhost.VirtualHostRegistry; import org.apache.qpid.server.virtualhost.plugins.QueueExistsException; -public final class VirtualHostAdapter extends AbstractAdapter implements VirtualHost, VirtualHostListener +public final class VirtualHostAdapter extends AbstractAdapter implements VirtualHost, VirtualHostListener, ReplicationGroupListener { private static final Logger LOGGER = Logger.getLogger(VirtualHostAdapter.class); @@ -111,6 +113,8 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual private final List<VirtualHostAlias> _aliases = new ArrayList<VirtualHostAlias>(); private StatisticsGatherer _brokerStatisticsGatherer; + private final List<ReplicationNode> _replicationNodes = new ArrayList<ReplicationNode>(); + public VirtualHostAdapter(UUID id, Map<String, Object> attributes, Broker broker, StatisticsGatherer brokerStatisticsGatherer, TaskExecutor taskExecutor) { super(id, null, MapValueConverter.convert(attributes, ATTRIBUTE_TYPES, false), taskExecutor, false); @@ -257,6 +261,15 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual } } + @Override + public Collection<ReplicationNode> getReplicationNodes() + { + synchronized (_replicationNodes) + { + return Collections.unmodifiableList(_replicationNodes); + } + } + public Exchange createExchange(Map<String, Object> attributes) throws AccessControlException, IllegalArgumentException @@ -1252,4 +1265,11 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual throw new AccessControlException("Setting of virtual host attributes is denied"); } } + + @Override + public void onReplicationNodeRecovered(ReplicationNode node) + { + _replicationNodes.add(node); + } + } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/replication/ReplicationGroupListener.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/replication/ReplicationGroupListener.java new file mode 100644 index 0000000000..0e26714173 --- /dev/null +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/replication/ReplicationGroupListener.java @@ -0,0 +1,59 @@ +/* + * + * 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.replication; + +import org.apache.qpid.server.model.ReplicationNode; + +public interface ReplicationGroupListener +{ + /** + * Fired when a remote replication node is added to a group. This event happens + * exactly once just after a new replication node is created. + */ + //void onReplicationNodeAddedToGroup(ReplicationNode node); + + /** + * Fired exactly once for each existing remote node. Used to inform the application + * on any existing nodes as it starts up for the first time. + */ + void onReplicationNodeRecovered(ReplicationNode node); + + /** + * Fired when a remote replication node is (permanently) removed from group. This event + * happens exactly once just after the existing replication node is deleted. + */ + //void onReplicationNodeRemovedFromGroup(ReplicationNode node); + + /** + * Fired when a remote replication node (that is already a member of the group) joins + * the group. This will typically occur when another replication node is started perhaps + * because the broker has been started. + */ + //void onReplicationNodeUp(); + + /** + * Fired when a remote replication node (that is already a member of the group) leaves + * the group. This will typically occur when another replication node is stopped perhaps + * because its broker has been stopped. + */ + //void onReplicationNodeDown(); + +} diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java index 2ebbedccd4..55c705c5ce 100755 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java @@ -25,7 +25,6 @@ import java.util.Map; import java.util.UUID; import java.util.concurrent.ScheduledFuture; import org.apache.qpid.AMQException; -import org.apache.qpid.AMQSecurityException; import org.apache.qpid.common.Closeable; import org.apache.qpid.server.configuration.VirtualHostConfiguration; import org.apache.qpid.server.connection.IConnectionRegistry; @@ -33,7 +32,6 @@ import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.plugin.ExchangeType; import org.apache.qpid.server.protocol.LinkRegistry; import org.apache.qpid.server.queue.AMQQueue; -import org.apache.qpid.server.queue.QueueRegistry; import org.apache.qpid.server.security.SecurityManager; import org.apache.qpid.server.stats.StatisticsGatherer; import org.apache.qpid.server.store.DurableConfigurationStore; |
