diff options
| author | Robert Gemmell <robbie@apache.org> | 2012-05-17 17:26:04 +0000 |
|---|---|---|
| committer | Robert Gemmell <robbie@apache.org> | 2012-05-17 17:26:04 +0000 |
| commit | f5d67044a9797c397764a7ac1aa1a1ed4aa893a3 (patch) | |
| tree | cf8a9cf6a5f741e31417ca4d32a6b708bb3b9fdd /qpid/java/broker/src/main | |
| parent | f523b9e510fc90ce3f7f7d7c2960f3bfee3d42df (diff) | |
| download | qpid-python-f5d67044a9797c397764a7ac1aa1a1ed4aa893a3.tar.gz | |
QPID-4006: add support for using BDB HA to form an active-passive cluster for persistent messaging
- Includes support for setting BDB configuration parameters via the store configuration, both for the existing store and the new HA variant.
- Removes the MessageStoreFactory and reverts store configuration to historical values.
Applied patch from Keith Wall, Andrew MacBean <andymacbean@gmail.com>, Oleksandr Rudyy <orudyy@gmail.com>, Philip Harvey <phil@philharveyonline.com>, and myself.
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1339728 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker/src/main')
11 files changed, 116 insertions, 164 deletions
diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/configuration/VirtualHostConfiguration.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/configuration/VirtualHostConfiguration.java index 5f472b6ddd..59c6926b76 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/configuration/VirtualHostConfiguration.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/configuration/VirtualHostConfiguration.java @@ -32,7 +32,6 @@ import org.apache.qpid.server.configuration.plugins.ConfigurationPlugin; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.registry.ApplicationRegistry; import org.apache.qpid.server.store.MemoryMessageStore; -import org.apache.qpid.server.store.MemoryMessageStoreFactory; import java.util.ArrayList; import java.util.Arrays; @@ -103,14 +102,14 @@ public class VirtualHostConfiguration extends ConfigurationPlugin return getConfig().subset("store"); } - public String getMessageStoreFactoryClass() + public String getMessageStoreClass() { - return getStringValue("store.factoryclass", MemoryMessageStoreFactory.class.getName()); + return getStringValue("store.class", MemoryMessageStore.class.getName()); } - public void setMessageStoreFactoryClass(String storeFactoryClass) + public void setMessageStoreClass(String storeFactoryClass) { - getConfig().setProperty("store.factoryclass", storeFactoryClass); + getConfig().setProperty("store.class", storeFactoryClass); } public List getExchanges() diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/Event.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/Event.java index 9b5ceef35f..c681126c11 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/Event.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/Event.java @@ -23,12 +23,20 @@ public enum Event { BEFORE_INIT, AFTER_INIT, + BEFORE_ACTIVATE, AFTER_ACTIVATE, + BEFORE_PASSIVATE, AFTER_PASSIVATE, + BEFORE_CLOSE, AFTER_CLOSE, + + BEFORE_QUIESCE, + AFTER_QUIESCE, + BEFORE_RESTART, + PERSISTENT_MESSAGE_SIZE_OVERFULL, - PERSISTENT_MESSAGE_SIZE_UNDERFULL, + PERSISTENT_MESSAGE_SIZE_UNDERFULL } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/EventManager.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/EventManager.java index 21ae3924b8..bf3de2611d 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/EventManager.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/EventManager.java @@ -24,9 +24,12 @@ import java.util.EnumMap; import java.util.List; import java.util.Map; +import org.apache.log4j.Logger; + public class EventManager { private Map<Event, List<EventListener>> _listeners = new EnumMap<Event, List<EventListener>> (Event.class); + private static final Logger _LOGGER = Logger.getLogger(EventManager.class); public synchronized void addEventListener(EventListener listener, Event... events) { @@ -46,6 +49,11 @@ public class EventManager { if (_listeners.containsKey(event)) { + if(_LOGGER.isDebugEnabled()) + { + _LOGGER.debug("Received event " + event); + } + for (EventListener listener : _listeners.get(event)) { listener.event(event); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MessageStoreFactory.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/HAMessageStore.java index a35db62b03..59483751ca 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MessageStoreFactory.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/HAMessageStore.java @@ -19,9 +19,11 @@ */ package org.apache.qpid.server.store; -public interface MessageStoreFactory +public interface HAMessageStore extends MessageStore { - MessageStore createMessageStore(); - - String getStoreClassName(); + /** + * Used to indicate that a store requires to make itself unavailable for read and read/write + * operations. + */ + void passivate(); } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MemoryMessageStore.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MemoryMessageStore.java index 59624b7a75..7b98b30860 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MemoryMessageStore.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MemoryMessageStore.java @@ -83,19 +83,19 @@ public class MemoryMessageStore extends NullMessageStore @Override public void configureConfigStore(String name, ConfigurationRecoveryHandler recoveryHandler, Configuration config) throws Exception { - _stateManager.attainState(State.CONFIGURING); + _stateManager.attainState(State.INITIALISING); } @Override public void configureMessageStore(String name, MessageStoreRecoveryHandler recoveryHandler, TransactionLogRecoveryHandler tlogRecoveryHandler, Configuration config) throws Exception { - _stateManager.attainState(State.CONFIGURED); + _stateManager.attainState(State.INITIALISED); } @Override public void activate() throws Exception { - _stateManager.attainState(State.RECOVERING); + _stateManager.attainState(State.ACTIVATING); _stateManager.attainState(State.ACTIVE); } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MemoryMessageStoreFactory.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MemoryMessageStoreFactory.java deleted file mode 100644 index 8724f102c6..0000000000 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/MemoryMessageStoreFactory.java +++ /dev/null @@ -1,37 +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.store; - -public class MemoryMessageStoreFactory implements MessageStoreFactory -{ - - @Override - public MessageStore createMessageStore() - { - return new MemoryMessageStore(); - } - - @Override - public String getStoreClassName() - { - return MemoryMessageStore.class.getSimpleName(); - } - -} diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/State.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/State.java index 7cbdede85e..2783637b2a 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/State.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/State.java @@ -20,19 +20,30 @@ */ package org.apache.qpid.server.store; +import org.apache.qpid.server.configuration.ConfiguredObject; + public enum State { - + /** The initial state of the store. In practice, the store immediately transitions to the subsequent states. */ INITIAL, - CONFIGURING, - CONFIGURED, - RECOVERING, + + INITIALISING, + /** + * The initial set-up of the store has completed. + * If the store is persistent, it has not yet loaded configuration for {@link ConfiguredObject}'s from disk. + * + * From the point of view of the user, the store is essentially stopped. + */ + INITIALISED, + + ACTIVATING, ACTIVE, - QUIESCING, - QUIESCED, + CLOSING, - CLOSED; - + CLOSED, + QUIESCING, + /** The virtual host (and implicitly also the store) has been manually paused by the user to allow configuration changes to take place */ + QUIESCED; }
\ No newline at end of file diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/StateManager.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/StateManager.java index 5998be5bb6..613b329beb 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/StateManager.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/StateManager.java @@ -24,6 +24,8 @@ package org.apache.qpid.server.store; import java.util.EnumMap; import java.util.Map; +import org.apache.qpid.server.store.StateManager.Transition; + public class StateManager { private State _state = State.INITIAL; @@ -70,16 +72,23 @@ public class StateManager } - public static final Transition CONFIGURE = new Transition(State.INITIAL, State.CONFIGURING, Event.BEFORE_INIT); - public static final Transition CONFIGURE_COMPLETE = new Transition(State.CONFIGURING, State.CONFIGURED, Event.AFTER_INIT); - public static final Transition RECOVER = new Transition(State.CONFIGURED, State.RECOVERING, Event.BEFORE_ACTIVATE); - public static final Transition ACTIVATE = new Transition(State.RECOVERING, State.ACTIVE, Event.AFTER_ACTIVATE); + public static final Transition INITIALISE = new Transition(State.INITIAL, State.INITIALISING, Event.BEFORE_INIT); + public static final Transition INITALISE_COMPLETE = new Transition(State.INITIALISING, State.INITIALISED, Event.AFTER_INIT); + + public static final Transition ACTIVATE = new Transition(State.INITIALISED, State.ACTIVATING, Event.BEFORE_ACTIVATE); + public static final Transition ACTIVATE_COMPLETE = new Transition(State.ACTIVATING, State.ACTIVE, Event.AFTER_ACTIVATE); + + public static final Transition CLOSE_INITIALISED = new Transition(State.INITIALISED, State.CLOSING, Event.BEFORE_CLOSE);; public static final Transition CLOSE_ACTIVE = new Transition(State.ACTIVE, State.CLOSING, Event.BEFORE_CLOSE); public static final Transition CLOSE_QUIESCED = new Transition(State.QUIESCED, State.CLOSING, Event.BEFORE_CLOSE); public static final Transition CLOSE_COMPLETE = new Transition(State.CLOSING, State.CLOSED, Event.AFTER_CLOSE); - public static final Transition QUIESCE = new Transition(State.ACTIVE, State.QUIESCING, Event.BEFORE_PASSIVATE); - public static final Transition QUIESCE_COMPLETE = new Transition(State.QUIESCING, State.QUIESCED, Event.BEFORE_PASSIVATE); - public static final Transition RESTART = new Transition(State.QUIESCED, State.RECOVERING, Event.BEFORE_ACTIVATE); + + public static final Transition PASSIVATE = new Transition(State.ACTIVE, State.INITIALISED, Event.BEFORE_PASSIVATE); + + public static final Transition QUIESCE = new Transition(State.ACTIVE, State.QUIESCING, Event.BEFORE_QUIESCE); + public static final Transition QUIESCE_COMPLETE = new Transition(State.QUIESCING, State.QUIESCED, Event.AFTER_QUIESCE); + + public static final Transition RESTART = new Transition(State.QUIESCED, State.ACTIVATING, Event.BEFORE_RESTART); public StateManager(final EventManager eventManager) @@ -105,16 +114,6 @@ public class StateManager return _state; } - public synchronized void stateTransition(final State current, final State desired) - { - if (_state != current) - { - throw new IllegalStateException("Cannot transition to the state: " + desired + "; need to be in state: " + current - + "; currently in state: " + _state); - } - attainState(desired); - } - public synchronized void attainState(State desired) { Transition transition = null; diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/derby/DerbyMessageStore.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/derby/DerbyMessageStore.java index de1ce1a9db..c065eb263b 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/derby/DerbyMessageStore.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/derby/DerbyMessageStore.java @@ -260,7 +260,7 @@ public class DerbyMessageStore implements MessageStore ConfigurationRecoveryHandler configRecoveryHandler, Configuration storeConfiguration) throws Exception { - _stateManager.attainState(State.CONFIGURING); + _stateManager.attainState(State.INITIALISING); _configRecoveryHandler = configRecoveryHandler; commonConfiguration(name, storeConfiguration); @@ -276,13 +276,13 @@ public class DerbyMessageStore implements MessageStore _tlogRecoveryHandler = tlogRecoveryHandler; _messageRecoveryHandler = recoveryHandler; - _stateManager.attainState(State.CONFIGURED); + _stateManager.attainState(State.INITIALISED); } @Override public void activate() throws Exception { - _stateManager.attainState(State.RECOVERING); + _stateManager.attainState(State.ACTIVATING); // this recovers durable exchanges, queues, and bindings recoverConfiguration(_configRecoveryHandler); @@ -716,7 +716,7 @@ public class DerbyMessageStore implements MessageStore public void close() throws Exception { _closed.getAndSet(true); - _stateManager.stateTransition(State.ACTIVE, State.CLOSING); + _stateManager.attainState(State.CLOSING); try { @@ -737,7 +737,7 @@ public class DerbyMessageStore implements MessageStore } } - _stateManager.stateTransition(State.CLOSING, State.CLOSED); + _stateManager.attainState(State.CLOSED); } @Override diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/derby/DerbyMessageStoreFactory.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/derby/DerbyMessageStoreFactory.java deleted file mode 100644 index 12d7f64a8d..0000000000 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/derby/DerbyMessageStoreFactory.java +++ /dev/null @@ -1,40 +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.store.derby; - -import org.apache.qpid.server.store.MessageStore; -import org.apache.qpid.server.store.MessageStoreFactory; - -public class DerbyMessageStoreFactory implements MessageStoreFactory -{ - - @Override - public MessageStore createMessageStore() - { - return new DerbyMessageStore(); - } - - @Override - public String getStoreClassName() - { - return DerbyMessageStore.class.getSimpleName(); - } - -} diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostImpl.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostImpl.java index b05025467d..5a14092930 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostImpl.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostImpl.java @@ -59,8 +59,8 @@ import org.apache.qpid.server.security.SecurityManager; import org.apache.qpid.server.stats.StatisticsCounter; import org.apache.qpid.server.store.Event; import org.apache.qpid.server.store.EventListener; +import org.apache.qpid.server.store.HAMessageStore; import org.apache.qpid.server.store.MessageStore; -import org.apache.qpid.server.store.MessageStoreFactory; import org.apache.qpid.server.store.OperationalLoggingListener; import org.apache.qpid.server.txn.DtxRegistry; import org.apache.qpid.server.virtualhost.plugins.VirtualHostPlugin; @@ -173,7 +173,7 @@ public class VirtualHostImpl implements VirtualHost, IConnectionRegistry.Registr _brokerMBean = new AMQBrokerManagerMBean(_virtualHostMBean); - _messageStore = initialiseMessageStore(hostConfig.getMessageStoreFactoryClass()); + _messageStore = initialiseMessageStore(hostConfig.getMessageStoreClass()); configureMessageStore(hostConfig); @@ -329,20 +329,19 @@ public class VirtualHostImpl implements VirtualHost, IConnectionRegistry.Registr } - private MessageStore initialiseMessageStore(final String messageStoreFactoryClass) throws Exception + private MessageStore initialiseMessageStore(final String messageStoreClass) throws Exception { - final Class<?> clazz = Class.forName(messageStoreFactoryClass); + final Class<?> clazz = Class.forName(messageStoreClass); final Object o = clazz.newInstance(); - if (!(o instanceof MessageStoreFactory)) + if (!(o instanceof MessageStore)) { - throw new ClassCastException("Message store factory class must implement " + MessageStoreFactory.class + + throw new ClassCastException("Message store factory class must implement " + MessageStore.class + ". Class " + clazz + " does not."); } - final MessageStoreFactory messageStoreFactory = (MessageStoreFactory) o; - final MessageStore messageStore = messageStoreFactory.createMessageStore(); - final MessageStoreLogSubject storeLogSubject = new MessageStoreLogSubject(this, messageStoreFactory.getStoreClassName()); + final MessageStore messageStore = (MessageStore) o; + final MessageStoreLogSubject storeLogSubject = new MessageStoreLogSubject(this, clazz.getSimpleName()); OperationalLoggingListener.listen(messageStore, storeLogSubject); messageStore.addEventListener(new BeforeActivationListener(), Event.BEFORE_ACTIVATE); @@ -366,7 +365,10 @@ public class VirtualHostImpl implements VirtualHost, IConnectionRegistry.Registr private void activateNonHAMessageStore() throws Exception { - _messageStore.activate(); + if (!(_messageStore instanceof HAMessageStore)) + { + _messageStore.activate(); + } } private void initialiseModel(VirtualHostConfiguration config) throws ConfigurationException, AMQException @@ -801,42 +803,42 @@ public class VirtualHostImpl implements VirtualHost, IConnectionRegistry.Registr } private final class BeforeActivationListener implements EventListener - { - @Override - public void event(Event event) - { - try - { - _exchangeRegistry.initialise(); - initialiseModel(_vhostConfig); - } - catch (Exception e) - { - throw new RuntimeException("Failed to initialise virtual host after state change", e); - } - } - } - - private final class AfterActivationListener implements EventListener - { - @Override - public void event(Event event) - { - initialiseHouseKeeping(_vhostConfig.getHousekeepingCheckPeriod()); - try - { - _brokerMBean.register(); - } - catch (JMException e) - { - throw new RuntimeException("Failed to register virtual host mbean for virtual host " + getName(), e); - } + { + @Override + public void event(Event event) + { + try + { + _exchangeRegistry.initialise(); + initialiseModel(_vhostConfig); + } + catch (Exception e) + { + throw new RuntimeException("Failed to initialise virtual host after state change", e); + } + } + } + + private final class AfterActivationListener implements EventListener + { + @Override + public void event(Event event) + { + initialiseHouseKeeping(_vhostConfig.getHousekeepingCheckPeriod()); + try + { + _brokerMBean.register(); + } + catch (JMException e) + { + throw new RuntimeException("Failed to register virtual host mbean for virtual host " + getName(), e); + } - _state = State.ACTIVE; - } - } + _state = State.ACTIVE; + } + } - public class BeforePassivationListener implements EventListener + private final class BeforePassivationListener implements EventListener { public void event(Event event) { |
