summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKeith Wall <kwall@apache.org>2013-12-13 18:50:15 +0000
committerKeith Wall <kwall@apache.org>2013-12-13 18:50:15 +0000
commit5f205440ea35de71c5d2a58bcb68afa1a5f32448 (patch)
tree66eb7ca7378e8303f171854d23893cfd2cb971f5
parentfda8441f2bc5c8cfcc1cd060f6a918d3ceab84bc (diff)
downloadqpid-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
-rw-r--r--qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/BDBHAVirtualHost.java8
-rw-r--r--qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacade.java94
-rw-r--r--qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeFactory.java19
-rw-r--r--qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNode.java220
-rw-r--r--qpid/java/bdbstore/src/main/java/org/apache/qpid/server/store/berkeleydb/replication/RemoteReplicationNodeFactory.java26
-rw-r--r--qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/BDBHAMessageStoreTest.java14
-rw-r--r--qpid/java/bdbstore/src/test/java/org/apache/qpid/server/store/berkeleydb/ReplicatedEnvironmentFacadeTest.java28
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/ReplicationNode.java114
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/UUIDGenerator.java5
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/VirtualHost.java1
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java22
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/replication/ReplicationGroupListener.java59
-rwxr-xr-xqpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java2
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;