[infinispan-commits] Infinispan SVN: r592 - in trunk: core/src/main/java/org/infinispan/commands and 22 other directories.
infinispan-commits at lists.jboss.org
infinispan-commits at lists.jboss.org
Mon Jul 20 11:35:57 EDT 2009
Author: mircea.markus
Date: 2009-07-20 11:35:56 -0400 (Mon, 20 Jul 2009)
New Revision: 592
Added:
trunk/core/src/main/java/org/infinispan/interceptors/DeadlockDetectingInterceptor.java
trunk/core/src/main/java/org/infinispan/marshall/exts/DeadlockDetectingGlobalTransactionExternalizer.java
trunk/core/src/main/java/org/infinispan/transaction/xa/DeadlockDetectingGlobalTransaction.java
trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransactionFactory.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectedException.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectingLockManager.java
trunk/core/src/test/java/org/infinispan/profiling/DeadlockDetectionPerformanceTest.java
trunk/core/src/test/java/org/infinispan/tx/DeadlockDetectionTest.java
Modified:
trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreIntegrationTest.java
trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreTest.java
trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java
trunk/core/src/main/java/org/infinispan/config/Configuration.java
trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java
trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorFactory.java
trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorNamedCacheFactory.java
trunk/core/src/main/java/org/infinispan/factories/EntryFactoryImpl.java
trunk/core/src/main/java/org/infinispan/factories/InterceptorChainFactory.java
trunk/core/src/main/java/org/infinispan/factories/LockManagerFactory.java
trunk/core/src/main/java/org/infinispan/interceptors/CacheMgmtInterceptor.java
trunk/core/src/main/java/org/infinispan/interceptors/DistributionInterceptor.java
trunk/core/src/main/java/org/infinispan/interceptors/LockingInterceptor.java
trunk/core/src/main/java/org/infinispan/interceptors/ReplicationInterceptor.java
trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java
trunk/core/src/main/java/org/infinispan/marshall/MarshallerImpl.java
trunk/core/src/main/java/org/infinispan/marshall/exts/GlobalTransactionExternalizer.java
trunk/core/src/main/java/org/infinispan/marshall/jboss/ConstantObjectTable.java
trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManager.java
trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java
trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java
trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransaction.java
trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java
trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManager.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManagerImpl.java
trunk/core/src/main/resources/config-samples/all.xml
trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd
trunk/core/src/test/java/org/infinispan/config/parsing/ConfigurationParserTest.java
trunk/core/src/test/java/org/infinispan/config/parsing/XmlFileParsingTest.java
trunk/core/src/test/java/org/infinispan/distribution/DefaultConsistentHashTest.java
trunk/core/src/test/java/org/infinispan/loaders/BaseCacheStoreTest.java
trunk/core/src/test/java/org/infinispan/loaders/decorators/ChainingCacheLoaderTest.java
trunk/core/src/test/java/org/infinispan/marshall/MarshallersTest.java
trunk/core/src/test/java/org/infinispan/marshall/jboss/JBossMarshallerTest.java
trunk/core/src/test/java/org/infinispan/tx/LocalModeTxTest.java
trunk/core/src/test/resources/configs/named-cache-test.xml
Log:
[ISPN-38] - (eager deadlock detection) - implementation
Modified: trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreIntegrationTest.java
===================================================================
--- trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreIntegrationTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreIntegrationTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -10,6 +10,7 @@
import org.infinispan.loaders.modifications.Store;
import org.infinispan.test.TestingUtil;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import org.testng.annotations.AfterTest;
import org.testng.annotations.BeforeTest;
import org.testng.annotations.Parameters;
@@ -28,6 +29,7 @@
public class BdbjeCacheStoreIntegrationTest extends BaseCacheStoreTest {
private String tmpDirectory;
+ private GlobalTransactionFactory gts = new GlobalTransactionFactory();
@BeforeTest
@Parameters({"basedir"})
@@ -61,7 +63,7 @@
mods.add(new Store(InternalEntryFactory.create("k1", "v1")));
mods.add(new Store(InternalEntryFactory.create("k2", "v2")));
mods.add(new Remove("k1"));
- GlobalTransaction tx = new GlobalTransaction(false);
+ GlobalTransaction tx = gts.newGlobalTransaction(null, false);
cs.prepare(mods, tx, false);
cs.commit(tx);
@@ -97,7 +99,7 @@
mods.add(new Store(InternalEntryFactory.create("k2", "v2")));
mods.add(new Remove("k1"));
mods.add(new Remove("old"));
- GlobalTransaction tx = new GlobalTransaction(false);
+ GlobalTransaction tx = gts.newGlobalTransaction(null, false);
cs.prepare(mods, tx, false);
cs.rollback(tx);
Modified: trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreTest.java
===================================================================
--- trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/cachestore/bdbje/src/test/java/org/infinispan/loaders/bdbje/BdbjeCacheStoreTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -19,6 +19,7 @@
import org.infinispan.marshall.Marshaller;
import org.infinispan.marshall.TestObjectStreamMarshaller;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import org.infinispan.util.ReflectionUtil;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
@@ -52,8 +53,9 @@
private PreparableTransactionRunner runner;
private CurrentTransaction currentTransaction;
+ private GlobalTransactionFactory gtf;
- private class MockBdbjeResourceFactory extends BdbjeResourceFactory {
+ private class MockBdbjeResourceFactory extends BdbjeResourceFactory {
@Override
public PreparableTransactionRunner createPreparableTransactionRunner(Environment env) {
@@ -272,7 +274,8 @@
@Test
public void testNoExceptionOnRollback() throws Exception {
start();
- GlobalTransaction tx = new GlobalTransaction(false);
+ gtf = new GlobalTransactionFactory();
+ GlobalTransaction tx = gtf.newGlobalTransaction(null, false);
replayAll();
cs.start();
cs.rollback(tx);
@@ -313,7 +316,7 @@
cs.start();
try {
txn = currentTransaction.beginTransaction(null);
- GlobalTransaction t = new GlobalTransaction(false);
+ GlobalTransaction t = gtf.newGlobalTransaction(null, false);
cs.prepare(Collections.singletonList(new Store(InternalEntryFactory.create("k", "v"))), t, false);
cs.commit(t);
assert false : "should have gotten an exception";
@@ -335,7 +338,7 @@
replayAll();
cs.start();
try {
- GlobalTransaction tx = new GlobalTransaction(false);
+ GlobalTransaction tx = gtf.newGlobalTransaction(null, false);
cs.prepare(Collections.singletonList(new Store(InternalEntryFactory.create("k", "v"))), tx, false);
assert false : "should have gotten an exception";
} catch (CacheLoaderException e) {
Modified: trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -49,7 +49,6 @@
import org.infinispan.factories.annotations.Inject;
import org.infinispan.factories.annotations.Start;
import org.infinispan.interceptors.InterceptorChain;
-import org.infinispan.loaders.CacheLoaderManager;
import org.infinispan.notifications.cachelistener.CacheNotifier;
import org.infinispan.transaction.xa.GlobalTransaction;
import org.infinispan.transaction.xa.TransactionTable;
@@ -67,7 +66,6 @@
private DataContainer dataContainer;
private CacheNotifier notifier;
private Cache cache;
- private CacheLoaderManager cacheLoaderManager;
private String cacheName;
// some stateless commands can be reused so that they aren't constructed again all the time.
@@ -82,14 +80,12 @@
@Inject
public void setupDependencies(DataContainer container, CacheNotifier notifier, Cache cache,
- InterceptorChain interceptorChain, CacheLoaderManager clManager,
- DistributionManager distributionManager, InvocationContextContainer icc,
- TransactionTable txTable) {
+ InterceptorChain interceptorChain, DistributionManager distributionManager,
+ InvocationContextContainer icc, TransactionTable txTable) {
this.dataContainer = container;
this.notifier = notifier;
this.cache = cache;
this.interceptorChain = interceptorChain;
- this.cacheLoaderManager = clManager;
this.distributionManager = distributionManager;
this.icc = icc;
this.txTable = txTable;
Modified: trunk/core/src/main/java/org/infinispan/config/Configuration.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/config/Configuration.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/config/Configuration.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -21,10 +21,6 @@
*/
package org.infinispan.config;
-import java.util.Collections;
-import java.util.List;
-import java.util.concurrent.TimeUnit;
-
import org.infinispan.config.parsing.ClusteringConfigReader;
import org.infinispan.config.parsing.CustomInterceptorConfigReader;
import org.infinispan.distribution.DefaultConsistentHash;
@@ -35,6 +31,10 @@
import org.infinispan.util.ReflectionUtil;
import org.infinispan.util.concurrent.IsolationLevel;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+
/**
* Encapsulates the configuration of a Cache.
*
@@ -61,14 +61,17 @@
@ConfigurationElement(name = "eviction", parent = "default", description = ""),
@ConfigurationElement(name = "expiration", parent = "default", description = ""),
@ConfigurationElement(name = "unsafe", parent = "default", description = ""),
- @ConfigurationElement(name = "customInterceptors", parent = "default",
+ @ConfigurationElement(name = "deadlockDetection", parent = "default", description = ""),
+ @ConfigurationElement(name = "customInterceptors", parent = "default",
customReader=CustomInterceptorConfigReader.class)
})
public class Configuration extends AbstractNamedCacheConfigurationBean {
private static final long serialVersionUID = 5553791890144997466L;
- private boolean useDeadlockDetection = false;
+ private boolean enableDeadlockDetection = false;
+ private long deadlockDetectionSpinDuration = 100;
+
// reference to a global configuration
private GlobalConfiguration globalConfiguration;
@@ -85,6 +88,28 @@
return fetchInMemoryState || (cacheLoaderManagerConfig != null && cacheLoaderManagerConfig.isFetchPersistentState());
}
+
+ public long getDeadlockDetectionSpinDuration() {
+ return deadlockDetectionSpinDuration;
+ }
+
+ @ConfigurationAttribute(name = "spinDuration", containingElement = "deadlockDetection")
+ public void setDeadlockDetectionSpinDuration(long eagerDeadlockSpinDuration) {
+ testImmutability("eagerDeadlockSpinDuration");
+ this.deadlockDetectionSpinDuration = eagerDeadlockSpinDuration;
+ }
+
+
+ public boolean isEnableDeadlockDetection() {
+ return enableDeadlockDetection;
+ }
+
+ @ConfigurationAttribute(name = "enabled", containingElement = "deadlockDetection")
+ public void setEnableDeadlockDetection(boolean useEagerDeadlockDetection) {
+ testImmutability("enableDeadlockDetection");
+ this.enableDeadlockDetection = useEagerDeadlockDetection;
+ }
+
public void setUseLockStriping(boolean useLockStriping) {
testImmutability("useLockStriping");
this.useLockStriping = useLockStriping;
@@ -115,15 +140,7 @@
public long getRehashRpcTimeout() {
return rehashRpcTimeout;
}
-
- public boolean isUseDeadlockDetection() {
- return useDeadlockDetection;
- }
- public void setUseDeadlockDetection(boolean useDeadlockDetection) {
- this.useDeadlockDetection = useDeadlockDetection;
- }
-
/**
* Cache replication mode.
Modified: trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -142,6 +142,7 @@
configureShutdown(getSingleElementInCoreNS("shutdown", globalElement), gc);
configureSerialization(getSingleElementInCoreNS("serialization", globalElement), gc);
configureGlobalJmxStatistics(getSingleElementInCoreNS("globalJmxStatistics", globalElement), gc);
+ configureGlobalJmxStatistics(getSingleElementInCoreNS("globalJmxStatistics", globalElement), gc);
}
}
@@ -161,9 +162,20 @@
configureCacheLoaders(getSingleElementInCoreNS("loaders", e), c);
configureCustomInterceptors(getSingleElementInCoreNS("customInterceptors", e), c);
configureUnsafe(getSingleElementInCoreNS("unsafe", e), c);
+ configureDeadlockDetection(getSingleElementInCoreNS("deadlockDetection", e), c);
return c;
}
+ void configureDeadlockDetection(Element element, Configuration c) {
+ if (element == null) return; //might me missing
+ String enabled = getAttributeValue(element, "enabled");
+ if (existsAttribute(enabled))
+ c.setEnableDeadlockDetection(getBoolean(enabled));
+ String spinDuration = getAttributeValue(element, "spinDuration");
+ if (existsAttribute(spinDuration))
+ c.setDeadlockDetectionSpinDuration(getLong(spinDuration));
+ }
+
private void assertInitialized() {
if (!initialized)
throw new ConfigurationException("Parser not initialized. Please invoke initialize() first, or use a constructor that initializes the parser.");
Modified: trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorFactory.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorFactory.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -7,6 +7,7 @@
import org.infinispan.factories.scopes.Scopes;
import org.infinispan.notifications.cachemanagerlistener.CacheManagerNotifier;
import org.infinispan.remoting.InboundInvocationHandler;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import org.infinispan.transaction.xa.TransactionTable;
import org.infinispan.util.Util;
@@ -17,7 +18,7 @@
* @author <a href="mailto:galder.zamarreno at jboss.com">Galder Zamarreno</a>
* @since 4.0
*/
- at DefaultFactoryFor(classes = {InboundInvocationHandler.class, CacheManagerNotifier.class, RemoteCommandFactory.class, TransactionTable.class})
+ at DefaultFactoryFor(classes = {InboundInvocationHandler.class, CacheManagerNotifier.class, RemoteCommandFactory.class, TransactionTable.class, GlobalTransactionFactory.class})
@Scope(Scopes.GLOBAL)
public class EmptyConstructorFactory extends AbstractComponentFactory implements AutoInstantiableFactory {
public <T> T construct(Class<T> componentType) {
Modified: trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorNamedCacheFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorNamedCacheFactory.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorNamedCacheFactory.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -52,7 +52,8 @@
if (componentType.isInterface()) {
Class componentImpl;
if (componentType.equals(Marshaller.class)) {
- componentImpl = VersionAwareMarshaller.class;
+ VersionAwareMarshaller versionAwareMarshaller = Util.getInstance(VersionAwareMarshaller.class);
+ return componentType.cast(versionAwareMarshaller);
} else if (componentType.equals(InvocationContextContainer.class)) {
componentImpl = InvocationContextContainerImpl.class;
} else {
Modified: trunk/core/src/main/java/org/infinispan/factories/EntryFactoryImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/EntryFactoryImpl.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/factories/EntryFactoryImpl.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -195,8 +195,8 @@
return true;
} else {
Object owner = lockManager.getOwner(key);
- throw new TimeoutException("Unable to acquire lock on key [" + key + "] after [" + getLockAcquisitionTimeout(ctx)
- + "] milliseconds for requestor [" + ctx.getLockOwner() + "]! Lock held by [" + owner + "]");
+ throw new TimeoutException("Unable to acquire lock on key [" + key + "] for requestor [" +
+ ctx.getLockOwner() + "]! Lock held by [" + owner + "]");
}
}
Modified: trunk/core/src/main/java/org/infinispan/factories/InterceptorChainFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/InterceptorChainFactory.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/factories/InterceptorChainFactory.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -22,6 +22,7 @@
package org.infinispan.factories;
+import org.infinispan.CacheException;
import org.infinispan.config.Configuration;
import org.infinispan.config.ConfigurationException;
import org.infinispan.config.CustomInterceptorConfig;
@@ -90,6 +91,10 @@
interceptorChain.appendIntereceptor(createInterceptor(NotificationInterceptor.class));
+ if (configuration.isEnableDeadlockDetection()) {
+ interceptorChain.appendIntereceptor(createInterceptor(DeadlockDetectingInterceptor.class));
+ }
+
switch (configuration.getCacheMode()) {
case REPL_SYNC:
case REPL_ASYNC:
@@ -168,6 +173,8 @@
public <T> T construct(Class<T> componentType) {
try {
return componentType.cast(buildInterceptorChain());
+ } catch (CacheException ce) {
+ throw ce;
}
catch (Exception e) {
throw new ConfigurationException("Unable to build interceptor chain", e);
Modified: trunk/core/src/main/java/org/infinispan/factories/LockManagerFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/LockManagerFactory.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/factories/LockManagerFactory.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -24,6 +24,8 @@
import org.infinispan.factories.annotations.DefaultFactoryFor;
import org.infinispan.util.concurrent.locks.LockManager;
import org.infinispan.util.concurrent.locks.LockManagerImpl;
+import org.infinispan.util.concurrent.locks.DeadlockDetectingLockManager;
+import org.infinispan.config.ConfigurationException;
/**
* // TODO: MANIK: Document this
@@ -32,8 +34,15 @@
* @since 4.0
*/
@DefaultFactoryFor(classes = LockManager.class)
-public class LockManagerFactory extends AbstractComponentFactory implements AutoInstantiableFactory {
+public class LockManagerFactory extends AbstractNamedCacheComponentFactory implements AutoInstantiableFactory {
public <T> T construct(Class<T> componentType) {
- return (T) new LockManagerImpl();
+ if (configuration.isEnableDeadlockDetection()) {
+ if (!configuration.getCacheMode().isSynchronous()) {
+ throw new ConfigurationException("Eager Dead lock detection can only be used for sync caches!");
+ }
+ return (T) new DeadlockDetectingLockManager();
+ } else {
+ return (T) new LockManagerImpl();
+ }
}
}
Modified: trunk/core/src/main/java/org/infinispan/interceptors/CacheMgmtInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/CacheMgmtInterceptor.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/interceptors/CacheMgmtInterceptor.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -98,6 +98,7 @@
}
@Override
+ //Map.put(key,value) :: oldValue
public Object visitPutKeyValueCommand(InvocationContext ctx, PutKeyValueCommand command) throws Throwable {
long t1 = System.currentTimeMillis();
Object retval = invokeNextInterceptor(ctx, command);
Added: trunk/core/src/main/java/org/infinispan/interceptors/DeadlockDetectingInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/DeadlockDetectingInterceptor.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/interceptors/DeadlockDetectingInterceptor.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,128 @@
+package org.infinispan.interceptors;
+
+import org.infinispan.commands.tx.PrepareCommand;
+import org.infinispan.commands.tx.RollbackCommand;
+import org.infinispan.commands.write.PutKeyValueCommand;
+import org.infinispan.context.InvocationContext;
+import org.infinispan.context.impl.TxInvocationContext;
+import org.infinispan.factories.annotations.Inject;
+import org.infinispan.factories.annotations.Start;
+import org.infinispan.interceptors.base.CommandInterceptor;
+import org.infinispan.remoting.transport.Address;
+import org.infinispan.transaction.xa.DeadlockDetectingGlobalTransaction;
+import org.infinispan.transaction.xa.TransactionTable;
+import org.infinispan.util.concurrent.locks.DeadlockDetectedException;
+import org.infinispan.util.concurrent.locks.LockManager;
+
+import javax.transaction.Transaction;
+import javax.transaction.TransactionManager;
+import java.util.Set;
+
+/**
+ * This interceptor populates the {@link org.infinispan.transaction.xa.DeadlockDetectingGlobalTransaction} with
+ * appropriate information needed in order to accomplish deadlock detection. It MUST process populate data before the
+ * replication takes place, so it will do all the tasks before calling {@link org.infinispan.interceptors.base.CommandInterceptor#invokeNextInterceptor(org.infinispan.context.InvocationContext,
+ * org.infinispan.commands.VisitableCommand)}
+ *
+ * @author Mircea.Markus at jboss.com
+ * @since 4.0
+ */
+public class DeadlockDetectingInterceptor extends CommandInterceptor {
+
+ private TransactionTable txTable;
+ private LockManager lockManager;
+ private TransactionManager txManager;
+
+ @Inject
+ public void init(TransactionTable txTable, LockManager lockManager, TransactionManager txManager) {
+ this.txTable = txTable;
+ this.lockManager = lockManager;
+ this.txManager = txManager;
+ }
+
+
+ /**
+ * Only does a sanity check.
+ */
+ @Start
+ public void start() {
+ if (!configuration.isEnableDeadlockDetection()) {
+ throw new IllegalStateException("This interceptor should not be present in the chain as deadlock detection is not used!");
+ }
+ }
+
+ @Override
+ public Object visitPutKeyValueCommand(InvocationContext ctx, PutKeyValueCommand command) throws Throwable {
+ if (ctx.isInTxScope()) {
+ DeadlockDetectingGlobalTransaction gtx = (DeadlockDetectingGlobalTransaction) ctx.getLockOwner();
+ gtx.setLockInterntion(command.getKey());
+ gtx.setProcessingThread(Thread.currentThread());
+ }
+ try {
+ return invokeNextInterceptor(ctx, command);
+ } catch (InterruptedException ie) {
+ if (ctx.isOriginLocal() && ctx.isInTxScope()) {
+ lockManager.releaseLocks(ctx);
+ Transaction transaction = txManager.getTransaction();
+ if (trace)
+ log.trace("Marking the transaction for rollback! : " + transaction);
+ if (transaction == null) {
+ throw new IllegalStateException("We're running in a local transaction, there MUST be one " +
+ "associated witht the local thread but none found! " + transaction);
+ }
+ transaction.setRollbackOnly();
+ throw new DeadlockDetectedException("Deadlock request was detected, tx " + transaction +
+ " was marked for rollback");
+ } else {
+ if (trace)
+ log.trace("Received an interrupt request, but we're not running within deadlock detection scenario, so passing it up the stack", ie);
+ throw ie;
+ }
+ }
+ }
+
+ @Override
+ public Object visitPrepareCommand(TxInvocationContext ctx, PrepareCommand command) throws Throwable {
+ DeadlockDetectingGlobalTransaction globalTransaction = (DeadlockDetectingGlobalTransaction) ctx.getGlobalTransaction();
+ globalTransaction.setProcessingThread(Thread.currentThread());
+ if (ctx.isOriginLocal()) {
+ if (configuration.getCacheMode().isDistributed()) {
+ Set<Address> transactionParticipants = ctx.getTransactionParticipants();
+ globalTransaction.setReplicatingTo(transactionParticipants);
+ } else {
+ globalTransaction.setReplicatingTo(null);
+ }
+ if (trace) log.trace("Deadlock detection information was added to " + globalTransaction);
+ }
+ try {
+ return invokeNextInterceptor(ctx, command);
+ } catch (Throwable dde) {
+ if (ctx.isOriginLocal()) {
+ globalTransaction.setMarkedForRollback(true);
+ boolean wasInterrupted = Thread.interrupted();
+ if (trace)
+ log.trace("Deadlock was detected on the remote side, marking the tx for rollback. Was this thread interrupted? " + wasInterrupted);
+ }
+ throw dde;
+ } finally {
+ if (!ctx.isOriginLocal()) {
+ if (!txTable.containRemoteTx(ctx.getGlobalTransaction())) {
+ if (trace) {
+ log.trace("While returning from prepare we determined that remote tx is no longer in the txTable. " +
+ "This means that a rollback was executed in between; releasing locks");
+ }
+ lockManager.releaseLocks(ctx);
+ }
+ }
+ }
+ }
+
+ @Override
+ public Object visitRollbackCommand(TxInvocationContext ctx, RollbackCommand command) throws Throwable {
+ if (!ctx.isOriginLocal()) {
+ DeadlockDetectingGlobalTransaction globalTransaction = (DeadlockDetectingGlobalTransaction) ctx.getGlobalTransaction();
+ globalTransaction.interruptProcessingThread();
+ }
+ return invokeNextInterceptor(ctx, command);
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/interceptors/DeadlockDetectingInterceptor.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Modified: trunk/core/src/main/java/org/infinispan/interceptors/DistributionInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/DistributionInterceptor.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/interceptors/DistributionInterceptor.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -175,7 +175,7 @@
return invokeNextInterceptor(ctx, command);
}
- // ---- TX boundard commands
+ // ---- TX boundary commands
@Override
public Object visitCommitCommand(TxInvocationContext ctx, CommitCommand command) throws Throwable {
if (ctx.isOriginLocal()) {
Modified: trunk/core/src/main/java/org/infinispan/interceptors/LockingInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/LockingInterceptor.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/interceptors/LockingInterceptor.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -45,6 +45,7 @@
import org.infinispan.transaction.xa.GlobalTransaction;
import org.infinispan.util.ReversibleOrderedSet;
import org.infinispan.util.concurrent.IsolationLevel;
+import org.infinispan.util.concurrent.TimeoutException;
import org.infinispan.util.concurrent.locks.LockManager;
import java.util.ArrayList;
@@ -83,7 +84,7 @@
return invokeNextInterceptor(ctx, command);
} finally {
if (ctx.isInTxScope()) {
- cleanupLocks(ctx, ctx.getLockOwner(), true);
+ cleanupLocks(ctx, true);
} else {
throw new IllegalStateException("Attempting to do a commit or rollback but there is no transactional context in scope. " + ctx);
}
@@ -96,7 +97,7 @@
return invokeNextInterceptor(ctx, command);
} finally {
if (ctx.isInTxScope()) {
- cleanupLocks(ctx, ctx.getLockOwner(), false);
+ cleanupLocks(ctx, false);
} else {
throw new IllegalStateException("Attempting to do a commit or rollback but there is no transactional context in scope. " + ctx);
}
@@ -107,9 +108,12 @@
public Object visitPrepareCommand(TxInvocationContext ctx, PrepareCommand command) throws Throwable {
try {
return invokeNextInterceptor(ctx, command);
+ } catch (TimeoutException te) {
+ cleanupLocks(ctx, false);
+ throw te;
} finally {
if (command.isOnePhaseCommit())
- cleanupLocks(ctx, ctx.getLockOwner(), true);
+ cleanupLocks(ctx, true);
}
}
@@ -235,7 +239,7 @@
private void doAfterCall(InvocationContext ctx) {
// for non-transactional stuff.
if (!ctx.isInTxScope()) {
- cleanupLocks(ctx, ctx.getLockOwner(), true);
+ cleanupLocks(ctx, true);
} else {
if (trace) log.trace("Transactional. Not cleaning up locks till the transaction ends.");
if (useReadCommitted) {
@@ -260,14 +264,12 @@
}
}
- private void cleanupLocks(InvocationContext ctx, Object owner, boolean commit) {
- // clean up.
- // unlocking needs to be done in reverse order.
- ReversibleOrderedSet<Map.Entry<Object, CacheEntry>> entries = ctx.getLookedUpEntries().entrySet();
- Iterator<Map.Entry<Object, CacheEntry>> it = entries.reverseIterator();
- if (trace) log.trace("Number of entries in context: {0}", entries.size());
-
+ private void cleanupLocks(InvocationContext ctx, boolean commit) {
if (commit) {
+ Object owner = ctx.getLockOwner();
+ ReversibleOrderedSet<Map.Entry<Object, CacheEntry>> entries = ctx.getLookedUpEntries().entrySet();
+ Iterator<Map.Entry<Object, CacheEntry>> it = entries.reverseIterator();
+ if (trace) log.trace("Number of entries in context: {0}", entries.size());
while (it.hasNext()) {
Map.Entry<Object, CacheEntry> e = it.next();
CacheEntry entry = e.getValue();
@@ -287,27 +289,11 @@
}
}
} else {
- while (it.hasNext()) {
- Map.Entry<Object, CacheEntry> e = it.next();
- CacheEntry entry = e.getValue();
- Object key = e.getKey();
- boolean needToUnlock = lockManager.possiblyLocked(entry);
- // could be null with read-committed
- if (entry != null && entry.isChanged()) entry.rollback();
- else {
- if (trace) log.trace("Entry for key {0} is null, not calling rollbackUpdate", key);
- }
- // and then unlock
- if (needToUnlock) {
- if (trace) log.trace("Releasing lock on [" + key + "] for owner " + owner);
- lockManager.unlock(key, owner);
- }
- }
+ lockManager.releaseLocks(ctx);
}
}
protected void commitEntry(InvocationContext ctx, CacheEntry entry) {
entry.commit(dataContainer);
}
-
}
Modified: trunk/core/src/main/java/org/infinispan/interceptors/ReplicationInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/ReplicationInterceptor.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/interceptors/ReplicationInterceptor.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -60,7 +60,6 @@
Object retVal = invokeNextInterceptor(ctx, command);
if (ctx.isOriginLocal() && command.hasModifications()) {
boolean async = configuration.getCacheMode() == Configuration.CacheMode.REPL_ASYNC;
- boolean useDeadlockDetection = configuration.isUseDeadlockDetection();
rpcManager.broadcastRpcCommand(command, !async, false);
}
return retVal;
Modified: trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -80,7 +80,7 @@
if (!command.isOnePhaseCommit()) {
transactionLog.logPrepare(command);
}
- if (getStatisticsEnabled()) prepares.incrementAndGet();
+ if (this.statsEnabled) prepares.incrementAndGet();
Object result = invokeNextInterceptor(ctx, command);
if (command.isOnePhaseCommit()) {
transactionLog.logOnePhaseCommit(ctx.getGlobalTransaction(), command.getModifications());
@@ -90,7 +90,7 @@
@Override
public Object visitCommitCommand(TxInvocationContext ctx, CommitCommand command) throws Throwable {
- if (getStatisticsEnabled()) commits.incrementAndGet();
+ if (this.statsEnabled) commits.incrementAndGet();
Object result = invokeNextInterceptor(ctx, command);
transactionLog.logCommit(command.getGlobalTransaction());
return result;
@@ -98,7 +98,7 @@
@Override
public Object visitRollbackCommand(TxInvocationContext ctx, RollbackCommand command) throws Throwable {
- if (getStatisticsEnabled()) rollbacks.incrementAndGet();
+ if (this.statsEnabled) rollbacks.incrementAndGet();
transactionLog.rollback(command.getGlobalTransaction());
return invokeNextInterceptor(ctx, command);
}
@@ -219,10 +219,6 @@
rollbacks.set(0);
}
- public boolean getStatisticsEnabled() {
- return this.statsEnabled;
- }
-
public void setStatisticsEnabled(boolean enabled) {
this.statsEnabled = enabled;
}
Modified: trunk/core/src/main/java/org/infinispan/marshall/MarshallerImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/marshall/MarshallerImpl.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/marshall/MarshallerImpl.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -26,17 +26,7 @@
import org.infinispan.commands.RemoteCommandFactory;
import org.infinispan.commands.ReplicableCommand;
import org.infinispan.commands.write.WriteCommand;
-import org.infinispan.container.entries.ImmortalCacheEntry;
-import org.infinispan.container.entries.ImmortalCacheValue;
-import org.infinispan.container.entries.InternalCacheEntry;
-import org.infinispan.container.entries.InternalCacheValue;
-import org.infinispan.container.entries.InternalEntryFactory;
-import org.infinispan.container.entries.MortalCacheEntry;
-import org.infinispan.container.entries.MortalCacheValue;
-import org.infinispan.container.entries.TransientCacheEntry;
-import org.infinispan.container.entries.TransientCacheValue;
-import org.infinispan.container.entries.TransientMortalCacheEntry;
-import org.infinispan.container.entries.TransientMortalCacheValue;
+import org.infinispan.container.entries.*;
import org.infinispan.io.ByteBuffer;
import org.infinispan.io.ExposedByteArrayOutputStream;
import org.infinispan.io.UnsignedNumeric;
@@ -49,8 +39,9 @@
import org.infinispan.remoting.responses.UnsuccessfulResponse;
import org.infinispan.remoting.transport.Address;
import org.infinispan.remoting.transport.jgroups.JGroupsAddress;
-import org.infinispan.transaction.xa.GlobalTransaction;
import org.infinispan.transaction.TransactionLog;
+import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import org.infinispan.util.FastCopyHashMap;
import org.infinispan.util.Immutables;
import org.infinispan.util.Util;
@@ -147,6 +138,8 @@
protected ClassLoader defaultClassLoader;
protected boolean useRefs = false;
+ protected GlobalTransactionFactory gtxFactory = new GlobalTransactionFactory();
+
public void init(ClassLoader defaultClassLoader, RemoteCommandFactory remoteCommandFactory) {
this.defaultClassLoader = defaultClassLoader;
this.remoteCommandFactory = remoteCommandFactory;
@@ -630,7 +623,7 @@
private GlobalTransaction unmarshallGlobalTransaction(ObjectInput in, UnmarshalledReferences refMap) throws IOException, ClassNotFoundException {
- GlobalTransaction gtx = new GlobalTransaction();
+ GlobalTransaction gtx = gtxFactory.instantiateGlobalTransaction();
long id = in.readLong();
Object address = unmarshallObject(in, refMap);
gtx.setId(id);
Added: trunk/core/src/main/java/org/infinispan/marshall/exts/DeadlockDetectingGlobalTransactionExternalizer.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/marshall/exts/DeadlockDetectingGlobalTransactionExternalizer.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/marshall/exts/DeadlockDetectingGlobalTransactionExternalizer.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,35 @@
+package org.infinispan.marshall.exts;
+
+import org.infinispan.transaction.xa.DeadlockDetectingGlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
+
+import java.io.IOException;
+import java.io.ObjectInput;
+import java.io.ObjectOutput;
+
+/**
+ * Externalizer for {@link DeadlockDetectingGlobalTransaction}.
+ *
+ * @author Mircea.Markus at jboss.com
+ */
+public class DeadlockDetectingGlobalTransactionExternalizer extends GlobalTransactionExternalizer {
+
+ public DeadlockDetectingGlobalTransactionExternalizer() {
+ super();
+ gtxFactory = new GlobalTransactionFactory(true);
+ }
+
+ @Override
+ public void writeObject(ObjectOutput output, Object subject) throws IOException {
+ super.writeObject(output, subject);
+ DeadlockDetectingGlobalTransaction ddGt = (DeadlockDetectingGlobalTransaction) subject;
+ output.writeLong(ddGt.getCoinToss());
+ }
+
+ @Override
+ public Object readObject(ObjectInput input) throws IOException, ClassNotFoundException {
+ DeadlockDetectingGlobalTransaction ddGt = (DeadlockDetectingGlobalTransaction) super.readObject(input);
+ ddGt.setCoinToss(input.readLong());
+ return ddGt;
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/marshall/exts/DeadlockDetectingGlobalTransactionExternalizer.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Modified: trunk/core/src/main/java/org/infinispan/marshall/exts/GlobalTransactionExternalizer.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/marshall/exts/GlobalTransactionExternalizer.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/marshall/exts/GlobalTransactionExternalizer.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -21,11 +21,10 @@
*/
package org.infinispan.marshall.exts;
- import net.jcip.annotations.Immutable;
-
import org.infinispan.marshall.Externalizer;
import org.infinispan.remoting.transport.Address;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import java.io.IOException;
import java.io.ObjectInput;
@@ -37,22 +36,23 @@
* @author Galder Zamarreño
* @since 4.0
*/
- at Immutable
public class GlobalTransactionExternalizer implements Externalizer {
+ protected GlobalTransactionFactory gtxFactory = new GlobalTransactionFactory();
+
+
public void writeObject(ObjectOutput output, Object subject) throws IOException {
GlobalTransaction gtx = (GlobalTransaction) subject;
output.writeLong(gtx.getId());
- output.writeObject(gtx.getAddress());
+ output.writeObject(gtx.getAddress());
}
public Object readObject(ObjectInput input) throws IOException, ClassNotFoundException {
- GlobalTransaction gtx = new GlobalTransaction();
+ GlobalTransaction gtx = gtxFactory.instantiateGlobalTransaction();
long id = input.readLong();
Object address = input.readObject();
gtx.setId(id);
gtx.setAddress((Address) address);
return gtx;
}
-
}
Modified: trunk/core/src/main/java/org/infinispan/marshall/jboss/ConstantObjectTable.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/marshall/jboss/ConstantObjectTable.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/marshall/jboss/ConstantObjectTable.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -53,36 +53,14 @@
import org.infinispan.loaders.bucket.Bucket;
import org.infinispan.marshall.Externalizer;
import org.infinispan.marshall.MarshalledValue;
-import org.infinispan.marshall.exts.ArrayListExternalizer;
-import org.infinispan.marshall.exts.BucketExternalizer;
-import org.infinispan.marshall.exts.DeltaAwareExternalizer;
-import org.infinispan.marshall.exts.ExceptionResponseExternalizer;
-import org.infinispan.marshall.exts.ExtendedResponseExternalizer;
-import org.infinispan.marshall.exts.GlobalTransactionExternalizer;
-import org.infinispan.marshall.exts.ImmortalCacheEntryExternalizer;
-import org.infinispan.marshall.exts.ImmortalCacheValueExternalizer;
-import org.infinispan.marshall.exts.ImmutableMapExternalizer;
-import org.infinispan.marshall.exts.JGroupsAddressExternalizer;
-import org.infinispan.marshall.exts.LinkedListExternalizer;
-import org.infinispan.marshall.exts.MapExternalizer;
-import org.infinispan.marshall.exts.MarshalledValueExternalizer;
-import org.infinispan.marshall.exts.MortalCacheEntryExternalizer;
-import org.infinispan.marshall.exts.MortalCacheValueExternalizer;
-import org.infinispan.marshall.exts.ReplicableCommandExternalizer;
-import org.infinispan.marshall.exts.SetExternalizer;
-import org.infinispan.marshall.exts.SingletonListExternalizer;
-import org.infinispan.marshall.exts.SuccessfulResponseExternalizer;
-import org.infinispan.marshall.exts.TransactionLogExternalizer;
-import org.infinispan.marshall.exts.TransientCacheEntryExternalizer;
-import org.infinispan.marshall.exts.TransientCacheValueExternalizer;
-import org.infinispan.marshall.exts.TransientMortalCacheEntryExternalizer;
-import org.infinispan.marshall.exts.TransientMortalCacheValueExternalizer;
+import org.infinispan.marshall.exts.*;
import org.infinispan.remoting.responses.ExceptionResponse;
import org.infinispan.remoting.responses.ExtendedResponse;
import org.infinispan.remoting.responses.RequestIgnoredResponse;
import org.infinispan.remoting.responses.SuccessfulResponse;
import org.infinispan.remoting.responses.UnsuccessfulResponse;
import org.infinispan.remoting.transport.jgroups.JGroupsAddress;
+import org.infinispan.transaction.xa.DeadlockDetectingGlobalTransaction;
import org.infinispan.transaction.xa.GlobalTransaction;
import org.infinispan.util.FastCopyHashMap;
import org.infinispan.util.Util;
@@ -122,6 +100,7 @@
static {
EXTERNALIZERS.put(GlobalTransaction.class.getName(), GlobalTransactionExternalizer.class.getName());
+ EXTERNALIZERS.put(DeadlockDetectingGlobalTransaction.class.getName(), DeadlockDetectingGlobalTransactionExternalizer.class.getName());
EXTERNALIZERS.put(JGroupsAddress.class.getName(), JGroupsAddressExternalizer.class.getName());
EXTERNALIZERS.put(ArrayList.class.getName(), ArrayListExternalizer.class.getName());
EXTERNALIZERS.put(LinkedList.class.getName(), LinkedListExternalizer.class.getName());
Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManager.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManager.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManager.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -53,7 +53,7 @@
* @return a list of responses from each member contacted.
* @throws Exception in the event of problems.
*/
- List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue, ResponseFilter responseFilter) throws Exception;
+ List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue, ResponseFilter responseFilter);
/**
* Invokes an RPC call on other caches in the cluster.
@@ -68,7 +68,7 @@
* @return a list of responses from each member contacted.
* @throws Exception in the event of problems.
*/
- List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue) throws Exception;
+ List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue);
/**
* Invokes an RPC call on other caches in the cluster.
Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -75,7 +75,7 @@
return !sync && replicationQueue != null && replicationQueue.isEnabled();
}
- public final List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue, ResponseFilter responseFilter) throws Exception {
+ public final List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue, ResponseFilter responseFilter) {
List<Address> members = t.getMembers();
if (members.size() < 2) {
if (log.isDebugEnabled())
@@ -87,7 +87,10 @@
if (isStatisticsEnabled()) replicationCount.incrementAndGet();
return result;
} catch (CacheException e) {
- if (log.isTraceEnabled()) log.trace("replication exception: ", e);
+ if (log.isTraceEnabled()) {
+ log.trace("replication exception: ", e);
+ }
+
if (isStatisticsEnabled()) replicationFailures.incrementAndGet();
throw e;
} catch (Throwable th) {
@@ -98,7 +101,7 @@
}
}
- public final List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue) throws Exception {
+ public final List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue) {
return invokeRemotely(recipients, rpcCommand, mode, timeout, usePriorityQueue, null);
}
@@ -206,17 +209,9 @@
rpc = cf.buildSingleRpcCommand(rpc);
}
List rsps;
- try {
rsps = invokeRemotely(recipients, rpc, getResponseMode(sync), timeout, usePriorityQueue);
if (trace) log.trace("responses=" + rsps);
if (sync) checkResponses(rsps);
- } catch (CacheException e) {
- log.error("Replication exception", e);
- throw e;
- } catch (Exception ex) {
- log.error("Unexpected exception", ex);
- throw new ReplicationException("Unexpected exception while replicating", ex);
- }
}
}
Modified: trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -310,7 +310,7 @@
if (mode.isAsynchronous()) return Collections.emptyList();// async case
if (trace)
- log.trace("Cache [{0}]: responses for command {1}:\n{2}", getAddress(), rpcCommand.getClass().getSimpleName(), rsps);
+ log.trace("Cache [{0}], is caller thread interupted? {3}: responses for command {1}:\n{2}", getAddress(), rpcCommand.getClass().getSimpleName(), rsps, Thread.currentThread().isInterrupted());
// short-circuit no-return-value calls.
if (rsps == null) return Collections.emptyList();
@@ -333,7 +333,7 @@
Exception e = ((ExceptionResponse) value).getException();
if (!(e instanceof ReplicationException)) {
// if we have any application-level exceptions make sure we throw them!!
- if (trace) log.trace("Recieved exception '{0}' from {1}", e, rsp.getSender());
+ if (trace) log.trace("Received exception from " + rsp.getSender(), e);
throw e;
}
}
Added: trunk/core/src/main/java/org/infinispan/transaction/xa/DeadlockDetectingGlobalTransaction.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/DeadlockDetectingGlobalTransaction.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/DeadlockDetectingGlobalTransaction.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,165 @@
+package org.infinispan.transaction.xa;
+
+import org.infinispan.remoting.transport.Address;
+import org.infinispan.util.logging.Log;
+import org.infinispan.util.logging.LogFactory;
+
+import java.util.Collections;
+import java.util.Set;
+
+/**
+ * This class is used when deadlock detection is enabled.
+ *
+ * @author Mircea.Markus at jboss.com
+ */
+public class DeadlockDetectingGlobalTransaction extends GlobalTransaction {
+
+ private static Log log = LogFactory.getLog(DeadlockDetectingGlobalTransaction.class);
+
+ public static final boolean trace = log.isTraceEnabled();
+
+ private Set<Address> replicatingTo = Collections.EMPTY_SET;
+
+ private volatile transient Thread processingThread;
+
+ private volatile long coinToss;
+
+ private volatile boolean isMarkedForRollback;
+
+ private transient volatile Object lockInterntion;
+
+
+ public DeadlockDetectingGlobalTransaction() {
+ }
+
+ DeadlockDetectingGlobalTransaction(Address addr, boolean remote) {
+ super(addr, remote);
+ }
+
+ DeadlockDetectingGlobalTransaction(boolean remote) {
+ super(null, remote);
+ }
+
+ /**
+ * Is this global transaction replicating to the given address?
+ */
+ public boolean isReplicatingTo(Address address) {
+ if (this.replicatingTo == null) {
+ return true;
+ } else {
+ return this.replicatingTo.contains(address);
+ }
+ }
+
+ /**
+ * On a node, this will set the thread that handles replicatin on the given node.
+ */
+ public void setProcessingThread(Thread replicationThread) {
+ if (trace) log.trace("Setting thread " + Thread.currentThread() + "on tx [" + this + "]");
+ this.processingThread = replicationThread;
+ }
+
+ /**
+ * Tries to interrupt the processing thread.
+ */
+ public synchronized void interruptProcessingThread() {
+ if (isMarkedForRollback) {
+ if (trace) log.trace("Not interrupting as tx is marked for rollback");
+ return;
+ }
+ if (processingThread == null) {
+ if(trace) log.trace("Processing thread is null, nothing to interrupt");
+ return;
+ }
+ if (trace) {
+ StackTraceElement[] stackTraceElements = processingThread.getStackTrace();
+ StringBuilder builder = new StringBuilder();
+ for (StackTraceElement stackTraceElement : stackTraceElements) {
+ builder.append(" ").append(stackTraceElement).append('\n');
+ }
+ log.trace("About to interrupt thread: " + processingThread + ". Thread's stack trace is: \n" + builder.toString());
+ }
+ this.processingThread.interrupt();
+ }
+
+ /**
+ * Sets the set og <b>Address</b> objects this node is replicating to.
+ */
+ public void setReplicatingTo(Set<Address> targets) {
+ this.replicatingTo = targets;
+ }
+
+ /**
+ * Based on the coin toss, determine whether this tx will continue working or this thread will be stopped.
+ */
+ public boolean thisWillInterrupt(DeadlockDetectingGlobalTransaction globalTransaction) {
+ return this.coinToss > globalTransaction.coinToss;
+ }
+
+ /**
+ * Sets the reandom number that defines the coin toss.
+ */
+ public void setCoinToss(long coinToss) {
+ this.coinToss = coinToss;
+ }
+
+ public long getCoinToss() {
+ return coinToss;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) return true;
+ if (!(o instanceof DeadlockDetectingGlobalTransaction)) return false;
+ if (!super.equals(o)) return false;
+
+ DeadlockDetectingGlobalTransaction that = (DeadlockDetectingGlobalTransaction) o;
+
+ if (coinToss != that.coinToss) return false;
+ if (replicatingTo != null ? !replicatingTo.equals(that.replicatingTo) : that.replicatingTo != null) return false;
+
+ return true;
+ }
+
+ @Override
+ public int hashCode() {
+ int result = super.hashCode();
+ result = 31 * result + (replicatingTo != null ? replicatingTo.hashCode() : 0);
+ result = 31 * result + (int) (coinToss ^ (coinToss >>> 32));
+ return result;
+ }
+
+ @Override
+ public String toString() {
+ return "DeadlockDetectingGlobalTransaction{" +
+ "replicatingTo=" + replicatingTo +
+ ", replicationThread=" + processingThread +
+ ", coinToss=" + coinToss +
+ "} " + super.toString();
+ }
+
+ /**
+ * Once marked for rollback, the call to {@link #interruptProcessingThread()} will be ignored.
+ */
+ public synchronized boolean isMarkedForRollback() {
+ return isMarkedForRollback;
+ }
+
+ public synchronized void setMarkedForRollback(boolean markedForRollback) {
+ isMarkedForRollback = markedForRollback;
+ }
+
+ /**
+ * Returns the key this transaction intends to lock.
+ */
+ public Object getLockInterntion() {
+ return lockInterntion;
+ }
+
+ /**
+ * Sets the lock this transaction intends to lock.
+ */
+ public void setLockInterntion(Object lockInterntion) {
+ this.lockInterntion = lockInterntion;
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/transaction/xa/DeadlockDetectingGlobalTransaction.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Modified: trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransaction.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransaction.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransaction.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -23,10 +23,10 @@
import org.infinispan.remoting.transport.Address;
-import java.io.Externalizable;
import java.io.IOException;
import java.io.ObjectInput;
import java.io.ObjectOutput;
+import java.io.Externalizable;
import java.util.concurrent.atomic.AtomicLong;
@@ -47,33 +47,23 @@
private long id = -1;
- private transient Address addr = null;
+ protected transient Address addr = null;
private transient int hash_code = -1; // in the worst case, hashCode() returns 0, then increases, so we're safe here
private transient boolean remote = false;
/**
* empty ctor used by externalization.
*/
- public GlobalTransaction() {
+ GlobalTransaction() {
}
- public GlobalTransaction(Address addr, boolean remote) {
+ GlobalTransaction(Address addr, boolean remote) {
this.id = sid.incrementAndGet();
this.addr = addr;
this.remote = remote;
}
- public GlobalTransaction(long id, Address addr) {
- this.id = id;
- this.addr = addr;
- this.remote = true;
- }
-
- public GlobalTransaction(boolean remote) {
- this(null, remote);
- }
-
- public Object getAddress() {
+ public Address getAddress() {
return addr;
}
Added: trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransactionFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransactionFactory.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransactionFactory.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,61 @@
+package org.infinispan.transaction.xa;
+
+import org.infinispan.config.Configuration;
+import org.infinispan.factories.annotations.Inject;
+import org.infinispan.factories.annotations.Start;
+import org.infinispan.remoting.transport.Address;
+
+import java.util.Random;
+
+/**
+ * Factory for GlobalTransaction/DadlockDetectingGlobalTransaction.
+ *
+ * @author Mircea.Markus at jboss.com
+ */
+public class GlobalTransactionFactory {
+
+ private boolean isEddEnabled = false;
+
+ /** this class is internally synchronized, so it can be shared between instances */
+ private final Random rnd = new Random();
+
+ private long generateRandomdId() {
+ return rnd.nextLong();
+ }
+
+
+ public GlobalTransactionFactory() {
+ }
+
+ public GlobalTransactionFactory(boolean eddEnabled) {
+ isEddEnabled = eddEnabled;
+ }
+
+ @Inject
+ public void init(Configuration configuration) {
+ isEddEnabled = configuration.isEnableDeadlockDetection();
+ }
+
+ @Start
+ public void start() {
+
+ }
+
+ public GlobalTransaction instantiateGlobalTransaction() {
+ if (isEddEnabled) {
+ return new DeadlockDetectingGlobalTransaction();
+ } else {
+ return new GlobalTransaction();
+ }
+ }
+
+ public GlobalTransaction newGlobalTransaction(Address addr, boolean remote) {
+ if (isEddEnabled) {
+ DeadlockDetectingGlobalTransaction globalTransaction = new DeadlockDetectingGlobalTransaction(addr, remote);
+ globalTransaction.setCoinToss(generateRandomdId());
+ return globalTransaction;
+ } else {
+ return new GlobalTransaction(addr, remote);
+ }
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/transaction/xa/GlobalTransactionFactory.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Modified: trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -40,17 +40,19 @@
private InterceptorChain invoker;
private CacheNotifier notifier;
private RpcManager rpcManager;
+ private GlobalTransactionFactory gtf;
@Inject
public void initialize(CommandsFactory commandsFactory, RpcManager rpcManager, Configuration configuration,
- InvocationContextContainer icc, InterceptorChain invoker, CacheNotifier notifier) {
+ InvocationContextContainer icc, InterceptorChain invoker, CacheNotifier notifier, GlobalTransactionFactory gtf) {
this.commandsFactory = commandsFactory;
this.rpcManager = rpcManager;
this.configuration = configuration;
this.icc = icc;
this.invoker = invoker;
this.notifier = notifier;
+ this.gtf = gtf;
}
@@ -91,7 +93,7 @@
TransactionXaAdapter current = localTransactions.get(transaction);
if (current == null) {
Address localAddress = rpcManager != null ? rpcManager.getTransport().getAddress() : null;
- GlobalTransaction tx = localAddress == null ? new GlobalTransaction(false) : new GlobalTransaction(localAddress, false);
+ GlobalTransaction tx = gtf.newGlobalTransaction(localAddress, false);
current = new TransactionXaAdapter(tx, icc, invoker, commandsFactory, configuration, this, transaction);
localTransactions.put(transaction, current);
try {
@@ -136,4 +138,8 @@
public TransactionXaAdapter getXaCacheAdapter(Transaction tx) {
return localTransactions.get(tx);
}
+
+ public boolean containRemoteTx(GlobalTransaction globalTransaction) {
+ return remoteTransactions.containsKey(globalTransaction);
+ }
}
Modified: trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -127,7 +127,7 @@
try {
invoker.invoke(ctx, rollbackCommand);
} catch (Throwable e) {
- log.error("Exception while ", e);
+ log.error("Exception while rollback", e);
throw new XAException(XAException.XA_HEURHAZ);
} finally {
txTable.removeLocalTransaction(transaction);
@@ -136,11 +136,11 @@
}
public void start(Xid xid, int i) throws XAException {
- if (trace) log.trace("start called");
+ if (trace) log.trace("start called on tx " + this.globalTx);
}
public void end(Xid xid, int i) throws XAException {
- if (trace) log.trace("end called");
+ if (trace) log.trace("end called on tx " + this.globalTx);
}
public void forget(Xid xid) throws XAException {
Added: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectedException.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectedException.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectedException.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,14 @@
+package org.infinispan.util.concurrent.locks;
+
+import org.infinispan.CacheException;
+
+/**
+ * Exception signaling detected deadlocks.
+ *
+ * @author Mircea.Markus at jboss.com
+ */
+public class DeadlockDetectedException extends CacheException {
+ public DeadlockDetectedException(String msg) {
+ super(msg);
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectedException.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Added: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectingLockManager.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectingLockManager.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectingLockManager.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,164 @@
+package org.infinispan.util.concurrent.locks;
+
+import org.infinispan.context.InvocationContext;
+import org.infinispan.context.impl.TxInvocationContext;
+import org.infinispan.factories.annotations.Start;
+import org.infinispan.jmx.annotations.MBean;
+import org.infinispan.jmx.annotations.ManagedAttribute;
+import org.infinispan.jmx.annotations.ManagedOperation;
+import org.infinispan.remoting.transport.Address;
+import org.infinispan.transaction.xa.DeadlockDetectingGlobalTransaction;
+import org.infinispan.util.logging.Log;
+import org.infinispan.util.logging.LogFactory;
+
+import static java.util.concurrent.TimeUnit.MILLISECONDS;
+import java.util.concurrent.atomic.AtomicLong;
+
+/**
+ * Lock manager in charge with processing deadlock detections.
+ *
+ * @author Mircea.Markus at jboss.com
+ */
+ at MBean(description = "Information about the number of deadlocks that were detected")
+public class DeadlockDetectingLockManager extends LockManagerImpl {
+
+ private static final Log log = LogFactory.getLog(DeadlockDetectingLockManager.class);
+
+ private volatile long spinDuration;
+
+ private volatile boolean exposeJmxStats;
+
+ private AtomicLong detectedDeadlocks = new AtomicLong(0);
+
+ private AtomicLong locallyInterruptedTransactions = new AtomicLong(0);
+
+ private AtomicLong overlapWithNotDeadlockAwareLockOwners = new AtomicLong(0);
+
+ @Start
+ public void init() {
+ spinDuration = configuration.getDeadlockDetectionSpinDuration();
+ exposeJmxStats = configuration.isExposeJmxStatistics();
+ }
+
+ public boolean lockAndRecord(Object key, InvocationContext ctx) throws InterruptedException {
+ long lockTimeout = getLockAcquisitionTimeout(ctx);
+ if (trace) log.trace("Attempting to lock {0} with acquisition timeout of {1} millis", key, lockTimeout);
+
+
+ if (ctx.isInTxScope()) {
+ if (trace) log.trace("Using early dead lock detection");
+ final long start = System.currentTimeMillis();
+ long now;
+ while ((now = System.currentTimeMillis()) < (start + lockTimeout)) {
+ if (lockContainer.acquireLock(key, spinDuration, MILLISECONDS)) {
+ if (trace) log.trace("successfully acquired lock on " + key + ", returning ...");
+ return true;
+ } else {
+ if (trace)
+ log.trace("Could not acquire lock on '" + key + "' as it is locked by '" + getOwner(key) + "', check for dead locks");
+ Object owner = getOwner(key);
+ if (!(owner instanceof DeadlockDetectingGlobalTransaction)) {
+ if (trace)
+ log.trace("Owner is not instance of DeadlockDetectingGlobalTransaction: " + owner + ", continuing ...");
+ if (exposeJmxStats) overlapWithNotDeadlockAwareLockOwners.incrementAndGet();
+ continue; //try to acquire lock again, for the rest of the time
+ }
+ DeadlockDetectingGlobalTransaction lockOwnerTx = (DeadlockDetectingGlobalTransaction) owner;
+ if (!ctx.isOriginLocal() && !lockOwnerTx.isRemote()) {
+ return remoteVsRemoteDld(key, ctx, lockTimeout, start, now, lockOwnerTx);
+ }
+ if (ctx.isOriginLocal() && !lockOwnerTx.isRemote()) {
+ localVsLocalDld(ctx, lockOwnerTx);
+ }
+ }
+ }
+ } else {
+ if (lockContainer.acquireLock(key, lockTimeout, MILLISECONDS)) {
+ return true;
+ }
+ }
+ // couldn't acquire lock!
+ return false;
+ }
+
+ private void localVsLocalDld(InvocationContext ctx, DeadlockDetectingGlobalTransaction lockOwnerTx) {
+ if (trace) log.trace("Looking for local vs local deadlocks");
+ DeadlockDetectingGlobalTransaction thisThreadsTx = (DeadlockDetectingGlobalTransaction) ctx.getLockOwner();
+ boolean weOwnLock = ownsLock(lockOwnerTx.getLockInterntion(), thisThreadsTx);
+ if (trace) {
+ log.trace("Other owner's intention is " + lockOwnerTx.getLockInterntion() + ". Do we(" + thisThreadsTx + ") own lock for it? " + weOwnLock + ". Lock owner is " + getOwner(lockOwnerTx.getLockInterntion()));
+ }
+ if (weOwnLock) {
+ boolean iShouldInterrupt = thisThreadsTx.thisWillInterrupt(lockOwnerTx);
+ if (trace)
+ log.trace("deadlock situation detected. Shall I interrupt?" + iShouldInterrupt );
+ if (iShouldInterrupt) {
+ lockOwnerTx.interruptProcessingThread();
+ if (exposeJmxStats) detectedDeadlocks.incrementAndGet();
+ }
+ }
+ }
+
+ private boolean remoteVsRemoteDld(Object key, InvocationContext ctx, long lockTimeout, long start, long now, DeadlockDetectingGlobalTransaction lockOwnerTx) throws InterruptedException {
+ TxInvocationContext remoteTxContext = (TxInvocationContext) ctx;
+ Address origin = remoteTxContext.getGlobalTransaction().getAddress();
+ DeadlockDetectingGlobalTransaction remoteGlobalTransaction = (DeadlockDetectingGlobalTransaction) ctx.getLockOwner();
+ boolean thisShouldInterrupt = remoteGlobalTransaction.thisWillInterrupt(lockOwnerTx);
+ if (trace) log.trace("Should I interrupt other transaction ? " + thisShouldInterrupt);
+ boolean isDeadLock = (configuration.getCacheMode().isReplicated() || lockOwnerTx.isReplicatingTo(origin)) && !lockOwnerTx.isRemote();
+ if (thisShouldInterrupt && isDeadLock) {
+ lockOwnerTx.interruptProcessingThread();
+ if (exposeJmxStats) {
+ detectedDeadlocks.incrementAndGet();
+ locallyInterruptedTransactions.incrementAndGet();
+ }
+ return lockForTheRemainingTime(key, lockTimeout, start, now);
+ } else if (!isDeadLock) {
+ return lockForTheRemainingTime(key, lockTimeout, start, now);
+ } else {
+ if (trace)
+ log.trace("Not trying to acquire lock anymore, as we're in deadlock and this will be rollback at origin");
+ if (exposeJmxStats) {
+ detectedDeadlocks.incrementAndGet();
+ }
+ remoteGlobalTransaction.setMarkedForRollback(true);
+ throw new DeadlockDetectedException("Deadlock situation detected on tx: " + remoteTxContext.getLockOwner());
+ }
+ }
+
+ private boolean lockForTheRemainingTime(Object key, long lockTimeout, long start, long now) throws InterruptedException {
+ long remainingLockingTime = (start + lockTimeout) - now;
+ if (remainingLockingTime < 0)
+ throw new IllegalStateException("No remaining time!!! The outer while condition MUST make sure this always stands true!");
+ if (trace) log.trace("trying to lock for the remaining time: " + remainingLockingTime + " millis ");
+ return lockContainer.acquireLock(key, remainingLockingTime, MILLISECONDS);
+ }
+
+ public void setExposeJmxStats(boolean exposeJmxStats) {
+ this.exposeJmxStats = exposeJmxStats;
+ }
+
+ @ManagedAttribute(description = "Number of situtations when we try to determine a deadlock and the other lock owner is e.g. a local tx. In this scenario we cannot run the deadlock detection mechanism")
+ public long getOverlapWithNotDeadlockAwareLockOwners() {
+ return overlapWithNotDeadlockAwareLockOwners.get();
+ }
+
+
+ @ManagedAttribute(description = "Number of locally originated transactions that were interrupted as a deadlock situation was detected")
+ public long getLocallyInterruptedTransactions() {
+ return locallyInterruptedTransactions.get();
+ }
+
+ @ManagedAttribute(description = "Total number of deadlocks detected")
+ public long getDetectedDeadlocks() {
+ return detectedDeadlocks.get();
+ }
+
+ @ManagedOperation(description = "Resets statistics gathered by this component")
+ public void resetStatistics() {
+ overlapWithNotDeadlockAwareLockOwners.set(0);
+ locallyInterruptedTransactions.set(0);
+ detectedDeadlocks.set(0);
+ }
+
+}
Property changes on: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/DeadlockDetectingLockManager.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Modified: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManager.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManager.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManager.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -104,4 +104,9 @@
* @return true if the entry *might* be locked, false if the entry definitely is *not* locked.
*/
boolean possiblyLocked(CacheEntry entry);
+
+ /**
+ * Cleanups the locks within the given context.
+ */
+ void releaseLocks(InvocationContext ctx);
}
Modified: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManagerImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManagerImpl.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManagerImpl.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -58,7 +58,7 @@
private TransactionManager transactionManager;
private InvocationContextContainer invocationContextContainer;
private static final Log log = LogFactory.getLog(LockManagerImpl.class);
- private static final boolean trace = log.isTraceEnabled();
+ protected static final boolean trace = log.isTraceEnabled();
@Inject
public void injectDependencies(Configuration configuration, TransactionManager transactionManager, InvocationContextContainer invocationContextContainer) {
@@ -85,7 +85,7 @@
return false;
}
- private long getLockAcquisitionTimeout(InvocationContext ctx) {
+ protected long getLockAcquisitionTimeout(InvocationContext ctx) {
return ctx.hasFlag(Flag.ZERO_LOCK_ACQUISITION_TIMEOUT) ?
0 : configuration.getLockAcquisitionTimeout();
}
@@ -143,6 +143,32 @@
return entry == null || entry.isChanged() || entry.isNull();
}
+ public void releaseLocks(InvocationContext ctx) {
+ Object owner = ctx.getLockOwner();
+ // clean up.
+ // unlocking needs to be done in reverse order.
+ ReversibleOrderedSet<Map.Entry<Object, CacheEntry>> entries = ctx.getLookedUpEntries().entrySet();
+ Iterator<Map.Entry<Object, CacheEntry>> it = entries.reverseIterator();
+ if (trace) log.trace("Number of entries in context: {0}", entries.size());
+
+ while (it.hasNext()) {
+ Map.Entry<Object, CacheEntry> e = it.next();
+ CacheEntry entry = e.getValue();
+ Object key = e.getKey();
+ boolean needToUnlock = possiblyLocked(entry);
+ // could be null with read-committed
+ if (entry != null && entry.isChanged()) entry.rollback();
+ else {
+ if (trace) log.trace("Entry for key {0} is null, not calling rollbackUpdate", key);
+ }
+ // and then unlock
+ if (needToUnlock) {
+ if (trace) log.trace("Releasing lock on [" + key + "] for owner " + owner);
+ unlock(key, owner);
+ }
+ }
+ }
+
@ManagedAttribute(writable = false, description = "The concurrency level that the MVCC Lock Manager has been configured with.")
public int getConcurrencyLevel() {
return configuration.getConcurrencyLevel();
Modified: trunk/core/src/main/resources/config-samples/all.xml
===================================================================
--- trunk/core/src/main/resources/config-samples/all.xml 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/resources/config-samples/all.xml 2009-07-20 15:35:56 UTC (rev 592)
@@ -107,6 +107,7 @@
-->
<!--<async useReplQueue="true" replQueueInterval="10000" replQueueMaxElements="500"/>-->
</clustering>
+
</default>
<!-- ************************************** -->
@@ -161,6 +162,8 @@
<async enabled="true" batchSize="1000" threadPoolSize="5"/>
</loader>
</loaders>
+
+ <deadlockDetection enabled="true" spinDuration="1000"/>
</namedCache>
Modified: trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd
===================================================================
--- trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd 2009-07-20 15:35:56 UTC (rev 592)
@@ -60,6 +60,7 @@
<xs:element name="loaders" type="tns:loadersType" minOccurs="0" maxOccurs="1"/>
<xs:element name="customInterceptors" type="tns:customInterceptorsType" minOccurs="0" maxOccurs="1"/>
<xs:element name="clustering" type="tns:clusteringType" minOccurs="0" maxOccurs="1"/>
+ <xs:element name="deadlockDetection" type="tns:deadlockDetectionType" minOccurs="0" maxOccurs="1"/>
<xs:element name="unsafe" type="tns:unsafeType" minOccurs="0" maxOccurs="1"/>
</xs:all>
<xs:attribute name="name" type="xs:string"/>
@@ -275,6 +276,11 @@
</xs:sequence>
</xs:complexType>
+ <xs:complexType name="deadlockDetectionType">
+ <xs:attribute name="enabled" type="tns:booleanType"/>
+ <xs:attribute name="spinDuration" type="xs:integer"/>
+ </xs:complexType>
+
<xs:complexType name="propertyType">
<xs:attribute name="name" type="xs:string"/>
<xs:attribute name="value" type="xs:string"/>
Modified: trunk/core/src/test/java/org/infinispan/config/parsing/ConfigurationParserTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/config/parsing/ConfigurationParserTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/config/parsing/ConfigurationParserTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -110,6 +110,18 @@
assert c.isInvocationBatchingEnabled();
}
+ public void testDeadlockDetection() throws Exception {
+ XmlConfigurationParserImpl parser = new XmlConfigurationParserImpl();
+ String xml = "<deadlockDetection enabled=\"true\" spinDuration=\"123\"/>";
+ Element e = XmlConfigHelper.stringToElement(xml);
+
+ Configuration c = new Configuration();
+ parser.configureDeadlockDetection(e, c);
+
+ assert c.isEnableDeadlockDetection();
+ assert c.getDeadlockDetectionSpinDuration() == 123;
+ }
+
public void testInvocationBatchingDefaults() throws Exception {
XmlConfigurationParserImpl parser = new XmlConfigurationParserImpl();
String xml = "<invocationBatching />";
Modified: trunk/core/src/test/java/org/infinispan/config/parsing/XmlFileParsingTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/config/parsing/XmlFileParsingTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/config/parsing/XmlFileParsingTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -154,6 +154,10 @@
assert c.getEvictionStrategy().equals(EvictionStrategy.FIFO);
assert c.getExpirationLifespan() == 60000;
assert c.getExpirationMaxIdle() == 1000;
+
+ c = namedCaches.get("withDeadlockDetection");
+ assert c.isEnableDeadlockDetection();
+ assert c.getDeadlockDetectionSpinDuration() == 1221;
}
private void testConfigurationMerging(XmlConfigurationParser parser) throws IOException {
Modified: trunk/core/src/test/java/org/infinispan/distribution/DefaultConsistentHashTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/distribution/DefaultConsistentHashTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/distribution/DefaultConsistentHashTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -107,4 +107,8 @@
public int hashCode() {
return addressNum;
}
+
+ public int compareTo(Object o) {
+ return this.addressNum - ((TestAddress)o).addressNum;
+ }
}
Modified: trunk/core/src/test/java/org/infinispan/loaders/BaseCacheStoreTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/loaders/BaseCacheStoreTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/loaders/BaseCacheStoreTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -13,6 +13,7 @@
import org.infinispan.marshall.Marshaller;
import org.infinispan.marshall.TestObjectStreamMarshaller;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import org.infinispan.util.Util;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
@@ -39,6 +40,8 @@
protected CacheStore cs;
+ protected GlobalTransactionFactory gtf = new GlobalTransactionFactory();
+
@BeforeMethod
public void setUp() throws Exception {
try {
@@ -204,7 +207,7 @@
mods.add(new Store(InternalEntryFactory.create("k1", "v1")));
mods.add(new Store(InternalEntryFactory.create("k2", "v2")));
mods.add(new Remove("k1"));
- GlobalTransaction tx = new GlobalTransaction(true);
+ GlobalTransaction tx = gtf.newGlobalTransaction(null, true);
cs.prepare(mods, tx, true);
assert cs.load("k2").getValue().equals("v2");
@@ -229,7 +232,7 @@
mods.add(new Store(InternalEntryFactory.create("k1", "v1")));
mods.add(new Store(InternalEntryFactory.create("k2", "v2")));
mods.add(new Remove("k1"));
- GlobalTransaction tx = new GlobalTransaction(false);
+ GlobalTransaction tx = gtf.newGlobalTransaction(null, false);
cs.prepare(mods, tx, false);
assert !cs.containsKey("k1");
@@ -271,7 +274,7 @@
mods.add(new Store(InternalEntryFactory.create("k2", "v2")));
mods.add(new Remove("k1"));
mods.add(new Remove("old"));
- GlobalTransaction tx = new GlobalTransaction(false);
+ GlobalTransaction tx = gtf.newGlobalTransaction(null, false);
cs.prepare(mods, tx, false);
assert !cs.containsKey("k1");
@@ -313,7 +316,7 @@
mods.add(new Store(InternalEntryFactory.create("k2", "v2")));
mods.add(new Remove("k1"));
mods.add(new Remove("old"));
- final GlobalTransaction tx = new GlobalTransaction(false);
+ final GlobalTransaction tx = gtf.newGlobalTransaction(null, false);
cs.prepare(mods, tx, false);
Thread t = new Thread(new Runnable() {
@@ -353,7 +356,7 @@
public void testCommitAndRollbackWithoutPrepare() throws CacheLoaderException {
cs.store(InternalEntryFactory.create("old", "old"));
- GlobalTransaction tx = new GlobalTransaction(false);
+ GlobalTransaction tx = gtf.newGlobalTransaction(null, false);
cs.commit(tx);
cs.store(InternalEntryFactory.create("old", "old"));
cs.rollback(tx);
Modified: trunk/core/src/test/java/org/infinispan/loaders/decorators/ChainingCacheLoaderTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/loaders/decorators/ChainingCacheLoaderTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/loaders/decorators/ChainingCacheLoaderTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -15,9 +15,9 @@
import org.infinispan.loaders.modifications.Store;
import org.infinispan.marshall.TestObjectStreamMarshaller;
import org.infinispan.transaction.xa.GlobalTransaction;
+import static org.testng.Assert.assertEquals;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.Test;
-import static org.testng.Assert.assertEquals;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
@@ -175,7 +175,7 @@
list.add(new Store(InternalEntryFactory.create("k5", "v5", lifespan)));
list.add(new Store(InternalEntryFactory.create("k6", "v6")));
list.add(new Remove("k6"));
- GlobalTransaction t = new GlobalTransaction(false);
+ GlobalTransaction t = gtf.newGlobalTransaction(null, false);
cs.prepare(list, t, true);
CacheStore[] allStores = new CacheStore[]{cs, store1, store2}; // for iteration
@@ -210,7 +210,7 @@
list.add(new Store(InternalEntryFactory.create("k5", "v5", lifespan)));
list.add(new Store(InternalEntryFactory.create("k6", "v6")));
list.add(new Remove("k6"));
- GlobalTransaction tx = new GlobalTransaction(false);
+ GlobalTransaction tx = gtf.newGlobalTransaction(null, false);
cs.prepare(list, tx, false);
CacheStore[] allStores = new CacheStore[]{cs, store1, store2}; // for iteration
Modified: trunk/core/src/test/java/org/infinispan/marshall/MarshallersTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/marshall/MarshallersTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/marshall/MarshallersTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -56,6 +56,7 @@
import org.infinispan.remoting.transport.jgroups.JGroupsAddress;
import org.infinispan.transaction.TransactionLog;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import org.infinispan.util.FastCopyHashMap;
import org.infinispan.util.Immutables;
import org.jgroups.stack.IpAddress;
@@ -77,6 +78,7 @@
public class MarshallersTest {
private final MarshallerImpl home = new MarshallerImpl();
+ private GlobalTransactionFactory gtf = new GlobalTransactionFactory();
private final JBossMarshaller jboss = new JBossMarshaller();
private final Marshaller[] marshallers = new Marshaller[] {home, jboss};
@@ -97,7 +99,7 @@
}
public void testGlobalTransactionMarshalling() throws Exception {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
checkEqualityAndSize(gtx);
}
@@ -105,7 +107,7 @@
List l1 = new ArrayList();
List l2 = new LinkedList();
for (int i = 0; i < 10; i++) {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
l1.add(gtx);
l2.add(gtx);
}
@@ -119,7 +121,7 @@
Map m3 = new HashMap();
Map<Integer, GlobalTransaction> m4 = new FastCopyHashMap<Integer, GlobalTransaction>();
for (int i = 0; i < 10; i++) {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
m1.put(1000 * i, gtx);
m2.put(1000 * i, gtx);
m4.put(1000 * i, gtx);
@@ -155,20 +157,20 @@
}
public void testMarshalledValueMarshalling() throws Exception {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
int bytesH = marshallAndAssertEquality(home, new MarshalledValue(gtx, true, home));
int bytesJ = marshallAndAssertEquality(jboss, new MarshalledValue(gtx, true, jboss));
assert bytesJ < bytesH : "JBoss Marshaller should write less bytes: bytesJBoss=" + bytesJ + ", bytesHome=" + bytesH;
}
public void testSingletonListMarshalling() throws Exception {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
List l = Collections.singletonList(gtx);
checkEqualityAndSize(l);
}
public void testTransactionLogMarshalling() throws Exception {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
PutKeyValueCommand command = new PutKeyValueCommand("k", "v", false, null, 0, 0);
TransactionLog.LogEntry entry = new TransactionLog.LogEntry(gtx, command);
@@ -282,14 +284,14 @@
Map m1 = new HashMap();
for (int i = 0; i < 10; i++) {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
m1.put(1000 * i, gtx);
}
PutMapCommand c10 = new PutMapCommand(m1, null, 0, 0);
checkEqualityAndSize(c10);
Address local = new JGroupsAddress(new IpAddress(12345));
- GlobalTransaction gtx = new GlobalTransaction(local, false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(local, false);
PrepareCommand c11 = new PrepareCommand(gtx, true, c5, c6, c8, c10);
checkEqualityAndSize(c11);
Modified: trunk/core/src/test/java/org/infinispan/marshall/jboss/JBossMarshallerTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/marshall/jboss/JBossMarshallerTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/marshall/jboss/JBossMarshallerTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -59,6 +59,7 @@
import org.infinispan.statetransfer.Person;
import org.infinispan.transaction.TransactionLog;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.GlobalTransactionFactory;
import org.infinispan.util.FastCopyHashMap;
import org.infinispan.util.Immutables;
import org.infinispan.util.concurrent.TimeoutException;
@@ -87,6 +88,8 @@
// private final JBossMarshaller marshaller = new JBossMarshaller();
private final VersionAwareMarshaller marshaller = new VersionAwareMarshaller();
+ private GlobalTransactionFactory gtf = new GlobalTransactionFactory();
+
@BeforeTest
public void setUp() {
marshaller.inject(Thread.currentThread().getContextClassLoader(), new RemoteCommandFactory());
@@ -105,7 +108,7 @@
public void testGlobalTransactionMarshalling() throws Exception {
JGroupsAddress jGroupsAddress = new JGroupsAddress(new IpAddress(12345));
- GlobalTransaction gtx = new GlobalTransaction(jGroupsAddress, false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(jGroupsAddress, false);
marshallAndAssertEquality(gtx);
}
@@ -114,7 +117,7 @@
List l2 = new LinkedList();
for (int i = 0; i < 10; i++) {
JGroupsAddress jGroupsAddress = new JGroupsAddress(new IpAddress(1000 * i));
- GlobalTransaction gtx = new GlobalTransaction(jGroupsAddress, false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(jGroupsAddress, false);
l1.add(gtx);
l2.add(gtx);
}
@@ -129,7 +132,7 @@
Map<Integer, GlobalTransaction> m4 = new FastCopyHashMap<Integer, GlobalTransaction>();
for (int i = 0; i < 10; i++) {
JGroupsAddress jGroupsAddress = new JGroupsAddress(new IpAddress(1000 * i));
- GlobalTransaction gtx = new GlobalTransaction(jGroupsAddress, false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(jGroupsAddress, false);
m1.put(1000 * i, gtx);
m2.put(1000 * i, gtx);
m4.put(1000 * i, gtx);
@@ -175,13 +178,13 @@
}
public void testSingletonListMarshalling() throws Exception {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
List l = Collections.singletonList(gtx);
marshallAndAssertEquality(l);
}
public void testTransactionLogMarshalling() throws Exception {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(12345)), false);
PutKeyValueCommand command = new PutKeyValueCommand("k", "v", false, null, 0, 0);
TransactionLog.LogEntry entry = new TransactionLog.LogEntry(gtx, command);
byte[] bytes = marshaller.objectToByteBuffer(entry);
@@ -257,14 +260,14 @@
Map m1 = new HashMap();
for (int i = 0; i < 10; i++) {
- GlobalTransaction gtx = new GlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(new JGroupsAddress(new IpAddress(1000 * i)), false);
m1.put(1000 * i, gtx);
}
PutMapCommand c10 = new PutMapCommand(m1, null, 0, 0);
marshallAndAssertEquality(c10);
Address local = new JGroupsAddress(new IpAddress(12345));
- GlobalTransaction gtx = new GlobalTransaction(local, false);
+ GlobalTransaction gtx = gtf.newGlobalTransaction(local, false);
PrepareCommand c11 = new PrepareCommand(gtx, true, c5, c6, c8, c10);
marshallAndAssertEquality(c11);
Added: trunk/core/src/test/java/org/infinispan/profiling/DeadlockDetectionPerformanceTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/profiling/DeadlockDetectionPerformanceTest.java (rev 0)
+++ trunk/core/src/test/java/org/infinispan/profiling/DeadlockDetectionPerformanceTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,206 @@
+package org.infinispan.profiling;
+
+import org.infinispan.Cache;
+import org.infinispan.config.Configuration;
+import org.infinispan.manager.CacheManager;
+import org.infinispan.test.TestingUtil;
+import org.infinispan.test.fwk.TestCacheManagerFactory;
+import org.infinispan.transaction.lookup.DummyTransactionManagerLookup;
+import org.testng.annotations.BeforeTest;
+import org.testng.annotations.Test;
+
+import javax.transaction.TransactionManager;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Random;
+import java.util.Set;
+import java.util.concurrent.CountDownLatch;
+
+/**
+ * Test for benchmarking the performance of deadlock detection code. It startes multiple threads that operate on the
+ * same
+ *
+ * @author Mircea.Markus at jboss.com
+ */
+ at Test (groups = "profiling", enabled = true, testName = "profiling.DeadlockDetectionPerformanceTest")
+public class DeadlockDetectionPerformanceTest {
+
+ public static final int KEY_POOL_SIZE = 10;
+
+ public static int TX_SIZE = 5;
+
+ public static final int THREAD_COUNT = 3;
+
+ public static final long BENCHMARK_DURATION = 60000;
+
+ public static boolean USE_DLD = true;
+
+ public static List<String> keyPool;
+
+ @BeforeTest
+ public void generateKeyPool() {
+ keyPool = new ArrayList<String>();
+ for (int i = 0; i < KEY_POOL_SIZE; i++) {
+ keyPool.add("key" + i);
+ }
+ }
+
+ @Test(invocationCount = 10, enabled = false)
+ public void testLocalDifferentTxSize() throws Exception {
+ USE_DLD = false;
+ for (int i = 2; i < KEY_POOL_SIZE; i++) {
+ TX_SIZE = i;
+ runLocalTest();
+ }
+ USE_DLD = true;
+ for (int i = 2; i < KEY_POOL_SIZE; i++) {
+ TX_SIZE = i;
+ runLocalTest();
+ }
+ }
+
+ @Test(invocationCount = 10, enabled = false)
+ public void testReplDifferentTxSize() throws Exception {
+ USE_DLD = false;
+ for (int i = 2; i < KEY_POOL_SIZE; i++) {
+ TX_SIZE = i;
+ runDistributedTest();
+ }
+ USE_DLD = true;
+ for (int i = 2; i < KEY_POOL_SIZE; i++) {
+ TX_SIZE = i;
+ runDistributedTest();
+ }
+ }
+
+ private void runDistributedTest() throws Exception {
+ CacheManager cm = null;
+ List<CacheManager> managers = new ArrayList<CacheManager>();
+ try {
+ CountDownLatch startLatch = new CountDownLatch(1);
+ List<ExecutorThread> executorThreads = new ArrayList<ExecutorThread>();
+ for (int i = 0; i < THREAD_COUNT; i++) {
+ cm = TestCacheManagerFactory.createClusteredCacheManager();
+ Configuration configuration = getConfiguration();
+ configuration.setCacheMode(Configuration.CacheMode.REPL_SYNC);
+ cm.defineCache("test", configuration);
+ Cache distCache = cm.getCache("test");
+ ExecutorThread executorThread = new ExecutorThread(startLatch, distCache);
+ executorThreads.add(executorThread);
+ managers.add(cm);
+ }
+ TestingUtil.blockUntilViewsReceived(10000, managers.toArray(new CacheManager[managers.size()]));
+ startLatch.countDown();
+ Thread.sleep(BENCHMARK_DURATION);
+ joinThreadsAndPrintResult(executorThreads);
+ } finally {
+ TestingUtil.killCacheManagers(managers);
+ }
+ }
+
+ private void runLocalTest() throws Exception {
+ CacheManager cm = TestCacheManagerFactory.createLocalCacheManager();
+ try {
+ Configuration configuration = getConfiguration();
+ cm.defineCache("test", configuration);
+ Cache localCache = cm.getCache("test");
+
+ CountDownLatch startLatch = new CountDownLatch(1);
+
+ List<ExecutorThread> executorThreads = new ArrayList<ExecutorThread>();
+ for (int i = 0; i < THREAD_COUNT; i++) {
+ ExecutorThread executorThread = new ExecutorThread(startLatch, localCache);
+ executorThreads.add(executorThread);
+ }
+ startLatch.countDown();
+ Thread.sleep(BENCHMARK_DURATION);
+ joinThreadsAndPrintResult(executorThreads);
+ } finally {
+ TestingUtil.killCacheManagers(cm);
+ }
+ }
+
+ private void joinThreadsAndPrintResult(List<ExecutorThread> executorThreads) throws InterruptedException {
+ int totalSuccess = 0;
+ int totalFailures = 0;
+ for (int i = 0; i < THREAD_COUNT; i++) {
+ ExecutorThread executorThread = executorThreads.get(i);
+ executorThread.join();
+ totalSuccess += executorThread.getSuccessfullTx();
+ totalFailures += executorThread.getFailedTx();
+ }
+ System.out.println("Use DDL? " + USE_DLD + " TX_SIZE = " + TX_SIZE + " totalSuccess = " + totalSuccess);
+ System.out.println("Use DDL? " + USE_DLD + " TX_SIZE = " + TX_SIZE + " totalFailures = " + totalFailures);
+ System.out.println("-------------------------------");
+ }
+
+ private Configuration getConfiguration() {
+ Configuration configuration = new Configuration();
+ configuration.setTransactionManagerLookupClass(DummyTransactionManagerLookup.class.getName());
+ configuration.setEnableDeadlockDetection(USE_DLD);
+ configuration.setUseLockStriping(false);
+ return configuration;
+ }
+
+ public static class ExecutorThread extends Thread {
+ private volatile CountDownLatch startLatch;
+ private volatile int successfullTx;
+ private volatile int failedTx;
+ private volatile Cache cache;
+ private volatile TransactionManager txManager;
+ static int TX_INDEX = 0;
+
+ public ExecutorThread(CountDownLatch startLatch, Cache cache) {
+ super("EXECUTOR-THREAD-" + TX_INDEX++);
+ this.startLatch = startLatch;
+ this.cache = cache;
+ txManager = TestingUtil.getTransactionManager(cache);
+ start();
+ }
+
+ @Override
+ public void run() {
+ long start = System.currentTimeMillis();
+ try {
+ startLatch.await();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ while ((start + BENCHMARK_DURATION) - System.currentTimeMillis() > 0) {
+ try {
+ txManager.begin();
+ List<String> keysToUpdate = getKeysPerTx();
+ for (String key : keysToUpdate) {
+ cache.put(key, "value");
+ }
+ txManager.commit();
+ successfullTx++;
+ } catch (Throwable e) {
+ failedTx++;
+ }
+ }
+ }
+
+ public int getFailedTx() {
+ return failedTx;
+ }
+
+ public int getSuccessfullTx() {
+ return successfullTx;
+ }
+ }
+
+ private static List<String> getKeysPerTx() {
+ Random rnd = new Random();
+ Set<String> result = new HashSet<String>();
+ while (result.size() < TX_SIZE) {
+ String key = keyPool.get(rnd.nextInt(KEY_POOL_SIZE));
+ result.add(key);
+ }
+ ArrayList resultList = new ArrayList(result);
+ Collections.shuffle(resultList);
+ return resultList;
+ }
+}
Property changes on: trunk/core/src/test/java/org/infinispan/profiling/DeadlockDetectionPerformanceTest.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Added: trunk/core/src/test/java/org/infinispan/tx/DeadlockDetectionTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/tx/DeadlockDetectionTest.java (rev 0)
+++ trunk/core/src/test/java/org/infinispan/tx/DeadlockDetectionTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -0,0 +1,558 @@
+package org.infinispan.tx;
+
+import org.infinispan.Cache;
+import org.infinispan.api.mvcc.LockAssert;
+import org.infinispan.commands.ReplicableCommand;
+import org.infinispan.config.Configuration;
+import org.infinispan.config.ConfigurationException;
+import org.infinispan.context.impl.NonTxInvocationContext;
+import org.infinispan.interceptors.DeadlockDetectingInterceptor;
+import org.infinispan.interceptors.InterceptorChain;
+import org.infinispan.manager.CacheManager;
+import org.infinispan.remoting.ReplicationException;
+import org.infinispan.remoting.responses.Response;
+import org.infinispan.remoting.rpc.ResponseFilter;
+import org.infinispan.remoting.rpc.ResponseMode;
+import org.infinispan.remoting.rpc.RpcManager;
+import org.infinispan.remoting.transport.Address;
+import org.infinispan.remoting.transport.Transport;
+import org.infinispan.statetransfer.StateTransferException;
+import org.infinispan.test.MultipleCacheManagersTest;
+import org.infinispan.test.TestingUtil;
+import org.infinispan.test.fwk.TestCacheManagerFactory;
+import org.infinispan.transaction.lookup.DummyTransactionManagerLookup;
+import org.infinispan.util.concurrent.NotifyingNotifiableFuture;
+import org.infinispan.util.concurrent.locks.DeadlockDetectingLockManager;
+import org.infinispan.util.concurrent.locks.LockManager;
+import org.infinispan.util.concurrent.locks.DeadlockDetectedException;
+import org.infinispan.util.logging.Log;
+import org.infinispan.util.logging.LogFactory;
+import static org.testng.Assert.assertEquals;
+import org.testng.annotations.AfterMethod;
+import org.testng.annotations.BeforeMethod;
+import org.testng.annotations.Test;
+
+import javax.transaction.TransactionManager;
+import javax.transaction.RollbackException;
+import java.util.List;
+import java.util.concurrent.ArrayBlockingQueue;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.CountDownLatch;
+
+/**
+ * Functional test for deadlock detection.
+ *
+ * @author Mircea.Markus at jboss.com
+ * <p/>
+ * TODO - test for deadlock on invalidation
+ * TODO - add test deadlock with distribution
+ */
+ at Test(testName = "tx.DeadlockDetectionTest", groups = "functional")
+public class DeadlockDetectionTest extends MultipleCacheManagersTest {
+
+ private ControlledRpcManager controlledRpcManager1;
+ private ControlledRpcManager controlledRpcManager2;
+ private CountDownLatch replicationLatch;
+ private ExecutorThread t1;
+ private ExecutorThread t2;
+ private DeadlockDetectingLockManager ddLm1;
+ private DeadlockDetectingLockManager ddLm2;
+
+
+ protected void createCacheManagers() throws Throwable {
+ Configuration config = getDefaultClusteredConfig(Configuration.CacheMode.REPL_SYNC);
+ config.setTransactionManagerLookupClass(DummyTransactionManagerLookup.class.getName());
+ config.setEnableDeadlockDetection(true);
+ config.setSyncCommitPhase(true);
+ config.setSyncRollbackPhase(true);
+ config.setUseLockStriping(false);
+ assert config.isEnableDeadlockDetection();
+ createClusteredCaches(2, "test", config);
+ assert config.isEnableDeadlockDetection();
+
+ assert cache(0, "test").getConfiguration().isEnableDeadlockDetection();
+ assert cache(1, "test").getConfiguration().isEnableDeadlockDetection();
+ assert !cache(0, "test").getConfiguration().isExposeJmxStatistics();
+ assert !cache(1, "test").getConfiguration().isExposeJmxStatistics();
+
+ ((DeadlockDetectingLockManager) TestingUtil.extractLockManager(cache(0, "test"))).setExposeJmxStats(true);
+ ((DeadlockDetectingLockManager) TestingUtil.extractLockManager(cache(1, "test"))).setExposeJmxStats(true);
+
+ RpcManager rpcManager1 = TestingUtil.extractComponent(cache(0, "test"), RpcManager.class);
+ RpcManager rpcManager2 = TestingUtil.extractComponent(cache(1, "test"), RpcManager.class);
+
+ controlledRpcManager1 = new ControlledRpcManager(rpcManager1);
+ controlledRpcManager2 = new ControlledRpcManager(rpcManager2);
+ TestingUtil.replaceComponent(cache(0, "test"), RpcManager.class, controlledRpcManager1, true);
+ TestingUtil.replaceComponent(cache(1, "test"), RpcManager.class, controlledRpcManager2, true);
+
+ assert TestingUtil.extractComponent(cache(0, "test"), RpcManager.class) instanceof ControlledRpcManager;
+ assert TestingUtil.extractComponent(cache(1, "test"), RpcManager.class) instanceof ControlledRpcManager;
+
+ ddLm1 = (DeadlockDetectingLockManager) TestingUtil.extractLockManager(cache(0, "test"));
+ ddLm2 = (DeadlockDetectingLockManager) TestingUtil.extractLockManager(cache(1, "test"));
+ }
+
+
+ @BeforeMethod
+ public void beforeMethod() {
+ t1 = new ExecutorThread(cache(0, "test"), 1);
+ t2 = new ExecutorThread(cache(1, "test"), 2);
+ replicationLatch = new CountDownLatch(1);
+ controlledRpcManager1.setReplicationLatch(replicationLatch);
+ controlledRpcManager2.setReplicationLatch(replicationLatch);
+ log.trace("_________________________ Here is beggins");
+ }
+
+ @AfterMethod
+ public void afterMethod() {
+ t1.stopThread();
+ t2.stopThread();
+ ((DeadlockDetectingLockManager) TestingUtil.extractLockManager(cache(0, "test"))).resetStatistics();
+ ((DeadlockDetectingLockManager) TestingUtil.extractLockManager(cache(1, "test"))).resetStatistics();
+ }
+
+ public void testDeadlockDetectionAndAsyncCaches() {
+ Configuration config = getDefaultClusteredConfig(Configuration.CacheMode.REPL_ASYNC);
+ config.setEnableDeadlockDetection(true);
+ config.setUseLockStriping(false);
+ CacheManager cm = TestCacheManagerFactory.createClusteredCacheManager();
+ cm.defineCache("test", config);
+ try {
+ cm.getCache("test");
+ assert false : "Exception expected";
+ } catch (ConfigurationException e) {
+ //expected
+ System.out.println("Error message is " + e.getMessage());
+ }
+ cm.stop();
+ }
+
+ public void testExpectedInnerStructure() {
+ LockManager lockManager = TestingUtil.extractComponent(cache(0, "test"), LockManager.class);
+ assert lockManager instanceof DeadlockDetectingLockManager;
+
+ InterceptorChain ic = TestingUtil.extractComponent(cache(0, "test"), InterceptorChain.class);
+ assert ic.containsInterceptorType(DeadlockDetectingInterceptor.class);
+ }
+
+ public void testDeadlockDetectedTwoTransactions() throws Exception {
+ t1.setKeyValue("key", "value1");
+ t2.setKeyValue("key", "value2");
+ assert OperationsResult.BEGGIN_TX_OK == t1.execute(Operations.BEGGIN_TX);
+ assert OperationsResult.BEGGIN_TX_OK == t2.execute(Operations.BEGGIN_TX);
+ System.out.println("After beggin");
+
+ t1.execute(Operations.PUT_KEY_VALUE);
+ t2.execute(Operations.PUT_KEY_VALUE);
+ System.out.println("After put key value");
+
+ t1.clearResponse();
+ t2.clearResponse();
+
+ t1.executeNoResponse(Operations.COMMIT_TX);
+ t2.executeNoResponse(Operations.COMMIT_TX);
+
+ System.out.println("Now replication is triggered");
+ replicationLatch.countDown();
+
+
+ Object t1Commit = t1.waitForResponse();
+ Object t2Commit = t2.waitForResponse();
+ System.out.println("After commit: " + t1Commit + ", " + t2Commit);
+
+ assert xor(t1Commit instanceof Exception, t2Commit instanceof Exception) : "only one thread must be failing " + t1Commit + "," + t2Commit;
+ System.out.println("t2Commit = " + t2Commit);
+ System.out.println("t1Commit = " + t1Commit);
+
+ if (t1Commit instanceof Exception) {
+ System.out.println("t1 rolled back");
+ Object o = cache(0, "test").get("key");
+ assert o != null;
+ assert o.equals("value2");
+ } else {
+ System.out.println("t2 rolled back");
+ Object o = cache(0, "test").get("key");
+ assert o != null;
+ assert o.equals("value1");
+ o = cache(1, "test").get("key");
+ assert o != null;
+ assert o.equals("value1");
+ }
+
+ assert ddLm1.getDetectedDeadlocks() + ddLm2.getDetectedDeadlocks() >= 1;
+
+ LockManager lm1 = TestingUtil.extractComponent(cache(0, "test"), LockManager.class);
+ assert !lm1.isLocked("key") : "It is locked by " + lm1.getOwner("key");
+ LockManager lm2 = TestingUtil.extractComponent(cache(1, "test"), LockManager.class);
+ assert !lm2.isLocked("key") : "It is locked by " + lm2.getOwner("key");
+ LockAssert.assertNoLocks(cache(0, "test"));
+ }
+
+ public void testLocalVsLocalTxDeadlock() {
+ CacheManager cm = null;
+ try {
+ cm = TestCacheManagerFactory.createLocalCacheManager();
+ Configuration configuration = new Configuration();
+ configuration.setTransactionManagerLookupClass(DummyTransactionManagerLookup.class.getName());
+ configuration.setEnableDeadlockDetection(true);
+ configuration.setUseLockStriping(false);
+ configuration.setExposeJmxStatistics(true);
+ cm.defineCache("test", configuration);
+ Cache localCache = cm.getCache("test");
+ DeadlockDetectingLockManager lockManager = (DeadlockDetectingLockManager) TestingUtil.extractLockManager(localCache);
+
+ ExecutorThread t1 = new ExecutorThread(localCache, 0);
+ ExecutorThread t2 = new ExecutorThread(localCache, 1);
+
+
+ assert OperationsResult.BEGGIN_TX_OK == t1.execute(Operations.BEGGIN_TX);
+ assert OperationsResult.BEGGIN_TX_OK == t2.execute(Operations.BEGGIN_TX);
+ System.out.println("After beggin");
+
+ t1.setKeyValue("k1", "value_1_t1");
+ t2.setKeyValue("k2", "value_2_t2");
+
+ assert OperationsResult.PUT_KEY_VALUE_OK == t1.execute(Operations.PUT_KEY_VALUE);
+ assert OperationsResult.PUT_KEY_VALUE_OK == t2.execute(Operations.PUT_KEY_VALUE);
+
+ System.out.println("After first PUT");
+ assert lockManager.isLocked("k1");
+ assert lockManager.isLocked("k2");
+
+
+ t1.setKeyValue("k2", "value_2_t1");
+ t2.setKeyValue("k1", "value_1_t2");
+ t1.executeNoResponse(Operations.PUT_KEY_VALUE);
+ t2.executeNoResponse(Operations.PUT_KEY_VALUE);
+
+ Object response1 = t1.waitForResponse();
+ Object response2 = t2.waitForResponse();
+
+ assert xor(response1 instanceof DeadlockDetectedException, response2 instanceof DeadlockDetectedException) : "expected one and only one exception: " + response1 + ", " + response2;
+ assert xor(response1 == OperationsResult.PUT_KEY_VALUE_OK, response2 == OperationsResult.PUT_KEY_VALUE_OK) : "expected one and only one exception: " + response1 + ", " + response2;
+
+ assert lockManager.isLocked("k1");
+ assert lockManager.isLocked("k2");
+ assert lockManager.getOwner("k1") == lockManager.getOwner("k2");
+
+ if (response1 instanceof Exception) {
+ assert OperationsResult.COMMIT_TX_OK == t2.execute(Operations.COMMIT_TX);
+ assertEquals("value_1_t2", localCache.get("k1"));
+ assertEquals("value_2_t2", localCache.get("k2"));
+ assert t1.execute(Operations.COMMIT_TX) instanceof RollbackException;
+ } else {
+ assert OperationsResult.COMMIT_TX_OK == t1.execute(Operations.COMMIT_TX);
+ assertEquals("value_1_t1", localCache.get("k1"));
+ assertEquals("value_2_t1", localCache.get("k2"));
+ assert t2.execute(Operations.COMMIT_TX) instanceof RollbackException;
+ }
+ assert lockManager.getNumberOfLocksHeld() == 0;
+ assertEquals(lockManager.getDetectedDeadlocks(), 1);
+ } finally {
+ TestingUtil.killCacheManagers(cm);
+ }
+ }
+
+
+ public void testDeadlockDetectedOneTx() throws Exception {
+ t1.setKeyValue("key", "value1");
+
+ LockManager lm2 = TestingUtil.extractComponent(cache(1, "test"), LockManager.class);
+ NonTxInvocationContext ctx = cache(1, "test").getAdvancedCache().getInvocationContextContainer().createNonTxInvocationContext();
+ lm2.lockAndRecord("key", ctx);
+ assert lm2.isLocked("key");
+
+
+ assert OperationsResult.BEGGIN_TX_OK == t1.execute(Operations.BEGGIN_TX) : "but received " + t1.lastResponse();
+ t1.execute(Operations.PUT_KEY_VALUE);
+
+ t1.clearResponse();
+ t1.executeNoResponse(Operations.COMMIT_TX);
+
+ replicationLatch.countDown();
+ System.out.println("Now replication is triggered");
+
+ t1.waitForResponse();
+
+
+ Object t1CommitRsp = t1.lastResponse();
+
+ assert t1CommitRsp instanceof Exception : "expected exception, received " + t1.lastResponse();
+
+ LockManager lm1 = TestingUtil.extractComponent(cache(0, "test"), LockManager.class);
+ assert !lm1.isLocked("key") : "It is locked by " + lm1.getOwner("key");
+
+ lm2.unlock("key", ctx.getLockOwner());
+ assert !lm2.isLocked("key");
+ assert !lm1.isLocked("key");
+ }
+
+ public void testLockReleasedWhileTryingToAcquire() throws Exception {
+ t1.setKeyValue("key", "value1");
+
+ LockManager lm2 = TestingUtil.extractComponent(cache(1, "test"), LockManager.class);
+ NonTxInvocationContext ctx = cache(1, "test").getAdvancedCache().getInvocationContextContainer().createNonTxInvocationContext();
+ lm2.lockAndRecord("key", ctx);
+ assert lm2.isLocked("key");
+
+
+ assert OperationsResult.BEGGIN_TX_OK == t1.execute(Operations.BEGGIN_TX) : "but received " + t1.lastResponse();
+ t1.execute(Operations.PUT_KEY_VALUE);
+
+ t1.clearResponse();
+ t1.executeNoResponse(Operations.COMMIT_TX);
+
+ replicationLatch.countDown();
+
+ Thread.sleep(3000); //just to make sure the remote tx thread managed to spin around for some times.
+ lm2.unlock("key", ctx.getLockOwner());
+
+ t1.waitForResponse();
+
+
+ Object t1CommitRsp = t1.lastResponse();
+
+ assert t1CommitRsp == OperationsResult.COMMIT_TX_OK : "expected true, received " + t1.lastResponse();
+
+ LockManager lm1 = TestingUtil.extractComponent(cache(0, "test"), LockManager.class);
+ assert !lm1.isLocked("key") : "It is locked by " + lm1.getOwner("key");
+
+ assert !lm2.isLocked("key");
+ assert !lm1.isLocked("key");
+ }
+
+ public static enum Operations {
+ BEGGIN_TX, COMMIT_TX, PUT_KEY_VALUE, STOP_THREAD
+ }
+
+ public static enum OperationsResult {
+ BEGGIN_TX_OK, COMMIT_TX_OK, PUT_KEY_VALUE_OK, STOP_THREAD_OK
+ }
+
+ public static final class ExecutorThread extends Thread {
+
+ private static Log log = LogFactory.getLog(ExecutorThread.class);
+
+ private Cache<Object, Object> cache;
+ private BlockingQueue<Object> toExecute = new ArrayBlockingQueue<Object>(1);
+ private volatile Object response;
+ private CountDownLatch responseLatch = new CountDownLatch(1);
+
+ private volatile Object key, value;
+
+ public void setKeyValue(Object key, Object value) {
+ this.key = key;
+ this.value = value;
+ }
+
+ public ExecutorThread(Cache<Object, Object> cache, int index) {
+ super("ExecutorThread-" + index);
+ this.cache = cache;
+ start();
+ }
+
+ public Object execute(Operations op) {
+ try {
+ responseLatch = new CountDownLatch(1);
+ toExecute.put(op);
+ responseLatch.await();
+ return response;
+ } catch (InterruptedException e) {
+ throw new RuntimeException("Unexpected", e);
+ }
+ }
+
+ public void executeNoResponse(Operations op) {
+ try {
+ responseLatch = null;
+ response = null;
+ toExecute.put(op);
+ } catch (InterruptedException e) {
+ throw new RuntimeException("Unexpected", e);
+ }
+ }
+
+ @Override
+ public void run() {
+ Operations operation;
+ boolean run = true;
+ while (run) {
+ try {
+ operation = (Operations) toExecute.take();
+ } catch (InterruptedException e) {
+ throw new RuntimeException(e);
+ }
+ System.out.println("about to process operation " + operation);
+ switch (operation) {
+ case BEGGIN_TX: {
+ TransactionManager txManager = TestingUtil.getTransactionManager(cache);
+ try {
+ txManager.begin();
+ setResponse(OperationsResult.BEGGIN_TX_OK);
+ } catch (Exception e) {
+ log.trace("Failure on beggining tx", e);
+ setResponse(e);
+ }
+ break;
+ }
+ case COMMIT_TX: {
+ TransactionManager txManager = TestingUtil.getTransactionManager(cache);
+ try {
+ txManager.commit();
+ setResponse(OperationsResult.COMMIT_TX_OK);
+ } catch (Exception e) {
+ log.trace("Exception while committing tx", e);
+ setResponse(e);
+ }
+ break;
+ }
+ case PUT_KEY_VALUE: {
+ try {
+ cache.put(key, value);
+ log.trace("Successfully exucuted putKeyValue(" + key + ", " + value + ")");
+ setResponse(OperationsResult.PUT_KEY_VALUE_OK);
+ } catch (Exception e) {
+ log.trace("Exception while executing putKeyValue(" + key + ", " + value + ")", e);
+ setResponse(e);
+ }
+ break;
+ }
+ case STOP_THREAD: {
+ System.out.println("Exiting...");
+ toExecute = null;
+ run = false;
+ break;
+ }
+ }
+ if (responseLatch != null) responseLatch.countDown();
+ }
+ setResponse("EXIT");
+ }
+
+ private void setResponse(Object e) {
+ log.trace("setResponse to " + e);
+ response = e;
+ }
+
+ public void stopThread() {
+ execute(Operations.STOP_THREAD);
+ while (!this.getState().equals(State.TERMINATED)) {
+ try {
+ Thread.sleep(50);
+ } catch (InterruptedException e) {
+ throw new IllegalStateException(e);
+ }
+ }
+ }
+
+ public Object lastResponse() {
+ return response;
+ }
+
+ public void clearResponse() {
+ response = null;
+ }
+
+ public Object waitForResponse() {
+ while (response == null) {
+ try {
+ Thread.sleep(50);
+ } catch (InterruptedException e) {
+ throw new RuntimeException(e);
+ }
+ }
+ return response;
+ }
+ }
+
+ private boolean xor(boolean b1, boolean b2) {
+ return (b1 || b2) && !(b1 && b2);
+ }
+
+ public static final class ControlledRpcManager implements RpcManager {
+
+ private volatile CountDownLatch replicationLatch;
+
+ public ControlledRpcManager(RpcManager realOne) {
+ this.realOne = realOne;
+ }
+
+ private RpcManager realOne;
+
+ public void setReplicationLatch(CountDownLatch replicationLatch) {
+ this.replicationLatch = replicationLatch;
+ }
+
+ public List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue, ResponseFilter responseFilter) {
+ return realOne.invokeRemotely(recipients, rpcCommand, mode, timeout, usePriorityQueue, responseFilter);
+ }
+
+ public List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue) {
+ return realOne.invokeRemotely(recipients, rpcCommand, mode, timeout, usePriorityQueue);
+ }
+
+ public List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout) throws Exception {
+ return realOne.invokeRemotely(recipients, rpcCommand, mode, timeout);
+ }
+
+ public void retrieveState(String cacheName, long timeout) throws StateTransferException {
+ realOne.retrieveState(cacheName, timeout);
+ }
+
+ public void broadcastRpcCommand(ReplicableCommand rpc, boolean sync) throws ReplicationException {
+ waitFirst();
+ realOne.broadcastRpcCommand(rpc, sync);
+ }
+
+ public void broadcastRpcCommand(ReplicableCommand rpc, boolean sync, boolean usePriorityQueue) throws ReplicationException {
+ waitFirst();
+ realOne.broadcastRpcCommand(rpc, sync, usePriorityQueue);
+ }
+
+ private void waitFirst() {
+ System.out.println(Thread.currentThread().getName() + " -- replication trigger called!");
+ try {
+ replicationLatch.await();
+ } catch (Exception e) {
+ throw new RuntimeException("Unexpected exception!", e);
+ }
+ }
+
+ public void broadcastRpcCommandInFuture(ReplicableCommand rpc, NotifyingNotifiableFuture<Object> future) {
+ realOne.broadcastRpcCommandInFuture(rpc, future);
+ }
+
+ public void broadcastRpcCommandInFuture(ReplicableCommand rpc, boolean usePriorityQueue, NotifyingNotifiableFuture<Object> future) {
+ realOne.broadcastRpcCommandInFuture(rpc, usePriorityQueue, future);
+ }
+
+ public void invokeRemotely(List<Address> recipients, ReplicableCommand rpc, boolean sync) throws ReplicationException {
+ realOne.invokeRemotely(recipients, rpc, sync);
+ }
+
+ public void invokeRemotely(List<Address> recipients, ReplicableCommand rpc, boolean sync, boolean usePriorityQueue) throws ReplicationException {
+ realOne.invokeRemotely(recipients, rpc, sync, usePriorityQueue);
+ }
+
+ public void invokeRemotelyInFuture(List<Address> recipients, ReplicableCommand rpc, NotifyingNotifiableFuture<Object> future) {
+ realOne.invokeRemotelyInFuture(recipients, rpc, future);
+ }
+
+ public void invokeRemotelyInFuture(List<Address> recipients, ReplicableCommand rpc, boolean usePriorityQueue, NotifyingNotifiableFuture<Object> future) {
+ realOne.invokeRemotelyInFuture(recipients, rpc, usePriorityQueue, future);
+ }
+
+ public void invokeRemotelyInFuture(List<Address> recipients, ReplicableCommand rpc, boolean usePriorityQueue, NotifyingNotifiableFuture<Object> future, long timeout) {
+ realOne.invokeRemotelyInFuture(recipients, rpc, usePriorityQueue, future, timeout);
+ }
+
+ public Transport getTransport() {
+ return realOne.getTransport();
+ }
+
+ public Address getCurrentStateTransferSource() {
+ return realOne.getCurrentStateTransferSource();
+ }
+ }
+}
Property changes on: trunk/core/src/test/java/org/infinispan/tx/DeadlockDetectionTest.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Modified: trunk/core/src/test/java/org/infinispan/tx/LocalModeTxTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/tx/LocalModeTxTest.java 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/java/org/infinispan/tx/LocalModeTxTest.java 2009-07-20 15:35:56 UTC (rev 592)
@@ -85,4 +85,5 @@
assert c.get("key").equals("value");
assert !c.isEmpty();
}
+
}
Modified: trunk/core/src/test/resources/configs/named-cache-test.xml
===================================================================
--- trunk/core/src/test/resources/configs/named-cache-test.xml 2009-07-20 14:32:07 UTC (rev 591)
+++ trunk/core/src/test/resources/configs/named-cache-test.xml 2009-07-20 15:35:56 UTC (rev 592)
@@ -127,8 +127,16 @@
</clustering>
<jmxStatistics enabled="false"/>
</namedCache>
+
+ <namedCache name="withDeadlockDetection">
+ <clustering>
+ <sync replTimeout="20000"/>
+ </clustering>
+ <jmxStatistics enabled="false"/>
+ <deadlockDetection enabled="true" spinDuration="1221"/>
+ </namedCache>
+
-
<namedCache name="cacheWithCustomInterceptors">
<!--
More information about the infinispan-commits
mailing list