summaryrefslogtreecommitdiff
path: root/qpid/java/broker-core/src
diff options
context:
space:
mode:
authorKeith Wall <kwall@apache.org = kwall = Keith Wall kwall@apache.org@apache.org>2014-04-14 08:54:19 +0000
committerKeith Wall <kwall@apache.org = kwall = Keith Wall kwall@apache.org@apache.org>2014-04-14 08:54:19 +0000
commitcde1072e86b57286594eb4fdb494576689aa8bca (patch)
tree8e0f378d16d5cf564f8ab0d2f93e5ec6f338621f /qpid/java/broker-core/src
parent981b8f5357355f842a523e4b50a1d5c711095a68 (diff)
downloadqpid-python-cde1072e86b57286594eb4fdb494576689aa8bca.tar.gz
QPID-5685: Store configuration version as an attribute of virtualhost within configuration store rather than within separate database/table
* ConfiguredObjectRecordHandler begin/end methods no longer take/return config version * DefaultUpgraderProvider uses the virtualhost record for the config version only and uses this to trigger the correct upgrade. Note this record is *not* recovered (yet). * BDB/SQL Upgraders migrate the config version from database/table to be the modelVersion attribute of a virtualhost entry. * BDB Upgrader tests (7 to 8). git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1587165 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker-core/src')
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/startup/BrokerStoreUpgrader.java7
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStore.java7
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandler.java13
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/MemoryConfigurationEntryStore.java9
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/model/VirtualHost.java2
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java119
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractMemoryMessageStore.java2
-rwxr-xr-xqpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfigurationRecoveryHandler.java6
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfiguredObjectRecordRecoveverAndUpgrader.java8
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/DurableConfigurationRecoverer.java59
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/JsonFileConfigStore.java9
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/UpgraderProvider.java2
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/handler/ConfiguredObjectRecordHandler.java7
-rw-r--r--qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java11
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java20
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/JsonFileConfigStoreTest.java25
-rw-r--r--qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/DurableConfigurationRecovererTest.java44
17 files changed, 161 insertions, 189 deletions
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/startup/BrokerStoreUpgrader.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/startup/BrokerStoreUpgrader.java
index e82a92bb83..c900ea047d 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/startup/BrokerStoreUpgrader.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/startup/BrokerStoreUpgrader.java
@@ -601,7 +601,6 @@ public class BrokerStoreUpgrader
private DurableConfigurationStoreUpgrader _upgrader;
private DurableConfigurationStore _store;
private final Map<UUID, ConfiguredObjectRecord> _records = new HashMap<UUID, ConfiguredObjectRecord>();
- private int _version;
private final SystemContext _systemContext;
private BrokerStoreRecoveryHandler(final SystemContext systemContext, DurableConfigurationStore store)
@@ -612,9 +611,8 @@ public class BrokerStoreUpgrader
@Override
- public void begin(final int configVersion)
+ public void begin()
{
- _version = configVersion;
}
@Override
@@ -625,7 +623,7 @@ public class BrokerStoreUpgrader
}
@Override
- public int end()
+ public void end()
{
String version = getCurrentVersion();
@@ -748,7 +746,6 @@ public class BrokerStoreUpgrader
});
- return _version;
}
private void applyRecursively(final ConfiguredObject<?> object, final Action<ConfiguredObject<?>> action)
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStore.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStore.java
index 59f248c9f5..d4409e666e 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStore.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/JsonConfigurationEntryStore.java
@@ -126,11 +126,9 @@ public class JsonConfigurationEntryStore extends MemoryConfigurationEntryStore
final Collection<ConfiguredObjectRecord> records = new ArrayList<ConfiguredObjectRecord>();
final ConfiguredObjectRecordHandler replayHandler = new ConfiguredObjectRecordHandler()
{
- private int _configVersion;
@Override
- public void begin(final int configVersion)
+ public void begin()
{
- _configVersion = configVersion;
}
@Override
@@ -141,9 +139,8 @@ public class JsonConfigurationEntryStore extends MemoryConfigurationEntryStore
}
@Override
- public int end()
+ public void end()
{
- return _configVersion;
}
};
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandler.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandler.java
index 88817be972..a17bbbff47 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandler.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandler.java
@@ -89,9 +89,8 @@ public class ManagementModeStoreHandler implements DurableConfigurationStore
private boolean _quiesceHttpPort = _options.getManagementModeHttpPortOverride() > 0;
@Override
- public void begin(final int configVersion)
+ public void begin()
{
- _version = configVersion;
}
@Override
@@ -160,9 +159,8 @@ public class ManagementModeStoreHandler implements DurableConfigurationStore
@Override
- public int end()
+ public void end()
{
- return _version;
}
};
@@ -186,7 +184,7 @@ public class ManagementModeStoreHandler implements DurableConfigurationStore
{
- recoveryHandler.begin(0);
+ recoveryHandler.begin();
for(ConfiguredObjectRecord record : _records.values())
{
@@ -366,7 +364,7 @@ public class ManagementModeStoreHandler implements DurableConfigurationStore
_store.visitConfiguredObjectRecords(new ConfiguredObjectRecordHandler()
{
@Override
- public void begin(final int configVersion)
+ public void begin()
{
}
@@ -428,9 +426,8 @@ public class ManagementModeStoreHandler implements DurableConfigurationStore
@Override
- public int end()
+ public void end()
{
- return 0;
}
});
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/MemoryConfigurationEntryStore.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/MemoryConfigurationEntryStore.java
index b3f91d7ee1..e6e4d0052b 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/MemoryConfigurationEntryStore.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/configuration/store/MemoryConfigurationEntryStore.java
@@ -124,11 +124,9 @@ public class MemoryConfigurationEntryStore implements ConfigurationEntryStore
final Collection<ConfiguredObjectRecord> records = new ArrayList<ConfiguredObjectRecord>();
final ConfiguredObjectRecordHandler replayHandler = new ConfiguredObjectRecordHandler()
{
- private int _configVersion;
@Override
- public void begin(final int configVersion)
+ public void begin()
{
- _configVersion = configVersion;
}
@Override
@@ -139,9 +137,8 @@ public class MemoryConfigurationEntryStore implements ConfigurationEntryStore
}
@Override
- public int end()
+ public void end()
{
- return _configVersion;
}
};
@@ -363,7 +360,7 @@ public class MemoryConfigurationEntryStore implements ConfigurationEntryStore
public void visitConfiguredObjectRecords(final ConfiguredObjectRecordHandler recoveryHandler) throws StoreException
{
- recoveryHandler.begin(0);
+ recoveryHandler.begin();
final Map<UUID,Map<String,UUID>> parentMap = new HashMap<UUID, Map<String, UUID>>();
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 33ccee8e71..d7b4868fb4 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
@@ -49,8 +49,6 @@ public interface VirtualHost<X extends VirtualHost<X, Q, E>, Q extends Queue<?>,
String CONFIGURATION_STORE_SETTINGS = "configurationStoreSettings";
String MESSAGE_STORE_SETTINGS = "messageStoreSettings";
- int CURRENT_CONFIG_VERSION = 5;
-
@ManagedAttribute
Collection<String> getSupportedExchangeTypes();
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java
index 6be5460d5f..6ff959fad8 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractJDBCMessageStore.java
@@ -48,6 +48,8 @@ import java.util.concurrent.atomic.AtomicLong;
import org.apache.log4j.Logger;
import org.apache.qpid.server.message.EnqueueableMessage;
import org.apache.qpid.server.model.ConfiguredObject;
+import org.apache.qpid.server.model.Model;
+import org.apache.qpid.server.model.UUIDGenerator;
import org.apache.qpid.server.plugin.MessageMetaDataType;
import org.apache.qpid.server.store.handler.ConfiguredObjectRecordHandler;
import org.apache.qpid.server.store.handler.DistributedTransactionHandler;
@@ -84,7 +86,7 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
private static final int DEFAULT_CONFIG_VERSION = 0;
- public static final Set<String> CONFIGURATION_STORE_TABLE_NAMES = new HashSet<String>(Arrays.asList(CONFIGURED_OBJECTS_TABLE_NAME, CONFIGURATION_VERSION_TABLE_NAME));
+ public static final Set<String> CONFIGURATION_STORE_TABLE_NAMES = new HashSet<String>(Arrays.asList(CONFIGURED_OBJECTS_TABLE_NAME, CONFIGURED_OBJECT_HIERARCHY_TABLE_NAME));
public static final Set<String> MESSAGE_STORE_TABLE_NAMES = new HashSet<String>(Arrays.asList(DB_VERSION_TABLE_NAME,
META_DATA_TABLE_NAME, MESSAGE_CONTENT_TABLE_NAME,
QUEUE_ENTRY_TABLE_NAME,
@@ -100,10 +102,8 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
private static final String UPDATE_DB_VERSION = "UPDATE " + DB_VERSION_TABLE_NAME + " SET version = ?";
- private static final String CREATE_CONFIG_VERSION_TABLE = "CREATE TABLE "+ CONFIGURATION_VERSION_TABLE_NAME + " ( version int not null )";
- private static final String INSERT_INTO_CONFIG_VERSION = "INSERT INTO "+ CONFIGURATION_VERSION_TABLE_NAME + " ( version ) VALUES ( ? )";
private static final String SELECT_FROM_CONFIG_VERSION = "SELECT version FROM " + CONFIGURATION_VERSION_TABLE_NAME;
- private static final String UPDATE_CONFIG_VERSION = "UPDATE " + CONFIGURATION_VERSION_TABLE_NAME + " SET version = ?";
+ private static final String DROP_CONFIG_VERSION_TABLE = "DROP TABLE "+ CONFIGURATION_VERSION_TABLE_NAME;
private static final String INSERT_INTO_QUEUE_ENTRY = "INSERT INTO " + QUEUE_ENTRY_TABLE_NAME + " (queue_id, message_id) values (?,?)";
@@ -230,16 +230,9 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
try
{
- int configVersion = getConfigVersion();
-
- handler.begin(configVersion);
+ handler.begin();
doVisitAllConfiguredObjectRecords(handler);
-
- int newConfigVersion = handler.end();
- if(newConfigVersion != configVersion)
- {
- setConfigVersion(newConfigVersion);
- }
+ handler.end();
}
catch (SQLException e)
{
@@ -470,6 +463,32 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
Connection connection = newConnection();
try
{
+ UUID virtualHostId = UUIDGenerator.generateVhostUUID(parent.getName());
+
+ String stringifiedConfigVersion = Model.MODEL_VERSION;
+ boolean tableExists = tableExists(CONFIGURATION_VERSION_TABLE_NAME, connection);
+ if(tableExists)
+ {
+ int configVersion = getConfigVersion(connection);
+ if (getLogger().isDebugEnabled())
+ {
+ getLogger().debug("Upgrader read existing config version " + configVersion);
+ }
+
+ stringifiedConfigVersion = "0." + configVersion;
+ }
+
+ Map<String, Object> virtualHostAttributes = new HashMap<String, Object>();
+ virtualHostAttributes.put("modelVersion", stringifiedConfigVersion);
+
+ ConfiguredObjectRecord configuredObject = new ConfiguredObjectRecordImpl(virtualHostId, "VirtualHost", virtualHostAttributes);
+ insertConfiguredObject(configuredObject, connection);
+
+ if (getLogger().isDebugEnabled())
+ {
+ getLogger().debug("Upgrader created VirtualHost configuration entry with config version " + stringifiedConfigVersion);
+ }
+
Map<UUID,Map<String,Object>> bindingsToUpdate = new HashMap<UUID, Map<String, Object>>();
List<UUID> others = new ArrayList<UUID>();
final ObjectMapper objectMapper = new ObjectMapper();
@@ -525,7 +544,7 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
{
stmt.setString(1, id.toString());
stmt.setString(2, "VirtualHost");
- stmt.setString(3, parent.getId().toString());
+ stmt.setString(3, virtualHostId.toString());
stmt.execute();
}
for(Map.Entry<UUID, Map<String,Object>> bindingEntry : bindingsToUpdate.entrySet())
@@ -586,6 +605,11 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
stmt.close();
}
connection.commit();
+
+ if (tableExists)
+ {
+ dropConfigVersionTable(connection);
+ }
}
finally
{
@@ -643,7 +667,6 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
{
Connection conn = newAutoCommitConnection();
- createConfigVersionTable(conn);
createConfiguredObjectsTable(conn);
createConfiguredObjectHierarchyTable(conn);
@@ -677,30 +700,19 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
}
}
- private void createConfigVersionTable(final Connection conn) throws SQLException
+ private void dropConfigVersionTable(final Connection conn) throws SQLException
{
if(!tableExists(CONFIGURATION_VERSION_TABLE_NAME, conn))
{
Statement stmt = conn.createStatement();
try
{
- stmt.execute(CREATE_CONFIG_VERSION_TABLE);
+ stmt.execute(DROP_CONFIG_VERSION_TABLE);
}
finally
{
stmt.close();
}
-
- PreparedStatement pstmt = conn.prepareStatement(INSERT_INTO_CONFIG_VERSION);
- try
- {
- pstmt.setInt(1, DEFAULT_CONFIG_VERSION);
- pstmt.execute();
- }
- finally
- {
- pstmt.close();
- }
}
}
@@ -872,63 +884,30 @@ abstract public class AbstractJDBCMessageStore implements MessageStore, DurableC
}
}
- private void setConfigVersion(int version) throws SQLException
- {
- Connection conn = newAutoCommitConnection();
- try
- {
-
- PreparedStatement stmt = conn.prepareStatement(UPDATE_CONFIG_VERSION);
- try
- {
- stmt.setInt(1, version);
- stmt.execute();
-
- }
- finally
- {
- stmt.close();
- }
- }
- finally
- {
- conn.close();
- }
- }
-
- private int getConfigVersion() throws SQLException
+ private int getConfigVersion(Connection conn) throws SQLException
{
- Connection conn = newAutoCommitConnection();
+ Statement stmt = conn.createStatement();
try
{
-
- Statement stmt = conn.createStatement();
+ ResultSet rs = stmt.executeQuery(SELECT_FROM_CONFIG_VERSION);
try
{
- ResultSet rs = stmt.executeQuery(SELECT_FROM_CONFIG_VERSION);
- try
- {
- if(rs.next())
- {
- return rs.getInt(1);
- }
- return DEFAULT_CONFIG_VERSION;
- }
- finally
+ if(rs.next())
{
- rs.close();
+ return rs.getInt(1);
}
-
+ return DEFAULT_CONFIG_VERSION;
}
finally
{
- stmt.close();
+ rs.close();
}
+
}
finally
{
- conn.close();
+ stmt.close();
}
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractMemoryMessageStore.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractMemoryMessageStore.java
index 99785c48a9..63413fc0de 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractMemoryMessageStore.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/AbstractMemoryMessageStore.java
@@ -261,7 +261,7 @@ abstract class AbstractMemoryMessageStore implements MessageStore, DurableConfig
@Override
public void visitConfiguredObjectRecords(ConfiguredObjectRecordHandler handler) throws StoreException
{
- handler.begin(VirtualHost.CURRENT_CONFIG_VERSION);
+ handler.begin();
for (ConfiguredObjectRecord record : _configuredObjectRecords.values())
{
if (!handler.handle(record))
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfigurationRecoveryHandler.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfigurationRecoveryHandler.java
index c8aef92a95..81d3f7f9da 100755
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfigurationRecoveryHandler.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfigurationRecoveryHandler.java
@@ -20,12 +20,10 @@
*/
package org.apache.qpid.server.store;
-import java.util.Map;
-import java.util.UUID;
public interface ConfigurationRecoveryHandler
{
- void beginConfigurationRecovery(DurableConfigurationStore store, int configVersion);
+ void beginConfigurationRecovery(DurableConfigurationStore store);
void configuredObject(ConfiguredObjectRecord object);
@@ -33,6 +31,6 @@ public interface ConfigurationRecoveryHandler
*
* @return the model version of the configuration
*/
- int completeConfigurationRecovery();
+ String completeConfigurationRecovery();
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfiguredObjectRecordRecoveverAndUpgrader.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfiguredObjectRecordRecoveverAndUpgrader.java
index 2cadb8ac4c..e9e6d88f36 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfiguredObjectRecordRecoveverAndUpgrader.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/ConfiguredObjectRecordRecoveverAndUpgrader.java
@@ -39,9 +39,9 @@ public class ConfiguredObjectRecordRecoveverAndUpgrader implements ConfiguredObj
}
@Override
- public void begin(int configVersion)
+ public void begin()
{
- _configRecoverer.beginConfigurationRecovery(_store, configVersion);
+ _configRecoverer.beginConfigurationRecovery(_store);
}
@Override
@@ -52,9 +52,9 @@ public class ConfiguredObjectRecordRecoveverAndUpgrader implements ConfiguredObj
}
@Override
- public int end()
+ public void end()
{
- return _configRecoverer.completeConfigurationRecovery();
+ _configRecoverer.completeConfigurationRecovery();
}
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/DurableConfigurationRecoverer.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/DurableConfigurationRecoverer.java
index 5975bf58b3..f8a8741a37 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/DurableConfigurationRecoverer.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/DurableConfigurationRecoverer.java
@@ -33,8 +33,7 @@ import org.apache.qpid.server.configuration.IllegalConfigurationException;
import org.apache.qpid.server.logging.EventLogger;
import org.apache.qpid.server.logging.messages.ConfigStoreMessages;
import org.apache.qpid.server.logging.subjects.MessageStoreLogSubject;
-
-import static org.apache.qpid.server.model.VirtualHost.CURRENT_CONFIG_VERSION;
+import org.apache.qpid.server.model.Model;
public class DurableConfigurationRecoverer implements ConfigurationRecoveryHandler
{
@@ -44,6 +43,7 @@ public class DurableConfigurationRecoverer implements ConfigurationRecoveryHandl
private final Map<String, Map<UUID, UnresolvedObject>> _unresolvedObjects =
new HashMap<String, Map<UUID, UnresolvedObject>>();
+ private final List<ConfiguredObjectRecord> _records = new ArrayList<ConfiguredObjectRecord>();
private final Map<String, Map<UUID, List<DependencyListener>>> _dependencyListeners =
new HashMap<String, Map<UUID, List<DependencyListener>>>();
@@ -69,19 +69,56 @@ public class DurableConfigurationRecoverer implements ConfigurationRecoveryHandl
}
@Override
- public void beginConfigurationRecovery(final DurableConfigurationStore store, final int configVersion)
+ public void beginConfigurationRecovery(final DurableConfigurationStore store)
{
_logSubject = new MessageStoreLogSubject(_name, store.getClass().getSimpleName());
_store = store;
- _upgrader = _upgraderProvider.getUpgrader(configVersion, this);
_eventLogger.message(_logSubject, ConfigStoreMessages.RECOVERY_START());
}
@Override
public void configuredObject(ConfiguredObjectRecord record)
{
- _upgrader.configuredObject(record);
+ _records.add(record);
+ }
+
+ @Override
+ public String completeConfigurationRecovery()
+ {
+ String configVersion = getConfigVersionFromRecords();
+
+ _upgrader = _upgraderProvider.getUpgrader(configVersion, this);
+
+ for (ConfiguredObjectRecord record : _records)
+ {
+ // We don't yet recover the VirtualHost record.
+ if (!"VirtualHost".equals(record.getType()))
+ {
+ _upgrader.configuredObject(record);
+ }
+ }
+ _upgrader.complete();
+ checkUnresolvedDependencies();
+ applyUpgrade();
+
+ _eventLogger.message(_logSubject, ConfigStoreMessages.RECOVERY_COMPLETE());
+ return Model.MODEL_VERSION;
+ }
+
+ private String getConfigVersionFromRecords()
+ {
+ String configVersion = Model.MODEL_VERSION;
+ for (ConfiguredObjectRecord record : _records)
+ {
+ if ("VirtualHost".equals(record.getType()))
+ {
+ configVersion = (String) record.getAttributes().get("modelVersion");
+ _logger.debug("Confifuration has config version : " + configVersion);
+ break;
+ }
+ }
+ return configVersion;
}
void onConfiguredObject(ConfiguredObjectRecord record)
@@ -94,23 +131,13 @@ public class DurableConfigurationRecoverer implements ConfigurationRecoveryHandl
recoverer.load(this, record);
}
+
private DurableConfiguredObjectRecoverer getRecoverer(final String type)
{
DurableConfiguredObjectRecoverer recoverer = _recoverers.get(type);
return recoverer;
}
- @Override
- public int completeConfigurationRecovery()
- {
- _upgrader.complete();
- checkUnresolvedDependencies();
- applyUpgrade();
-
- _eventLogger.message(_logSubject, ConfigStoreMessages.RECOVERY_COMPLETE());
- return CURRENT_CONFIG_VERSION;
- }
-
private void applyUpgrade()
{
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/JsonFileConfigStore.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/JsonFileConfigStore.java
index a5ace16cfa..1cf1d1a36d 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/JsonFileConfigStore.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/JsonFileConfigStore.java
@@ -100,7 +100,7 @@ public class JsonFileConfigStore implements DurableConfigurationStore
@Override
public void visitConfiguredObjectRecords(ConfiguredObjectRecordHandler handler)
{
- handler.begin(_configVersion);
+ handler.begin();
List<ConfiguredObjectRecord> records = new ArrayList<ConfiguredObjectRecord>(_objectsById.values());
for(ConfiguredObjectRecord record : records)
{
@@ -110,12 +110,7 @@ public class JsonFileConfigStore implements DurableConfigurationStore
break;
}
}
- int oldConfigVersion = _configVersion;
- _configVersion = handler.end();
- if(oldConfigVersion != _configVersion)
- {
- save();
- }
+ handler.end();
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/UpgraderProvider.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/UpgraderProvider.java
index c2ea0745ff..9785be78a6 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/UpgraderProvider.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/UpgraderProvider.java
@@ -22,5 +22,5 @@ package org.apache.qpid.server.store;
public interface UpgraderProvider
{
- DurableConfigurationStoreUpgrader getUpgrader(int configVersion, DurableConfigurationRecoverer recoverer);
+ DurableConfigurationStoreUpgrader getUpgrader(String configVersion, DurableConfigurationRecoverer recoverer);
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/handler/ConfiguredObjectRecordHandler.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/handler/ConfiguredObjectRecordHandler.java
index 747c735ff1..f251a442c7 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/handler/ConfiguredObjectRecordHandler.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/store/handler/ConfiguredObjectRecordHandler.java
@@ -24,8 +24,7 @@ import org.apache.qpid.server.store.ConfiguredObjectRecord;
public interface ConfiguredObjectRecordHandler
{
- // TODO configVersion argument will be removed.
- void begin(int configVersion);
+ void begin();
/**
* Handles the given record.
@@ -35,7 +34,5 @@ public interface ConfiguredObjectRecordHandler
*/
boolean handle(ConfiguredObjectRecord record);
- //TODO: return should be void
- // temporarily returning new config version
- int end();
+ void end();
}
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java
index 8c3ebaf9be..fe2f867e41 100644
--- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java
+++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java
@@ -20,8 +20,6 @@
*/
package org.apache.qpid.server.virtualhost;
-import static org.apache.qpid.server.model.VirtualHost.CURRENT_CONFIG_VERSION;
-
import java.util.Collections;
import java.util.HashMap;
import java.util.LinkedHashMap;
@@ -34,6 +32,7 @@ import org.apache.qpid.server.exchange.TopicExchange;
import org.apache.qpid.server.filter.FilterSupport;
import org.apache.qpid.server.model.Binding;
import org.apache.qpid.server.model.Exchange;
+import org.apache.qpid.server.model.Model;
import org.apache.qpid.server.model.Queue;
import org.apache.qpid.server.model.UUIDGenerator;
import org.apache.qpid.server.queue.QueueArgumentsConverter;
@@ -76,14 +75,16 @@ public class DefaultUpgraderProvider implements UpgraderProvider
_defaultExchangeIds = Collections.unmodifiableMap(defaultExchangeIds);
}
- public DurableConfigurationStoreUpgrader getUpgrader(final int configVersion, DurableConfigurationRecoverer recoverer)
+ public DurableConfigurationStoreUpgrader getUpgrader(final String configVersion, DurableConfigurationRecoverer recoverer)
{
if (LOGGER.isDebugEnabled())
{
LOGGER.debug("Getting upgrader for configVersion: " + configVersion);
}
DurableConfigurationStoreUpgrader currentUpgrader = null;
- switch(configVersion)
+
+ int conigVersionAsInteger = Integer.parseInt(configVersion.replace(".", ""));
+ switch(conigVersionAsInteger)
{
case 0:
currentUpgrader = addUpgrader(currentUpgrader, new Version0Upgrader());
@@ -95,7 +96,7 @@ public class DefaultUpgraderProvider implements UpgraderProvider
currentUpgrader = addUpgrader(currentUpgrader, new Version3Upgrader());
case 4:
currentUpgrader = addUpgrader(currentUpgrader, new Version4Upgrader());
- case CURRENT_CONFIG_VERSION:
+ case (Model.MODEL_MAJOR_VERSION * 10) + Model.MODEL_MINOR_VERSION:
currentUpgrader = addUpgrader(currentUpgrader, new NullUpgrader(recoverer));
break;
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java
index 0fe116570a..7859a4110b 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/configuration/store/ManagementModeStoreHandlerTest.java
@@ -471,9 +471,8 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
private int _version;
@Override
- public void begin(final int configVersion)
+ public void begin()
{
- _version = configVersion;
}
@Override
@@ -488,9 +487,8 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
}
@Override
- public int end()
+ public void end()
{
- return _version;
}
public ConfiguredObjectRecord getBrokerRecord()
@@ -503,7 +501,6 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
{
private final UUID _id;
private ConfiguredObjectRecord _foundRecord;
- private int _version;
private RecordFinder(final UUID id)
{
@@ -511,9 +508,8 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
}
@Override
- public void begin(final int configVersion)
+ public void begin()
{
- _version = configVersion;
}
@Override
@@ -528,9 +524,8 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
}
@Override
- public int end()
+ public void end()
{
- return _version;
}
public ConfiguredObjectRecord getFoundRecord()
@@ -543,7 +538,6 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
{
private final Collection<UUID> _childIds = new HashSet<UUID>();
private final ConfiguredObjectRecord _parent;
- private int _version;
private ChildFinder(final ConfiguredObjectRecord parent)
{
@@ -551,9 +545,8 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
}
@Override
- public void begin(final int configVersion)
+ public void begin()
{
- _version = configVersion;
}
@Override
@@ -575,9 +568,8 @@ public class ManagementModeStoreHandlerTest extends QpidTestCase
}
@Override
- public int end()
+ public void end()
{
- return _version;
}
public Collection<UUID> getChildIds()
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/JsonFileConfigStoreTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/JsonFileConfigStoreTest.java
index 6907898a6c..cdaab22fed 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/JsonFileConfigStoreTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/store/JsonFileConfigStoreTest.java
@@ -121,35 +121,14 @@ public class JsonFileConfigStoreTest extends QpidTestCase
{
_store.openConfigurationStore(_virtualHost, _configurationStoreSettings);
_store.visitConfiguredObjectRecords(_handler);
+
InOrder inorder = inOrder(_handler);
- inorder.verify(_handler).begin(eq(0));
+ inorder.verify(_handler).begin();
inorder.verify(_handler,never()).handle(any(ConfiguredObjectRecord.class));
inorder.verify(_handler).end();
_store.closeConfigurationStore();
}
- public void testUpdatedConfigVersionIsRetained() throws Exception
- {
- final int NEW_CONFIG_VERSION = 42;
- when(_handler.end()).thenReturn(NEW_CONFIG_VERSION);
-
- _store.openConfigurationStore(_virtualHost, _configurationStoreSettings);
- _store.visitConfiguredObjectRecords(_handler);
- _store.closeConfigurationStore();
-
- _store.openConfigurationStore(_virtualHost, _configurationStoreSettings);
- _store.visitConfiguredObjectRecords(_handler);
- InOrder inorder = inOrder(_handler);
-
- // first time the config version should be the initial version - 0
- inorder.verify(_handler).begin(eq(0));
-
- // second time the config version should be the updated version
- inorder.verify(_handler).begin(eq(NEW_CONFIG_VERSION));
-
- _store.closeConfigurationStore();
- }
-
public void testCreateObject() throws Exception
{
_store.openConfigurationStore(_virtualHost, _configurationStoreSettings);
diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/DurableConfigurationRecovererTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/DurableConfigurationRecovererTest.java
index 5d5856bf4a..e6b57d8039 100644
--- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/DurableConfigurationRecovererTest.java
+++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/virtualhost/DurableConfigurationRecovererTest.java
@@ -41,8 +41,10 @@ import org.apache.qpid.server.exchange.FanoutExchange;
import org.apache.qpid.server.exchange.HeadersExchange;
import org.apache.qpid.server.exchange.TopicExchange;
import org.apache.qpid.server.model.Binding;
+import org.apache.qpid.server.model.Model;
import org.apache.qpid.server.model.Queue;
import org.apache.qpid.server.model.UUIDGenerator;
+import org.apache.qpid.server.model.VirtualHost;
import org.apache.qpid.server.plugin.ExchangeType;
import org.apache.qpid.server.queue.AMQQueue;
import org.apache.qpid.server.queue.QueueFactory;
@@ -63,11 +65,11 @@ import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
-import static org.apache.qpid.server.model.VirtualHost.CURRENT_CONFIG_VERSION;
public class DurableConfigurationRecovererTest extends QpidTestCase
{
private static final String VIRTUAL_HOST_NAME = "test";
+ private static final UUID VIRTUAL_HOST_ID = UUID.randomUUID();
private static final UUID QUEUE_ID = new UUID(0,0);
private static final UUID TOPIC_EXCHANGE_ID = UUIDGenerator.generateExchangeUUID(TopicExchange.TYPE.getDefaultExchangeName(), VIRTUAL_HOST_NAME);
private static final UUID DIRECT_EXCHANGE_ID = UUIDGenerator.generateExchangeUUID(DirectExchange.TYPE.getDefaultExchangeName(), VIRTUAL_HOST_NAME);
@@ -205,19 +207,22 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
public void testUpgradeEmptyStore() throws Exception
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 0);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
assertEquals("Did not upgrade to the expected version",
- CURRENT_CONFIG_VERSION,
+ Model.MODEL_VERSION,
_durableConfigurationRecoverer.completeConfigurationRecovery());
}
public void testUpgradeNewerStoreFails() throws Exception
{
+ String bumpedModelVersion = Model.MODEL_MAJOR_VERSION + "." + (Model.MODEL_MINOR_VERSION + 1);
try
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, CURRENT_CONFIG_VERSION + 1);
- _durableConfigurationRecoverer.completeConfigurationRecovery();
- fail("Should not be able to start when config model is newer than current");
+
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
+ _durableConfigurationRecoverer.configuredObject(getVirtualHostModelRecord(bumpedModelVersion));
+ String newVersion = _durableConfigurationRecoverer.completeConfigurationRecovery();
+ fail("Should not be able to start when config model is newer than current. Actually upgraded to " + newVersion);
}
catch (IllegalStateException e)
{
@@ -225,10 +230,20 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
}
}
+ private ConfiguredObjectRecordImpl getVirtualHostModelRecord(
+ String modelVersion)
+ {
+ ConfiguredObjectRecordImpl virtualHostRecord = new ConfiguredObjectRecordImpl(VIRTUAL_HOST_ID,
+ VirtualHost.class.getSimpleName(),
+ Collections.<String,Object>singletonMap("modelVersion", modelVersion));
+ return virtualHostRecord;
+ }
+
public void testUpgradeRemovesBindingsToNonTopicExchanges() throws Exception
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 0);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
+ _durableConfigurationRecoverer.configuredObject(getVirtualHostModelRecord("0.0"));
_durableConfigurationRecoverer.configuredObject(new ConfiguredObjectRecordImpl(new UUID(1, 0),
"org.apache.qpid.server.model.Binding",
@@ -252,7 +267,8 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
public void testUpgradeOnlyRemovesSelectorBindings() throws Exception
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 0);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
+ _durableConfigurationRecoverer.configuredObject(getVirtualHostModelRecord("0.0"));
_durableConfigurationRecoverer.configuredObject(new ConfiguredObjectRecordImpl(new UUID(1, 0),
"org.apache.qpid.server.model.Binding",
@@ -336,7 +352,8 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
public void testUpgradeKeepsBindingsToTopicExchanges() throws Exception
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 0);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
+ _durableConfigurationRecoverer.configuredObject(getVirtualHostModelRecord("0.0"));
_durableConfigurationRecoverer.configuredObject(new ConfiguredObjectRecordImpl(new UUID(1, 0),
"org.apache.qpid.server.model.Binding",
@@ -358,7 +375,8 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
public void testUpgradeDoesNotRecur() throws Exception
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 2);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
+ _durableConfigurationRecoverer.configuredObject(getVirtualHostModelRecord("0.0"));
_durableConfigurationRecoverer.configuredObject(new ConfiguredObjectRecordImpl(new UUID(1, 0),
"Binding",
@@ -375,7 +393,7 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
public void testFailsWithUnresolvedObjects()
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 2);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
_durableConfigurationRecoverer.configuredObject(new ConfiguredObjectRecordImpl(new UUID(1, 0),
@@ -400,7 +418,7 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
public void testFailsWithUnknownObjectType()
{
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 2);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
try
@@ -468,7 +486,7 @@ public class DurableConfigurationRecovererTest extends QpidTestCase
}
});
- _durableConfigurationRecoverer.beginConfigurationRecovery(_store, 2);
+ _durableConfigurationRecoverer.beginConfigurationRecovery(_store);
_durableConfigurationRecoverer.configuredObject(new ConfiguredObjectRecordImpl(queueId, Queue.class.getSimpleName(),
createQueue("testQueue", exchangeId)));