[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