[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