[infinispan-commits] Infinispan SVN: r1353 - in trunk/core/src: test/java/org/infinispan/loaders/decorators and 1 other directory.

infinispan-commits at lists.jboss.org infinispan-commits at lists.jboss.org
Fri Jan 8 05:41:22 EST 2010


Author: galder.zamarreno at jboss.com
Date: 2010-01-08 05:41:22 -0500 (Fri, 08 Jan 2010)
New Revision: 1353

Modified:
   trunk/core/src/main/java/org/infinispan/loaders/decorators/AsyncStore.java
   trunk/core/src/test/java/org/infinispan/loaders/decorators/AsyncTest.java
Log:
[ISPN-325] (Race condition in AsyncStore) Fixed.

Modified: trunk/core/src/main/java/org/infinispan/loaders/decorators/AsyncStore.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/loaders/decorators/AsyncStore.java	2010-01-07 19:26:12 UTC (rev 1352)
+++ trunk/core/src/main/java/org/infinispan/loaders/decorators/AsyncStore.java	2010-01-08 10:41:22 UTC (rev 1353)
@@ -14,10 +14,12 @@
 import org.infinispan.loaders.modifications.PurgeExpired;
 import org.infinispan.loaders.modifications.Remove;
 import org.infinispan.loaders.modifications.Store;
+import org.infinispan.util.concurrent.locks.containers.ReentrantPerEntryLockContainer;
 import org.infinispan.util.logging.Log;
 import org.infinispan.util.logging.LogFactory;
 
 import java.util.ArrayList;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
@@ -77,6 +79,7 @@
    private final Lock write = mapLock.writeLock();
    private int concurrencyLevel;
    @GuardedBy("mapLock") private ConcurrentMap<Object, Modification> state;
+   private ReleaseAllLockContainer lockContainer;
    
    public AsyncStore(CacheStore delegate, AsyncStoreConfig asyncStoreConfig) {
       super(delegate);
@@ -87,6 +90,7 @@
    public void init(CacheLoaderConfig config, Cache cache, Marshaller m) throws CacheLoaderException {
       super.init(config, cache, m);
       concurrencyLevel = cache == null || cache.getConfiguration() == null ? 16 : cache.getConfiguration().getConcurrencyLevel();
+      lockContainer = new ReleaseAllLockContainer(concurrencyLevel);
    }
 
    @Override
@@ -167,13 +171,13 @@
                super.purgeExpired();
                break;
          }
-      }      
+      }
    }
-   
+
    protected Runnable createAsyncProcessor() {
       return new AsyncProcessor();
    }
