mirror of https://github.com/apache/activemq.git
Fix for the failing QueueMasterSlaveTestUsingSharedFileTest on windows. The lock on the AsyncDataManager was not being retried.
git-svn-id: https://svn.apache.org/repos/asf/activemq/trunk@636037 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
parent
75d43c027f
commit
67d38fafd1
|
@ -113,11 +113,7 @@ public class AsyncDataManager {
|
|||
}
|
||||
|
||||
started = true;
|
||||
directory.mkdirs();
|
||||
synchronized (this) {
|
||||
controlFile = new ControlFile(new File(directory, filePrefix + "control"), CONTROL_RECORD_MAX_LENGTH);
|
||||
controlFile.lock();
|
||||
}
|
||||
lock();
|
||||
|
||||
ByteSequence sequence = controlFile.load();
|
||||
if (sequence != null && sequence.getLength() > 0) {
|
||||
|
@ -194,6 +190,16 @@ public class AsyncDataManager {
|
|||
Scheduler.executePeriodically(cleanupTask, 1000 * 30);
|
||||
}
|
||||
|
||||
public void lock() throws IOException {
|
||||
synchronized (this) {
|
||||
if( controlFile == null ) {
|
||||
directory.mkdirs();
|
||||
controlFile = new ControlFile(new File(directory, filePrefix + "control"), CONTROL_RECORD_MAX_LENGTH);
|
||||
}
|
||||
controlFile.lock();
|
||||
}
|
||||
}
|
||||
|
||||
protected Location recoveryCheck(DataFile dataFile, Location location) throws IOException {
|
||||
if (location == null) {
|
||||
location = new Location();
|
||||
|
|
|
@ -30,9 +30,9 @@ import java.util.concurrent.CopyOnWriteArraySet;
|
|||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
|
||||
import org.apache.activeio.journal.Journal;
|
||||
import org.apache.activeio.journal.active.JournalLockedException;
|
||||
import org.apache.activemq.broker.BrokerService;
|
||||
import org.apache.activemq.broker.BrokerServiceAware;
|
||||
import org.apache.activemq.broker.ConnectionContext;
|
||||
|
@ -90,6 +90,7 @@ public class AMQPersistenceAdapter implements PersistenceAdapter, UsageListener,
|
|||
private static final String PROPERTY_PREFIX = "org.apache.activemq.store.amq";
|
||||
private static final boolean BROKEN_FILE_LOCK;
|
||||
private static final boolean DISABLE_LOCKING;
|
||||
private static final int JOURNAL_LOCKED_WAIT_DELAY = 10 * 1000;
|
||||
|
||||
private AsyncDataManager asyncDataManager;
|
||||
private ReferenceStoreAdapter referenceStoreAdapter;
|
||||
|
@ -125,6 +126,7 @@ public class AMQPersistenceAdapter implements PersistenceAdapter, UsageListener,
|
|||
private RandomAccessFile lockFile;
|
||||
private FileLock lock;
|
||||
private boolean disableLocking = DISABLE_LOCKING;
|
||||
private boolean failIfJournalIsLocked;
|
||||
|
||||
public String getBrokerName() {
|
||||
return this.brokerName;
|
||||
|
@ -186,6 +188,24 @@ public class AMQPersistenceAdapter implements PersistenceAdapter, UsageListener,
|
|||
if (taskRunnerFactory == null) {
|
||||
taskRunnerFactory = createTaskRunnerFactory();
|
||||
}
|
||||
|
||||
if (failIfJournalIsLocked) {
|
||||
asyncDataManager.lock();
|
||||
} else {
|
||||
while (true) {
|
||||
try {
|
||||
asyncDataManager.lock();
|
||||
break;
|
||||
} catch (IOException e) {
|
||||
LOG.info("Journal is locked... waiting " + (JOURNAL_LOCKED_WAIT_DELAY / 1000) + " seconds for the journal to be unlocked.");
|
||||
try {
|
||||
Thread.sleep(JOURNAL_LOCKED_WAIT_DELAY);
|
||||
} catch (InterruptedException e1) {
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
asyncDataManager.start();
|
||||
if (deleteAllMessages) {
|
||||
asyncDataManager.delete();
|
||||
|
|
|
@ -17,6 +17,10 @@
|
|||
package org.apache.activemq.broker.ft;
|
||||
|
||||
import java.io.File;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import org.apache.activemq.ActiveMQConnectionFactory;
|
||||
import org.apache.activemq.JmsTopicSendReceiveWithTwoConnectionsTest;
|
||||
import org.apache.activemq.broker.BrokerService;
|
||||
|
@ -32,7 +36,8 @@ public class QueueMasterSlaveTest extends JmsTopicSendReceiveWithTwoConnectionsT
|
|||
private static final transient Log LOG = LogFactory.getLog(QueueMasterSlaveTest.class);
|
||||
|
||||
protected BrokerService master;
|
||||
protected BrokerService slave;
|
||||
protected AtomicReference<BrokerService> slave = new AtomicReference<BrokerService>();
|
||||
protected CountDownLatch slaveStarted = new CountDownLatch(1);
|
||||
protected int inflightMessageCount;
|
||||
protected int failureCount = 50;
|
||||
protected String uriString = "failover://(tcp://localhost:62001,tcp://localhost:62002)?randomize=false";
|
||||
|
@ -62,7 +67,12 @@ public class QueueMasterSlaveTest extends JmsTopicSendReceiveWithTwoConnectionsT
|
|||
|
||||
protected void tearDown() throws Exception {
|
||||
super.tearDown();
|
||||
slave.stop();
|
||||
|
||||
slaveStarted.await(5, TimeUnit.SECONDS);
|
||||
BrokerService brokerService = slave.get();
|
||||
if( brokerService!=null ) {
|
||||
brokerService.stop();
|
||||
}
|
||||
master.stop();
|
||||
}
|
||||
|
||||
|
@ -93,7 +103,9 @@ public class QueueMasterSlaveTest extends JmsTopicSendReceiveWithTwoConnectionsT
|
|||
protected void createSlave() throws Exception {
|
||||
BrokerFactoryBean brokerFactory = new BrokerFactoryBean(new ClassPathResource(getSlaveXml()));
|
||||
brokerFactory.afterPropertiesSet();
|
||||
slave = brokerFactory.getBroker();
|
||||
slave.start();
|
||||
BrokerService broker = brokerFactory.getBroker();
|
||||
broker.start();
|
||||
slave.set(broker);
|
||||
slaveStarted.countDown();
|
||||
}
|
||||
}
|
||||
|
|
|
@ -16,10 +16,6 @@
|
|||
*/
|
||||
package org.apache.activemq.broker.ft;
|
||||
|
||||
import org.apache.activemq.xbean.BrokerFactoryBean;
|
||||
import org.springframework.core.io.ClassPathResource;
|
||||
|
||||
|
||||
public class QueueMasterSlaveTestUsingSharedFileTest extends
|
||||
QueueMasterSlaveTest {
|
||||
|
||||
|
@ -32,16 +28,14 @@ public class QueueMasterSlaveTestUsingSharedFileTest extends
|
|||
}
|
||||
|
||||
protected void createSlave() throws Exception {
|
||||
// Start the Brokers async since starting them up could be a blocking operation..
|
||||
new Thread(new Runnable() {
|
||||
|
||||
public void run() {
|
||||
try {
|
||||
QueueMasterSlaveTestUsingSharedFileTest.super.createSlave();
|
||||
} catch (Exception e) {
|
||||
|
||||
e.printStackTrace();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}).start();
|
||||
|
|
Loading…
Reference in New Issue