-   
+
    private void enqueue(Object key, Modification mod) {
       try {
          if (stopped.get()) {
@@ -248,7 +252,8 @@
     */
    class AsyncProcessor implements Runnable {
       private ConcurrentMap<Object, Modification> swap = newStateMap();
-      
+      private final Set<Object> lockedKeys = new HashSet<Object>();
+
       public void run() {
          while (!Thread.interrupted()) {
             try {
@@ -272,23 +277,43 @@
       void run0() throws InterruptedException {
          if (trace) log.trace("Checking for modifications");
          boolean unlock = false;
+         
          try {
             acquireLock(write);
             unlock = true;
             swap = state;
             state = newStateMap();
+
+            // This needs doing within the WL section, because if a key is in use, we need to put it back in the state
+            // map for later processing and we don't wanna do it in such way that we override a newer value that might 
+            // have been enqueued by a user thread.
+            for (Object key : swap.keySet()) {
+               boolean acquired = lockContainer.acquireLock(key, 0, TimeUnit.NANOSECONDS);
+               if (trace) log.trace("Lock for key {0} was acquired={1}", key, acquired);
+               if (!acquired) {
+                  Modification prev = swap.remove(key);
+                  state.put(key, prev);
+               } else {
+                  lockedKeys.add(key);
+               }
+            }
          } finally {
             if (unlock) write.unlock();
          }
-         
-         int size = swap.size();
-         if (size == 0) 
-            awaitNotEmpty();
-         else 
-            decrementAndGet(size);
 
-         if (trace) log.trace("Calling put(List) with {0} modifications", size);
-         put(swap);
+         try {
+            int size = swap.size();
+            if (size == 0) 
+               awaitNotEmpty();
+            else 
+               decrementAndGet(size);
+
+            if (trace) log.trace("Apply {0} modifications", size);
+            put(swap);
+         } finally {
+            lockContainer.releaseLocks(lockedKeys);
+            lockedKeys.clear();
+         }
       }
       
       void put(ConcurrentMap<Object, Modification> mods) {
@@ -304,4 +329,17 @@
    private ConcurrentMap<Object, Modification> newStateMap() {
       return new ConcurrentHashMap<Object, Modification>(64, 0.75f, concurrencyLevel);
    }
+
+   private static class ReleaseAllLockContainer extends ReentrantPerEntryLockContainer {
+      private ReleaseAllLockContainer(int concurrencyLevel) {
+         super(concurrencyLevel);
+      }
+
+      void releaseLocks(Set<Object> keys) {
+         for (Object key : keys) {
+            if (trace) log.trace("Release lock for key {0}", key);
+            releaseLock(key);
+         }
+      }
+   }
 }

Modified: trunk/core/src/test/java/org/infinispan/loaders/decorators/AsyncTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/loaders/decorators/AsyncTest.java	2010-01-07 19:26:12 UTC (rev 1352)
+++ trunk/core/src/test/java/org/infinispan/loaders/decorators/AsyncTest.java	2010-01-08 10:41:22 UTC (rev 1353)
@@ -4,7 +4,9 @@
 import org.infinispan.container.entries.InternalCacheEntry;
 import org.infinispan.container.entries.InternalEntryFactory;
 import org.infinispan.loaders.CacheLoaderException;
+import org.infinispan.loaders.CacheStore;
 import org.infinispan.loaders.dummy.DummyInMemoryCacheStore;
+import org.infinispan.loaders.modifications.Modification;
 import org.infinispan.test.AbstractInfinispanTest;
 import org.infinispan.test.TestingUtil;
 import org.infinispan.util.logging.Log;
@@ -13,7 +15,11 @@
 import org.testng.annotations.BeforeTest;
 import org.testng.annotations.Test;
 
+import java.lang.reflect.Method;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ExecutorService;
+import java.util.concurrent.TimeUnit;
 
 @Test(groups = "unit", testName = "loaders.decorators.AsyncTest")
 public class AsyncTest extends AbstractInfinispanTest {
@@ -91,7 +97,27 @@
       store.start();
       doTestRemove(number, key);
    }
-   
+
+   public void testThreadSafetyWritingDiffValuesForKey(Method m) throws Exception {
+      final String key = "k1";
+      final CountDownLatch v1Latch = new CountDownLatch(1);
+      final CountDownLatch v2Latch = new CountDownLatch(1);
+      final CountDownLatch endLatch = new CountDownLatch(1);
+      DummyInMemoryCacheStore underlying = new DummyInMemoryCacheStore();
+      store = new MockAsyncStore(key, v1Latch, v2Latch, endLatch, underlying, asyncConfig);
+      dummyCfg = new DummyInMemoryCacheStore.Cfg();
+      dummyCfg.setStore(m.getName());
+      store.init(dummyCfg, null, null);
+      store.start();
+      
+      store.store(InternalEntryFactory.create(key, "v1"));
+      v2Latch.await();
+      store.store(InternalEntryFactory.create(key, "v2"));
+      endLatch.await();
+
+      assert store.load(key).getValue().equals("v2");
+   }
+
    private void doTestPut(int number, String key, String value) throws Exception {
       for (int i = 0; i < number; i++) store.store(InternalEntryFactory.create(key + i, value + i));
       
@@ -182,4 +208,40 @@
       }
    }
 
+   class MockAsyncStore extends AsyncStore {
+      volatile boolean block = true;
+      final CountDownLatch v1Latch;
+      final CountDownLatch v2Latch;
+      final CountDownLatch endLatch;
+      final Object key;
+
+      MockAsyncStore(Object key, CountDownLatch v1Latch, CountDownLatch v2Latch, CountDownLatch endLatch, 
+               CacheStore delegate, AsyncStoreConfig asyncStoreConfig) {
+         super(delegate, asyncStoreConfig);
+         this.v1Latch = v1Latch;
+         this.v2Latch = v2Latch;
+         this.endLatch = endLatch;
+         this.key = key;
+      }
+
+      @Override
+      protected void applyModificationsSync(ConcurrentMap<Object, Modification> mods) throws CacheLoaderException {
+         if (mods.get(key) != null && block) {
+            if (log.isTraceEnabled()) log.trace("Wait for v1 latch");
+            try {
+               v2Latch.countDown();
+               block = false;
+               v1Latch.await(2, TimeUnit.SECONDS);
+            } catch (InterruptedException e) {
+            }
+            super.applyModificationsSync(mods);
+         } else if (mods.get(key) != null && !block) {
+            if (log.isTraceEnabled()) log.trace("Do v2 modification and unleash v1 latch");
+            super.applyModificationsSync(mods);
+            v1Latch.countDown();
+            endLatch.countDown();
+         }
+      }
+      
+   };
 }



More information about the infinispan-commits mailing list